Skip to main content

solana_pubsub_client/nonblocking/
pubsub_client.rs

1//! A client for subscribing to messages from the RPC server.
2//!
3//! The [`PubsubClient`] implements [Solana WebSocket event
4//! subscriptions][spec].
5//!
6//! [spec]: https://solana.com/docs/rpc/websocket
7//!
8//! This is a nonblocking (async) API. For a blocking API use the synchronous
9//! client in [`crate::pubsub_client`].
10//!
11//! A single `PubsubClient` client may be used to subscribe to many events via
12//! subscription methods like [`PubsubClient::account_subscribe`]. These methods
13//! return a [`PubsubClientResult`] of a pair, the first element being a
14//! [`BoxStream`] of subscription-specific [`RpcResponse`]s, the second being an
15//! unsubscribe closure, an asynchronous function that can be called and
16//! `await`ed to unsubscribe.
17//!
18//! Note that `BoxStream` contains an immutable reference to the `PubsubClient`
19//! that created it. This makes `BoxStream` not `Send`, forcing it to stay in
20//! the same task as its `PubsubClient`. `PubsubClient` though is `Send` and
21//! `Sync`, and can be shared between tasks by putting it in an `Arc`. Thus
22//! one viable pattern to creating multiple subscriptions is:
23//!
24//! - create an `Arc<PubsubClient>`
25//! - spawn one task for each subscription, sharing the `PubsubClient`.
26//! - in each task:
27//!   - create a subscription
28//!   - send the `UnsubscribeFn` to another task to handle shutdown
29//!   - loop while receiving messages from the subscription
30//!
31//! This pattern is illustrated in the example below.
32//!
33//! By default the [`block_subscribe`] and [`vote_subscribe`] events are
34//! disabled on RPC nodes. They can be enabled by passing
35//! `--rpc-pubsub-enable-block-subscription` and
36//! `--rpc-pubsub-enable-vote-subscription` to `agave-validator`. When these
37//! methods are disabled, the RPC server will return a "Method not found" error
38//! message.
39//!
40//! [`block_subscribe`]: https://docs.rs/solana-rpc/latest/solana_rpc/rpc_pubsub/trait.RpcSolPubSub.html#tymethod.block_subscribe
41//! [`vote_subscribe`]: https://docs.rs/solana-rpc/latest/solana_rpc/rpc_pubsub/trait.RpcSolPubSub.html#tymethod.vote_subscribe
42//!
43//! # Examples
44//!
45//! Demo two async `PubsubClient` subscriptions with clean shutdown.
46//!
47//! This spawns a task for each subscription type, each of which subscribes and
48//! sends back a ready message and an unsubscribe channel (closure), then loops
49//! on printing messages. The main task then waits for user input before
50//! unsubscribing and waiting on the tasks.
51//!
52//! ```
53//! use anyhow::Result;
54//! use futures_util::StreamExt;
55//! use solana_pubsub_client::nonblocking::pubsub_client::PubsubClient;
56//! use std::sync::Arc;
57//! use tokio::io::AsyncReadExt;
58//! use tokio::sync::mpsc::unbounded_channel;
59//!
60//! pub async fn watch_subscriptions(
61//!     websocket_url: &str,
62//! ) -> Result<()> {
63//!
64//!     // Subscription tasks will send a ready signal when they have subscribed.
65//!     let (ready_sender, mut ready_receiver) = unbounded_channel::<()>();
66//!
67//!     // Channel to receive unsubscribe channels (actually closures).
68//!     // These receive a pair of `(Box<dyn FnOnce() -> BoxFuture<'static, ()> + Send>), &'static str)`,
69//!     // where the first is a closure to call to unsubscribe, the second is the subscription name.
70//!     let (unsubscribe_sender, mut unsubscribe_receiver) = unbounded_channel::<(_, &'static str)>();
71//!
72//!     // The `PubsubClient` must be `Arc`ed to share it across tasks.
73//!     let pubsub_client = Arc::new(PubsubClient::new(websocket_url).await?);
74//!
75//!     let mut join_handles = vec![];
76//!
77//!     join_handles.push(("slot", tokio::spawn({
78//!         // Clone things we need before moving their clones into the `async move` block.
79//!         //
80//!         // The subscriptions have to be made from the tasks that will receive the subscription messages,
81//!         // because the subscription streams hold a reference to the `PubsubClient`.
82//!         // Otherwise we would just subscribe on the main task and send the receivers out to other tasks.
83//!
84//!         let ready_sender = ready_sender.clone();
85//!         let unsubscribe_sender = unsubscribe_sender.clone();
86//!         let pubsub_client = Arc::clone(&pubsub_client);
87//!         async move {
88//!             let (mut slot_notifications, slot_unsubscribe) =
89//!                 pubsub_client.slot_subscribe().await?;
90//!
91//!             // With the subscription started,
92//!             // send a signal back to the main task for synchronization.
93//!             ready_sender.send(()).expect("channel");
94//!
95//!             // Send the unsubscribe closure back to the main task.
96//!             unsubscribe_sender.send((slot_unsubscribe, "slot"))
97//!                 .map_err(|e| format!("{}", e)).expect("channel");
98//!
99//!             // Drop senders so that the channels can close.
100//!             // The main task will receive until channels are closed.
101//!             drop((ready_sender, unsubscribe_sender));
102//!
103//!             // Do something with the subscribed messages.
104//!             // This loop will end once the main task unsubscribes.
105//!             while let Some(slot_info) = slot_notifications.next().await {
106//!                 println!("------------------------------------------------------------");
107//!                 println!("slot pubsub result: {:?}", slot_info);
108//!             }
109//!
110//!             // This type hint is necessary to allow the `async move` block to use `?`.
111//!             Ok::<_, anyhow::Error>(())
112//!         }
113//!     })));
114//!
115//!     join_handles.push(("root", tokio::spawn({
116//!         let ready_sender = ready_sender.clone();
117//!         let unsubscribe_sender = unsubscribe_sender.clone();
118//!         let pubsub_client = Arc::clone(&pubsub_client);
119//!         async move {
120//!             let (mut root_notifications, root_unsubscribe) =
121//!                 pubsub_client.root_subscribe().await?;
122//!
123//!             ready_sender.send(()).expect("channel");
124//!             unsubscribe_sender.send((root_unsubscribe, "root"))
125//!                 .map_err(|e| format!("{}", e)).expect("channel");
126//!             drop((ready_sender, unsubscribe_sender));
127//!
128//!             while let Some(root) = root_notifications.next().await {
129//!                 println!("------------------------------------------------------------");
130//!                 println!("root pubsub result: {:?}", root);
131//!             }
132//!
133//!             Ok::<_, anyhow::Error>(())
134//!         }
135//!     })));
136//!
137//!     // Drop these senders so that the channels can close
138//!     // and their receivers return `None` below.
139//!     drop(ready_sender);
140//!     drop(unsubscribe_sender);
141//!
142//!     // Wait until all subscribers are ready before proceeding with application logic.
143//!     while let Some(_) = ready_receiver.recv().await { }
144//!
145//!     // Do application logic here.
146//!
147//!     // Wait for input or some application-specific shutdown condition.
148//!     tokio::io::stdin().read_u8().await?;
149//!
150//!     // Unsubscribe from everything, which will shutdown all the tasks.
151//!     while let Some((unsubscribe, name)) = unsubscribe_receiver.recv().await {
152//!         println!("unsubscribing from {}", name);
153//!         unsubscribe().await
154//!     }
155//!
156//!     // Wait for the tasks.
157//!     for (name, handle) in join_handles {
158//!         println!("waiting on task {}", name);
159//!         if let Ok(Err(e)) = handle.await {
160//!             println!("task {} failed: {}", name, e);
161//!         }
162//!     }
163//!
164//!     Ok(())
165//! }
166//! # Ok::<(), anyhow::Error>(())
167//! ```
168
169use {
170    futures_util::{
171        future::{BoxFuture, FutureExt, ready},
172        sink::SinkExt,
173        stream::{BoxStream, StreamExt},
174    },
175    log::*,
176    serde::de::DeserializeOwned,
177    serde_json::{Map, Value, json},
178    solana_account_decoder_client_types::UiAccount,
179    solana_clock::Slot,
180    solana_pubkey::Pubkey,
181    solana_rpc_client_types::{
182        config::{
183            RpcAccountInfoConfig, RpcBlockSubscribeConfig, RpcBlockSubscribeFilter,
184            RpcProgramAccountsConfig, RpcSignatureSubscribeConfig, RpcTransactionLogsConfig,
185            RpcTransactionLogsFilter,
186        },
187        error_object::RpcErrorObject,
188        response::{
189            Response as RpcResponse, RpcBlockUpdate, RpcKeyedAccount, RpcLogsResponse,
190            RpcSignatureResult, RpcVote, SlotInfo, SlotUpdate,
191        },
192    },
193    solana_signature::Signature,
194    std::collections::BTreeMap,
195    thiserror::Error,
196    tokio::{
197        net::TcpStream,
198        sync::{mpsc, oneshot},
199        task::JoinHandle,
200        time::{Duration, sleep},
201    },
202    tokio_stream::wrappers::UnboundedReceiverStream,
203    tokio_tungstenite::{
204        MaybeTlsStream, WebSocketStream, connect_async,
205        tungstenite::{
206            Message,
207            protocol::frame::{CloseFrame, coding::CloseCode},
208        },
209    },
210    tungstenite::{
211        Bytes,
212        client::IntoClientRequest,
213        http::{StatusCode, header},
214    },
215};
216
217pub type PubsubClientResult<T = ()> = Result<T, PubsubClientError>;
218
219#[derive(Debug, Error)]
220pub enum PubsubClientError {
221    #[error("url parse error")]
222    UrlParseError(#[from] url::ParseError),
223
224    #[error("unable to connect to server")]
225    ConnectionError(Box<tokio_tungstenite::tungstenite::Error>),
226
227    #[error("websocket error")]
228    WsError(#[from] Box<tokio_tungstenite::tungstenite::Error>),
229
230    #[error("connection closed (({0})")]
231    ConnectionClosed(String),
232
233    #[error("json parse error")]
234    JsonParseError(#[from] serde_json::error::Error),
235
236    #[error("subscribe failed: {reason}")]
237    SubscribeFailed { reason: String, message: String },
238
239    #[error("unexpected message format: {0}")]
240    UnexpectedMessageError(String),
241
242    #[error("request failed: {reason}")]
243    RequestFailed { reason: String, message: String },
244
245    #[error("request error: {0}")]
246    RequestError(String),
247
248    #[error("could not find subscription id: {0}")]
249    UnexpectedSubscriptionResponse(String),
250
251    #[error("could not find node version: {0}")]
252    UnexpectedGetVersionResponse(String),
253}
254
255type UnsubscribeFn = Box<dyn FnOnce() -> BoxFuture<'static, ()> + Send>;
256type SubscribeResponseMsg =
257    Result<(mpsc::UnboundedReceiver<Value>, UnsubscribeFn), PubsubClientError>;
258type SubscribeRequestMsg = (String, Value, oneshot::Sender<SubscribeResponseMsg>);
259type SubscribeResult<'a, T> = PubsubClientResult<(BoxStream<'a, T>, UnsubscribeFn)>;
260type RequestMsg = (
261    String,
262    Value,
263    oneshot::Sender<Result<Value, PubsubClientError>>,
264);
265
266/// A client for subscribing to messages from the RPC server.
267///
268/// See the [module documentation][self].
269#[derive(Debug)]
270pub struct PubsubClient {
271    subscribe_sender: mpsc::UnboundedSender<SubscribeRequestMsg>,
272    _request_sender: mpsc::UnboundedSender<RequestMsg>,
273    shutdown_sender: oneshot::Sender<()>,
274    ws: JoinHandle<PubsubClientResult>,
275}
276
277async fn connect_with_retry<R: IntoClientRequest>(
278    request: R,
279) -> Result<WebSocketStream<MaybeTlsStream<TcpStream>>, Box<tungstenite::Error>> {
280    let mut connection_retries = 5;
281    let client_request = request.into_client_request().map_err(Box::new)?;
282    loop {
283        let result = connect_async(client_request.clone())
284            .await
285            .map(|(socket, _)| socket);
286        if let Err(tungstenite::Error::Http(response)) = &result
287            && response.status() == StatusCode::TOO_MANY_REQUESTS
288            && connection_retries > 0
289        {
290            let mut duration = Duration::from_millis(500);
291            if let Some(retry_after) = response.headers().get(header::RETRY_AFTER)
292                && let Ok(retry_after) = retry_after.to_str()
293                && let Ok(retry_after) = retry_after.parse::<u64>()
294                && retry_after < 120
295            {
296                duration = Duration::from_secs(retry_after);
297            }
298
299            connection_retries -= 1;
300            debug!(
301                "Too many requests: server responded with {response:?}, {connection_retries} \
302                 retries left, pausing for {duration:?}"
303            );
304
305            sleep(duration).await;
306            continue;
307        }
308        return result.map_err(Box::new);
309    }
310}
311
312impl PubsubClient {
313    pub async fn new<R: IntoClientRequest>(request: R) -> PubsubClientResult<Self> {
314        let client_request = request.into_client_request().map_err(Box::new)?;
315        let ws = connect_with_retry(client_request)
316            .await
317            .map_err(PubsubClientError::ConnectionError)?;
318
319        let (subscribe_sender, subscribe_receiver) = mpsc::unbounded_channel();
320        let (_request_sender, request_receiver) = mpsc::unbounded_channel();
321        let (shutdown_sender, shutdown_receiver) = oneshot::channel();
322
323        #[allow(clippy::used_underscore_binding)]
324        Ok(Self {
325            subscribe_sender,
326            _request_sender,
327            shutdown_sender,
328            ws: tokio::spawn(PubsubClient::run_ws(
329                ws,
330                subscribe_receiver,
331                request_receiver,
332                shutdown_receiver,
333            )),
334        })
335    }
336
337    pub async fn shutdown(self) -> PubsubClientResult {
338        let _ = self.shutdown_sender.send(());
339        self.ws.await.unwrap() // WS future should not be cancelled or panicked
340    }
341
342    async fn subscribe<'a, T>(&self, operation: &str, params: Value) -> SubscribeResult<'a, T>
343    where
344        T: DeserializeOwned + Send + 'a,
345    {
346        let (response_sender, response_receiver) = oneshot::channel();
347        self.subscribe_sender
348            .send((operation.to_string(), params, response_sender))
349            .map_err(|err| PubsubClientError::ConnectionClosed(err.to_string()))?;
350
351        let (notifications, unsubscribe) = response_receiver
352            .await
353            .map_err(|err| PubsubClientError::ConnectionClosed(err.to_string()))??;
354        Ok((
355            UnboundedReceiverStream::new(notifications)
356                .filter_map(|value| ready(serde_json::from_value::<T>(value).ok()))
357                .boxed(),
358            unsubscribe,
359        ))
360    }
361
362    /// Subscribe to account events.
363    ///
364    /// Receives messages of type [`UiAccount`] when an account's lamports or data changes.
365    ///
366    /// # RPC Reference
367    ///
368    /// This method corresponds directly to the [`accountSubscribe`] RPC method.
369    ///
370    /// [`accountSubscribe`]: https://solana.com/docs/rpc/websocket#accountsubscribe
371    pub async fn account_subscribe(
372        &self,
373        pubkey: &Pubkey,
374        config: Option<RpcAccountInfoConfig>,
375    ) -> SubscribeResult<'_, RpcResponse<UiAccount>> {
376        let params = json!([pubkey.to_string(), config]);
377        self.subscribe("account", params).await
378    }
379
380    /// Subscribe to block events.
381    ///
382    /// Receives messages of type [`RpcBlockUpdate`] when a block is confirmed or finalized.
383    ///
384    /// This method is disabled by default. It can be enabled by passing
385    /// `--rpc-pubsub-enable-block-subscription` to `agave-validator`.
386    ///
387    /// # RPC Reference
388    ///
389    /// This method corresponds directly to the [`blockSubscribe`] RPC method.
390    ///
391    /// [`blockSubscribe`]: https://solana.com/docs/rpc/websocket#blocksubscribe
392    pub async fn block_subscribe(
393        &self,
394        filter: RpcBlockSubscribeFilter,
395        config: Option<RpcBlockSubscribeConfig>,
396    ) -> SubscribeResult<'_, RpcResponse<RpcBlockUpdate>> {
397        self.subscribe("block", json!([filter, config])).await
398    }
399
400    /// Subscribe to transaction log events.
401    ///
402    /// Receives messages of type [`RpcLogsResponse`] when a transaction is committed.
403    ///
404    /// # RPC Reference
405    ///
406    /// This method corresponds directly to the [`logsSubscribe`] RPC method.
407    ///
408    /// [`logsSubscribe`]: https://solana.com/docs/rpc/websocket#logssubscribe
409    pub async fn logs_subscribe(
410        &self,
411        filter: RpcTransactionLogsFilter,
412        config: RpcTransactionLogsConfig,
413    ) -> SubscribeResult<'_, RpcResponse<RpcLogsResponse>> {
414        self.subscribe("logs", json!([filter, config])).await
415    }
416
417    /// Subscribe to program account events.
418    ///
419    /// Receives messages of type [`RpcKeyedAccount`] when an account owned
420    /// by the given program changes.
421    ///
422    /// # RPC Reference
423    ///
424    /// This method corresponds directly to the [`programSubscribe`] RPC method.
425    ///
426    /// [`programSubscribe`]: https://solana.com/docs/rpc/websocket#programsubscribe
427    pub async fn program_subscribe(
428        &self,
429        pubkey: &Pubkey,
430        config: Option<RpcProgramAccountsConfig>,
431    ) -> SubscribeResult<'_, RpcResponse<RpcKeyedAccount>> {
432        let params = json!([pubkey.to_string(), config]);
433        self.subscribe("program", params).await
434    }
435
436    /// Subscribe to vote events.
437    ///
438    /// Receives messages of type [`RpcVote`] when a new vote is observed. These
439    /// votes are observed prior to confirmation and may never be confirmed.
440    ///
441    /// This method is disabled by default. It can be enabled by passing
442    /// `--rpc-pubsub-enable-vote-subscription` to `agave-validator`.
443    ///
444    /// # RPC Reference
445    ///
446    /// This method corresponds directly to the [`voteSubscribe`] RPC method.
447    ///
448    /// [`voteSubscribe`]: https://solana.com/docs/rpc/websocket#votesubscribe
449    pub async fn vote_subscribe(&self) -> SubscribeResult<'_, RpcVote> {
450        self.subscribe("vote", json!([])).await
451    }
452
453    /// Subscribe to root events.
454    ///
455    /// Receives messages of type [`Slot`] when a new [root] is set by the
456    /// validator.
457    ///
458    /// [root]: https://solana.com/docs/terminology#root
459    ///
460    /// # RPC Reference
461    ///
462    /// This method corresponds directly to the [`rootSubscribe`] RPC method.
463    ///
464    /// [`rootSubscribe`]: https://solana.com/docs/rpc/websocket#rootsubscribe
465    pub async fn root_subscribe(&self) -> SubscribeResult<'_, Slot> {
466        self.subscribe("root", json!([])).await
467    }
468
469    /// Subscribe to transaction confirmation events.
470    ///
471    /// Receives messages of type [`RpcSignatureResult`] when a transaction
472    /// with the given signature is committed.
473    ///
474    /// This is a subscription to a single notification. It is automatically
475    /// cancelled by the server once the notification is sent.
476    ///
477    /// # RPC Reference
478    ///
479    /// This method corresponds directly to the [`signatureSubscribe`] RPC method.
480    ///
481    /// [`signatureSubscribe`]: https://solana.com/docs/rpc/websocket#signaturesubscribe
482    pub async fn signature_subscribe(
483        &self,
484        signature: &Signature,
485        config: Option<RpcSignatureSubscribeConfig>,
486    ) -> SubscribeResult<'_, RpcResponse<RpcSignatureResult>> {
487        let params = json!([signature.to_string(), config]);
488        self.subscribe("signature", params).await
489    }
490
491    /// Subscribe to slot events.
492    ///
493    /// Receives messages of type [`SlotInfo`] when processing of a slot begins.
494    ///
495    /// # RPC Reference
496    ///
497    /// This method corresponds directly to the [`slotSubscribe`] RPC method.
498    ///
499    /// [`slotSubscribe`]: https://solana.com/docs/rpc/websocket#slotsubscribe
500    pub async fn slot_subscribe(&self) -> SubscribeResult<'_, SlotInfo> {
501        self.subscribe("slot", json!([])).await
502    }
503
504    /// Subscribe to slot update events.
505    ///
506    /// Receives messages of type [`SlotUpdate`] when various updates to a slot occur.
507    ///
508    /// Note that this method operates differently than other subscriptions:
509    /// instead of sending the message to a receiver on a channel, it accepts a
510    /// `handler` callback that processes the message directly. This processing
511    /// occurs on another thread.
512    ///
513    /// # RPC Reference
514    ///
515    /// This method corresponds directly to the [`slotUpdatesSubscribe`] RPC method.
516    ///
517    /// [`slotUpdatesSubscribe`]: https://solana.com/docs/rpc/websocket#slotsupdatessubscribe
518    pub async fn slot_updates_subscribe(&self) -> SubscribeResult<'_, SlotUpdate> {
519        self.subscribe("slotsUpdates", json!([])).await
520    }
521
522    async fn run_ws(
523        mut ws: WebSocketStream<MaybeTlsStream<TcpStream>>,
524        mut subscribe_receiver: mpsc::UnboundedReceiver<SubscribeRequestMsg>,
525        mut request_receiver: mpsc::UnboundedReceiver<RequestMsg>,
526        mut shutdown_receiver: oneshot::Receiver<()>,
527    ) -> PubsubClientResult {
528        let mut request_id: u64 = 0;
529
530        let mut requests_subscribe = BTreeMap::new();
531        let mut requests_unsubscribe = BTreeMap::<u64, oneshot::Sender<()>>::new();
532        let mut other_requests = BTreeMap::new();
533        let mut subscriptions = BTreeMap::new();
534        let (unsubscribe_sender, mut unsubscribe_receiver) = mpsc::unbounded_channel();
535
536        loop {
537            tokio::select! {
538                // Send close on shutdown signal
539                _ = (&mut shutdown_receiver) => {
540                    let frame = CloseFrame { code: CloseCode::Normal, reason: "".into() };
541                    ws.send(Message::Close(Some(frame))).await.map_err(Box::new)?;
542                    ws.flush().await.map_err(Box::new)?;
543                    break;
544                },
545                // Send `Message::Ping` each 10s if no any other communication
546                () = sleep(Duration::from_secs(10)) => {
547                    ws.send(Message::Ping(Bytes::new())).await.map_err(Box::new)?;
548                },
549                // Read message for subscribe
550                Some((operation, params, response_sender)) = subscribe_receiver.recv() => {
551                    request_id += 1;
552                    let method = format!("{operation}Subscribe");
553                    let text = json!({"jsonrpc":"2.0","id":request_id,"method":method,"params":params}).to_string();
554                    ws.send(Message::Text(text.into())).await.map_err(Box::new)?;
555                    requests_subscribe.insert(request_id, (operation, response_sender));
556                },
557                // Read message for unsubscribe
558                Some((operation, sid, response_sender)) = unsubscribe_receiver.recv() => {
559                    subscriptions.remove(&sid);
560                    request_id += 1;
561                    let method = format!("{operation}Unsubscribe");
562                    let text = json!({"jsonrpc":"2.0","id":request_id,"method":method,"params":[sid]}).to_string();
563                    ws.send(Message::Text(text.into())).await.map_err(Box::new)?;
564                    requests_unsubscribe.insert(request_id, response_sender);
565                },
566                // Read message for other requests
567                Some((method, params, response_sender)) = request_receiver.recv() => {
568                    request_id += 1;
569                    let text = json!({"jsonrpc":"2.0","id":request_id,"method":method,"params":params}).to_string();
570                    ws.send(Message::Text(text.into())).await.map_err(Box::new)?;
571                    other_requests.insert(request_id, response_sender);
572                }
573                // Read incoming WebSocket message
574                next_msg = ws.next() => {
575                    let msg = match next_msg {
576                        Some(msg) => msg.map_err(Box::new)?,
577                        None => break,
578                    };
579                    trace!("ws.next(): {msg:?}");
580
581                    // Get text from the message
582                    let text = match msg {
583                        Message::Text(text) => text,
584                        Message::Binary(_data) => continue, // Ignore
585                        Message::Ping(data) => {
586                            ws.send(Message::Pong(data)).await.map_err(Box::new)?;
587                            continue
588                        },
589                        Message::Pong(_data) => continue,
590                        Message::Close(_frame) => break,
591                        Message::Frame(_frame) => continue,
592                    };
593
594
595                    let mut json: Map<String, Value> = serde_json::from_str(&text)?;
596
597                    // Subscribe/Unsubscribe response, example:
598                    // `{"jsonrpc":"2.0","result":5308752,"id":1}`
599                    if let Some(id) = json.get("id") {
600                        let id = id.as_u64().ok_or_else(|| {
601                            PubsubClientError::SubscribeFailed { reason: "invalid `id` field".into(), message: text.to_string() }
602                        })?;
603
604                        let err = json.get("error").map(|error_object| {
605                            match serde_json::from_value::<RpcErrorObject>(error_object.clone()) {
606                                Ok(rpc_error_object) => {
607                                    format!("{} ({})",  rpc_error_object.message, rpc_error_object.code)
608                                }
609                                Err(err) => format!(
610                                    "Failed to deserialize RPC error response: {} [{}]",
611                                    serde_json::to_string(error_object).unwrap(),
612                                    err
613                                )
614                            }
615                        });
616
617                        if let Some(response_sender) = other_requests.remove(&id) {
618                            match err {
619                                Some(reason) => {
620                                    let _ = response_sender.send(Err(PubsubClientError::RequestFailed { reason, message: text.to_string()}));
621                                },
622                                None => {
623                                    let json_result = json.get("result").ok_or_else(|| {
624                                        PubsubClientError::RequestFailed { reason: "missing `result` field".into(), message: text.to_string() }
625                                    })?;
626                                    if response_sender.send(Ok(json_result.clone())).is_err() {
627                                        break;
628                                    }
629                                }
630                            }
631                        } else if let Some(response_sender) = requests_unsubscribe.remove(&id) {
632                            let _ = response_sender.send(()); // do not care if receiver is closed
633                        } else if let Some((operation, response_sender)) = requests_subscribe.remove(&id) {
634                            match err {
635                                Some(reason) => {
636                                    let _ = response_sender.send(Err(PubsubClientError::SubscribeFailed { reason, message: text.to_string()}));
637                                },
638                                None => {
639                                    // Subscribe Id
640                                    let sid = json.get("result").and_then(Value::as_u64).ok_or_else(|| {
641                                        PubsubClientError::SubscribeFailed { reason: "invalid `result` field".into(), message: text.to_string() }
642                                    })?;
643
644                                    // Create notifications channel and unsubscribe function
645                                    let (notifications_sender, notifications_receiver) = mpsc::unbounded_channel();
646                                    let unsubscribe_sender = unsubscribe_sender.clone();
647                                    let unsubscribe = Box::new(move || async move {
648                                        let (response_sender, response_receiver) = oneshot::channel();
649                                        // do nothing if ws already closed
650                                        if unsubscribe_sender.send((operation, sid, response_sender)).is_ok() {
651                                            let _ = response_receiver.await; // channel can be closed only if ws is closed
652                                        }
653                                    }.boxed());
654
655                                    if response_sender.send(Ok((notifications_receiver, unsubscribe))).is_err() {
656                                        break;
657                                    }
658                                    subscriptions.insert(sid, notifications_sender);
659                                }
660                            }
661                        } else {
662                            error!("Unknown request id: {id}");
663                            break;
664                        }
665                        continue;
666                    }
667
668                    // Notification, example:
669                    // `{"jsonrpc":"2.0","method":"logsNotification","params":{"result":{...},"subscription":3114862}}`
670                    if let Some(Value::Object(params)) = json.get_mut("params")
671                        && let Some(sid) = params.get("subscription").and_then(Value::as_u64)
672                    {
673                        let mut unsubscribe_required = false;
674
675                        if let Some(notifications_sender) = subscriptions.get(&sid) {
676                            if let Some(result) = params.remove("result")
677                                && notifications_sender.send(result).is_err()
678                            {
679                                unsubscribe_required = true;
680                            }
681                        } else {
682                            unsubscribe_required = true;
683                        }
684
685                        if unsubscribe_required
686                            && let Some(Value::String(method)) = json.remove("method")
687                            && let Some(operation) = method.strip_suffix("Notification")
688                        {
689                            let (response_sender, _response_receiver) = oneshot::channel();
690                            let _ = unsubscribe_sender
691                                .send((operation.to_string(), sid, response_sender));
692                        }
693                    }
694                }
695            }
696        }
697
698        Ok(())
699    }
700}
701
702#[cfg(test)]
703mod tests {
704    // see client-test/test/client.rs
705}