Skip to main content

reth_rpc_eth_api/helpers/
subscriptions.rs

1//! Streams subscriptions providers for `eth_subscribe`.
2
3use crate::{EthApiTypes, RpcConvert, RpcLog, RpcNodeCore, RpcReceipt};
4use alloy_consensus::{transaction::TxHashRef, BlockHeader, TxReceipt};
5use alloy_rpc_types_eth::{pubsub::TransactionReceiptsParams, Filter};
6use futures::StreamExt;
7use reth_chain_state::CanonStateSubscriptions;
8use reth_primitives_traits::TransactionMeta;
9use reth_rpc_convert::{transaction::ConvertReceiptInput, RpcHeader};
10use reth_rpc_eth_types::logs_utils;
11use tracing::error;
12
13/// Provides streams subscriptions for `eth_subscribe`.
14///
15/// Override the default methods to inject additional data sources (e.g. flashblocks).
16pub trait EthSubscriptions:
17    RpcNodeCore + EthApiTypes<RpcConvert: RpcConvert<Primitives = Self::Primitives>>
18{
19    /// Returns a stream that yields matching logs from canonical chain updates.
20    fn log_stream(
21        &self,
22        filter: Filter,
23    ) -> impl futures::Stream<Item = RpcLog<Self::NetworkTypes>> + Send + Unpin {
24        let converter = self.converter();
25        self.provider().canonical_state_stream().flat_map(move |canon_state| {
26            let reverted_chains = canon_state.reverted();
27            let committed_chain = canon_state.committed();
28            let reverted = reverted_chains.iter().flat_map(|chain| {
29                chain.blocks_and_receipts().map(|(block, receipts)| (block, receipts, true))
30            });
31            let committed = committed_chain
32                .blocks_and_receipts()
33                .map(|(block, receipts)| (block, receipts, false));
34            let mut all_logs = Vec::new();
35
36            for (block, receipts, removed) in reverted.chain(committed) {
37                let result = logs_utils::matching_block_logs_with_tx_hashes(
38                    converter,
39                    &filter,
40                    block.sealed_header(),
41                    block
42                        .transactions_recovered()
43                        .zip(receipts.iter())
44                        .map(|(tx, receipt)| (*tx.tx_hash(), receipt)),
45                    removed,
46                );
47                match result {
48                    Ok(logs) => all_logs.extend(logs),
49                    Err(err) => {
50                        error!(target = "rpc", %err, "Failed to convert logs");
51                    }
52                }
53            }
54
55            futures::stream::iter(all_logs)
56        })
57    }
58
59    /// Returns a stream that yields new block headers from canonical chain updates.
60    fn header_stream(
61        &self,
62    ) -> impl futures::Stream<Item = RpcHeader<Self::NetworkTypes>> + Send + Unpin {
63        let converter = self.converter();
64        self.provider().canonical_state_stream().flat_map(move |new_chain| {
65            let headers = new_chain
66                .committed()
67                .blocks_iter()
68                .filter_map(|block| {
69                    match converter
70                        .convert_header(block.clone_sealed_header(), Some(block.rlp_length()))
71                    {
72                        Ok(header) => Some(header),
73                        Err(err) => {
74                            error!(target = "rpc", %err, "Failed to convert header");
75                            None
76                        }
77                    }
78                })
79                .collect::<Vec<_>>();
80            futures::stream::iter(headers)
81        })
82    }
83
84    /// Returns a stream that yields matching transaction receipts from canonical chain updates.
85    fn transaction_receipts_stream(
86        &self,
87        filter: TransactionReceiptsParams,
88    ) -> impl futures::Stream<Item = Vec<RpcReceipt<Self::NetworkTypes>>> + Send + Unpin {
89        let converter = self.converter();
90        self.provider().canonical_state_stream().flat_map(move |new_chain| {
91            let results: Vec<_> = new_chain
92                .committed()
93                .blocks_and_receipts()
94                .filter_map(|(block, receipts)| {
95                    let block_hash = block.hash();
96                    let block_number = block.number();
97                    let base_fee = block.base_fee_per_gas();
98                    let excess_blob_gas = block.excess_blob_gas();
99                    let timestamp = block.timestamp();
100
101                    let mut gas_used: u64 = 0;
102                    let mut next_log_index: usize = 0;
103
104                    let inputs: Vec<_> = block
105                        .transactions_recovered()
106                        .zip(receipts.iter())
107                        .enumerate()
108                        .filter_map(|(idx, (tx, receipt))| {
109                            let gas_used_before = gas_used;
110                            let next_log_index_before = next_log_index;
111                            let cumulative_gas_used = receipt.cumulative_gas_used();
112
113                            gas_used = cumulative_gas_used;
114                            next_log_index += receipt.logs().len();
115
116                            let matches = match &filter.transaction_hashes {
117                                Some(hashes) if !hashes.is_empty() => hashes.contains(tx.tx_hash()),
118                                _ => true,
119                            };
120
121                            matches.then(|| ConvertReceiptInput {
122                                tx,
123                                gas_used: cumulative_gas_used - gas_used_before,
124                                next_log_index: next_log_index_before,
125                                meta: TransactionMeta {
126                                    tx_hash: *tx.tx_hash(),
127                                    index: idx as u64,
128                                    block_hash,
129                                    block_number,
130                                    base_fee,
131                                    excess_blob_gas,
132                                    timestamp,
133                                },
134                                receipt: receipt.clone(),
135                            })
136                        })
137                        .collect();
138
139                    if inputs.is_empty() {
140                        return None;
141                    }
142
143                    match converter.convert_receipts_with_block(inputs, block.sealed_block()) {
144                        Ok(rpc_receipts) => Some(rpc_receipts),
145                        Err(err) => {
146                            error!(target = "rpc", %err, "Failed to convert receipts");
147                            None
148                        }
149                    }
150                })
151                .collect();
152
153            futures::stream::iter(results)
154        })
155    }
156}