Skip to main content

reth_rpc_eth_types/
logs_utils.rs

1//! Helper functions for `reth_rpc_eth_api::EthFilterApiServer` implementation.
2//!
3//! Log parsing for building filter.
4
5use crate::EthApiError;
6use alloy_consensus::{transaction::TxHashRef, BlockHeader, TxReceipt};
7use alloy_primitives::TxHash;
8use alloy_rpc_types_eth::{Filter, Log};
9use jsonrpsee_types::ErrorObject;
10use reth_chainspec::ChainInfo;
11use reth_errors::ProviderError;
12use reth_primitives_traits::{
13    BlockBody, NodePrimitives, RecoveredBlock, SealedHeaderFor, SignedTransaction,
14};
15use reth_rpc_convert::{RpcConvert, RpcLog};
16use reth_storage_api::{BlockReader, ProviderBlock};
17use std::sync::Arc;
18use thiserror::Error;
19
20/// Returns all matching and converted logs of a block's receipts when the transaction hashes are
21/// known.
22pub fn matching_block_logs_with_tx_hashes<'a, I, R, C>(
23    converter: &C,
24    filter: &Filter,
25    header: &SealedHeaderFor<C::Primitives>,
26    tx_hashes_and_receipts: I,
27    removed: bool,
28) -> Result<Vec<RpcLog<C::Network>>, C::Error>
29where
30    I: IntoIterator<Item = (TxHash, &'a R)>,
31    R: TxReceipt<Log = alloy_primitives::Log> + 'a,
32    C: RpcConvert<Primitives: NodePrimitives<Receipt = R>>,
33{
34    let block_num_hash = header.num_hash();
35    if !filter.matches_block(&block_num_hash) {
36        return Ok(vec![])
37    }
38
39    // An empty address/topic set matches any bloom, so this can only skip blocks that have no
40    // matching logs.
41    if !filter.matches_bloom(header.logs_bloom()) {
42        return Ok(vec![])
43    }
44
45    let mut all_logs = Vec::new();
46    // Tracks the index of a log in the entire block.
47    let mut log_index: u64 = 0;
48
49    // Iterate over transaction hashes and receipts and append matching logs.
50    for (receipt_idx, (tx_hash, receipt)) in tx_hashes_and_receipts.into_iter().enumerate() {
51        for log in receipt.logs() {
52            if filter.matches(log) {
53                let log = Log {
54                    inner: log.clone(),
55                    block_hash: Some(block_num_hash.hash),
56                    block_number: Some(block_num_hash.number),
57                    transaction_hash: Some(tx_hash),
58                    // The transaction and receipt index is always the same.
59                    transaction_index: Some(receipt_idx as u64),
60                    log_index: Some(log_index),
61                    removed,
62                    block_timestamp: Some(header.timestamp()),
63                };
64                all_logs.push(converter.convert_log(log, receipt, header)?);
65            }
66            log_index += 1;
67        }
68    }
69    Ok(all_logs)
70}
71
72/// Helper enum to fetch a transaction either from a block or from the provider.
73#[derive(Debug)]
74pub enum ProviderOrBlock<'a, P: BlockReader> {
75    /// Provider
76    Provider(&'a P),
77    /// [`RecoveredBlock`]
78    Block(Arc<RecoveredBlock<ProviderBlock<P>>>),
79}
80
81/// Appends all matching and converted logs of a block's receipts.
82/// If the log matches, look up the corresponding transaction hash.
83pub fn append_matching_block_logs<P, C>(
84    all_logs: &mut Vec<RpcLog<C::Network>>,
85    converter: &C,
86    provider_or_block: ProviderOrBlock<'_, P>,
87    filter: &Filter,
88    header: &SealedHeaderFor<C::Primitives>,
89    receipts: &[P::Receipt],
90    removed: bool,
91) -> Result<(), EthApiError>
92where
93    P: BlockReader<Transaction: SignedTransaction>,
94    C: RpcConvert<Primitives: NodePrimitives<Block = ProviderBlock<P>, Receipt = P::Receipt>>,
95{
96    let block_num_hash = header.num_hash();
97    if !filter.matches_block(&block_num_hash) {
98        return Ok(());
99    }
100
101    // Tracks the index of a log in the entire block.
102    let mut log_index: u64 = 0;
103
104    // Lazy loaded number of the first transaction in the block.
105    // This is useful for blocks with multiple matching logs because it
106    // prevents re-querying the block body indices.
107    let mut loaded_first_tx_num = None;
108
109    // Iterate over receipts and append matching logs.
110    for (receipt_idx, receipt) in receipts.iter().enumerate() {
111        // The transaction hash of the current receipt.
112        let mut transaction_hash = None;
113
114        for log in receipt.logs() {
115            if filter.matches(log) {
116                // if this is the first match in the receipt's logs, look up the transaction hash
117                if transaction_hash.is_none() {
118                    transaction_hash = match &provider_or_block {
119                        ProviderOrBlock::Block(block) => {
120                            block.body().transactions().get(receipt_idx).map(|t| *t.tx_hash())
121                        }
122                        ProviderOrBlock::Provider(provider) => {
123                            let first_tx_num = match loaded_first_tx_num {
124                                Some(num) => num,
125                                None => {
126                                    let block_body_indices = provider
127                                        .block_body_indices(block_num_hash.number)?
128                                        .ok_or(ProviderError::BlockBodyIndicesNotFound(
129                                            block_num_hash.number,
130                                        ))?;
131                                    loaded_first_tx_num = Some(block_body_indices.first_tx_num);
132                                    block_body_indices.first_tx_num
133                                }
134                            };
135
136                            // This is safe because Transactions and Receipts have the same
137                            // keys.
138                            let transaction_id = first_tx_num + receipt_idx as u64;
139                            let transaction =
140                                provider.transaction_by_id(transaction_id)?.ok_or_else(|| {
141                                    ProviderError::TransactionNotFound(transaction_id.into())
142                                })?;
143
144                            Some(*transaction.tx_hash())
145                        }
146                    };
147                }
148
149                let log = Log {
150                    inner: log.clone(),
151                    block_hash: Some(block_num_hash.hash),
152                    block_number: Some(block_num_hash.number),
153                    transaction_hash,
154                    // The transaction and receipt index is always the same.
155                    transaction_index: Some(receipt_idx as u64),
156                    log_index: Some(log_index),
157                    removed,
158                    block_timestamp: Some(header.timestamp()),
159                };
160                let log = converter
161                    .convert_log(log, receipt, header)
162                    .map_err(|err| EthApiError::other(Into::<ErrorObject<'static>>::into(err)))?;
163                all_logs.push(log);
164            }
165            log_index += 1;
166        }
167    }
168
169    Ok(())
170}
171
172/// Computes the block range based on the filter range and current block numbers.
173///
174/// Returns an error for invalid ranges rather than silently clamping values.
175pub fn get_filter_block_range(
176    from_block: Option<u64>,
177    to_block: Option<u64>,
178    start_block: u64,
179    info: ChainInfo,
180) -> Result<(u64, u64), FilterBlockRangeError> {
181    let from_block_number = from_block.unwrap_or(start_block);
182    let to_block_number = to_block.unwrap_or(info.best_number);
183
184    // from > to is an invalid range
185    if from_block_number > to_block_number {
186        return Err(FilterBlockRangeError::InvalidBlockRange);
187    }
188
189    // we cannot query blocks that don't exist yet
190    if to_block_number > info.best_number {
191        return Err(FilterBlockRangeError::BlockRangeExceedsHead {
192            requested: to_block_number,
193            head: info.best_number,
194        });
195    }
196
197    Ok((from_block_number, to_block_number))
198}
199
200/// Errors for filter block range validation.
201///
202/// See also <https://github.com/ethereum/go-ethereum/blob/master/eth/filters/filter.go#L224-L230>.
203#[derive(Debug, Clone, Copy, PartialEq, Eq, Error)]
204pub enum FilterBlockRangeError {
205    /// `from_block > to_block`
206    #[error("invalid block range params")]
207    InvalidBlockRange,
208    /// Block range extends beyond current head
209    #[error("block range extends beyond current head block: requested {requested}, head {head}")]
210    BlockRangeExceedsHead {
211        /// The requested `toBlock` number
212        requested: u64,
213        /// The current head block number
214        head: u64,
215    },
216}
217
218#[cfg(test)]
219mod tests {
220    use alloy_rpc_types_eth::Filter;
221
222    use super::*;
223
224    #[test]
225    fn test_log_range_from_and_to() {
226        let from = 14000000u64;
227        let to = 14000100u64;
228        let info = ChainInfo { best_number: 15000000, ..Default::default() };
229        let range = get_filter_block_range(Some(from), Some(to), info.best_number, info).unwrap();
230        assert_eq!(range, (from, to));
231    }
232
233    #[test]
234    fn test_log_range_from() {
235        let from = 14000000u64;
236        let info = ChainInfo { best_number: 15000000, ..Default::default() };
237        let range = get_filter_block_range(Some(from), None, 0, info).unwrap();
238        assert_eq!(range, (from, info.best_number));
239    }
240
241    #[test]
242    fn test_log_range_to() {
243        let to = 14000000u64;
244        let start_block = 0u64;
245        let info = ChainInfo { best_number: 15000000, ..Default::default() };
246        let range = get_filter_block_range(None, Some(to), start_block, info).unwrap();
247        assert_eq!(range, (start_block, to));
248    }
249
250    #[test]
251    fn test_log_range_higher_error() {
252        // Range extends beyond head -> should error instead of clamping
253        let from = 15000001u64;
254        let to = 15000002u64;
255        let info = ChainInfo { best_number: 15000000, ..Default::default() };
256        let err = get_filter_block_range(Some(from), Some(to), info.best_number, info).unwrap_err();
257        assert_eq!(
258            err,
259            FilterBlockRangeError::BlockRangeExceedsHead { requested: to, head: info.best_number }
260        );
261    }
262
263    #[test]
264    fn test_log_range_to_below_start_error() {
265        // to_block < start_block, default from -> invalid range
266        let to = 14000000u64;
267        let info = ChainInfo { best_number: 15000000, ..Default::default() };
268        let err = get_filter_block_range(None, Some(to), info.best_number, info).unwrap_err();
269        assert_eq!(err, FilterBlockRangeError::InvalidBlockRange);
270    }
271
272    #[test]
273    fn test_log_range_empty() {
274        let info = ChainInfo { best_number: 15000000, ..Default::default() };
275        let range = get_filter_block_range(None, None, info.best_number, info).unwrap();
276
277        // no range given -> head
278        assert_eq!(range, (info.best_number, info.best_number));
279    }
280
281    #[test]
282    fn test_invalid_block_range_error() {
283        let from = 100;
284        let to = 50;
285        let info = ChainInfo { best_number: 150, ..Default::default() };
286        let err = get_filter_block_range(Some(from), Some(to), 0, info).unwrap_err();
287        assert_eq!(err, FilterBlockRangeError::InvalidBlockRange);
288    }
289
290    #[test]
291    fn test_block_range_exceeds_head_error() {
292        let from = 100;
293        let to = 200;
294        let info = ChainInfo { best_number: 150, ..Default::default() };
295        let err = get_filter_block_range(Some(from), Some(to), 0, info).unwrap_err();
296        assert_eq!(
297            err,
298            FilterBlockRangeError::BlockRangeExceedsHead { requested: to, head: info.best_number }
299        );
300    }
301
302    #[test]
303    fn parse_log_from_only() {
304        let s = r#"{"fromBlock":"0xf47a42","address":["0x7de93682b9b5d80d45cd371f7a14f74d49b0914c","0x0f00392fcb466c0e4e4310d81b941e07b4d5a079","0xebf67ab8cff336d3f609127e8bbf8bd6dd93cd81"],"topics":["0x0559884fd3a460db3073b7fc896cc77986f16e378210ded43186175bf646fc5f"]}"#;
305        let filter: Filter = serde_json::from_str(s).unwrap();
306
307        assert_eq!(filter.get_from_block(), Some(16022082));
308        assert!(filter.get_to_block().is_none());
309
310        let best_number = 17229427;
311        let info = ChainInfo { best_number, ..Default::default() };
312
313        let (from_block, to_block) = filter.block_option.as_range();
314
315        let start_block = info.best_number;
316
317        let (from_block_number, to_block_number) = get_filter_block_range(
318            from_block.and_then(alloy_rpc_types_eth::BlockNumberOrTag::as_number),
319            to_block.and_then(alloy_rpc_types_eth::BlockNumberOrTag::as_number),
320            start_block,
321            info,
322        )
323        .unwrap();
324        assert_eq!(from_block_number, 16022082);
325        assert_eq!(to_block_number, best_number);
326    }
327}