Skip to main content

reth_rpc/eth/
filter.rs

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