reth_rpc_eth_api/helpers/
subscriptions.rs1use 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
13pub trait EthSubscriptions:
17 RpcNodeCore + EthApiTypes<RpcConvert: RpcConvert<Primitives = Self::Primitives>>
18{
19 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 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 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}