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
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 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}