Skip to main content

reth_rpc/eth/
pubsub.rs

1//! `eth_` `PubSub` RPC handler implementation
2
3use std::sync::Arc;
4
5use alloy_primitives::TxHash;
6use alloy_rpc_types_eth::{
7    pubsub::{
8        Params, PubSubSyncStatus, SubscriptionKind, SyncStatusMetadata, TransactionReceiptsParams,
9    },
10    Filter,
11};
12use futures::StreamExt;
13use jsonrpsee::{
14    server::SubscriptionMessage, types::ErrorObject, PendingSubscriptionSink, SubscriptionSink,
15};
16use reth_chain_state::CanonStateSubscriptions;
17use reth_network_api::NetworkInfo;
18use reth_rpc_convert::RpcHeader;
19use reth_rpc_eth_api::{
20    helpers::EthSubscriptions, pubsub::EthPubSubApiServer, RpcConvert, RpcLog, RpcNodeCore,
21    RpcTransaction,
22};
23use reth_rpc_server_types::result::{internal_rpc_err, invalid_params_rpc_err};
24use reth_storage_api::BlockNumReader;
25use reth_tasks::Runtime;
26use reth_transaction_pool::{NewTransactionEvent, TransactionPool};
27use serde::Serialize;
28use tokio_stream::{
29    wrappers::{BroadcastStream, ReceiverStream},
30    Stream,
31};
32use tracing::error;
33
34/// `Eth` pubsub RPC implementation.
35///
36/// This handles `eth_subscribe` RPC calls.
37#[derive(Clone)]
38pub struct EthPubSub<Eth> {
39    /// All nested fields bundled together.
40    inner: Arc<EthPubSubInner<Eth>>,
41}
42
43// === impl EthPubSub ===
44
45impl<Eth> EthPubSub<Eth> {
46    /// Creates a new, shareable instance.
47    pub fn new(eth_api: Eth, subscription_task_spawner: Runtime) -> Self {
48        let inner = EthPubSubInner { eth_api, subscription_task_spawner };
49        Self { inner: Arc::new(inner) }
50    }
51}
52
53impl<Eth> EthPubSub<Eth>
54where
55    Eth: EthSubscriptions,
56{
57    /// Returns the current sync status for the `syncing` subscription
58    pub fn sync_status(&self, is_syncing: bool) -> PubSubSyncStatus {
59        self.inner.sync_status(is_syncing)
60    }
61
62    /// Returns a stream that yields all transaction hashes emitted by the txpool.
63    pub fn pending_transaction_hashes_stream(&self) -> impl Stream<Item = TxHash> {
64        self.inner.pending_transaction_hashes_stream()
65    }
66
67    /// Returns a stream that yields all transactions emitted by the txpool.
68    pub fn full_pending_transaction_stream(
69        &self,
70    ) -> impl Stream<Item = NewTransactionEvent<<Eth::Pool as TransactionPool>::Transaction>> {
71        self.inner.full_pending_transaction_stream()
72    }
73
74    /// Returns a stream that yields new block headers.
75    pub fn new_headers_stream(&self) -> impl Stream<Item = RpcHeader<Eth::NetworkTypes>> {
76        self.inner.eth_api.header_stream()
77    }
78
79    /// Returns a stream that yields matching logs.
80    pub fn log_stream(&self, filter: Filter) -> impl Stream<Item = RpcLog<Eth::NetworkTypes>> {
81        self.inner.eth_api.log_stream(filter)
82    }
83
84    /// The actual handler for an accepted [`EthPubSub::subscribe`] call.
85    pub async fn handle_accepted(
86        &self,
87        accepted_sink: SubscriptionSink,
88        kind: SubscriptionKind,
89        params: Option<Params>,
90    ) -> Result<(), ErrorObject<'static>> {
91        #[allow(unreachable_patterns)]
92        match kind {
93            SubscriptionKind::NewHeads => {
94                pipe_from_stream(accepted_sink, self.new_headers_stream()).await
95            }
96            SubscriptionKind::Logs => {
97                // if no params are provided, used default filter params
98                let filter = match params {
99                    Some(Params::Logs(filter)) => *filter,
100                    Some(Params::Bool(_)) => {
101                        return Err(invalid_params_rpc_err("Invalid params for logs"))
102                    }
103                    _ => Default::default(),
104                };
105                pipe_from_stream(accepted_sink, self.log_stream(filter)).await
106            }
107            SubscriptionKind::NewPendingTransactions => {
108                if let Some(params) = params {
109                    match params {
110                        Params::Bool(true) => {
111                            // full transaction objects requested
112                            let stream = self.full_pending_transaction_stream().filter_map(|tx| {
113                                let tx_value = match self
114                                    .inner
115                                    .eth_api
116                                    .converter()
117                                    .fill_pending(tx.transaction.to_consensus())
118                                {
119                                    Ok(tx) => Some(tx),
120                                    Err(err) => {
121                                        error!(target = "rpc",
122                                            %err,
123                                            "Failed to fill transaction with block context"
124                                        );
125                                        None
126                                    }
127                                };
128                                std::future::ready(tx_value)
129                            });
130                            return pipe_from_stream(accepted_sink, stream).await
131                        }
132                        Params::Bool(false) | Params::None => {
133                            // only hashes requested
134                        }
135                        _ => {
136                            return Err(invalid_params_rpc_err(
137                                "Invalid params for newPendingTransactions",
138                            ))
139                        }
140                    }
141                }
142
143                pipe_from_stream(accepted_sink, self.pending_transaction_hashes_stream()).await
144            }
145            SubscriptionKind::Syncing => {
146                // get new block subscription
147                let mut canon_state = BroadcastStream::new(
148                    self.inner.eth_api.provider().subscribe_to_canonical_state(),
149                );
150                // get current sync status
151                let mut initial_sync_status = self.inner.eth_api.network().is_syncing();
152                let current_sub_res = self.sync_status(initial_sync_status);
153
154                // send the current status immediately
155                let msg = SubscriptionMessage::new(
156                    accepted_sink.method_name(),
157                    accepted_sink.subscription_id(),
158                    &current_sub_res,
159                )
160                .map_err(SubscriptionSerializeError::new)?;
161
162                if accepted_sink.send(msg).await.is_err() {
163                    return Ok(())
164                }
165
166                loop {
167                    // Sends only happen when the sync status changes, so a failed send cannot be
168                    // relied on to detect a closed subscription.
169                    tokio::select! {
170                        _ = accepted_sink.closed() => break,
171                        maybe_event = canon_state.next() => {
172                            if maybe_event.is_none() {
173                                break
174                            }
175                        }
176                    }
177
178                    let current_syncing = self.inner.eth_api.network().is_syncing();
179                    // Only send a new response if the sync status has changed
180                    if current_syncing != initial_sync_status {
181                        // Update the sync status on each new block
182                        initial_sync_status = current_syncing;
183
184                        // send a new message now that the status changed
185                        let sync_status = self.sync_status(current_syncing);
186                        let msg = SubscriptionMessage::new(
187                            accepted_sink.method_name(),
188                            accepted_sink.subscription_id(),
189                            &sync_status,
190                        )
191                        .map_err(SubscriptionSerializeError::new)?;
192
193                        if accepted_sink.send(msg).await.is_err() {
194                            break
195                        }
196                    }
197                }
198
199                Ok(())
200            }
201            SubscriptionKind::TransactionReceipts => {
202                let filter = match params {
203                    Some(Params::TransactionReceipts(filter)) => filter,
204                    None | Some(Params::None) => TransactionReceiptsParams::default(),
205                    _ => {
206                        return Err(invalid_params_rpc_err("Invalid params for transactionReceipts"))
207                    }
208                };
209
210                pipe_from_stream(
211                    accepted_sink,
212                    self.inner.eth_api.transaction_receipts_stream(filter),
213                )
214                .await
215            }
216            _ => Err(invalid_params_rpc_err("Unsupported subscription kind")),
217        }
218    }
219}
220
221#[async_trait::async_trait]
222impl<Eth> EthPubSubApiServer<RpcTransaction<Eth::NetworkTypes>> for EthPubSub<Eth>
223where
224    Eth: EthSubscriptions,
225{
226    /// Handler for `eth_subscribe`
227    async fn subscribe(
228        &self,
229        pending: PendingSubscriptionSink,
230        kind: SubscriptionKind,
231        params: Option<Params>,
232    ) -> jsonrpsee::core::SubscriptionResult {
233        let sink = pending.accept().await?;
234        let pubsub = self.clone();
235        self.inner.subscription_task_spawner.spawn_task(async move {
236            let _ = pubsub.handle_accepted(sink, kind, params).await;
237        });
238
239        Ok(())
240    }
241}
242
243/// Helper to convert a serde error into an [`ErrorObject`]
244#[derive(Debug, thiserror::Error)]
245#[error("Failed to serialize subscription item: {0}")]
246pub struct SubscriptionSerializeError(#[from] serde_json::Error);
247
248impl SubscriptionSerializeError {
249    const fn new(err: serde_json::Error) -> Self {
250        Self(err)
251    }
252}
253
254impl From<SubscriptionSerializeError> for ErrorObject<'static> {
255    fn from(value: SubscriptionSerializeError) -> Self {
256        internal_rpc_err(value.to_string())
257    }
258}
259
260/// Pipes all stream items to the subscription sink.
261async fn pipe_from_stream<T, St>(
262    sink: SubscriptionSink,
263    mut stream: St,
264) -> Result<(), ErrorObject<'static>>
265where
266    St: Stream<Item = T> + Unpin,
267    T: Serialize,
268{
269    loop {
270        tokio::select! {
271            _ = sink.closed() => {
272                // connection dropped
273                break Ok(())
274            },
275            maybe_item = stream.next() => {
276                let item = match maybe_item {
277                    Some(item) => item,
278                    None => {
279                        // stream ended
280                        break  Ok(())
281                    },
282                };
283                let msg = SubscriptionMessage::new(
284                    sink.method_name(),
285                    sink.subscription_id(),
286                    &item
287                ).map_err(SubscriptionSerializeError::new)?;
288
289                if sink.send(msg).await.is_err() {
290                    break Ok(());
291                }
292            }
293        }
294    }
295}
296
297impl<Eth> std::fmt::Debug for EthPubSub<Eth> {
298    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
299        f.debug_struct("EthPubSub").finish_non_exhaustive()
300    }
301}
302
303/// Container type `EthPubSub`
304#[derive(Clone)]
305struct EthPubSubInner<EthApi> {
306    /// The `eth` API.
307    eth_api: EthApi,
308    /// The type that's used to spawn subscription tasks.
309    subscription_task_spawner: Runtime,
310}
311
312// == impl EthPubSubInner ===
313
314impl<Eth> EthPubSubInner<Eth>
315where
316    Eth: RpcNodeCore<Provider: BlockNumReader>,
317{
318    /// Returns the current sync status for the `syncing` subscription
319    fn sync_status(&self, is_syncing: bool) -> PubSubSyncStatus {
320        if is_syncing {
321            let current_block = self
322                .eth_api
323                .provider()
324                .chain_info()
325                .map(|info| info.best_number)
326                .unwrap_or_default();
327            PubSubSyncStatus::Detailed(SyncStatusMetadata {
328                syncing: true,
329                starting_block: 0,
330                current_block,
331                highest_block: Some(current_block),
332            })
333        } else {
334            PubSubSyncStatus::Simple(false)
335        }
336    }
337}
338
339impl<Eth> EthPubSubInner<Eth>
340where
341    Eth: RpcNodeCore<Pool: TransactionPool>,
342{
343    /// Returns a stream that yields all transaction hashes emitted by the txpool.
344    fn pending_transaction_hashes_stream(&self) -> impl Stream<Item = TxHash> {
345        ReceiverStream::new(self.eth_api.pool().pending_transactions_listener())
346    }
347
348    /// Returns a stream that yields all transactions emitted by the txpool.
349    fn full_pending_transaction_stream(
350        &self,
351    ) -> impl Stream<Item = NewTransactionEvent<<Eth::Pool as TransactionPool>::Transaction>> {
352        self.eth_api.pool().new_pending_pool_transactions_listener()
353    }
354}