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