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