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.convert_header(block.clone_sealed_header(), block.rlp_length())
70                    {
71                        Ok(header) => Some(header),
72                        Err(err) => {
73                            error!(target = "rpc", %err, "Failed to convert header");
74                            None
75                        }
76                    }
77                })
78                .collect::<Vec<_>>();
79            futures::stream::iter(headers)
80        })
81    }
82
83    /// Returns a stream that yields matching transaction receipts from canonical chain updates.
84    fn transaction_receipts_stream(
85        &self,
86        filter: TransactionReceiptsParams,
87    ) -> impl futures::Stream<Item = Vec<RpcReceipt<Self::NetworkTypes>>> + Send + Unpin {
88        let converter = self.converter();
89        self.provider().canonical_state_stream().flat_map(move |new_chain| {
90            let results: Vec<_> = new_chain
91                .committed()
92                .blocks_and_receipts()
93                .filter_map(|(block, receipts)| {
94                    let block_hash = block.hash();
95                    let block_number = block.number();
96                    let base_fee = block.base_fee_per_gas();
97                    let excess_blob_gas = block.excess_blob_gas();
98                    let timestamp = block.timestamp();
99
100                    let mut gas_used: u64 = 0;
101                    let mut next_log_index: usize = 0;
102
103                    let inputs: Vec<_> = block
104                        .transactions_recovered()
105                        .zip(receipts.iter())
106                        .enumerate()
107                        .filter_map(|(idx, (tx, receipt))| {
108                            let gas_used_before = gas_used;
109                            let next_log_index_before = next_log_index;
110                            let cumulative_gas_used = receipt.cumulative_gas_used();
111
112                            gas_used = cumulative_gas_used;
113                            next_log_index += receipt.logs().len();
114
115                            let matches = match &filter.transaction_hashes {
116                                Some(hashes) if !hashes.is_empty() => hashes.contains(tx.tx_hash()),
117                                _ => true,
118                            };
119
120                            matches.then(|| ConvertReceiptInput {
121                                tx,
122                                gas_used: cumulative_gas_used - gas_used_before,
123                                next_log_index: next_log_index_before,
124                                meta: TransactionMeta {
125                                    tx_hash: *tx.tx_hash(),
126                                    index: idx as u64,
127                                    block_hash,
128                                    block_number,
129                                    base_fee,
130                                    excess_blob_gas,
131                                    timestamp,
132                                },
133                                receipt: receipt.clone(),
134                            })
135                        })
136                        .collect();
137
138                    if inputs.is_empty() {
139                        return None;
140                    }
141
142                    match converter.convert_receipts_with_block(inputs, block.sealed_block()) {
143                        Ok(rpc_receipts) => Some(rpc_receipts),
144                        Err(err) => {
145                            error!(target = "rpc", %err, "Failed to convert receipts");
146                            None
147                        }
148                    }
149                })
150                .collect();
151
152            futures::stream::iter(results)
153        })
154    }
155}