Skip to main content

reth_rpc/eth/
filter.rs

1//! `eth_` `Filter` RPC handler implementation
2
3use alloy_consensus::BlockHeader;
4use alloy_eips::BlockNumberOrTag;
5use alloy_primitives::{Sealable, TxHash};
6use alloy_rpc_types_eth::{
7    Filter, FilterBlockOption, FilterChanges, FilterId, PendingTransactionFilterKind,
8};
9use async_trait::async_trait;
10use futures::{
11    future::TryFutureExt,
12    stream::{FuturesOrdered, StreamExt},
13    Future,
14};
15use itertools::Itertools;
16use jsonrpsee::{core::RpcResult, server::IdProvider};
17use reth_errors::ProviderError;
18use reth_primitives_traits::{NodePrimitives, SealedHeader};
19use reth_rpc_eth_api::{
20    helpers::{EthBlocks, LoadReceipt},
21    EngineEthFilter, EthApiTypes, EthFilterApiServer, FullEthApiTypes, QueryLimits, RpcConvert,
22    RpcLog, RpcNodeCoreExt, RpcTransaction,
23};
24use reth_rpc_eth_types::{
25    logs_utils::{self, append_matching_block_logs, ProviderOrBlock},
26    EthApiError, EthFilterConfig, EthStateCache, EthSubscriptionIdProvider,
27};
28use reth_rpc_server_types::{result::rpc_error_with_code, ToRpcResult};
29use reth_storage_api::{
30    BlockHashReader, BlockIdReader, BlockNumReader, BlockReader, HeaderProvider, ProviderBlock,
31    ProviderReceipt, ReceiptProvider,
32};
33use reth_tasks::Runtime;
34use reth_transaction_pool::{NewSubpoolTransactionStream, PoolTransaction, TransactionPool};
35use std::{
36    collections::{HashMap, VecDeque},
37    fmt,
38    iter::{Peekable, StepBy},
39    ops::RangeInclusive,
40    pin::Pin,
41    sync::Arc,
42    time::{Duration, Instant},
43};
44use tokio::{
45    sync::{mpsc::Receiver, oneshot, Mutex},
46    time::MissedTickBehavior,
47};
48use tracing::{debug, error, trace};
49
50impl<Eth> EngineEthFilter<RpcLog<Eth::NetworkTypes>> for EthFilter<Eth>
51where
52    Eth: FullEthApiTypes
53        + RpcNodeCoreExt<Provider: BlockIdReader>
54        + LoadReceipt
55        + EthBlocks
56        + 'static,
57{
58    /// Returns logs matching given filter object, no query limits
59    fn logs(
60        &self,
61        filter: Filter,
62        limits: QueryLimits,
63    ) -> impl Future<Output = RpcResult<Vec<RpcLog<Eth::NetworkTypes>>>> + Send {
64        trace!(target: "rpc::eth", "Serving eth_getLogs");
65        self.logs_for_filter(filter, limits).map_err(|e| e.into())
66    }
67}
68
69/// Threshold for deciding between cached and range mode processing
70const CACHED_MODE_BLOCK_THRESHOLD: u64 = 250;
71
72/// Threshold for bloom filter matches that triggers reduced caching
73const HIGH_BLOOM_MATCH_THRESHOLD: usize = 20;
74
75/// Threshold for bloom filter matches that triggers moderately reduced caching
76const MODERATE_BLOOM_MATCH_THRESHOLD: usize = 10;
77
78/// Minimum block count to apply bloom filter match adjustments
79const BLOOM_ADJUSTMENT_MIN_BLOCKS: u64 = 100;
80
81/// The maximum number of headers we read at once when handling a range filter.
82const MAX_HEADERS_RANGE: u64 = 1_000; // with ~530bytes per header this is ~500kb
83
84// Cached mode is only reachable for ranges that fit into a single header window, which is what
85// keeps the mode decision independent of how a range is split.
86const _: () = assert!(CACHED_MODE_BLOCK_THRESHOLD <= MAX_HEADERS_RANGE);
87
88/// Minimum number of bloom matching blocks in a header window for fetching their receipts in
89/// parallel
90const PARALLEL_PROCESSING_THRESHOLD: usize = 1000;
91
92/// Default concurrency for parallel processing
93const DEFAULT_PARALLEL_CONCURRENCY: usize = 4;
94
95/// Maximum number of blocks whose receipts the parallel fetching holds in memory at once
96const MAX_PARALLEL_BATCH_SIZE: usize = 256;
97
98/// `Eth` filter RPC implementation.
99///
100/// This type handles `eth_` rpc requests related to filters (`eth_getLogs`).
101pub struct EthFilter<Eth: EthApiTypes> {
102    /// All nested fields bundled together
103    inner: Arc<EthFilterInner<Eth>>,
104}
105
106impl<Eth> Clone for EthFilter<Eth>
107where
108    Eth: EthApiTypes,
109{
110    fn clone(&self) -> Self {
111        Self { inner: self.inner.clone() }
112    }
113}
114
115impl<Eth> EthFilter<Eth>
116where
117    Eth: EthApiTypes + 'static,
118{
119    /// Creates a new, shareable instance.
120    ///
121    /// This uses the given pool to get notified about new transactions, the provider to interact
122    /// with the blockchain, the cache to fetch cacheable data, like the logs.
123    ///
124    /// See also [`EthFilterConfig`].
125    ///
126    /// This also spawns a task that periodically clears stale filters.
127    ///
128    /// # Create a new instance with [`EthApi`](crate::EthApi)
129    ///
130    /// ```no_run
131    /// use reth_evm_ethereum::EthEvmConfig;
132    /// use reth_network_api::noop::NoopNetwork;
133    /// use reth_provider::noop::NoopProvider;
134    /// use reth_rpc::{EthApi, EthFilter};
135    /// use reth_tasks::Runtime;
136    /// use reth_transaction_pool::noop::NoopTransactionPool;
137    /// let eth_api = EthApi::builder(
138    ///     NoopProvider::default(),
139    ///     NoopTransactionPool::default(),
140    ///     NoopNetwork::default(),
141    ///     EthEvmConfig::mainnet(),
142    /// )
143    /// .build();
144    /// let filter = EthFilter::new(eth_api, Default::default(), Runtime::test());
145    /// ```
146    pub fn new(eth_api: Eth, config: EthFilterConfig, task_spawner: Runtime) -> Self {
147        let EthFilterConfig { max_blocks_per_filter, max_logs_per_response, stale_filter_ttl } =
148            config;
149        let inner = EthFilterInner {
150            eth_api,
151            active_filters: ActiveFilters::new(),
152            id_provider: Arc::new(EthSubscriptionIdProvider::default()),
153            max_headers_range: MAX_HEADERS_RANGE,
154            task_spawner,
155            stale_filter_ttl,
156            query_limits: QueryLimits { max_blocks_per_filter, max_logs_per_response },
157        };
158
159        let eth_filter = Self { inner: Arc::new(inner) };
160
161        let this = eth_filter.clone();
162        eth_filter.inner.task_spawner.spawn_critical_task(
163            "eth-filters_stale-filters-clean",
164            async move {
165                this.watch_and_clear_stale_filters().await;
166            },
167        );
168
169        eth_filter
170    }
171
172    /// Returns all currently active filters
173    pub fn active_filters(&self) -> &ActiveFilters<RpcTransaction<Eth::NetworkTypes>> {
174        &self.inner.active_filters
175    }
176
177    /// Endless future that [`Self::clear_stale_filters`] every `stale_filter_ttl` interval.
178    /// Nonetheless, this endless future frees the thread at every await point.
179    async fn watch_and_clear_stale_filters(&self) {
180        let mut interval = tokio::time::interval_at(
181            tokio::time::Instant::now() + self.inner.stale_filter_ttl,
182            self.inner.stale_filter_ttl,
183        );
184        interval.set_missed_tick_behavior(MissedTickBehavior::Delay);
185        loop {
186            interval.tick().await;
187            self.clear_stale_filters(Instant::now()).await;
188        }
189    }
190
191    /// Clears all filters that have not been polled for longer than the configured
192    /// `stale_filter_ttl` at the given instant.
193    pub async fn clear_stale_filters(&self, now: Instant) {
194        trace!(target: "rpc::eth", "clear stale filters");
195        let mut filters = self.active_filters().inner.lock().await;
196        filters.retain(|id, filter| {
197            let is_valid = (now - filter.last_poll_timestamp) < self.inner.stale_filter_ttl;
198
199            if !is_valid {
200                trace!(target: "rpc::eth", "evict filter with id: {:?}", id);
201            }
202
203            is_valid
204        });
205        filters.shrink_to_fit();
206    }
207}
208
209impl<Eth> EthFilter<Eth>
210where
211    Eth: FullEthApiTypes<Provider: BlockReader + BlockIdReader>
212        + RpcNodeCoreExt
213        + LoadReceipt
214        + EthBlocks
215        + 'static,
216{
217    /// Access the underlying provider.
218    fn provider(&self) -> &Eth::Provider {
219        self.inner.eth_api.provider()
220    }
221
222    /// Access the underlying pool.
223    fn pool(&self) -> &Eth::Pool {
224        self.inner.eth_api.pool()
225    }
226
227    /// Returns all the filter changes for the given id, if any
228    pub async fn filter_changes(
229        &self,
230        id: FilterId,
231    ) -> Result<
232        FilterChanges<RpcTransaction<Eth::NetworkTypes>, RpcLog<Eth::NetworkTypes>>,
233        EthFilterError,
234    > {
235        let info = self.provider().chain_info()?;
236        let best_number = info.best_number;
237
238        // start_block is the block from which we should start fetching changes, the next block from
239        // the last time changes were polled, in other words the best block at last poll + 1
240        let (start_block, kind) = {
241            let mut filters = self.inner.active_filters.inner.lock().await;
242            let filter = filters.get_mut(&id).ok_or(EthFilterError::FilterNotFound(id))?;
243
244            if filter.block > best_number {
245                // no new blocks since the last poll
246                return Ok(FilterChanges::Empty)
247            }
248
249            // update filter
250            // we fetch all changes from [filter.block..best_block], so we advance the filter's
251            // block to `best_block +1`, the next from which we should start fetching changes again
252            let mut block = best_number + 1;
253            std::mem::swap(&mut filter.block, &mut block);
254            filter.last_poll_timestamp = Instant::now();
255
256            (block, filter.kind.clone())
257        };
258
259        match kind {
260            FilterKind::PendingTransaction(filter) => Ok(match filter.drain().await {
261                FilterChanges::Empty => FilterChanges::Empty,
262                FilterChanges::Hashes(hashes) => FilterChanges::Hashes(hashes),
263                FilterChanges::Transactions(transactions) => {
264                    FilterChanges::Transactions(transactions)
265                }
266                FilterChanges::Logs(_) => unreachable!("pending transaction filter returned logs"),
267            }),
268            FilterKind::Block => {
269                // Note: we need to fetch the block hashes from inclusive range
270                // [start_block..best_block]
271                let end_block = best_number + 1;
272                let block_hashes =
273                    self.provider().canonical_hashes_range(start_block, end_block).map_err(
274                        |_| EthApiError::HeaderRangeNotFound(start_block.into(), end_block.into()),
275                    )?;
276                Ok(FilterChanges::Hashes(block_hashes))
277            }
278            FilterKind::Log(filter) => {
279                let (from_block_number, to_block_number) = match filter.block_option {
280                    FilterBlockOption::Range { from_block, to_block } => {
281                        let from = from_block
282                            .map(|num| self.provider().convert_block_number(num))
283                            .transpose()?
284                            .flatten();
285                        let to = to_block
286                            .map(|num| self.provider().convert_block_number(num))
287                            .transpose()?
288                            .flatten();
289                        logs_utils::get_filter_block_range(from, to, start_block, info)?
290                    }
291                    FilterBlockOption::AtBlockHash(block_hash) => {
292                        // blockHash is equivalent to fromBlock = toBlock = the block number with
293                        // hash blockHash
294                        // get_logs_in_block_range is inclusive
295                        let block_number = self
296                            .provider()
297                            .block_number(block_hash)?
298                            .ok_or(ProviderError::HeaderNotFound(block_hash.into()))?;
299                        (block_number, block_number)
300                    }
301                };
302                let logs = self
303                    .inner
304                    .clone()
305                    .get_logs_in_block_range(
306                        *filter,
307                        from_block_number,
308                        to_block_number,
309                        self.inner.query_limits,
310                    )
311                    .await?;
312                Ok(FilterChanges::Logs(logs))
313            }
314        }
315    }
316
317    /// Returns an array of all logs matching filter with given id.
318    ///
319    /// Returns an error if no matching log filter exists.
320    ///
321    /// Handler for `eth_getFilterLogs`
322    pub async fn filter_logs(
323        &self,
324        id: FilterId,
325    ) -> Result<Vec<RpcLog<Eth::NetworkTypes>>, EthFilterError> {
326        let filter = {
327            let mut filters = self.inner.active_filters.inner.lock().await;
328            let filter =
329                filters.get_mut(&id).ok_or_else(|| EthFilterError::FilterNotFound(id.clone()))?;
330            if let FilterKind::Log(ref inner_filter) = filter.kind {
331                filter.last_poll_timestamp = Instant::now();
332                *inner_filter.clone()
333            } else {
334                // Not a log filter
335                return Err(EthFilterError::FilterNotFound(id))
336            }
337        };
338
339        self.logs_for_filter(filter, self.inner.query_limits).await
340    }
341
342    /// Returns logs matching given filter object.
343    async fn logs_for_filter(
344        &self,
345        filter: Filter,
346        limits: QueryLimits,
347    ) -> Result<Vec<RpcLog<Eth::NetworkTypes>>, EthFilterError> {
348        self.inner.clone().logs_for_filter(filter, limits).await
349    }
350}
351
352#[async_trait]
353impl<Eth> EthFilterApiServer<RpcTransaction<Eth::NetworkTypes>, RpcLog<Eth::NetworkTypes>>
354    for EthFilter<Eth>
355where
356    Eth: FullEthApiTypes + RpcNodeCoreExt + LoadReceipt + EthBlocks + 'static,
357{
358    /// Handler for `eth_newFilter`
359    async fn new_filter(&self, filter: Filter) -> RpcResult<FilterId> {
360        trace!(target: "rpc::eth", "Serving eth_newFilter");
361        self.inner
362            .install_filter(FilterKind::<RpcTransaction<Eth::NetworkTypes>>::Log(Box::new(filter)))
363            .await
364    }
365
366    /// Handler for `eth_newBlockFilter`
367    async fn new_block_filter(&self) -> RpcResult<FilterId> {
368        trace!(target: "rpc::eth", "Serving eth_newBlockFilter");
369        self.inner.install_filter(FilterKind::<RpcTransaction<Eth::NetworkTypes>>::Block).await
370    }
371
372    /// Handler for `eth_newPendingTransactionFilter`
373    async fn new_pending_transaction_filter(
374        &self,
375        kind: Option<PendingTransactionFilterKind>,
376    ) -> RpcResult<FilterId> {
377        trace!(target: "rpc::eth", "Serving eth_newPendingTransactionFilter");
378
379        let transaction_kind = match kind.unwrap_or_default() {
380            PendingTransactionFilterKind::Hashes => {
381                let receiver = self.pool().pending_transactions_listener();
382                let pending_txs_receiver = PendingTransactionsReceiver::new(receiver);
383                FilterKind::PendingTransaction(PendingTransactionKind::Hashes(pending_txs_receiver))
384            }
385            PendingTransactionFilterKind::Full => {
386                let stream = self.pool().new_pending_pool_transactions_listener();
387                let full_txs_receiver = FullTransactionsReceiver::new(
388                    stream,
389                    dyn_clone::clone(self.inner.eth_api.converter()),
390                );
391                FilterKind::PendingTransaction(PendingTransactionKind::FullTransaction(Arc::new(
392                    full_txs_receiver,
393                )))
394            }
395        };
396
397        // Install the filter and propagate any errors
398        self.inner.install_filter(transaction_kind).await
399    }
400
401    /// Handler for `eth_getFilterChanges`
402    async fn filter_changes(
403        &self,
404        id: FilterId,
405    ) -> RpcResult<FilterChanges<RpcTransaction<Eth::NetworkTypes>, RpcLog<Eth::NetworkTypes>>>
406    {
407        trace!(target: "rpc::eth", "Serving eth_getFilterChanges");
408        Ok(Self::filter_changes(self, id).await?)
409    }
410
411    /// Returns an array of all logs matching filter with given id.
412    ///
413    /// Returns an error if no matching log filter exists.
414    ///
415    /// Handler for `eth_getFilterLogs`
416    async fn filter_logs(&self, id: FilterId) -> RpcResult<Vec<RpcLog<Eth::NetworkTypes>>> {
417        trace!(target: "rpc::eth", "Serving eth_getFilterLogs");
418        Ok(Self::filter_logs(self, id).await?)
419    }
420
421    /// Handler for `eth_uninstallFilter`
422    async fn uninstall_filter(&self, id: FilterId) -> RpcResult<bool> {
423        trace!(target: "rpc::eth", "Serving eth_uninstallFilter");
424        let mut filters = self.inner.active_filters.inner.lock().await;
425        if filters.remove(&id).is_some() {
426            trace!(target: "rpc::eth::filter", ?id, "uninstalled filter");
427            Ok(true)
428        } else {
429            Ok(false)
430        }
431    }
432
433    /// Returns logs matching given filter object.
434    ///
435    /// Handler for `eth_getLogs`
436    async fn logs(&self, filter: Filter) -> RpcResult<Vec<RpcLog<Eth::NetworkTypes>>> {
437        trace!(target: "rpc::eth", "Serving eth_getLogs");
438        Ok(self.logs_for_filter(filter, self.inner.query_limits).await?)
439    }
440}
441
442impl<Eth> std::fmt::Debug for EthFilter<Eth>
443where
444    Eth: EthApiTypes,
445{
446    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
447        f.debug_struct("EthFilter").finish_non_exhaustive()
448    }
449}
450
451/// Container type `EthFilter`
452#[derive(Debug)]
453struct EthFilterInner<Eth: EthApiTypes> {
454    /// Inner `eth` API implementation.
455    eth_api: Eth,
456    /// All currently installed filters.
457    active_filters: ActiveFilters<RpcTransaction<Eth::NetworkTypes>>,
458    /// Provides ids to identify filters
459    id_provider: Arc<dyn IdProvider>,
460    /// limits for logs queries
461    query_limits: QueryLimits,
462    /// maximum number of headers to read at once for range filter
463    max_headers_range: u64,
464    /// The type that can spawn tasks.
465    task_spawner: Runtime,
466    /// Duration since the last filter poll, after which the filter is considered stale
467    stale_filter_ttl: Duration,
468}
469
470impl<Eth> EthFilterInner<Eth>
471where
472    Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
473        + EthApiTypes<NetworkTypes: reth_rpc_eth_api::types::RpcTypes>
474        + LoadReceipt
475        + EthBlocks
476        + 'static,
477{
478    /// Access the underlying provider.
479    fn provider(&self) -> &Eth::Provider {
480        self.eth_api.provider()
481    }
482
483    /// Access the underlying [`EthStateCache`].
484    fn eth_cache(&self) -> &EthStateCache<Eth::Primitives> {
485        self.eth_api.cache()
486    }
487
488    /// Returns logs matching given filter object.
489    async fn logs_for_filter(
490        self: Arc<Self>,
491        filter: Filter,
492        limits: QueryLimits,
493    ) -> Result<Vec<RpcLog<Eth::NetworkTypes>>, EthFilterError> {
494        match filter.block_option {
495            FilterBlockOption::AtBlockHash(block_hash) => {
496                // First try to get cached block and receipts, as it's likely they're already cached
497                let Some((receipts, maybe_block)) =
498                    self.eth_cache().get_receipts_and_maybe_block(block_hash).await?
499                else {
500                    // the block itself may still exist with its receipts pruned
501                    return Err(match self.provider().block_number(block_hash)? {
502                        Some(number) => {
503                            let earliest_available = self.provider().earliest_block_number()?;
504                            if number < earliest_available {
505                                EthApiError::PrunedHistoryUnavailable {
506                                    requested: number,
507                                    earliest_available,
508                                }
509                                .into()
510                            } else {
511                                EthFilterError::ReceiptsUnavailable(number)
512                            }
513                        }
514                        None => ProviderError::HeaderNotFound(block_hash.into()).into(),
515                    })
516                };
517
518                let header = if let Some(block) = &maybe_block {
519                    block.clone_sealed_header()
520                } else {
521                    let header = self
522                        .provider()
523                        .header_by_hash_or_number(block_hash.into())?
524                        .ok_or_else(|| ProviderError::HeaderNotFound(block_hash.into()))?;
525                    SealedHeader::new(header, block_hash)
526                };
527
528                // Check if the block has been pruned (EIP-4444)
529                let earliest_block = self.provider().earliest_block_number()?;
530                if header.number() < earliest_block {
531                    return Err(EthApiError::PrunedHistoryUnavailable {
532                        requested: header.number(),
533                        earliest_available: earliest_block,
534                    }
535                    .into());
536                }
537
538                if !filter.matches_bloom(header.logs_bloom()) {
539                    return Ok(Vec::new())
540                }
541
542                let mut all_logs = Vec::new();
543                append_matching_block_logs(
544                    &mut all_logs,
545                    self.eth_api.converter(),
546                    maybe_block
547                        .map(ProviderOrBlock::Block)
548                        .unwrap_or_else(|| ProviderOrBlock::Provider(self.provider())),
549                    &filter,
550                    &header,
551                    &receipts,
552                    false,
553                )?;
554                Ok(all_logs)
555            }
556            FilterBlockOption::Range { from_block, to_block } => {
557                // Handle special case where from block is pending
558                if from_block.is_some_and(|b| b.is_pending()) {
559                    let to_block = to_block.unwrap_or(BlockNumberOrTag::Pending);
560                    if !(to_block.is_pending() || to_block.is_number()) {
561                        // always empty range
562                        return Ok(Vec::new());
563                    }
564                    // Try to get pending block and receipts
565                    if let Ok(Some(pending_block)) = self.eth_api.local_pending_block().await {
566                        if let BlockNumberOrTag::Number(to_block) = to_block &&
567                            to_block < pending_block.block.number()
568                        {
569                            // this block range is empty based on the user input
570                            return Ok(Vec::new());
571                        }
572
573                        let info = self.provider().chain_info()?;
574                        if pending_block.block.number() > info.best_number {
575                            // only consider the pending block if it is ahead of the chain
576                            let mut all_logs = Vec::new();
577                            let header = pending_block.block.clone_sealed_header();
578                            append_matching_block_logs(
579                                &mut all_logs,
580                                self.eth_api.converter(),
581                                ProviderOrBlock::<Eth::Provider>::Block(pending_block.block),
582                                &filter,
583                                &header,
584                                &pending_block.receipts,
585                                false, // removed = false for pending blocks
586                            )?;
587                            return Ok(all_logs)
588                        }
589                    }
590                }
591
592                let info = self.provider().chain_info()?;
593                let start_block = info.best_number;
594                // Without a pending block to serve, a `pending` bound resolves to the head on both
595                // ends instead of to whatever payload the engine currently holds
596                let from = from_block
597                    .filter(|num| !num.is_pending())
598                    .map(|num| self.provider().convert_block_number(num))
599                    .transpose()?
600                    .flatten();
601                let to = to_block
602                    .filter(|num| !num.is_pending())
603                    .map(|num| self.provider().convert_block_number(num))
604                    .transpose()?
605                    .flatten();
606
607                // Return error if toBlock exceeds current head
608                if let Some(t) = to &&
609                    t > info.best_number
610                {
611                    return Err(EthFilterError::BlockRangeExceedsHead {
612                        requested: t,
613                        head: info.best_number,
614                    });
615                }
616
617                let (from_block_number, to_block_number) =
618                    logs_utils::get_filter_block_range(from, to, start_block, info)?;
619
620                // Check if the requested range overlaps with pruned history (EIP-4444)
621                let earliest_block = self.provider().earliest_block_number()?;
622                if from_block_number < earliest_block {
623                    return Err(EthApiError::PrunedHistoryUnavailable {
624                        requested: from_block_number,
625                        earliest_available: earliest_block,
626                    }
627                    .into());
628                }
629
630                self.get_logs_in_block_range(filter, from_block_number, to_block_number, limits)
631                    .await
632            }
633        }
634    }
635
636    /// Installs a new filter and returns the new identifier.
637    async fn install_filter(
638        &self,
639        kind: FilterKind<RpcTransaction<Eth::NetworkTypes>>,
640    ) -> RpcResult<FilterId> {
641        let last_poll_block_number = self.provider().best_block_number().to_rpc_result()?;
642        let subscription_id = self.id_provider.next_id();
643
644        let id = match subscription_id {
645            jsonrpsee_types::SubscriptionId::Num(n) => FilterId::Num(n),
646            jsonrpsee_types::SubscriptionId::Str(s) => FilterId::Str(s.into_owned()),
647        };
648        let mut filters = self.active_filters.inner.lock().await;
649        filters.insert(
650            id.clone(),
651            ActiveFilter {
652                block: last_poll_block_number,
653                last_poll_timestamp: Instant::now(),
654                kind,
655            },
656        );
657        Ok(id)
658    }
659
660    /// Returns all logs in the given _inclusive_ range that match the filter
661    ///
662    /// Returns an error if:
663    ///  - underlying database error
664    ///  - amount of matches exceeds configured limit
665    async fn get_logs_in_block_range(
666        self: Arc<Self>,
667        filter: Filter,
668        from_block: u64,
669        to_block: u64,
670        limits: QueryLimits,
671    ) -> Result<Vec<RpcLog<Eth::NetworkTypes>>, EthFilterError> {
672        trace!(target: "rpc::eth::filter", from=from_block, to=to_block, ?filter, "finding logs in range");
673
674        // perform boundary checks first
675        if to_block < from_block {
676            return Err(EthFilterError::InvalidBlockRangeParams)
677        }
678
679        if let Some(max_blocks_per_filter) =
680            limits.max_blocks_per_filter.filter(|limit| to_block - from_block > *limit)
681        {
682            return Err(EthFilterError::QueryExceedsMaxBlocks(max_blocks_per_filter))
683        }
684
685        // The scan occupies a blocking thread until it completes, so it shares the budget for
686        // blocking IO requests with `eth_call` and friends instead of pinning an unbounded number
687        // of pool threads.
688        let permit = self
689            .eth_api
690            .acquire_owned_blocking_io()
691            .await
692            .map_err(|_| EthFilterError::InternalError)?;
693
694        let (mut tx, rx) = oneshot::channel();
695        let this = self.clone();
696        self.task_spawner.spawn_blocking_task(async move {
697            let _permit = permit;
698            let fut = this.get_logs_in_block_range_inner(&filter, from_block, to_block, limits);
699            tokio::pin!(fut);
700            let res = tokio::select! {
701                // Range scans perform blocking reads before their first yield.
702                biased;
703                _ = tx.closed() => None,
704                res = &mut fut => Some(res),
705            };
706            if let Some(res) = res {
707                let _ = tx.send(res);
708            }
709        });
710
711        rx.await.map_err(|_| EthFilterError::InternalError)?
712    }
713
714    /// Returns all logs in the given _inclusive_ range that match the filter
715    ///
716    /// Note: This function uses a mix of blocking db operations for fetching indices and header
717    /// ranges and utilizes the rpc cache for optimistically fetching receipts and blocks.
718    /// This function is considered blocking and should thus be spawned on a blocking task.
719    ///
720    /// Returns an error if:
721    ///  - underlying database error
722    async fn get_logs_in_block_range_inner(
723        self: Arc<Self>,
724        filter: &Filter,
725        from_block: u64,
726        to_block: u64,
727        limits: QueryLimits,
728    ) -> Result<Vec<RpcLog<Eth::NetworkTypes>>, EthFilterError> {
729        let mut all_logs = Vec::new();
730
731        // get current chain tip to determine processing mode
732        let chain_tip = self.provider().best_block_number()?;
733
734        // Scan the range window by window so that receipts are fetched while headers are still
735        // being read: the log limit can end the query after the first window, memory is bounded by
736        // one window, and a cancelled query stops at the next window.
737        for (from, to) in
738            BlockRangeInclusiveIter::new(from_block..=to_block, self.max_headers_range)
739        {
740            // reading headers is blocking, this gives the cancellation check a chance to run
741            tokio::task::yield_now().await;
742
743            // collect the headers of this window that match the bloom filter
744            let mut matching_headers = Vec::new();
745            let headers = self.provider().headers_range(from..=to)?;
746
747            let mut headers_iter = headers.into_iter().peekable();
748
749            while let Some(header) = headers_iter.next() {
750                if !filter.matches_bloom(header.logs_bloom()) {
751                    continue
752                }
753
754                let current_number = header.number();
755
756                let block_hash = match headers_iter.peek() {
757                    Some(next_header) if next_header.number() == current_number + 1 => {
758                        // Headers are consecutive, use the more efficient parent_hash
759                        next_header.parent_hash()
760                    }
761                    _ => {
762                        // Headers not consecutive or last header, calculate hash
763                        header.hash_slow()
764                    }
765                };
766
767                matching_headers.push(SealedHeader::new(header, block_hash));
768            }
769
770            // initialize the appropriate range mode based on collected headers
771            let mut range_mode = RangeMode::new(
772                self.clone(),
773                matching_headers,
774                from_block,
775                to_block,
776                self.max_headers_range,
777                chain_tip,
778            );
779
780            // iterate through the range mode to get receipts and blocks
781            while let Some(ReceiptBlockResult { receipts, recovered_block, header }) =
782                range_mode.next().await?
783            {
784                let num_hash = header.num_hash();
785                append_matching_block_logs(
786                    &mut all_logs,
787                    self.eth_api.converter(),
788                    recovered_block
789                        .map(ProviderOrBlock::Block)
790                        .unwrap_or_else(|| ProviderOrBlock::Provider(self.provider())),
791                    filter,
792                    &header,
793                    &receipts,
794                    false,
795                )?;
796
797                // size check but only if range is multiple blocks, so we always return all
798                // logs of a single block
799                let is_multi_block_range = from_block != to_block;
800                if let Some(max_logs_per_response) = limits.max_logs_per_response &&
801                    is_multi_block_range &&
802                    all_logs.len() > max_logs_per_response
803                {
804                    let retry_to_block = if num_hash.number == from_block {
805                        from_block
806                    } else {
807                        num_hash.number - 1
808                    };
809
810                    debug!(
811                        target: "rpc::eth::filter",
812                        logs_found = all_logs.len(),
813                        max_logs_per_response,
814                        from_block,
815                        to_block = retry_to_block,
816                        "Query exceeded max logs per response limit"
817                    );
818                    return Err(EthFilterError::QueryExceedsMaxResults {
819                        max_logs: max_logs_per_response,
820                        from_block,
821                        to_block: retry_to_block,
822                    });
823                }
824            }
825        }
826
827        Ok(all_logs)
828    }
829}
830
831/// All active filters
832#[derive(Debug, Clone, Default)]
833pub struct ActiveFilters<T> {
834    inner: Arc<Mutex<HashMap<FilterId, ActiveFilter<T>>>>,
835}
836
837impl<T> ActiveFilters<T> {
838    /// Returns an empty instance.
839    pub fn new() -> Self {
840        Self { inner: Arc::new(Mutex::new(HashMap::default())) }
841    }
842
843    /// Returns `true` if a filter with the given id exists.
844    pub async fn contains(&self, id: &FilterId) -> bool {
845        self.inner.lock().await.contains_key(id)
846    }
847
848    /// Returns the number of currently active filters.
849    pub async fn len(&self) -> usize {
850        self.inner.lock().await.len()
851    }
852
853    /// Returns `true` if there are no active filters.
854    pub async fn is_empty(&self) -> bool {
855        self.inner.lock().await.is_empty()
856    }
857
858    /// Returns all active filter ids.
859    pub async fn ids(&self) -> Vec<FilterId> {
860        self.inner.lock().await.keys().cloned().collect()
861    }
862}
863
864/// An installed filter
865#[derive(Debug)]
866struct ActiveFilter<T> {
867    /// At which block the filter was polled last.
868    block: u64,
869    /// Last time this filter was polled.
870    last_poll_timestamp: Instant,
871    /// What kind of filter it is.
872    kind: FilterKind<T>,
873}
874
875/// A receiver for pending transactions that returns all new transactions since the last poll.
876#[derive(Debug, Clone)]
877struct PendingTransactionsReceiver {
878    txs_receiver: Arc<Mutex<Receiver<TxHash>>>,
879}
880
881impl PendingTransactionsReceiver {
882    fn new(receiver: Receiver<TxHash>) -> Self {
883        Self { txs_receiver: Arc::new(Mutex::new(receiver)) }
884    }
885
886    /// Returns all new pending transactions received since the last poll.
887    async fn drain<T>(&self) -> FilterChanges<T> {
888        let mut pending_txs = Vec::new();
889        let mut prepared_stream = self.txs_receiver.lock().await;
890
891        while let Ok(tx_hash) = prepared_stream.try_recv() {
892            pending_txs.push(tx_hash);
893        }
894
895        // Convert the vector of hashes into FilterChanges::Hashes
896        FilterChanges::Hashes(pending_txs)
897    }
898}
899
900/// A structure to manage and provide access to a stream of full transaction details.
901#[derive(Debug, Clone)]
902struct FullTransactionsReceiver<T: PoolTransaction, TxCompat> {
903    txs_stream: Arc<Mutex<NewSubpoolTransactionStream<T>>>,
904    converter: TxCompat,
905}
906
907impl<T, TxCompat> FullTransactionsReceiver<T, TxCompat>
908where
909    T: PoolTransaction + 'static,
910    TxCompat: RpcConvert<Primitives: NodePrimitives<SignedTx = T::Consensus>>,
911{
912    /// Creates a new `FullTransactionsReceiver` encapsulating the provided transaction stream.
913    fn new(stream: NewSubpoolTransactionStream<T>, converter: TxCompat) -> Self {
914        Self { txs_stream: Arc::new(Mutex::new(stream)), converter }
915    }
916
917    /// Returns all new pending transactions received since the last poll.
918    async fn drain(&self) -> FilterChanges<RpcTransaction<TxCompat::Network>> {
919        let mut pending_txs = Vec::new();
920        let mut prepared_stream = self.txs_stream.lock().await;
921
922        while let Ok(tx) = prepared_stream.try_recv() {
923            match self.converter.fill_pending(tx.transaction.to_consensus()) {
924                Ok(tx) => pending_txs.push(tx),
925                Err(err) => {
926                    error!(target: "rpc",
927                        %err,
928                        "Failed to fill txn with block context"
929                    );
930                }
931            }
932        }
933        FilterChanges::Transactions(pending_txs)
934    }
935}
936
937/// Helper trait for [`FullTransactionsReceiver`] to erase the `Transaction` type.
938#[async_trait]
939trait FullTransactionsFilter<T>: fmt::Debug + Send + Sync + Unpin + 'static {
940    async fn drain(&self) -> FilterChanges<T>;
941}
942
943#[async_trait]
944impl<T, TxCompat> FullTransactionsFilter<RpcTransaction<TxCompat::Network>>
945    for FullTransactionsReceiver<T, TxCompat>
946where
947    T: PoolTransaction + 'static,
948    TxCompat: RpcConvert<Primitives: NodePrimitives<SignedTx = T::Consensus>> + 'static,
949{
950    async fn drain(&self) -> FilterChanges<RpcTransaction<TxCompat::Network>> {
951        Self::drain(self).await
952    }
953}
954
955/// Represents the kind of pending transaction data that can be retrieved.
956///
957/// This enum differentiates between two kinds of pending transaction data:
958/// - Just the transaction hashes.
959/// - Full transaction details.
960#[derive(Debug, Clone)]
961enum PendingTransactionKind<T> {
962    Hashes(PendingTransactionsReceiver),
963    FullTransaction(Arc<dyn FullTransactionsFilter<T>>),
964}
965
966impl<T: 'static> PendingTransactionKind<T> {
967    async fn drain(&self) -> FilterChanges<T> {
968        match self {
969            Self::Hashes(receiver) => receiver.drain().await,
970            Self::FullTransaction(receiver) => receiver.drain().await,
971        }
972    }
973}
974
975#[derive(Clone, Debug)]
976enum FilterKind<T> {
977    Log(Box<Filter>),
978    Block,
979    PendingTransaction(PendingTransactionKind<T>),
980}
981
982/// An iterator that yields _inclusive_ block ranges of a given step size
983#[derive(Debug)]
984struct BlockRangeInclusiveIter {
985    iter: StepBy<RangeInclusive<u64>>,
986    step: u64,
987    end: u64,
988}
989
990impl BlockRangeInclusiveIter {
991    fn new(range: RangeInclusive<u64>, step: u64) -> Self {
992        Self { end: *range.end(), iter: range.step_by(step as usize + 1), step }
993    }
994}
995
996impl Iterator for BlockRangeInclusiveIter {
997    type Item = (u64, u64);
998
999    fn next(&mut self) -> Option<Self::Item> {
1000        let start = self.iter.next()?;
1001        let end = (start + self.step).min(self.end);
1002        if start > end {
1003            return None
1004        }
1005        Some((start, end))
1006    }
1007}
1008
1009/// Errors that can occur in the handler implementation
1010#[derive(Debug, thiserror::Error)]
1011pub enum EthFilterError {
1012    /// Filter not found.
1013    #[error("filter not found")]
1014    FilterNotFound(FilterId),
1015    /// Invalid block range.
1016    #[error("invalid block range params")]
1017    InvalidBlockRangeParams,
1018    /// Block range extends beyond current head.
1019    #[error("block range extends beyond current head block: requested {requested}, head {head}")]
1020    BlockRangeExceedsHead {
1021        /// The requested `toBlock` number
1022        requested: u64,
1023        /// The current head block number
1024        head: u64,
1025    },
1026    /// Query scope is too broad.
1027    #[error("query exceeds max block range {0}")]
1028    QueryExceedsMaxBlocks(u64),
1029    /// Receipts of a block the filter matched are gone, most likely pruned.
1030    #[error("pruned history unavailable")]
1031    ReceiptsUnavailable(u64),
1032    /// Query result is too large.
1033    #[error("query exceeds max results {max_logs}, retry with the range {from_block}-{to_block}")]
1034    QueryExceedsMaxResults {
1035        /// Maximum number of logs allowed per response
1036        max_logs: usize,
1037        /// Start block of the suggested retry range
1038        from_block: u64,
1039        /// End block of the suggested retry range (last successfully processed block)
1040        to_block: u64,
1041    },
1042    /// Error serving request in `eth_` namespace.
1043    #[error(transparent)]
1044    EthAPIError(#[from] EthApiError),
1045    /// Error thrown when a spawned task failed to deliver a response.
1046    #[error("internal filter error")]
1047    InternalError,
1048}
1049
1050impl From<EthFilterError> for jsonrpsee::types::error::ErrorObject<'static> {
1051    fn from(err: EthFilterError) -> Self {
1052        match err {
1053            // geth and Nethermind answer -32000 for unknown filter ids
1054            EthFilterError::FilterNotFound(_) => rpc_error_with_code(
1055                jsonrpsee::types::error::CALL_EXECUTION_FAILED_CODE,
1056                "filter not found",
1057            ),
1058            err @ EthFilterError::InternalError => {
1059                rpc_error_with_code(jsonrpsee::types::error::INTERNAL_ERROR_CODE, err.to_string())
1060            }
1061            EthFilterError::EthAPIError(err) => err.into(),
1062            err @ EthFilterError::ReceiptsUnavailable(_) => {
1063                rpc_error_with_code(4444, err.to_string())
1064            }
1065            err @ (EthFilterError::InvalidBlockRangeParams |
1066            EthFilterError::QueryExceedsMaxBlocks(_) |
1067            EthFilterError::QueryExceedsMaxResults { .. } |
1068            EthFilterError::BlockRangeExceedsHead { .. }) => {
1069                rpc_error_with_code(jsonrpsee::types::error::INVALID_PARAMS_CODE, err.to_string())
1070            }
1071        }
1072    }
1073}
1074
1075impl From<ProviderError> for EthFilterError {
1076    fn from(err: ProviderError) -> Self {
1077        Self::EthAPIError(err.into())
1078    }
1079}
1080
1081impl From<logs_utils::FilterBlockRangeError> for EthFilterError {
1082    fn from(err: logs_utils::FilterBlockRangeError) -> Self {
1083        match err {
1084            logs_utils::FilterBlockRangeError::InvalidBlockRange => Self::InvalidBlockRangeParams,
1085            logs_utils::FilterBlockRangeError::BlockRangeExceedsHead { requested, head } => {
1086                Self::BlockRangeExceedsHead { requested, head }
1087            }
1088        }
1089    }
1090}
1091
1092/// Helper type for the common pattern of returning receipts, block and the original header that is
1093/// a match for the filter.
1094struct ReceiptBlockResult<P>
1095where
1096    P: ReceiptProvider + BlockReader,
1097{
1098    /// We always need the entire receipts for the matching block.
1099    receipts: Arc<Vec<ProviderReceipt<P>>>,
1100    /// Block can be optional and we can fetch it lazily when needed.
1101    recovered_block: Option<Arc<reth_primitives_traits::RecoveredBlock<ProviderBlock<P>>>>,
1102    /// The header of the block.
1103    header: SealedHeader<<P as HeaderProvider>::Header>,
1104}
1105
1106/// Represents different modes for processing block ranges when filtering logs
1107enum RangeMode<
1108    Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1109        + EthApiTypes
1110        + LoadReceipt
1111        + EthBlocks
1112        + 'static,
1113> {
1114    /// Use cache-based processing for recent blocks
1115    Cached(CachedMode<Eth>),
1116    /// Use range-based processing for older blocks
1117    Range(RangeBlockMode<Eth>),
1118}
1119
1120impl<
1121        Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1122            + EthApiTypes
1123            + LoadReceipt
1124            + EthBlocks
1125            + 'static,
1126    > RangeMode<Eth>
1127{
1128    /// Creates a new `RangeMode`.
1129    fn new(
1130        filter_inner: Arc<EthFilterInner<Eth>>,
1131        sealed_headers: Vec<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>,
1132        from_block: u64,
1133        to_block: u64,
1134        max_headers_range: u64,
1135        chain_tip: u64,
1136    ) -> Self {
1137        let block_count = to_block - from_block + 1;
1138        let distance_from_tip = chain_tip.saturating_sub(to_block);
1139
1140        // Determine if we should use cached mode based on range characteristics
1141        let use_cached_mode =
1142            Self::should_use_cached_mode(&sealed_headers, block_count, distance_from_tip);
1143
1144        // Fetching receipts in parallel is only worth the extra tasks when most of the window has
1145        // them to read; the sequential path serves the rest from the receipt cache where it can
1146        let parallel = sealed_headers.len() >= PARALLEL_PROCESSING_THRESHOLD;
1147
1148        if use_cached_mode && !sealed_headers.is_empty() {
1149            Self::Cached(CachedMode { filter_inner, headers_iter: sealed_headers.into_iter() })
1150        } else {
1151            Self::Range(RangeBlockMode {
1152                filter_inner,
1153                iter: sealed_headers.into_iter().peekable(),
1154                next: VecDeque::new(),
1155                max_range: (max_headers_range as usize).min(MAX_PARALLEL_BATCH_SIZE),
1156                parallel,
1157                pending_tasks: FuturesOrdered::new(),
1158            })
1159        }
1160    }
1161
1162    /// Determines whether to use cached mode based on bloom filter matches and range size
1163    const fn should_use_cached_mode(
1164        headers: &[SealedHeader<<Eth::Provider as HeaderProvider>::Header>],
1165        block_count: u64,
1166        distance_from_tip: u64,
1167    ) -> bool {
1168        // Headers are already filtered by bloom, so count equals length
1169        let bloom_matches = headers.len();
1170
1171        // Calculate adjusted threshold based on bloom matches
1172        let adjusted_threshold = Self::calculate_adjusted_threshold(block_count, bloom_matches);
1173
1174        block_count <= adjusted_threshold && distance_from_tip <= adjusted_threshold
1175    }
1176
1177    /// Calculates the adjusted cache threshold based on bloom filter matches
1178    const fn calculate_adjusted_threshold(block_count: u64, bloom_matches: usize) -> u64 {
1179        // Only apply adjustments for larger ranges
1180        if block_count <= BLOOM_ADJUSTMENT_MIN_BLOCKS {
1181            return CACHED_MODE_BLOCK_THRESHOLD;
1182        }
1183
1184        match bloom_matches {
1185            n if n > HIGH_BLOOM_MATCH_THRESHOLD => CACHED_MODE_BLOCK_THRESHOLD / 2,
1186            n if n > MODERATE_BLOOM_MATCH_THRESHOLD => (CACHED_MODE_BLOCK_THRESHOLD * 3) / 4,
1187            _ => CACHED_MODE_BLOCK_THRESHOLD,
1188        }
1189    }
1190
1191    /// Gets the next (receipts, `maybe_block`, header, `block_hash`) tuple.
1192    async fn next(&mut self) -> Result<Option<ReceiptBlockResult<Eth::Provider>>, EthFilterError> {
1193        match self {
1194            Self::Cached(cached) => cached.next().await,
1195            Self::Range(range) => range.next().await,
1196        }
1197    }
1198}
1199
1200/// Mode for processing blocks using cache optimization for recent blocks
1201struct CachedMode<
1202    Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1203        + EthApiTypes
1204        + LoadReceipt
1205        + EthBlocks
1206        + 'static,
1207> {
1208    filter_inner: Arc<EthFilterInner<Eth>>,
1209    headers_iter: std::vec::IntoIter<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>,
1210}
1211
1212impl<
1213        Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1214            + EthApiTypes
1215            + LoadReceipt
1216            + EthBlocks
1217            + 'static,
1218    > CachedMode<Eth>
1219{
1220    async fn next(&mut self) -> Result<Option<ReceiptBlockResult<Eth::Provider>>, EthFilterError> {
1221        let Some(header) = self.headers_iter.next() else { return Ok(None) };
1222
1223        // Use get_receipts_and_maybe_block which has automatic fallback to provider
1224        let Some((receipts, maybe_block)) =
1225            self.filter_inner.eth_cache().get_receipts_and_maybe_block(header.hash()).await?
1226        else {
1227            return Err(EthFilterError::ReceiptsUnavailable(header.number()))
1228        };
1229
1230        Ok(Some(ReceiptBlockResult { receipts, recovered_block: maybe_block, header }))
1231    }
1232}
1233
1234/// Type alias for parallel receipt fetching task futures used in `RangeBlockMode`
1235type ReceiptFetchFuture<P> =
1236    Pin<Box<dyn Future<Output = Result<Vec<ReceiptBlockResult<P>>, EthFilterError>> + Send>>;
1237
1238/// Mode for processing blocks using range queries for older blocks
1239struct RangeBlockMode<
1240    Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1241        + EthApiTypes
1242        + LoadReceipt
1243        + EthBlocks
1244        + 'static,
1245> {
1246    filter_inner: Arc<EthFilterInner<Eth>>,
1247    iter: Peekable<std::vec::IntoIter<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>>,
1248    next: VecDeque<ReceiptBlockResult<Eth::Provider>>,
1249    /// Maximum number of consecutive blocks fetched by one batch of parallel tasks
1250    max_range: usize,
1251    /// Whether receipts are fetched in parallel batches of consecutive blocks
1252    parallel: bool,
1253    // Stream of ongoing receipt fetching tasks
1254    pending_tasks: FuturesOrdered<ReceiptFetchFuture<Eth::Provider>>,
1255}
1256
1257impl<
1258        Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1259            + EthApiTypes
1260            + LoadReceipt
1261            + EthBlocks
1262            + 'static,
1263    > RangeBlockMode<Eth>
1264{
1265    async fn next(&mut self) -> Result<Option<ReceiptBlockResult<Eth::Provider>>, EthFilterError> {
1266        loop {
1267            // First, try to return any already processed result from buffer
1268            if let Some(result) = self.next.pop_front() {
1269                return Ok(Some(result));
1270            }
1271
1272            // Try to get a completed task result if there are pending tasks
1273            if let Some(task_result) = self.pending_tasks.next().await {
1274                self.next.extend(task_result?);
1275                continue;
1276            }
1277
1278            // No pending tasks - try to generate more work
1279            let Some(next_header) = self.iter.next() else {
1280                // No more headers to process
1281                return Ok(None);
1282            };
1283
1284            if !self.parallel {
1285                // Process the block on its own so that only one block's receipts are held
1286                if let Some(result) = self.process_small_range(vec![next_header]).await? {
1287                    return Ok(Some(result));
1288                }
1289                // Continue loop to check for more work
1290                continue;
1291            }
1292
1293            let mut range_headers = Vec::with_capacity(self.max_range.min(self.iter.len() + 1));
1294            range_headers.push(next_header);
1295
1296            // Collect consecutive blocks up to max_range size
1297            while range_headers.len() < self.max_range {
1298                let Some(peeked) = self.iter.peek() else { break };
1299                let Some(last_header) = range_headers.last() else { break };
1300
1301                let expected_next = last_header.number() + 1;
1302                if peeked.number() != expected_next {
1303                    trace!(
1304                        target: "rpc::eth::filter",
1305                        last_block = last_header.number(),
1306                        next_block = peeked.number(),
1307                        expected = expected_next,
1308                        range_size = range_headers.len(),
1309                        "Non-consecutive block detected, stopping range collection"
1310                    );
1311                    break; // Non-consecutive block, stop here
1312                }
1313
1314                let Some(next_header) = self.iter.next() else { break };
1315                range_headers.push(next_header);
1316            }
1317
1318            self.spawn_parallel_tasks(range_headers);
1319            // Continue loop to await the spawned tasks
1320        }
1321    }
1322
1323    /// Process a small range of headers sequentially
1324    ///
1325    /// This is used for ranges below [`PARALLEL_PROCESSING_THRESHOLD`], one block at a time.
1326    async fn process_small_range(
1327        &mut self,
1328        range_headers: Vec<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>,
1329    ) -> Result<Option<ReceiptBlockResult<Eth::Provider>>, EthFilterError> {
1330        // Process each header individually to avoid queuing for all receipts
1331        for header in range_headers {
1332            // First check if already cached to avoid unnecessary provider calls
1333            let (maybe_block, maybe_receipts) = self
1334                .filter_inner
1335                .eth_cache()
1336                .maybe_cached_block_and_receipts(header.hash())
1337                .await?;
1338
1339            let receipts = match maybe_receipts {
1340                Some(receipts) => receipts,
1341                None => {
1342                    // Not cached - fetch directly from provider
1343                    match self.filter_inner.provider().receipts_by_block(header.hash().into())? {
1344                        Some(receipts) => Arc::new(receipts),
1345                        None => return Err(EthFilterError::ReceiptsUnavailable(header.number())),
1346                    }
1347                }
1348            };
1349
1350            if !receipts.is_empty() {
1351                self.next.push_back(ReceiptBlockResult {
1352                    receipts,
1353                    recovered_block: maybe_block,
1354                    header,
1355                });
1356            }
1357        }
1358
1359        Ok(self.next.pop_front())
1360    }
1361
1362    /// Spawn parallel tasks for processing a large range of headers
1363    ///
1364    /// This is used for ranges of at least [`PARALLEL_PROCESSING_THRESHOLD`] blocks.
1365    fn spawn_parallel_tasks(
1366        &mut self,
1367        range_headers: Vec<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>,
1368    ) {
1369        // Split headers into chunks
1370        let chunk_size = range_headers.len().div_ceil(DEFAULT_PARALLEL_CONCURRENCY).max(1);
1371        let header_chunks = range_headers
1372            .into_iter()
1373            .chunks(chunk_size)
1374            .into_iter()
1375            .map(|chunk| chunk.collect::<Vec<_>>())
1376            .collect::<Vec<_>>();
1377
1378        // Spawn each chunk as a separate task directly into the FuturesOrdered stream
1379        for chunk_headers in header_chunks {
1380            let filter_inner = self.filter_inner.clone();
1381            let fetch = move || Self::fetch_chunk_receipts(&filter_inner, chunk_headers);
1382
1383            // A parallel task occupies an additional blocking thread, so it needs its own share
1384            // of the blocking IO budget. Without one the chunk is fetched on this task instead.
1385            let chunk_task: ReceiptFetchFuture<Eth::Provider> = match self
1386                .filter_inner
1387                .eth_api
1388                .blocking_io_task_guard()
1389                .clone()
1390                .try_acquire_owned()
1391            {
1392                Ok(permit) => Box::pin(async move {
1393                    let chunk_task = tokio::task::spawn_blocking(move || {
1394                        let _permit = permit;
1395                        fetch()
1396                    });
1397
1398                    // Await the blocking task and handle the result
1399                    match chunk_task.await {
1400                        Ok(chunk_results) => chunk_results,
1401                        Err(join_err) => {
1402                            trace!(target: "rpc::eth::filter", error = ?join_err, "Task join error");
1403                            Err(EthFilterError::InternalError)
1404                        }
1405                    }
1406                }),
1407                Err(_) => Box::pin(async move { fetch() }),
1408            };
1409
1410            self.pending_tasks.push_back(chunk_task);
1411        }
1412    }
1413
1414    /// Fetches the receipts of the given blocks from the provider.
1415    fn fetch_chunk_receipts(
1416        filter_inner: &EthFilterInner<Eth>,
1417        chunk_headers: Vec<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>,
1418    ) -> Result<Vec<ReceiptBlockResult<Eth::Provider>>, EthFilterError> {
1419        let mut chunk_results = Vec::with_capacity(chunk_headers.len());
1420
1421        for header in chunk_headers {
1422            // Fetch directly from provider - RangeMode is used for older blocks
1423            // unlikely to be cached
1424            let receipts = match filter_inner.provider().receipts_by_block(header.hash().into())? {
1425                Some(receipts) => Arc::new(receipts),
1426                None => return Err(EthFilterError::ReceiptsUnavailable(header.number())),
1427            };
1428
1429            if !receipts.is_empty() {
1430                chunk_results.push(ReceiptBlockResult { receipts, recovered_block: None, header });
1431            }
1432        }
1433
1434        Ok(chunk_results)
1435    }
1436}
1437
1438#[cfg(test)]
1439mod tests {
1440    use super::*;
1441    use crate::{eth::EthApi, EthApiBuilder};
1442    use alloy_network::Ethereum;
1443    use alloy_primitives::FixedBytes;
1444    use rand::Rng;
1445    use reth_chainspec::{ChainSpec, ChainSpecProvider};
1446    use reth_ethereum_primitives::TxType;
1447    use reth_evm_ethereum::EthEvmConfig;
1448    use reth_network_api::noop::NoopNetwork;
1449    use reth_provider::test_utils::MockEthProvider;
1450    use reth_rpc_convert::RpcConverter;
1451    use reth_rpc_eth_api::node::RpcNodeCoreAdapter;
1452    use reth_rpc_eth_types::receipt::EthReceiptConverter;
1453    use reth_tasks::Runtime;
1454    use reth_testing_utils::generators;
1455    use reth_transaction_pool::test_utils::{testing_pool, TestPool};
1456    use std::{collections::VecDeque, sync::Arc};
1457
1458    #[test]
1459    fn receipts_unavailable_error_matches_geth() {
1460        let err: jsonrpsee::types::error::ErrorObject<'static> =
1461            EthFilterError::ReceiptsUnavailable(100).into();
1462        assert_eq!(err.code(), 4444);
1463        assert_eq!(err.message(), "pruned history unavailable");
1464    }
1465
1466    #[test]
1467    fn test_block_range_iter() {
1468        let mut rng = generators::rng();
1469
1470        let start = rng.random::<u32>() as u64;
1471        let end = start.saturating_add(rng.random::<u32>() as u64);
1472        let step = rng.random::<u16>() as u64;
1473        let range = start..=end;
1474        let mut iter = BlockRangeInclusiveIter::new(range.clone(), step);
1475        let (from, mut end) = iter.next().unwrap();
1476        assert_eq!(from, start);
1477        assert_eq!(end, (from + step).min(*range.end()));
1478
1479        for (next_from, next_end) in iter {
1480            // ensure range starts with previous end + 1
1481            assert_eq!(next_from, end + 1);
1482            end = next_end;
1483        }
1484
1485        assert_eq!(end, *range.end());
1486    }
1487
1488    // Helper function to create a test EthApi instance
1489    #[expect(clippy::type_complexity)]
1490    fn build_test_eth_api(
1491        provider: MockEthProvider,
1492    ) -> EthApi<
1493        RpcNodeCoreAdapter<MockEthProvider, TestPool, NoopNetwork, EthEvmConfig>,
1494        RpcConverter<Ethereum, EthEvmConfig, EthReceiptConverter<ChainSpec>>,
1495    > {
1496        EthApiBuilder::new(
1497            provider.clone(),
1498            testing_pool(),
1499            NoopNetwork::default(),
1500            EthEvmConfig::new(provider.chain_spec()),
1501        )
1502        .build()
1503    }
1504
1505    #[tokio::test]
1506    async fn test_logs_for_filter_from_block_beyond_head() {
1507        let provider = MockEthProvider::default();
1508        provider.add_header(FixedBytes::random(), alloy_consensus::Header::default());
1509        let eth_api = build_test_eth_api(provider);
1510
1511        let eth_filter =
1512            super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1513
1514        let filter = Filter::new().from_block(100u64).to_block(BlockNumberOrTag::Latest);
1515        let result = eth_filter.inner.clone().logs_for_filter(filter, QueryLimits::default()).await;
1516        assert!(matches!(result, Err(EthFilterError::InvalidBlockRangeParams)), "{result:?}");
1517    }
1518
1519    #[tokio::test]
1520    async fn test_range_block_mode_empty_range() {
1521        let provider = MockEthProvider::default();
1522        let eth_api = build_test_eth_api(provider);
1523
1524        let eth_filter =
1525            super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1526        let filter_inner = eth_filter.inner;
1527
1528        let headers = vec![];
1529        let max_range = 100;
1530
1531        let mut range_mode = RangeBlockMode {
1532            filter_inner,
1533            iter: headers.into_iter().peekable(),
1534            next: VecDeque::new(),
1535            max_range,
1536            parallel: false,
1537
1538            pending_tasks: FuturesOrdered::new(),
1539        };
1540
1541        let result = range_mode.next().await;
1542        assert!(result.is_ok());
1543        assert!(result.unwrap().is_none());
1544    }
1545
1546    #[tokio::test]
1547    async fn test_range_block_mode_queued_results_priority() {
1548        let provider = MockEthProvider::default();
1549
1550        let headers = vec![
1551            SealedHeader::new(
1552                alloy_consensus::Header { number: 100, ..Default::default() },
1553                FixedBytes::random(),
1554            ),
1555            SealedHeader::new(
1556                alloy_consensus::Header { number: 101, ..Default::default() },
1557                FixedBytes::random(),
1558            ),
1559        ];
1560        for header in &headers {
1561            provider.add_header(header.hash(), header.header().clone());
1562            provider.add_receipts(header.number(), vec![]);
1563        }
1564
1565        let eth_api = build_test_eth_api(provider);
1566
1567        let eth_filter =
1568            super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1569        let filter_inner = eth_filter.inner;
1570
1571        // create specific mock results to test ordering
1572        let expected_block_hash_1 = FixedBytes::from([1u8; 32]);
1573        let expected_block_hash_2 = FixedBytes::from([2u8; 32]);
1574
1575        // create mock receipts to test receipt handling
1576        let mock_receipt_1 = reth_ethereum_primitives::Receipt {
1577            tx_type: TxType::Legacy,
1578            cumulative_gas_used: 100_000,
1579            logs: vec![],
1580            success: true,
1581        };
1582        let mock_receipt_2 = reth_ethereum_primitives::Receipt {
1583            tx_type: TxType::Eip1559,
1584            cumulative_gas_used: 200_000,
1585            logs: vec![],
1586            success: true,
1587        };
1588        let mock_receipt_3 = reth_ethereum_primitives::Receipt {
1589            tx_type: TxType::Eip2930,
1590            cumulative_gas_used: 150_000,
1591            logs: vec![],
1592            success: false, // Different success status
1593        };
1594
1595        let mock_result_1 = ReceiptBlockResult {
1596            receipts: Arc::new(vec![mock_receipt_1.clone(), mock_receipt_2.clone()]),
1597            recovered_block: None,
1598            header: SealedHeader::new(
1599                alloy_consensus::Header { number: 42, ..Default::default() },
1600                expected_block_hash_1,
1601            ),
1602        };
1603
1604        let mock_result_2 = ReceiptBlockResult {
1605            receipts: Arc::new(vec![mock_receipt_3.clone()]),
1606            recovered_block: None,
1607            header: SealedHeader::new(
1608                alloy_consensus::Header { number: 43, ..Default::default() },
1609                expected_block_hash_2,
1610            ),
1611        };
1612
1613        let mut range_mode = RangeBlockMode {
1614            filter_inner,
1615            iter: headers.into_iter().peekable(),
1616            next: VecDeque::from([mock_result_1, mock_result_2]), // Queue two results
1617            max_range: 100,
1618            parallel: false,
1619
1620            pending_tasks: FuturesOrdered::new(),
1621        };
1622
1623        // first call should return the first queued result (FIFO order)
1624        let result1 = range_mode.next().await;
1625        assert!(result1.is_ok());
1626        let receipt_result1 = result1.unwrap().unwrap();
1627        assert_eq!(receipt_result1.header.hash(), expected_block_hash_1);
1628        assert_eq!(receipt_result1.header.number, 42);
1629
1630        // verify receipts
1631        assert_eq!(receipt_result1.receipts.len(), 2);
1632        assert_eq!(receipt_result1.receipts[0].tx_type, mock_receipt_1.tx_type);
1633        assert_eq!(
1634            receipt_result1.receipts[0].cumulative_gas_used,
1635            mock_receipt_1.cumulative_gas_used
1636        );
1637        assert_eq!(receipt_result1.receipts[0].success, mock_receipt_1.success);
1638        assert_eq!(receipt_result1.receipts[1].tx_type, mock_receipt_2.tx_type);
1639        assert_eq!(
1640            receipt_result1.receipts[1].cumulative_gas_used,
1641            mock_receipt_2.cumulative_gas_used
1642        );
1643        assert_eq!(receipt_result1.receipts[1].success, mock_receipt_2.success);
1644
1645        // second call should return the second queued result
1646        let result2 = range_mode.next().await;
1647        assert!(result2.is_ok());
1648        let receipt_result2 = result2.unwrap().unwrap();
1649        assert_eq!(receipt_result2.header.hash(), expected_block_hash_2);
1650        assert_eq!(receipt_result2.header.number, 43);
1651
1652        // verify receipts
1653        assert_eq!(receipt_result2.receipts.len(), 1);
1654        assert_eq!(receipt_result2.receipts[0].tx_type, mock_receipt_3.tx_type);
1655        assert_eq!(
1656            receipt_result2.receipts[0].cumulative_gas_used,
1657            mock_receipt_3.cumulative_gas_used
1658        );
1659        assert_eq!(receipt_result2.receipts[0].success, mock_receipt_3.success);
1660
1661        // queue should now be empty
1662        assert!(range_mode.next.is_empty());
1663
1664        let result3 = range_mode.next().await;
1665        assert!(result3.is_ok());
1666    }
1667
1668    #[tokio::test]
1669    async fn test_range_block_mode_single_block_missing_receipts() {
1670        let provider = MockEthProvider::default();
1671        let eth_api = build_test_eth_api(provider);
1672
1673        let eth_filter =
1674            super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1675        let filter_inner = eth_filter.inner;
1676
1677        let headers = vec![SealedHeader::new(
1678            alloy_consensus::Header { number: 100, ..Default::default() },
1679            FixedBytes::random(),
1680        )];
1681
1682        let mut range_mode = RangeBlockMode {
1683            filter_inner,
1684            iter: headers.into_iter().peekable(),
1685            next: VecDeque::new(),
1686            max_range: 100,
1687            parallel: false,
1688
1689            pending_tasks: FuturesOrdered::new(),
1690        };
1691
1692        // a block whose header matched the filter but whose receipts are gone must not be
1693        // silently skipped
1694        let Err(err) = range_mode.next().await else { panic!("missing receipts must be an error") };
1695        assert!(matches!(err, EthFilterError::ReceiptsUnavailable(100)), "{err:?}");
1696    }
1697
1698    #[tokio::test]
1699    async fn test_range_block_mode_provider_receipts() {
1700        let provider = MockEthProvider::default();
1701
1702        let header_1 = alloy_consensus::Header { number: 100, ..Default::default() };
1703        let header_2 = alloy_consensus::Header { number: 101, ..Default::default() };
1704        let header_3 = alloy_consensus::Header { number: 102, ..Default::default() };
1705
1706        let block_hash_1 = FixedBytes::random();
1707        let block_hash_2 = FixedBytes::random();
1708        let block_hash_3 = FixedBytes::random();
1709
1710        provider.add_header(block_hash_1, header_1.clone());
1711        provider.add_header(block_hash_2, header_2.clone());
1712        provider.add_header(block_hash_3, header_3.clone());
1713
1714        // create mock receipts to test provider fetching with mock logs
1715        let mock_log = alloy_primitives::Log {
1716            address: alloy_primitives::Address::ZERO,
1717            data: alloy_primitives::LogData::new_unchecked(vec![], alloy_primitives::Bytes::new()),
1718        };
1719
1720        let receipt_100_1 = reth_ethereum_primitives::Receipt {
1721            tx_type: TxType::Legacy,
1722            cumulative_gas_used: 21_000,
1723            logs: vec![mock_log.clone()],
1724            success: true,
1725        };
1726        let receipt_100_2 = reth_ethereum_primitives::Receipt {
1727            tx_type: TxType::Eip1559,
1728            cumulative_gas_used: 42_000,
1729            logs: vec![mock_log.clone()],
1730            success: true,
1731        };
1732        let receipt_101_1 = reth_ethereum_primitives::Receipt {
1733            tx_type: TxType::Eip2930,
1734            cumulative_gas_used: 30_000,
1735            logs: vec![mock_log.clone()],
1736            success: false,
1737        };
1738
1739        provider.add_receipts(100, vec![receipt_100_1.clone(), receipt_100_2.clone()]);
1740        provider.add_receipts(101, vec![receipt_101_1.clone()]);
1741        // a block without transactions, which a provider reports as an empty list
1742        provider.add_receipts(102, vec![]);
1743
1744        let eth_api = build_test_eth_api(provider);
1745
1746        let eth_filter =
1747            super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1748        let filter_inner = eth_filter.inner;
1749
1750        let headers = vec![
1751            SealedHeader::new(header_1, block_hash_1),
1752            SealedHeader::new(header_2, block_hash_2),
1753            SealedHeader::new(header_3, block_hash_3),
1754        ];
1755
1756        let mut range_mode = RangeBlockMode {
1757            filter_inner,
1758            iter: headers.into_iter().peekable(),
1759            next: VecDeque::new(),
1760            max_range: 3, // include the 3 blocks in the first queried results
1761            parallel: false,
1762
1763            pending_tasks: FuturesOrdered::new(),
1764        };
1765
1766        // first call should fetch receipts from provider and return first block with receipts
1767        let result = range_mode.next().await;
1768        assert!(result.is_ok());
1769        let receipt_result = result.unwrap().unwrap();
1770
1771        assert_eq!(receipt_result.header.hash(), block_hash_1);
1772        assert_eq!(receipt_result.header.number, 100);
1773        assert_eq!(receipt_result.receipts.len(), 2);
1774
1775        // verify receipts
1776        assert_eq!(receipt_result.receipts[0].tx_type, receipt_100_1.tx_type);
1777        assert_eq!(
1778            receipt_result.receipts[0].cumulative_gas_used,
1779            receipt_100_1.cumulative_gas_used
1780        );
1781        assert_eq!(receipt_result.receipts[0].success, receipt_100_1.success);
1782
1783        assert_eq!(receipt_result.receipts[1].tx_type, receipt_100_2.tx_type);
1784        assert_eq!(
1785            receipt_result.receipts[1].cumulative_gas_used,
1786            receipt_100_2.cumulative_gas_used
1787        );
1788        assert_eq!(receipt_result.receipts[1].success, receipt_100_2.success);
1789
1790        // second call should return the second block with receipts
1791        let result2 = range_mode.next().await;
1792        assert!(result2.is_ok());
1793        let receipt_result2 = result2.unwrap().unwrap();
1794
1795        assert_eq!(receipt_result2.header.hash(), block_hash_2);
1796        assert_eq!(receipt_result2.header.number, 101);
1797        assert_eq!(receipt_result2.receipts.len(), 1);
1798
1799        // verify receipts
1800        assert_eq!(receipt_result2.receipts[0].tx_type, receipt_101_1.tx_type);
1801        assert_eq!(
1802            receipt_result2.receipts[0].cumulative_gas_used,
1803            receipt_101_1.cumulative_gas_used
1804        );
1805        assert_eq!(receipt_result2.receipts[0].success, receipt_101_1.success);
1806
1807        // third call should return None since no more blocks with receipts
1808        let result3 = range_mode.next().await;
1809        assert!(result3.is_ok());
1810        assert!(result3.unwrap().is_none());
1811    }
1812
1813    #[tokio::test]
1814    async fn test_range_block_mode_iterator_exhaustion() {
1815        let provider = MockEthProvider::default();
1816
1817        let header_100 = alloy_consensus::Header { number: 100, ..Default::default() };
1818        let header_101 = alloy_consensus::Header { number: 101, ..Default::default() };
1819
1820        let block_hash_100 = FixedBytes::random();
1821        let block_hash_101 = FixedBytes::random();
1822
1823        // Associate headers with hashes first
1824        provider.add_header(block_hash_100, header_100.clone());
1825        provider.add_header(block_hash_101, header_101.clone());
1826
1827        // Add mock receipts so headers are actually processed
1828        let mock_receipt = reth_ethereum_primitives::Receipt {
1829            tx_type: TxType::Legacy,
1830            cumulative_gas_used: 21_000,
1831            logs: vec![],
1832            success: true,
1833        };
1834        provider.add_receipts(100, vec![mock_receipt.clone()]);
1835        provider.add_receipts(101, vec![mock_receipt.clone()]);
1836
1837        let eth_api = build_test_eth_api(provider);
1838
1839        let eth_filter =
1840            super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1841        let filter_inner = eth_filter.inner;
1842
1843        let headers = vec![
1844            SealedHeader::new(header_100, block_hash_100),
1845            SealedHeader::new(header_101, block_hash_101),
1846        ];
1847
1848        let mut range_mode = RangeBlockMode {
1849            filter_inner,
1850            iter: headers.into_iter().peekable(),
1851            next: VecDeque::new(),
1852            max_range: 1,
1853            parallel: false,
1854
1855            pending_tasks: FuturesOrdered::new(),
1856        };
1857
1858        let result1 = range_mode.next().await;
1859        assert!(result1.is_ok());
1860        assert!(result1.unwrap().is_some()); // Should have processed block 100
1861
1862        assert!(range_mode.iter.peek().is_some()); // Should still have block 101
1863
1864        let result2 = range_mode.next().await;
1865        assert!(result2.is_ok());
1866        assert!(result2.unwrap().is_some()); // Should have processed block 101
1867
1868        // now iterator should be exhausted
1869        assert!(range_mode.iter.peek().is_none());
1870
1871        // further calls should return None
1872        let result3 = range_mode.next().await;
1873        assert!(result3.is_ok());
1874        assert!(result3.unwrap().is_none());
1875    }
1876
1877    #[tokio::test]
1878    async fn test_cached_mode_with_mock_receipts() {
1879        // create test data
1880        let test_hash = FixedBytes::from([42u8; 32]);
1881        let test_block_number = 100u64;
1882        let test_header = SealedHeader::new(
1883            alloy_consensus::Header {
1884                number: test_block_number,
1885                gas_used: 50_000,
1886                ..Default::default()
1887            },
1888            test_hash,
1889        );
1890
1891        // add a mock receipt to the provider with a mock log
1892        let mock_log = alloy_primitives::Log {
1893            address: alloy_primitives::Address::ZERO,
1894            data: alloy_primitives::LogData::new_unchecked(vec![], alloy_primitives::Bytes::new()),
1895        };
1896
1897        let mock_receipt = reth_ethereum_primitives::Receipt {
1898            tx_type: TxType::Legacy,
1899            cumulative_gas_used: 21_000,
1900            logs: vec![mock_log],
1901            success: true,
1902        };
1903
1904        let provider = MockEthProvider::default();
1905        provider.add_header(test_hash, test_header.header().clone());
1906        provider.add_receipts(test_block_number, vec![mock_receipt.clone()]);
1907
1908        let eth_api = build_test_eth_api(provider);
1909        let eth_filter =
1910            super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1911        let filter_inner = eth_filter.inner;
1912
1913        let headers = vec![test_header.clone()];
1914
1915        let mut cached_mode = CachedMode { filter_inner, headers_iter: headers.into_iter() };
1916
1917        // should find the receipt from provider fallback (cache will be empty)
1918        let result = cached_mode.next().await.expect("next should succeed");
1919        let receipt_block_result = result.expect("should have receipt result");
1920        assert_eq!(receipt_block_result.header.hash(), test_hash);
1921        assert_eq!(receipt_block_result.header.number, test_block_number);
1922        assert_eq!(receipt_block_result.receipts.len(), 1);
1923        assert_eq!(receipt_block_result.receipts[0].tx_type, mock_receipt.tx_type);
1924        assert_eq!(
1925            receipt_block_result.receipts[0].cumulative_gas_used,
1926            mock_receipt.cumulative_gas_used
1927        );
1928        assert_eq!(receipt_block_result.receipts[0].success, mock_receipt.success);
1929
1930        // iterator should be exhausted
1931        let result2 = cached_mode.next().await;
1932        assert!(result2.is_ok());
1933        assert!(result2.unwrap().is_none());
1934    }
1935
1936    #[tokio::test]
1937    async fn test_cached_mode_empty_headers() {
1938        let provider = MockEthProvider::default();
1939        let eth_api = build_test_eth_api(provider);
1940
1941        let eth_filter =
1942            super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1943        let filter_inner = eth_filter.inner;
1944
1945        let headers: Vec<SealedHeader<alloy_consensus::Header>> = vec![];
1946
1947        let mut cached_mode = CachedMode { filter_inner, headers_iter: headers.into_iter() };
1948
1949        // should immediately return None for empty headers
1950        let result = cached_mode.next().await.expect("next should succeed");
1951        assert!(result.is_none());
1952    }
1953
1954    #[tokio::test]
1955    async fn test_log_limit_retry_range_excludes_overflow_block() {
1956        let provider = MockEthProvider::default();
1957
1958        use alloy_consensus::TxLegacy;
1959        use reth_db_api::models::StoredBlockBodyIndices;
1960        use reth_ethereum_primitives::{TransactionSigned, TxType};
1961
1962        let tx_inner = TxLegacy {
1963            chain_id: Some(1),
1964            nonce: 0,
1965            gas_price: 21_000,
1966            gas_limit: 21_000,
1967            to: alloy_primitives::TxKind::Call(alloy_primitives::Address::ZERO),
1968            value: alloy_primitives::U256::ZERO,
1969            input: alloy_primitives::Bytes::new(),
1970        };
1971        let signature = alloy_primitives::Signature::test_signature();
1972        let tx = TransactionSigned::new_unhashed(tx_inner.into(), signature);
1973
1974        let mock_log = alloy_primitives::Log {
1975            address: alloy_primitives::Address::ZERO,
1976            data: alloy_primitives::LogData::new_unchecked(vec![], alloy_primitives::Bytes::new()),
1977        };
1978
1979        let receipt = reth_ethereum_primitives::Receipt {
1980            tx_type: TxType::Legacy,
1981            cumulative_gas_used: 21_000,
1982            logs: vec![mock_log],
1983            success: true,
1984        };
1985
1986        let mut prev_hash = alloy_primitives::B256::default();
1987        for (idx, block_number) in (100u64..=102).enumerate() {
1988            let header = alloy_consensus::Header {
1989                number: block_number,
1990                parent_hash: prev_hash,
1991                logs_bloom: alloy_primitives::Bloom::from([1u8; 256]),
1992                ..Default::default()
1993            };
1994            let hash = header.hash_slow();
1995            prev_hash = hash;
1996
1997            let block = reth_ethereum_primitives::Block {
1998                header,
1999                body: reth_ethereum_primitives::BlockBody {
2000                    transactions: vec![tx.clone()],
2001                    ..Default::default()
2002                },
2003            };
2004            provider.add_block(hash, block);
2005            provider.add_receipts(block_number, vec![receipt.clone()]);
2006            provider.add_block_body_indices(
2007                block_number,
2008                StoredBlockBodyIndices { first_tx_num: idx as u64, tx_count: 1 },
2009            );
2010        }
2011
2012        let eth_api = build_test_eth_api(provider);
2013        let eth_filter = EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
2014        let err = eth_filter
2015            .inner
2016            .clone()
2017            .get_logs_in_block_range(
2018                Filter::default(),
2019                100,
2020                102,
2021                QueryLimits { max_blocks_per_filter: None, max_logs_per_response: Some(2) },
2022            )
2023            .await
2024            .expect_err("range should exceed max logs");
2025
2026        let EthFilterError::QueryExceedsMaxResults { max_logs, from_block, to_block } = err else {
2027            panic!("unexpected error: {err:?}");
2028        };
2029
2030        assert_eq!(max_logs, 2);
2031        assert_eq!(from_block, 100);
2032        assert_eq!(to_block, 101);
2033    }
2034
2035    #[tokio::test]
2036    async fn test_non_consecutive_headers_after_bloom_filter() {
2037        let provider = MockEthProvider::default();
2038
2039        // Create 4 headers where only blocks 100 and 102 will match bloom filter
2040        let mut expected_hashes = vec![];
2041        let mut prev_hash = alloy_primitives::B256::default();
2042
2043        // Create a transaction for blocks that will have receipts
2044        use alloy_consensus::TxLegacy;
2045        use reth_ethereum_primitives::{TransactionSigned, TxType};
2046
2047        let tx_inner = TxLegacy {
2048            chain_id: Some(1),
2049            nonce: 0,
2050            gas_price: 21_000,
2051            gas_limit: 21_000,
2052            to: alloy_primitives::TxKind::Call(alloy_primitives::Address::ZERO),
2053            value: alloy_primitives::U256::ZERO,
2054            input: alloy_primitives::Bytes::new(),
2055        };
2056        let signature = alloy_primitives::Signature::test_signature();
2057        let tx = TransactionSigned::new_unhashed(tx_inner.into(), signature);
2058
2059        for i in 100u64..=103 {
2060            let header = alloy_consensus::Header {
2061                number: i,
2062                parent_hash: prev_hash,
2063                // Set bloom to match filter only for blocks 100 and 102
2064                logs_bloom: if i == 100 || i == 102 {
2065                    alloy_primitives::Bloom::from([1u8; 256])
2066                } else {
2067                    alloy_primitives::Bloom::default()
2068                },
2069                ..Default::default()
2070            };
2071
2072            let hash = header.hash_slow();
2073            expected_hashes.push(hash);
2074            prev_hash = hash;
2075
2076            // Add transaction to blocks that will have receipts (100 and 102)
2077            let transactions = if i == 100 || i == 102 { vec![tx.clone()] } else { vec![] };
2078
2079            let block = reth_ethereum_primitives::Block {
2080                header,
2081                body: reth_ethereum_primitives::BlockBody { transactions, ..Default::default() },
2082            };
2083            provider.add_block(hash, block);
2084        }
2085
2086        // Add receipts with logs only to blocks that match bloom
2087        let mock_log = alloy_primitives::Log {
2088            address: alloy_primitives::Address::ZERO,
2089            data: alloy_primitives::LogData::new_unchecked(vec![], alloy_primitives::Bytes::new()),
2090        };
2091
2092        let receipt = reth_ethereum_primitives::Receipt {
2093            tx_type: TxType::Legacy,
2094            cumulative_gas_used: 21_000,
2095            logs: vec![mock_log],
2096            success: true,
2097        };
2098
2099        provider.add_receipts(100, vec![receipt.clone()]);
2100        provider.add_receipts(101, vec![]);
2101        provider.add_receipts(102, vec![receipt.clone()]);
2102        provider.add_receipts(103, vec![]);
2103
2104        // Add block body indices for each block so receipts can be fetched
2105        use reth_db_api::models::StoredBlockBodyIndices;
2106        provider
2107            .add_block_body_indices(100, StoredBlockBodyIndices { first_tx_num: 0, tx_count: 1 });
2108        provider
2109            .add_block_body_indices(101, StoredBlockBodyIndices { first_tx_num: 1, tx_count: 0 });
2110        provider
2111            .add_block_body_indices(102, StoredBlockBodyIndices { first_tx_num: 1, tx_count: 1 });
2112        provider
2113            .add_block_body_indices(103, StoredBlockBodyIndices { first_tx_num: 2, tx_count: 0 });
2114
2115        let eth_api = build_test_eth_api(provider);
2116        let eth_filter = EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
2117
2118        // Use default filter which will match any non-empty bloom
2119        let filter = Filter::default();
2120
2121        // Get logs in the range - this will trigger the bloom filtering
2122        let logs = eth_filter
2123            .inner
2124            .clone()
2125            .get_logs_in_block_range(filter, 100, 103, QueryLimits::default())
2126            .await
2127            .expect("should succeed");
2128
2129        // We should get logs from blocks 100 and 102 only (bloom filtered)
2130        assert_eq!(logs.len(), 2);
2131
2132        assert_eq!(logs[0].block_number, Some(100));
2133        assert_eq!(logs[1].block_number, Some(102));
2134
2135        // Each block hash should be the hash of its own header, not derived from any other header
2136        assert_eq!(logs[0].block_hash, Some(expected_hashes[0])); // block 100
2137        assert_eq!(logs[1].block_hash, Some(expected_hashes[2])); // block 102
2138    }
2139
2140    #[tokio::test]
2141    async fn test_range_scan_waits_for_blocking_io_permit() {
2142        use reth_rpc_eth_api::helpers::SpawnBlocking;
2143
2144        let provider = MockEthProvider::default();
2145        let header = alloy_consensus::Header::default();
2146        provider.add_header(header.hash_slow(), header);
2147        provider.add_receipts(0, vec![]);
2148        let eth_api = build_test_eth_api(provider);
2149
2150        // take every permit so the scan has to wait for one
2151        let guard = eth_api.blocking_io_task_guard().clone();
2152        let permits =
2153            guard.clone().acquire_many_owned(guard.available_permits() as u32).await.unwrap();
2154
2155        let eth_filter = EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
2156        let scan = eth_filter.inner.clone().get_logs_in_block_range(
2157            Filter::default(),
2158            0,
2159            0,
2160            QueryLimits::default(),
2161        );
2162        tokio::pin!(scan);
2163        assert!(
2164            tokio::time::timeout(Duration::from_millis(100), &mut scan).await.is_err(),
2165            "scan must wait for a blocking IO permit"
2166        );
2167
2168        drop(permits);
2169        assert!(scan.await.unwrap().is_empty());
2170    }
2171
2172    #[tokio::test]
2173    async fn test_range_scan_across_windows() {
2174        use alloy_consensus::TxLegacy;
2175        use alloy_primitives::{Address, Bloom, Bytes, Log, LogData, Signature};
2176        use reth_db_api::models::StoredBlockBodyIndices;
2177        use reth_ethereum_primitives::{Block, BlockBody, Receipt, TransactionSigned};
2178
2179        let provider = MockEthProvider::default();
2180        let tx = TransactionSigned::new_unhashed(
2181            TxLegacy {
2182                chain_id: Some(1),
2183                gas_price: 21_000,
2184                gas_limit: 21_000,
2185                ..Default::default()
2186            }
2187            .into(),
2188            Signature::test_signature(),
2189        );
2190        let receipt = Receipt {
2191            tx_type: TxType::Legacy,
2192            cumulative_gas_used: 21_000,
2193            logs: vec![Log {
2194                address: Address::ZERO,
2195                data: LogData::new_unchecked(vec![], Bytes::new()),
2196            }],
2197            success: true,
2198        };
2199
2200        // an empty filter matches every header, so every window is fetched in parallel while
2201        // only three blocks have logs to return
2202        let matching_blocks = [5u64, 1_200, 2_400];
2203        let mut parent_hash = FixedBytes::default();
2204        for number in 0..=2_500u64 {
2205            let matches = matching_blocks.contains(&number);
2206            let header = alloy_consensus::Header {
2207                number,
2208                parent_hash,
2209                logs_bloom: if matches { Bloom::from([1u8; 256]) } else { Bloom::default() },
2210                ..Default::default()
2211            };
2212            parent_hash = header.hash_slow();
2213            let transactions = if matches { vec![tx.clone()] } else { Vec::new() };
2214            provider.add_block(
2215                parent_hash,
2216                Block { header, body: BlockBody { transactions, ..Default::default() } },
2217            );
2218            if matches {
2219                let tx_num = matching_blocks.iter().position(|b| *b == number).unwrap() as u64;
2220                provider.add_receipts(number, vec![receipt.clone()]);
2221                provider.add_block_body_indices(
2222                    number,
2223                    StoredBlockBodyIndices { first_tx_num: tx_num, tx_count: 1 },
2224                );
2225            } else {
2226                provider.add_receipts(number, vec![]);
2227            }
2228        }
2229
2230        let eth_filter = EthFilter::new(
2231            build_test_eth_api(provider),
2232            EthFilterConfig::default(),
2233            Runtime::test(),
2234        );
2235        let logs = eth_filter
2236            .inner
2237            .clone()
2238            .get_logs_in_block_range(Filter::default(), 0, 2_500, QueryLimits::default())
2239            .await
2240            .unwrap();
2241        let blocks = logs.iter().map(|log| log.block_number.unwrap()).collect::<Vec<_>>();
2242        assert_eq!(blocks, matching_blocks);
2243
2244        // the log limit ends the scan in the window that exceeds it and names the last block that
2245        // fit as the range to retry with
2246        let err = eth_filter
2247            .inner
2248            .clone()
2249            .get_logs_in_block_range(
2250                Filter::default(),
2251                0,
2252                2_500,
2253                QueryLimits { max_blocks_per_filter: None, max_logs_per_response: Some(1) },
2254            )
2255            .await
2256            .unwrap_err();
2257        assert!(
2258            matches!(
2259                err,
2260                EthFilterError::QueryExceedsMaxResults {
2261                    max_logs: 1,
2262                    from_block: 0,
2263                    to_block: 1_199
2264                }
2265            ),
2266            "{err:?}"
2267        );
2268    }
2269
2270    #[tokio::test]
2271    async fn test_logs_for_filter_pending_to_block_ends_at_head() {
2272        let provider = MockEthProvider::default();
2273        let header = alloy_consensus::Header { number: 2, ..Default::default() };
2274        let hash = header.hash_slow();
2275        provider.add_header(hash, header);
2276        provider.add_receipts(2, vec![]);
2277        // the engine holds a payload that is not canonical yet
2278        provider.set_pending_block_num_hash(Some(alloy_eips::BlockNumHash::new(3, hash)));
2279
2280        let eth_filter = EthFilter::new(
2281            build_test_eth_api(provider),
2282            EthFilterConfig::default(),
2283            Runtime::test(),
2284        );
2285        for filter in [
2286            Filter::new().from_block(0u64).to_block(BlockNumberOrTag::Pending),
2287            Filter::new().select(BlockNumberOrTag::Pending..),
2288        ] {
2289            let logs = eth_filter
2290                .inner
2291                .clone()
2292                .logs_for_filter(filter, QueryLimits::default())
2293                .await
2294                .unwrap();
2295            assert!(logs.is_empty());
2296        }
2297    }
2298
2299    #[tokio::test]
2300    async fn test_logs_for_filter_over_pruned_receipts() {
2301        let provider = MockEthProvider::default();
2302        let mut parent_hash = FixedBytes::default();
2303        for number in 0..=2u64 {
2304            let header = alloy_consensus::Header {
2305                number,
2306                parent_hash,
2307                logs_bloom: alloy_primitives::Bloom::from([1u8; 256]),
2308                ..Default::default()
2309            };
2310            parent_hash = header.hash_slow();
2311            provider.add_block(
2312                parent_hash,
2313                reth_ethereum_primitives::Block { header, body: Default::default() },
2314            );
2315            // the receipts of block 1 were pruned
2316            if number != 1 {
2317                provider.add_receipts(number, vec![]);
2318            }
2319        }
2320
2321        let eth_filter = EthFilter::new(
2322            build_test_eth_api(provider),
2323            EthFilterConfig::default(),
2324            Runtime::test(),
2325        );
2326
2327        // a range that reaches the pruned block is rejected instead of served incompletely
2328        let err = eth_filter
2329            .inner
2330            .clone()
2331            .logs_for_filter(Filter::new().from_block(0u64).to_block(2u64), QueryLimits::default())
2332            .await
2333            .unwrap_err();
2334        assert!(matches!(err, EthFilterError::ReceiptsUnavailable(1)), "{err:?}");
2335
2336        // a range that avoids it is served
2337        let logs = eth_filter
2338            .inner
2339            .clone()
2340            .logs_for_filter(Filter::new().from_block(2u64).to_block(2u64), QueryLimits::default())
2341            .await
2342            .unwrap();
2343        assert!(logs.is_empty());
2344    }
2345}