1use 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#[derive(Clone)]
38pub struct EthPubSub<Eth> {
39 inner: Arc<EthPubSubInner<Eth>>,
41}
42
43impl<Eth> EthPubSub<Eth> {
46 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 pub fn sync_status(&self, is_syncing: bool) -> PubSubSyncStatus {
59 self.inner.sync_status(is_syncing)
60 }
61
62 pub fn pending_transaction_hashes_stream(&self) -> impl Stream<Item = TxHash> {
64 self.inner.pending_transaction_hashes_stream()
65 }
66
67 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 pub fn new_headers_stream(&self) -> impl Stream<Item = RpcHeader<Eth::NetworkTypes>> {
76 self.inner.eth_api.header_stream()
77 }
78
79 pub fn log_stream(&self, filter: Filter) -> impl Stream<Item = RpcLog<Eth::NetworkTypes>> {
81 self.inner.eth_api.log_stream(filter)
82 }
83
84 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 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 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 }
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 let mut canon_state = BroadcastStream::new(
148 self.inner.eth_api.provider().subscribe_to_canonical_state(),
149 );
150 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 let msg = SubscriptionMessage::new(
156 accepted_sink.method_name(),
157 accepted_sink.subscription_id(),
158 ¤t_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 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 if current_syncing != initial_sync_status {
181 initial_sync_status = current_syncing;
183
184 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 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#[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
260async 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 break Ok(())
274 },
275 maybe_item = stream.next() => {
276 let item = match maybe_item {
277 Some(item) => item,
278 None => {
279 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#[derive(Clone)]
305struct EthPubSubInner<EthApi> {
306 eth_api: EthApi,
308 subscription_task_spawner: Runtime,
310}
311
312impl<Eth> EthPubSubInner<Eth>
315where
316 Eth: RpcNodeCore<Provider: BlockNumReader>,
317{
318 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 fn pending_transaction_hashes_stream(&self) -> impl Stream<Item = TxHash> {
345 ReceiverStream::new(self.eth_api.pool().pending_transactions_listener())
346 }
347
348 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}