1use 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 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
69const CACHED_MODE_BLOCK_THRESHOLD: u64 = 250;
71
72const HIGH_BLOOM_MATCH_THRESHOLD: usize = 20;
74
75const MODERATE_BLOOM_MATCH_THRESHOLD: usize = 10;
77
78const BLOOM_ADJUSTMENT_MIN_BLOCKS: u64 = 100;
80
81const MAX_HEADERS_RANGE: u64 = 1_000; const _: () = assert!(CACHED_MODE_BLOCK_THRESHOLD <= MAX_HEADERS_RANGE);
87
88const PARALLEL_PROCESSING_THRESHOLD: usize = 1000;
91
92const DEFAULT_PARALLEL_CONCURRENCY: usize = 4;
94
95const MAX_PARALLEL_BATCH_SIZE: usize = 256;
97
98pub struct EthFilter<Eth: EthApiTypes> {
102 inner: Arc<EthFilterInner<Eth>>,
104}
105
106impl<Eth> Clone for EthFilter<Eth>
107where
108 Eth: EthApiTypes,
109{
110 fn clone(&self) -> Self {
111 Self { inner: self.inner.clone() }
112 }
113}
114
115impl<Eth> EthFilter<Eth>
116where
117 Eth: EthApiTypes + 'static,
118{
119 pub fn new(eth_api: Eth, config: EthFilterConfig, task_spawner: Runtime) -> Self {
147 let EthFilterConfig { max_blocks_per_filter, max_logs_per_response, stale_filter_ttl } =
148 config;
149 let inner = EthFilterInner {
150 eth_api,
151 active_filters: ActiveFilters::new(),
152 id_provider: Arc::new(EthSubscriptionIdProvider::default()),
153 max_headers_range: MAX_HEADERS_RANGE,
154 task_spawner,
155 stale_filter_ttl,
156 query_limits: QueryLimits { max_blocks_per_filter, max_logs_per_response },
157 };
158
159 let eth_filter = Self { inner: Arc::new(inner) };
160
161 let this = eth_filter.clone();
162 eth_filter.inner.task_spawner.spawn_critical_task(
163 "eth-filters_stale-filters-clean",
164 async move {
165 this.watch_and_clear_stale_filters().await;
166 },
167 );
168
169 eth_filter
170 }
171
172 pub fn active_filters(&self) -> &ActiveFilters<RpcTransaction<Eth::NetworkTypes>> {
174 &self.inner.active_filters
175 }
176
177 async fn watch_and_clear_stale_filters(&self) {
180 let mut interval = tokio::time::interval_at(
181 tokio::time::Instant::now() + self.inner.stale_filter_ttl,
182 self.inner.stale_filter_ttl,
183 );
184 interval.set_missed_tick_behavior(MissedTickBehavior::Delay);
185 loop {
186 interval.tick().await;
187 self.clear_stale_filters(Instant::now()).await;
188 }
189 }
190
191 pub async fn clear_stale_filters(&self, now: Instant) {
194 trace!(target: "rpc::eth", "clear stale filters");
195 let mut filters = self.active_filters().inner.lock().await;
196 filters.retain(|id, filter| {
197 let is_valid = (now - filter.last_poll_timestamp) < self.inner.stale_filter_ttl;
198
199 if !is_valid {
200 trace!(target: "rpc::eth", "evict filter with id: {:?}", id);
201 }
202
203 is_valid
204 });
205 filters.shrink_to_fit();
206 }
207}
208
209impl<Eth> EthFilter<Eth>
210where
211 Eth: FullEthApiTypes<Provider: BlockReader + BlockIdReader>
212 + RpcNodeCoreExt
213 + LoadReceipt
214 + EthBlocks
215 + 'static,
216{
217 fn provider(&self) -> &Eth::Provider {
219 self.inner.eth_api.provider()
220 }
221
222 fn pool(&self) -> &Eth::Pool {
224 self.inner.eth_api.pool()
225 }
226
227 pub async fn filter_changes(
229 &self,
230 id: FilterId,
231 ) -> Result<
232 FilterChanges<RpcTransaction<Eth::NetworkTypes>, RpcLog<Eth::NetworkTypes>>,
233 EthFilterError,
234 > {
235 let info = self.provider().chain_info()?;
236 let best_number = info.best_number;
237
238 let (start_block, kind) = {
241 let mut filters = self.inner.active_filters.inner.lock().await;
242 let filter = filters.get_mut(&id).ok_or(EthFilterError::FilterNotFound(id))?;
243
244 if filter.block > best_number {
245 return Ok(FilterChanges::Empty)
247 }
248
249 let mut block = best_number + 1;
253 std::mem::swap(&mut filter.block, &mut block);
254 filter.last_poll_timestamp = Instant::now();
255
256 (block, filter.kind.clone())
257 };
258
259 match kind {
260 FilterKind::PendingTransaction(filter) => Ok(match filter.drain().await {
261 FilterChanges::Empty => FilterChanges::Empty,
262 FilterChanges::Hashes(hashes) => FilterChanges::Hashes(hashes),
263 FilterChanges::Transactions(transactions) => {
264 FilterChanges::Transactions(transactions)
265 }
266 FilterChanges::Logs(_) => unreachable!("pending transaction filter returned logs"),
267 }),
268 FilterKind::Block => {
269 let end_block = best_number + 1;
272 let block_hashes =
273 self.provider().canonical_hashes_range(start_block, end_block).map_err(
274 |_| EthApiError::HeaderRangeNotFound(start_block.into(), end_block.into()),
275 )?;
276 Ok(FilterChanges::Hashes(block_hashes))
277 }
278 FilterKind::Log(filter) => {
279 let (from_block_number, to_block_number) = match filter.block_option {
280 FilterBlockOption::Range { from_block, to_block } => {
281 let from = from_block
282 .map(|num| self.provider().convert_block_number(num))
283 .transpose()?
284 .flatten();
285 let to = to_block
286 .map(|num| self.provider().convert_block_number(num))
287 .transpose()?
288 .flatten();
289 logs_utils::get_filter_block_range(from, to, start_block, info)?
290 }
291 FilterBlockOption::AtBlockHash(block_hash) => {
292 let block_number = self
296 .provider()
297 .block_number(block_hash)?
298 .ok_or(ProviderError::HeaderNotFound(block_hash.into()))?;
299 (block_number, block_number)
300 }
301 };
302 let logs = self
303 .inner
304 .clone()
305 .get_logs_in_block_range(
306 *filter,
307 from_block_number,
308 to_block_number,
309 self.inner.query_limits,
310 )
311 .await?;
312 Ok(FilterChanges::Logs(logs))
313 }
314 }
315 }
316
317 pub async fn filter_logs(
323 &self,
324 id: FilterId,
325 ) -> Result<Vec<RpcLog<Eth::NetworkTypes>>, EthFilterError> {
326 let filter = {
327 let mut filters = self.inner.active_filters.inner.lock().await;
328 let filter =
329 filters.get_mut(&id).ok_or_else(|| EthFilterError::FilterNotFound(id.clone()))?;
330 if let FilterKind::Log(ref inner_filter) = filter.kind {
331 filter.last_poll_timestamp = Instant::now();
332 *inner_filter.clone()
333 } else {
334 return Err(EthFilterError::FilterNotFound(id))
336 }
337 };
338
339 self.logs_for_filter(filter, self.inner.query_limits).await
340 }
341
342 async fn logs_for_filter(
344 &self,
345 filter: Filter,
346 limits: QueryLimits,
347 ) -> Result<Vec<RpcLog<Eth::NetworkTypes>>, EthFilterError> {
348 self.inner.clone().logs_for_filter(filter, limits).await
349 }
350}
351
352#[async_trait]
353impl<Eth> EthFilterApiServer<RpcTransaction<Eth::NetworkTypes>, RpcLog<Eth::NetworkTypes>>
354 for EthFilter<Eth>
355where
356 Eth: FullEthApiTypes + RpcNodeCoreExt + LoadReceipt + EthBlocks + 'static,
357{
358 async fn new_filter(&self, filter: Filter) -> RpcResult<FilterId> {
360 trace!(target: "rpc::eth", "Serving eth_newFilter");
361 self.inner
362 .install_filter(FilterKind::<RpcTransaction<Eth::NetworkTypes>>::Log(Box::new(filter)))
363 .await
364 }
365
366 async fn new_block_filter(&self) -> RpcResult<FilterId> {
368 trace!(target: "rpc::eth", "Serving eth_newBlockFilter");
369 self.inner.install_filter(FilterKind::<RpcTransaction<Eth::NetworkTypes>>::Block).await
370 }
371
372 async fn new_pending_transaction_filter(
374 &self,
375 kind: Option<PendingTransactionFilterKind>,
376 ) -> RpcResult<FilterId> {
377 trace!(target: "rpc::eth", "Serving eth_newPendingTransactionFilter");
378
379 let transaction_kind = match kind.unwrap_or_default() {
380 PendingTransactionFilterKind::Hashes => {
381 let receiver = self.pool().pending_transactions_listener();
382 let pending_txs_receiver = PendingTransactionsReceiver::new(receiver);
383 FilterKind::PendingTransaction(PendingTransactionKind::Hashes(pending_txs_receiver))
384 }
385 PendingTransactionFilterKind::Full => {
386 let stream = self.pool().new_pending_pool_transactions_listener();
387 let full_txs_receiver = FullTransactionsReceiver::new(
388 stream,
389 dyn_clone::clone(self.inner.eth_api.converter()),
390 );
391 FilterKind::PendingTransaction(PendingTransactionKind::FullTransaction(Arc::new(
392 full_txs_receiver,
393 )))
394 }
395 };
396
397 self.inner.install_filter(transaction_kind).await
399 }
400
401 async fn filter_changes(
403 &self,
404 id: FilterId,
405 ) -> RpcResult<FilterChanges<RpcTransaction<Eth::NetworkTypes>, RpcLog<Eth::NetworkTypes>>>
406 {
407 trace!(target: "rpc::eth", "Serving eth_getFilterChanges");
408 Ok(Self::filter_changes(self, id).await?)
409 }
410
411 async fn filter_logs(&self, id: FilterId) -> RpcResult<Vec<RpcLog<Eth::NetworkTypes>>> {
417 trace!(target: "rpc::eth", "Serving eth_getFilterLogs");
418 Ok(Self::filter_logs(self, id).await?)
419 }
420
421 async fn uninstall_filter(&self, id: FilterId) -> RpcResult<bool> {
423 trace!(target: "rpc::eth", "Serving eth_uninstallFilter");
424 let mut filters = self.inner.active_filters.inner.lock().await;
425 if filters.remove(&id).is_some() {
426 trace!(target: "rpc::eth::filter", ?id, "uninstalled filter");
427 Ok(true)
428 } else {
429 Ok(false)
430 }
431 }
432
433 async fn logs(&self, filter: Filter) -> RpcResult<Vec<RpcLog<Eth::NetworkTypes>>> {
437 trace!(target: "rpc::eth", "Serving eth_getLogs");
438 Ok(self.logs_for_filter(filter, self.inner.query_limits).await?)
439 }
440}
441
442impl<Eth> std::fmt::Debug for EthFilter<Eth>
443where
444 Eth: EthApiTypes,
445{
446 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
447 f.debug_struct("EthFilter").finish_non_exhaustive()
448 }
449}
450
451#[derive(Debug)]
453struct EthFilterInner<Eth: EthApiTypes> {
454 eth_api: Eth,
456 active_filters: ActiveFilters<RpcTransaction<Eth::NetworkTypes>>,
458 id_provider: Arc<dyn IdProvider>,
460 query_limits: QueryLimits,
462 max_headers_range: u64,
464 task_spawner: Runtime,
466 stale_filter_ttl: Duration,
468}
469
470impl<Eth> EthFilterInner<Eth>
471where
472 Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
473 + EthApiTypes<NetworkTypes: reth_rpc_eth_api::types::RpcTypes>
474 + LoadReceipt
475 + EthBlocks
476 + 'static,
477{
478 fn provider(&self) -> &Eth::Provider {
480 self.eth_api.provider()
481 }
482
483 fn eth_cache(&self) -> &EthStateCache<Eth::Primitives> {
485 self.eth_api.cache()
486 }
487
488 async fn logs_for_filter(
490 self: Arc<Self>,
491 filter: Filter,
492 limits: QueryLimits,
493 ) -> Result<Vec<RpcLog<Eth::NetworkTypes>>, EthFilterError> {
494 match filter.block_option {
495 FilterBlockOption::AtBlockHash(block_hash) => {
496 let Some((receipts, maybe_block)) =
498 self.eth_cache().get_receipts_and_maybe_block(block_hash).await?
499 else {
500 return Err(match self.provider().block_number(block_hash)? {
502 Some(number) => {
503 let earliest_available = self.provider().earliest_block_number()?;
504 if number < earliest_available {
505 EthApiError::PrunedHistoryUnavailable {
506 requested: number,
507 earliest_available,
508 }
509 .into()
510 } else {
511 EthFilterError::ReceiptsUnavailable(number)
512 }
513 }
514 None => ProviderError::HeaderNotFound(block_hash.into()).into(),
515 })
516 };
517
518 let header = if let Some(block) = &maybe_block {
519 block.clone_sealed_header()
520 } else {
521 let header = self
522 .provider()
523 .header_by_hash_or_number(block_hash.into())?
524 .ok_or_else(|| ProviderError::HeaderNotFound(block_hash.into()))?;
525 SealedHeader::new(header, block_hash)
526 };
527
528 let earliest_block = self.provider().earliest_block_number()?;
530 if header.number() < earliest_block {
531 return Err(EthApiError::PrunedHistoryUnavailable {
532 requested: header.number(),
533 earliest_available: earliest_block,
534 }
535 .into());
536 }
537
538 if !filter.matches_bloom(header.logs_bloom()) {
539 return Ok(Vec::new())
540 }
541
542 let mut all_logs = Vec::new();
543 append_matching_block_logs(
544 &mut all_logs,
545 self.eth_api.converter(),
546 maybe_block
547 .map(ProviderOrBlock::Block)
548 .unwrap_or_else(|| ProviderOrBlock::Provider(self.provider())),
549 &filter,
550 &header,
551 &receipts,
552 false,
553 )?;
554 Ok(all_logs)
555 }
556 FilterBlockOption::Range { from_block, to_block } => {
557 if from_block.is_some_and(|b| b.is_pending()) {
559 let to_block = to_block.unwrap_or(BlockNumberOrTag::Pending);
560 if !(to_block.is_pending() || to_block.is_number()) {
561 return Ok(Vec::new());
563 }
564 if let Ok(Some(pending_block)) = self.eth_api.local_pending_block().await {
566 if let BlockNumberOrTag::Number(to_block) = to_block &&
567 to_block < pending_block.block.number()
568 {
569 return Ok(Vec::new());
571 }
572
573 let info = self.provider().chain_info()?;
574 if pending_block.block.number() > info.best_number {
575 let mut all_logs = Vec::new();
577 let header = pending_block.block.clone_sealed_header();
578 append_matching_block_logs(
579 &mut all_logs,
580 self.eth_api.converter(),
581 ProviderOrBlock::<Eth::Provider>::Block(pending_block.block),
582 &filter,
583 &header,
584 &pending_block.receipts,
585 false, )?;
587 return Ok(all_logs)
588 }
589 }
590 }
591
592 let info = self.provider().chain_info()?;
593 let start_block = info.best_number;
594 let from = from_block
597 .filter(|num| !num.is_pending())
598 .map(|num| self.provider().convert_block_number(num))
599 .transpose()?
600 .flatten();
601 let to = to_block
602 .filter(|num| !num.is_pending())
603 .map(|num| self.provider().convert_block_number(num))
604 .transpose()?
605 .flatten();
606
607 if let Some(t) = to &&
609 t > info.best_number
610 {
611 return Err(EthFilterError::BlockRangeExceedsHead {
612 requested: t,
613 head: info.best_number,
614 });
615 }
616
617 let (from_block_number, to_block_number) =
618 logs_utils::get_filter_block_range(from, to, start_block, info)?;
619
620 let earliest_block = self.provider().earliest_block_number()?;
622 if from_block_number < earliest_block {
623 return Err(EthApiError::PrunedHistoryUnavailable {
624 requested: from_block_number,
625 earliest_available: earliest_block,
626 }
627 .into());
628 }
629
630 self.get_logs_in_block_range(filter, from_block_number, to_block_number, limits)
631 .await
632 }
633 }
634 }
635
636 async fn install_filter(
638 &self,
639 kind: FilterKind<RpcTransaction<Eth::NetworkTypes>>,
640 ) -> RpcResult<FilterId> {
641 let last_poll_block_number = self.provider().best_block_number().to_rpc_result()?;
642 let subscription_id = self.id_provider.next_id();
643
644 let id = match subscription_id {
645 jsonrpsee_types::SubscriptionId::Num(n) => FilterId::Num(n),
646 jsonrpsee_types::SubscriptionId::Str(s) => FilterId::Str(s.into_owned()),
647 };
648 let mut filters = self.active_filters.inner.lock().await;
649 filters.insert(
650 id.clone(),
651 ActiveFilter {
652 block: last_poll_block_number,
653 last_poll_timestamp: Instant::now(),
654 kind,
655 },
656 );
657 Ok(id)
658 }
659
660 async fn get_logs_in_block_range(
666 self: Arc<Self>,
667 filter: Filter,
668 from_block: u64,
669 to_block: u64,
670 limits: QueryLimits,
671 ) -> Result<Vec<RpcLog<Eth::NetworkTypes>>, EthFilterError> {
672 trace!(target: "rpc::eth::filter", from=from_block, to=to_block, ?filter, "finding logs in range");
673
674 if to_block < from_block {
676 return Err(EthFilterError::InvalidBlockRangeParams)
677 }
678
679 if let Some(max_blocks_per_filter) =
680 limits.max_blocks_per_filter.filter(|limit| to_block - from_block > *limit)
681 {
682 return Err(EthFilterError::QueryExceedsMaxBlocks(max_blocks_per_filter))
683 }
684
685 let permit = self
689 .eth_api
690 .acquire_owned_blocking_io()
691 .await
692 .map_err(|_| EthFilterError::InternalError)?;
693
694 let (mut tx, rx) = oneshot::channel();
695 let this = self.clone();
696 self.task_spawner.spawn_blocking_task(async move {
697 let _permit = permit;
698 let fut = this.get_logs_in_block_range_inner(&filter, from_block, to_block, limits);
699 tokio::pin!(fut);
700 let res = tokio::select! {
701 biased;
703 _ = tx.closed() => None,
704 res = &mut fut => Some(res),
705 };
706 if let Some(res) = res {
707 let _ = tx.send(res);
708 }
709 });
710
711 rx.await.map_err(|_| EthFilterError::InternalError)?
712 }
713
714 async fn get_logs_in_block_range_inner(
723 self: Arc<Self>,
724 filter: &Filter,
725 from_block: u64,
726 to_block: u64,
727 limits: QueryLimits,
728 ) -> Result<Vec<RpcLog<Eth::NetworkTypes>>, EthFilterError> {
729 let mut all_logs = Vec::new();
730
731 let chain_tip = self.provider().best_block_number()?;
733
734 for (from, to) in
738 BlockRangeInclusiveIter::new(from_block..=to_block, self.max_headers_range)
739 {
740 tokio::task::yield_now().await;
742
743 let mut matching_headers = Vec::new();
745 let headers = self.provider().headers_range(from..=to)?;
746
747 let mut headers_iter = headers.into_iter().peekable();
748
749 while let Some(header) = headers_iter.next() {
750 if !filter.matches_bloom(header.logs_bloom()) {
751 continue
752 }
753
754 let current_number = header.number();
755
756 let block_hash = match headers_iter.peek() {
757 Some(next_header) if next_header.number() == current_number + 1 => {
758 next_header.parent_hash()
760 }
761 _ => {
762 header.hash_slow()
764 }
765 };
766
767 matching_headers.push(SealedHeader::new(header, block_hash));
768 }
769
770 let mut range_mode = RangeMode::new(
772 self.clone(),
773 matching_headers,
774 from_block,
775 to_block,
776 self.max_headers_range,
777 chain_tip,
778 );
779
780 while let Some(ReceiptBlockResult { receipts, recovered_block, header }) =
782 range_mode.next().await?
783 {
784 let num_hash = header.num_hash();
785 append_matching_block_logs(
786 &mut all_logs,
787 self.eth_api.converter(),
788 recovered_block
789 .map(ProviderOrBlock::Block)
790 .unwrap_or_else(|| ProviderOrBlock::Provider(self.provider())),
791 filter,
792 &header,
793 &receipts,
794 false,
795 )?;
796
797 let is_multi_block_range = from_block != to_block;
800 if let Some(max_logs_per_response) = limits.max_logs_per_response &&
801 is_multi_block_range &&
802 all_logs.len() > max_logs_per_response
803 {
804 let retry_to_block = if num_hash.number == from_block {
805 from_block
806 } else {
807 num_hash.number - 1
808 };
809
810 debug!(
811 target: "rpc::eth::filter",
812 logs_found = all_logs.len(),
813 max_logs_per_response,
814 from_block,
815 to_block = retry_to_block,
816 "Query exceeded max logs per response limit"
817 );
818 return Err(EthFilterError::QueryExceedsMaxResults {
819 max_logs: max_logs_per_response,
820 from_block,
821 to_block: retry_to_block,
822 });
823 }
824 }
825 }
826
827 Ok(all_logs)
828 }
829}
830
831#[derive(Debug, Clone, Default)]
833pub struct ActiveFilters<T> {
834 inner: Arc<Mutex<HashMap<FilterId, ActiveFilter<T>>>>,
835}
836
837impl<T> ActiveFilters<T> {
838 pub fn new() -> Self {
840 Self { inner: Arc::new(Mutex::new(HashMap::default())) }
841 }
842
843 pub async fn contains(&self, id: &FilterId) -> bool {
845 self.inner.lock().await.contains_key(id)
846 }
847
848 pub async fn len(&self) -> usize {
850 self.inner.lock().await.len()
851 }
852
853 pub async fn is_empty(&self) -> bool {
855 self.inner.lock().await.is_empty()
856 }
857
858 pub async fn ids(&self) -> Vec<FilterId> {
860 self.inner.lock().await.keys().cloned().collect()
861 }
862}
863
864#[derive(Debug)]
866struct ActiveFilter<T> {
867 block: u64,
869 last_poll_timestamp: Instant,
871 kind: FilterKind<T>,
873}
874
875#[derive(Debug, Clone)]
877struct PendingTransactionsReceiver {
878 txs_receiver: Arc<Mutex<Receiver<TxHash>>>,
879}
880
881impl PendingTransactionsReceiver {
882 fn new(receiver: Receiver<TxHash>) -> Self {
883 Self { txs_receiver: Arc::new(Mutex::new(receiver)) }
884 }
885
886 async fn drain<T>(&self) -> FilterChanges<T> {
888 let mut pending_txs = Vec::new();
889 let mut prepared_stream = self.txs_receiver.lock().await;
890
891 while let Ok(tx_hash) = prepared_stream.try_recv() {
892 pending_txs.push(tx_hash);
893 }
894
895 FilterChanges::Hashes(pending_txs)
897 }
898}
899
900#[derive(Debug, Clone)]
902struct FullTransactionsReceiver<T: PoolTransaction, TxCompat> {
903 txs_stream: Arc<Mutex<NewSubpoolTransactionStream<T>>>,
904 converter: TxCompat,
905}
906
907impl<T, TxCompat> FullTransactionsReceiver<T, TxCompat>
908where
909 T: PoolTransaction + 'static,
910 TxCompat: RpcConvert<Primitives: NodePrimitives<SignedTx = T::Consensus>>,
911{
912 fn new(stream: NewSubpoolTransactionStream<T>, converter: TxCompat) -> Self {
914 Self { txs_stream: Arc::new(Mutex::new(stream)), converter }
915 }
916
917 async fn drain(&self) -> FilterChanges<RpcTransaction<TxCompat::Network>> {
919 let mut pending_txs = Vec::new();
920 let mut prepared_stream = self.txs_stream.lock().await;
921
922 while let Ok(tx) = prepared_stream.try_recv() {
923 match self.converter.fill_pending(tx.transaction.to_consensus()) {
924 Ok(tx) => pending_txs.push(tx),
925 Err(err) => {
926 error!(target: "rpc",
927 %err,
928 "Failed to fill txn with block context"
929 );
930 }
931 }
932 }
933 FilterChanges::Transactions(pending_txs)
934 }
935}
936
937#[async_trait]
939trait FullTransactionsFilter<T>: fmt::Debug + Send + Sync + Unpin + 'static {
940 async fn drain(&self) -> FilterChanges<T>;
941}
942
943#[async_trait]
944impl<T, TxCompat> FullTransactionsFilter<RpcTransaction<TxCompat::Network>>
945 for FullTransactionsReceiver<T, TxCompat>
946where
947 T: PoolTransaction + 'static,
948 TxCompat: RpcConvert<Primitives: NodePrimitives<SignedTx = T::Consensus>> + 'static,
949{
950 async fn drain(&self) -> FilterChanges<RpcTransaction<TxCompat::Network>> {
951 Self::drain(self).await
952 }
953}
954
955#[derive(Debug, Clone)]
961enum PendingTransactionKind<T> {
962 Hashes(PendingTransactionsReceiver),
963 FullTransaction(Arc<dyn FullTransactionsFilter<T>>),
964}
965
966impl<T: 'static> PendingTransactionKind<T> {
967 async fn drain(&self) -> FilterChanges<T> {
968 match self {
969 Self::Hashes(receiver) => receiver.drain().await,
970 Self::FullTransaction(receiver) => receiver.drain().await,
971 }
972 }
973}
974
975#[derive(Clone, Debug)]
976enum FilterKind<T> {
977 Log(Box<Filter>),
978 Block,
979 PendingTransaction(PendingTransactionKind<T>),
980}
981
982#[derive(Debug)]
984struct BlockRangeInclusiveIter {
985 iter: StepBy<RangeInclusive<u64>>,
986 step: u64,
987 end: u64,
988}
989
990impl BlockRangeInclusiveIter {
991 fn new(range: RangeInclusive<u64>, step: u64) -> Self {
992 Self { end: *range.end(), iter: range.step_by(step as usize + 1), step }
993 }
994}
995
996impl Iterator for BlockRangeInclusiveIter {
997 type Item = (u64, u64);
998
999 fn next(&mut self) -> Option<Self::Item> {
1000 let start = self.iter.next()?;
1001 let end = (start + self.step).min(self.end);
1002 if start > end {
1003 return None
1004 }
1005 Some((start, end))
1006 }
1007}
1008
1009#[derive(Debug, thiserror::Error)]
1011pub enum EthFilterError {
1012 #[error("filter not found")]
1014 FilterNotFound(FilterId),
1015 #[error("invalid block range params")]
1017 InvalidBlockRangeParams,
1018 #[error("block range extends beyond current head block: requested {requested}, head {head}")]
1020 BlockRangeExceedsHead {
1021 requested: u64,
1023 head: u64,
1025 },
1026 #[error("query exceeds max block range {0}")]
1028 QueryExceedsMaxBlocks(u64),
1029 #[error("pruned history unavailable")]
1031 ReceiptsUnavailable(u64),
1032 #[error("query exceeds max results {max_logs}, retry with the range {from_block}-{to_block}")]
1034 QueryExceedsMaxResults {
1035 max_logs: usize,
1037 from_block: u64,
1039 to_block: u64,
1041 },
1042 #[error(transparent)]
1044 EthAPIError(#[from] EthApiError),
1045 #[error("internal filter error")]
1047 InternalError,
1048}
1049
1050impl From<EthFilterError> for jsonrpsee::types::error::ErrorObject<'static> {
1051 fn from(err: EthFilterError) -> Self {
1052 match err {
1053 EthFilterError::FilterNotFound(_) => rpc_error_with_code(
1055 jsonrpsee::types::error::CALL_EXECUTION_FAILED_CODE,
1056 "filter not found",
1057 ),
1058 err @ EthFilterError::InternalError => {
1059 rpc_error_with_code(jsonrpsee::types::error::INTERNAL_ERROR_CODE, err.to_string())
1060 }
1061 EthFilterError::EthAPIError(err) => err.into(),
1062 err @ EthFilterError::ReceiptsUnavailable(_) => {
1063 rpc_error_with_code(4444, err.to_string())
1064 }
1065 err @ (EthFilterError::InvalidBlockRangeParams |
1066 EthFilterError::QueryExceedsMaxBlocks(_) |
1067 EthFilterError::QueryExceedsMaxResults { .. } |
1068 EthFilterError::BlockRangeExceedsHead { .. }) => {
1069 rpc_error_with_code(jsonrpsee::types::error::INVALID_PARAMS_CODE, err.to_string())
1070 }
1071 }
1072 }
1073}
1074
1075impl From<ProviderError> for EthFilterError {
1076 fn from(err: ProviderError) -> Self {
1077 Self::EthAPIError(err.into())
1078 }
1079}
1080
1081impl From<logs_utils::FilterBlockRangeError> for EthFilterError {
1082 fn from(err: logs_utils::FilterBlockRangeError) -> Self {
1083 match err {
1084 logs_utils::FilterBlockRangeError::InvalidBlockRange => Self::InvalidBlockRangeParams,
1085 logs_utils::FilterBlockRangeError::BlockRangeExceedsHead { requested, head } => {
1086 Self::BlockRangeExceedsHead { requested, head }
1087 }
1088 }
1089 }
1090}
1091
1092struct ReceiptBlockResult<P>
1095where
1096 P: ReceiptProvider + BlockReader,
1097{
1098 receipts: Arc<Vec<ProviderReceipt<P>>>,
1100 recovered_block: Option<Arc<reth_primitives_traits::RecoveredBlock<ProviderBlock<P>>>>,
1102 header: SealedHeader<<P as HeaderProvider>::Header>,
1104}
1105
1106enum RangeMode<
1108 Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1109 + EthApiTypes
1110 + LoadReceipt
1111 + EthBlocks
1112 + 'static,
1113> {
1114 Cached(CachedMode<Eth>),
1116 Range(RangeBlockMode<Eth>),
1118}
1119
1120impl<
1121 Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1122 + EthApiTypes
1123 + LoadReceipt
1124 + EthBlocks
1125 + 'static,
1126 > RangeMode<Eth>
1127{
1128 fn new(
1130 filter_inner: Arc<EthFilterInner<Eth>>,
1131 sealed_headers: Vec<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>,
1132 from_block: u64,
1133 to_block: u64,
1134 max_headers_range: u64,
1135 chain_tip: u64,
1136 ) -> Self {
1137 let block_count = to_block - from_block + 1;
1138 let distance_from_tip = chain_tip.saturating_sub(to_block);
1139
1140 let use_cached_mode =
1142 Self::should_use_cached_mode(&sealed_headers, block_count, distance_from_tip);
1143
1144 let parallel = sealed_headers.len() >= PARALLEL_PROCESSING_THRESHOLD;
1147
1148 if use_cached_mode && !sealed_headers.is_empty() {
1149 Self::Cached(CachedMode { filter_inner, headers_iter: sealed_headers.into_iter() })
1150 } else {
1151 Self::Range(RangeBlockMode {
1152 filter_inner,
1153 iter: sealed_headers.into_iter().peekable(),
1154 next: VecDeque::new(),
1155 max_range: (max_headers_range as usize).min(MAX_PARALLEL_BATCH_SIZE),
1156 parallel,
1157 pending_tasks: FuturesOrdered::new(),
1158 })
1159 }
1160 }
1161
1162 const fn should_use_cached_mode(
1164 headers: &[SealedHeader<<Eth::Provider as HeaderProvider>::Header>],
1165 block_count: u64,
1166 distance_from_tip: u64,
1167 ) -> bool {
1168 let bloom_matches = headers.len();
1170
1171 let adjusted_threshold = Self::calculate_adjusted_threshold(block_count, bloom_matches);
1173
1174 block_count <= adjusted_threshold && distance_from_tip <= adjusted_threshold
1175 }
1176
1177 const fn calculate_adjusted_threshold(block_count: u64, bloom_matches: usize) -> u64 {
1179 if block_count <= BLOOM_ADJUSTMENT_MIN_BLOCKS {
1181 return CACHED_MODE_BLOCK_THRESHOLD;
1182 }
1183
1184 match bloom_matches {
1185 n if n > HIGH_BLOOM_MATCH_THRESHOLD => CACHED_MODE_BLOCK_THRESHOLD / 2,
1186 n if n > MODERATE_BLOOM_MATCH_THRESHOLD => (CACHED_MODE_BLOCK_THRESHOLD * 3) / 4,
1187 _ => CACHED_MODE_BLOCK_THRESHOLD,
1188 }
1189 }
1190
1191 async fn next(&mut self) -> Result<Option<ReceiptBlockResult<Eth::Provider>>, EthFilterError> {
1193 match self {
1194 Self::Cached(cached) => cached.next().await,
1195 Self::Range(range) => range.next().await,
1196 }
1197 }
1198}
1199
1200struct CachedMode<
1202 Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1203 + EthApiTypes
1204 + LoadReceipt
1205 + EthBlocks
1206 + 'static,
1207> {
1208 filter_inner: Arc<EthFilterInner<Eth>>,
1209 headers_iter: std::vec::IntoIter<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>,
1210}
1211
1212impl<
1213 Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1214 + EthApiTypes
1215 + LoadReceipt
1216 + EthBlocks
1217 + 'static,
1218 > CachedMode<Eth>
1219{
1220 async fn next(&mut self) -> Result<Option<ReceiptBlockResult<Eth::Provider>>, EthFilterError> {
1221 let Some(header) = self.headers_iter.next() else { return Ok(None) };
1222
1223 let Some((receipts, maybe_block)) =
1225 self.filter_inner.eth_cache().get_receipts_and_maybe_block(header.hash()).await?
1226 else {
1227 return Err(EthFilterError::ReceiptsUnavailable(header.number()))
1228 };
1229
1230 Ok(Some(ReceiptBlockResult { receipts, recovered_block: maybe_block, header }))
1231 }
1232}
1233
1234type ReceiptFetchFuture<P> =
1236 Pin<Box<dyn Future<Output = Result<Vec<ReceiptBlockResult<P>>, EthFilterError>> + Send>>;
1237
1238struct RangeBlockMode<
1240 Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1241 + EthApiTypes
1242 + LoadReceipt
1243 + EthBlocks
1244 + 'static,
1245> {
1246 filter_inner: Arc<EthFilterInner<Eth>>,
1247 iter: Peekable<std::vec::IntoIter<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>>,
1248 next: VecDeque<ReceiptBlockResult<Eth::Provider>>,
1249 max_range: usize,
1251 parallel: bool,
1253 pending_tasks: FuturesOrdered<ReceiptFetchFuture<Eth::Provider>>,
1255}
1256
1257impl<
1258 Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1259 + EthApiTypes
1260 + LoadReceipt
1261 + EthBlocks
1262 + 'static,
1263 > RangeBlockMode<Eth>
1264{
1265 async fn next(&mut self) -> Result<Option<ReceiptBlockResult<Eth::Provider>>, EthFilterError> {
1266 loop {
1267 if let Some(result) = self.next.pop_front() {
1269 return Ok(Some(result));
1270 }
1271
1272 if let Some(task_result) = self.pending_tasks.next().await {
1274 self.next.extend(task_result?);
1275 continue;
1276 }
1277
1278 let Some(next_header) = self.iter.next() else {
1280 return Ok(None);
1282 };
1283
1284 if !self.parallel {
1285 if let Some(result) = self.process_small_range(vec![next_header]).await? {
1287 return Ok(Some(result));
1288 }
1289 continue;
1291 }
1292
1293 let mut range_headers = Vec::with_capacity(self.max_range.min(self.iter.len() + 1));
1294 range_headers.push(next_header);
1295
1296 while range_headers.len() < self.max_range {
1298 let Some(peeked) = self.iter.peek() else { break };
1299 let Some(last_header) = range_headers.last() else { break };
1300
1301 let expected_next = last_header.number() + 1;
1302 if peeked.number() != expected_next {
1303 trace!(
1304 target: "rpc::eth::filter",
1305 last_block = last_header.number(),
1306 next_block = peeked.number(),
1307 expected = expected_next,
1308 range_size = range_headers.len(),
1309 "Non-consecutive block detected, stopping range collection"
1310 );
1311 break; }
1313
1314 let Some(next_header) = self.iter.next() else { break };
1315 range_headers.push(next_header);
1316 }
1317
1318 self.spawn_parallel_tasks(range_headers);
1319 }
1321 }
1322
1323 async fn process_small_range(
1327 &mut self,
1328 range_headers: Vec<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>,
1329 ) -> Result<Option<ReceiptBlockResult<Eth::Provider>>, EthFilterError> {
1330 for header in range_headers {
1332 let (maybe_block, maybe_receipts) = self
1334 .filter_inner
1335 .eth_cache()
1336 .maybe_cached_block_and_receipts(header.hash())
1337 .await?;
1338
1339 let receipts = match maybe_receipts {
1340 Some(receipts) => receipts,
1341 None => {
1342 match self.filter_inner.provider().receipts_by_block(header.hash().into())? {
1344 Some(receipts) => Arc::new(receipts),
1345 None => return Err(EthFilterError::ReceiptsUnavailable(header.number())),
1346 }
1347 }
1348 };
1349
1350 if !receipts.is_empty() {
1351 self.next.push_back(ReceiptBlockResult {
1352 receipts,
1353 recovered_block: maybe_block,
1354 header,
1355 });
1356 }
1357 }
1358
1359 Ok(self.next.pop_front())
1360 }
1361
1362 fn spawn_parallel_tasks(
1366 &mut self,
1367 range_headers: Vec<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>,
1368 ) {
1369 let chunk_size = range_headers.len().div_ceil(DEFAULT_PARALLEL_CONCURRENCY).max(1);
1371 let header_chunks = range_headers
1372 .into_iter()
1373 .chunks(chunk_size)
1374 .into_iter()
1375 .map(|chunk| chunk.collect::<Vec<_>>())
1376 .collect::<Vec<_>>();
1377
1378 for chunk_headers in header_chunks {
1380 let filter_inner = self.filter_inner.clone();
1381 let fetch = move || Self::fetch_chunk_receipts(&filter_inner, chunk_headers);
1382
1383 let chunk_task: ReceiptFetchFuture<Eth::Provider> = match self
1386 .filter_inner
1387 .eth_api
1388 .blocking_io_task_guard()
1389 .clone()
1390 .try_acquire_owned()
1391 {
1392 Ok(permit) => Box::pin(async move {
1393 let chunk_task = tokio::task::spawn_blocking(move || {
1394 let _permit = permit;
1395 fetch()
1396 });
1397
1398 match chunk_task.await {
1400 Ok(chunk_results) => chunk_results,
1401 Err(join_err) => {
1402 trace!(target: "rpc::eth::filter", error = ?join_err, "Task join error");
1403 Err(EthFilterError::InternalError)
1404 }
1405 }
1406 }),
1407 Err(_) => Box::pin(async move { fetch() }),
1408 };
1409
1410 self.pending_tasks.push_back(chunk_task);
1411 }
1412 }
1413
1414 fn fetch_chunk_receipts(
1416 filter_inner: &EthFilterInner<Eth>,
1417 chunk_headers: Vec<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>,
1418 ) -> Result<Vec<ReceiptBlockResult<Eth::Provider>>, EthFilterError> {
1419 let mut chunk_results = Vec::with_capacity(chunk_headers.len());
1420
1421 for header in chunk_headers {
1422 let receipts = match filter_inner.provider().receipts_by_block(header.hash().into())? {
1425 Some(receipts) => Arc::new(receipts),
1426 None => return Err(EthFilterError::ReceiptsUnavailable(header.number())),
1427 };
1428
1429 if !receipts.is_empty() {
1430 chunk_results.push(ReceiptBlockResult { receipts, recovered_block: None, header });
1431 }
1432 }
1433
1434 Ok(chunk_results)
1435 }
1436}
1437
1438#[cfg(test)]
1439mod tests {
1440 use super::*;
1441 use crate::{eth::EthApi, EthApiBuilder};
1442 use alloy_network::Ethereum;
1443 use alloy_primitives::FixedBytes;
1444 use rand::Rng;
1445 use reth_chainspec::{ChainSpec, ChainSpecProvider};
1446 use reth_ethereum_primitives::TxType;
1447 use reth_evm_ethereum::EthEvmConfig;
1448 use reth_network_api::noop::NoopNetwork;
1449 use reth_provider::test_utils::MockEthProvider;
1450 use reth_rpc_convert::RpcConverter;
1451 use reth_rpc_eth_api::node::RpcNodeCoreAdapter;
1452 use reth_rpc_eth_types::receipt::EthReceiptConverter;
1453 use reth_tasks::Runtime;
1454 use reth_testing_utils::generators;
1455 use reth_transaction_pool::test_utils::{testing_pool, TestPool};
1456 use std::{collections::VecDeque, sync::Arc};
1457
1458 #[test]
1459 fn receipts_unavailable_error_matches_geth() {
1460 let err: jsonrpsee::types::error::ErrorObject<'static> =
1461 EthFilterError::ReceiptsUnavailable(100).into();
1462 assert_eq!(err.code(), 4444);
1463 assert_eq!(err.message(), "pruned history unavailable");
1464 }
1465
1466 #[test]
1467 fn test_block_range_iter() {
1468 let mut rng = generators::rng();
1469
1470 let start = rng.random::<u32>() as u64;
1471 let end = start.saturating_add(rng.random::<u32>() as u64);
1472 let step = rng.random::<u16>() as u64;
1473 let range = start..=end;
1474 let mut iter = BlockRangeInclusiveIter::new(range.clone(), step);
1475 let (from, mut end) = iter.next().unwrap();
1476 assert_eq!(from, start);
1477 assert_eq!(end, (from + step).min(*range.end()));
1478
1479 for (next_from, next_end) in iter {
1480 assert_eq!(next_from, end + 1);
1482 end = next_end;
1483 }
1484
1485 assert_eq!(end, *range.end());
1486 }
1487
1488 #[expect(clippy::type_complexity)]
1490 fn build_test_eth_api(
1491 provider: MockEthProvider,
1492 ) -> EthApi<
1493 RpcNodeCoreAdapter<MockEthProvider, TestPool, NoopNetwork, EthEvmConfig>,
1494 RpcConverter<Ethereum, EthEvmConfig, EthReceiptConverter<ChainSpec>>,
1495 > {
1496 EthApiBuilder::new(
1497 provider.clone(),
1498 testing_pool(),
1499 NoopNetwork::default(),
1500 EthEvmConfig::new(provider.chain_spec()),
1501 )
1502 .build()
1503 }
1504
1505 #[tokio::test]
1506 async fn test_logs_for_filter_from_block_beyond_head() {
1507 let provider = MockEthProvider::default();
1508 provider.add_header(FixedBytes::random(), alloy_consensus::Header::default());
1509 let eth_api = build_test_eth_api(provider);
1510
1511 let eth_filter =
1512 super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1513
1514 let filter = Filter::new().from_block(100u64).to_block(BlockNumberOrTag::Latest);
1515 let result = eth_filter.inner.clone().logs_for_filter(filter, QueryLimits::default()).await;
1516 assert!(matches!(result, Err(EthFilterError::InvalidBlockRangeParams)), "{result:?}");
1517 }
1518
1519 #[tokio::test]
1520 async fn test_range_block_mode_empty_range() {
1521 let provider = MockEthProvider::default();
1522 let eth_api = build_test_eth_api(provider);
1523
1524 let eth_filter =
1525 super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1526 let filter_inner = eth_filter.inner;
1527
1528 let headers = vec![];
1529 let max_range = 100;
1530
1531 let mut range_mode = RangeBlockMode {
1532 filter_inner,
1533 iter: headers.into_iter().peekable(),
1534 next: VecDeque::new(),
1535 max_range,
1536 parallel: false,
1537
1538 pending_tasks: FuturesOrdered::new(),
1539 };
1540
1541 let result = range_mode.next().await;
1542 assert!(result.is_ok());
1543 assert!(result.unwrap().is_none());
1544 }
1545
1546 #[tokio::test]
1547 async fn test_range_block_mode_queued_results_priority() {
1548 let provider = MockEthProvider::default();
1549
1550 let headers = vec![
1551 SealedHeader::new(
1552 alloy_consensus::Header { number: 100, ..Default::default() },
1553 FixedBytes::random(),
1554 ),
1555 SealedHeader::new(
1556 alloy_consensus::Header { number: 101, ..Default::default() },
1557 FixedBytes::random(),
1558 ),
1559 ];
1560 for header in &headers {
1561 provider.add_header(header.hash(), header.header().clone());
1562 provider.add_receipts(header.number(), vec![]);
1563 }
1564
1565 let eth_api = build_test_eth_api(provider);
1566
1567 let eth_filter =
1568 super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1569 let filter_inner = eth_filter.inner;
1570
1571 let expected_block_hash_1 = FixedBytes::from([1u8; 32]);
1573 let expected_block_hash_2 = FixedBytes::from([2u8; 32]);
1574
1575 let mock_receipt_1 = reth_ethereum_primitives::Receipt {
1577 tx_type: TxType::Legacy,
1578 cumulative_gas_used: 100_000,
1579 logs: vec![],
1580 success: true,
1581 };
1582 let mock_receipt_2 = reth_ethereum_primitives::Receipt {
1583 tx_type: TxType::Eip1559,
1584 cumulative_gas_used: 200_000,
1585 logs: vec![],
1586 success: true,
1587 };
1588 let mock_receipt_3 = reth_ethereum_primitives::Receipt {
1589 tx_type: TxType::Eip2930,
1590 cumulative_gas_used: 150_000,
1591 logs: vec![],
1592 success: false, };
1594
1595 let mock_result_1 = ReceiptBlockResult {
1596 receipts: Arc::new(vec![mock_receipt_1.clone(), mock_receipt_2.clone()]),
1597 recovered_block: None,
1598 header: SealedHeader::new(
1599 alloy_consensus::Header { number: 42, ..Default::default() },
1600 expected_block_hash_1,
1601 ),
1602 };
1603
1604 let mock_result_2 = ReceiptBlockResult {
1605 receipts: Arc::new(vec![mock_receipt_3.clone()]),
1606 recovered_block: None,
1607 header: SealedHeader::new(
1608 alloy_consensus::Header { number: 43, ..Default::default() },
1609 expected_block_hash_2,
1610 ),
1611 };
1612
1613 let mut range_mode = RangeBlockMode {
1614 filter_inner,
1615 iter: headers.into_iter().peekable(),
1616 next: VecDeque::from([mock_result_1, mock_result_2]), max_range: 100,
1618 parallel: false,
1619
1620 pending_tasks: FuturesOrdered::new(),
1621 };
1622
1623 let result1 = range_mode.next().await;
1625 assert!(result1.is_ok());
1626 let receipt_result1 = result1.unwrap().unwrap();
1627 assert_eq!(receipt_result1.header.hash(), expected_block_hash_1);
1628 assert_eq!(receipt_result1.header.number, 42);
1629
1630 assert_eq!(receipt_result1.receipts.len(), 2);
1632 assert_eq!(receipt_result1.receipts[0].tx_type, mock_receipt_1.tx_type);
1633 assert_eq!(
1634 receipt_result1.receipts[0].cumulative_gas_used,
1635 mock_receipt_1.cumulative_gas_used
1636 );
1637 assert_eq!(receipt_result1.receipts[0].success, mock_receipt_1.success);
1638 assert_eq!(receipt_result1.receipts[1].tx_type, mock_receipt_2.tx_type);
1639 assert_eq!(
1640 receipt_result1.receipts[1].cumulative_gas_used,
1641 mock_receipt_2.cumulative_gas_used
1642 );
1643 assert_eq!(receipt_result1.receipts[1].success, mock_receipt_2.success);
1644
1645 let result2 = range_mode.next().await;
1647 assert!(result2.is_ok());
1648 let receipt_result2 = result2.unwrap().unwrap();
1649 assert_eq!(receipt_result2.header.hash(), expected_block_hash_2);
1650 assert_eq!(receipt_result2.header.number, 43);
1651
1652 assert_eq!(receipt_result2.receipts.len(), 1);
1654 assert_eq!(receipt_result2.receipts[0].tx_type, mock_receipt_3.tx_type);
1655 assert_eq!(
1656 receipt_result2.receipts[0].cumulative_gas_used,
1657 mock_receipt_3.cumulative_gas_used
1658 );
1659 assert_eq!(receipt_result2.receipts[0].success, mock_receipt_3.success);
1660
1661 assert!(range_mode.next.is_empty());
1663
1664 let result3 = range_mode.next().await;
1665 assert!(result3.is_ok());
1666 }
1667
1668 #[tokio::test]
1669 async fn test_range_block_mode_single_block_missing_receipts() {
1670 let provider = MockEthProvider::default();
1671 let eth_api = build_test_eth_api(provider);
1672
1673 let eth_filter =
1674 super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1675 let filter_inner = eth_filter.inner;
1676
1677 let headers = vec![SealedHeader::new(
1678 alloy_consensus::Header { number: 100, ..Default::default() },
1679 FixedBytes::random(),
1680 )];
1681
1682 let mut range_mode = RangeBlockMode {
1683 filter_inner,
1684 iter: headers.into_iter().peekable(),
1685 next: VecDeque::new(),
1686 max_range: 100,
1687 parallel: false,
1688
1689 pending_tasks: FuturesOrdered::new(),
1690 };
1691
1692 let Err(err) = range_mode.next().await else { panic!("missing receipts must be an error") };
1695 assert!(matches!(err, EthFilterError::ReceiptsUnavailable(100)), "{err:?}");
1696 }
1697
1698 #[tokio::test]
1699 async fn test_range_block_mode_provider_receipts() {
1700 let provider = MockEthProvider::default();
1701
1702 let header_1 = alloy_consensus::Header { number: 100, ..Default::default() };
1703 let header_2 = alloy_consensus::Header { number: 101, ..Default::default() };
1704 let header_3 = alloy_consensus::Header { number: 102, ..Default::default() };
1705
1706 let block_hash_1 = FixedBytes::random();
1707 let block_hash_2 = FixedBytes::random();
1708 let block_hash_3 = FixedBytes::random();
1709
1710 provider.add_header(block_hash_1, header_1.clone());
1711 provider.add_header(block_hash_2, header_2.clone());
1712 provider.add_header(block_hash_3, header_3.clone());
1713
1714 let mock_log = alloy_primitives::Log {
1716 address: alloy_primitives::Address::ZERO,
1717 data: alloy_primitives::LogData::new_unchecked(vec![], alloy_primitives::Bytes::new()),
1718 };
1719
1720 let receipt_100_1 = reth_ethereum_primitives::Receipt {
1721 tx_type: TxType::Legacy,
1722 cumulative_gas_used: 21_000,
1723 logs: vec![mock_log.clone()],
1724 success: true,
1725 };
1726 let receipt_100_2 = reth_ethereum_primitives::Receipt {
1727 tx_type: TxType::Eip1559,
1728 cumulative_gas_used: 42_000,
1729 logs: vec![mock_log.clone()],
1730 success: true,
1731 };
1732 let receipt_101_1 = reth_ethereum_primitives::Receipt {
1733 tx_type: TxType::Eip2930,
1734 cumulative_gas_used: 30_000,
1735 logs: vec![mock_log.clone()],
1736 success: false,
1737 };
1738
1739 provider.add_receipts(100, vec![receipt_100_1.clone(), receipt_100_2.clone()]);
1740 provider.add_receipts(101, vec![receipt_101_1.clone()]);
1741 provider.add_receipts(102, vec![]);
1743
1744 let eth_api = build_test_eth_api(provider);
1745
1746 let eth_filter =
1747 super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1748 let filter_inner = eth_filter.inner;
1749
1750 let headers = vec![
1751 SealedHeader::new(header_1, block_hash_1),
1752 SealedHeader::new(header_2, block_hash_2),
1753 SealedHeader::new(header_3, block_hash_3),
1754 ];
1755
1756 let mut range_mode = RangeBlockMode {
1757 filter_inner,
1758 iter: headers.into_iter().peekable(),
1759 next: VecDeque::new(),
1760 max_range: 3, parallel: false,
1762
1763 pending_tasks: FuturesOrdered::new(),
1764 };
1765
1766 let result = range_mode.next().await;
1768 assert!(result.is_ok());
1769 let receipt_result = result.unwrap().unwrap();
1770
1771 assert_eq!(receipt_result.header.hash(), block_hash_1);
1772 assert_eq!(receipt_result.header.number, 100);
1773 assert_eq!(receipt_result.receipts.len(), 2);
1774
1775 assert_eq!(receipt_result.receipts[0].tx_type, receipt_100_1.tx_type);
1777 assert_eq!(
1778 receipt_result.receipts[0].cumulative_gas_used,
1779 receipt_100_1.cumulative_gas_used
1780 );
1781 assert_eq!(receipt_result.receipts[0].success, receipt_100_1.success);
1782
1783 assert_eq!(receipt_result.receipts[1].tx_type, receipt_100_2.tx_type);
1784 assert_eq!(
1785 receipt_result.receipts[1].cumulative_gas_used,
1786 receipt_100_2.cumulative_gas_used
1787 );
1788 assert_eq!(receipt_result.receipts[1].success, receipt_100_2.success);
1789
1790 let result2 = range_mode.next().await;
1792 assert!(result2.is_ok());
1793 let receipt_result2 = result2.unwrap().unwrap();
1794
1795 assert_eq!(receipt_result2.header.hash(), block_hash_2);
1796 assert_eq!(receipt_result2.header.number, 101);
1797 assert_eq!(receipt_result2.receipts.len(), 1);
1798
1799 assert_eq!(receipt_result2.receipts[0].tx_type, receipt_101_1.tx_type);
1801 assert_eq!(
1802 receipt_result2.receipts[0].cumulative_gas_used,
1803 receipt_101_1.cumulative_gas_used
1804 );
1805 assert_eq!(receipt_result2.receipts[0].success, receipt_101_1.success);
1806
1807 let result3 = range_mode.next().await;
1809 assert!(result3.is_ok());
1810 assert!(result3.unwrap().is_none());
1811 }
1812
1813 #[tokio::test]
1814 async fn test_range_block_mode_iterator_exhaustion() {
1815 let provider = MockEthProvider::default();
1816
1817 let header_100 = alloy_consensus::Header { number: 100, ..Default::default() };
1818 let header_101 = alloy_consensus::Header { number: 101, ..Default::default() };
1819
1820 let block_hash_100 = FixedBytes::random();
1821 let block_hash_101 = FixedBytes::random();
1822
1823 provider.add_header(block_hash_100, header_100.clone());
1825 provider.add_header(block_hash_101, header_101.clone());
1826
1827 let mock_receipt = reth_ethereum_primitives::Receipt {
1829 tx_type: TxType::Legacy,
1830 cumulative_gas_used: 21_000,
1831 logs: vec![],
1832 success: true,
1833 };
1834 provider.add_receipts(100, vec![mock_receipt.clone()]);
1835 provider.add_receipts(101, vec![mock_receipt.clone()]);
1836
1837 let eth_api = build_test_eth_api(provider);
1838
1839 let eth_filter =
1840 super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1841 let filter_inner = eth_filter.inner;
1842
1843 let headers = vec![
1844 SealedHeader::new(header_100, block_hash_100),
1845 SealedHeader::new(header_101, block_hash_101),
1846 ];
1847
1848 let mut range_mode = RangeBlockMode {
1849 filter_inner,
1850 iter: headers.into_iter().peekable(),
1851 next: VecDeque::new(),
1852 max_range: 1,
1853 parallel: false,
1854
1855 pending_tasks: FuturesOrdered::new(),
1856 };
1857
1858 let result1 = range_mode.next().await;
1859 assert!(result1.is_ok());
1860 assert!(result1.unwrap().is_some()); assert!(range_mode.iter.peek().is_some()); let result2 = range_mode.next().await;
1865 assert!(result2.is_ok());
1866 assert!(result2.unwrap().is_some()); assert!(range_mode.iter.peek().is_none());
1870
1871 let result3 = range_mode.next().await;
1873 assert!(result3.is_ok());
1874 assert!(result3.unwrap().is_none());
1875 }
1876
1877 #[tokio::test]
1878 async fn test_cached_mode_with_mock_receipts() {
1879 let test_hash = FixedBytes::from([42u8; 32]);
1881 let test_block_number = 100u64;
1882 let test_header = SealedHeader::new(
1883 alloy_consensus::Header {
1884 number: test_block_number,
1885 gas_used: 50_000,
1886 ..Default::default()
1887 },
1888 test_hash,
1889 );
1890
1891 let mock_log = alloy_primitives::Log {
1893 address: alloy_primitives::Address::ZERO,
1894 data: alloy_primitives::LogData::new_unchecked(vec![], alloy_primitives::Bytes::new()),
1895 };
1896
1897 let mock_receipt = reth_ethereum_primitives::Receipt {
1898 tx_type: TxType::Legacy,
1899 cumulative_gas_used: 21_000,
1900 logs: vec![mock_log],
1901 success: true,
1902 };
1903
1904 let provider = MockEthProvider::default();
1905 provider.add_header(test_hash, test_header.header().clone());
1906 provider.add_receipts(test_block_number, vec![mock_receipt.clone()]);
1907
1908 let eth_api = build_test_eth_api(provider);
1909 let eth_filter =
1910 super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1911 let filter_inner = eth_filter.inner;
1912
1913 let headers = vec![test_header.clone()];
1914
1915 let mut cached_mode = CachedMode { filter_inner, headers_iter: headers.into_iter() };
1916
1917 let result = cached_mode.next().await.expect("next should succeed");
1919 let receipt_block_result = result.expect("should have receipt result");
1920 assert_eq!(receipt_block_result.header.hash(), test_hash);
1921 assert_eq!(receipt_block_result.header.number, test_block_number);
1922 assert_eq!(receipt_block_result.receipts.len(), 1);
1923 assert_eq!(receipt_block_result.receipts[0].tx_type, mock_receipt.tx_type);
1924 assert_eq!(
1925 receipt_block_result.receipts[0].cumulative_gas_used,
1926 mock_receipt.cumulative_gas_used
1927 );
1928 assert_eq!(receipt_block_result.receipts[0].success, mock_receipt.success);
1929
1930 let result2 = cached_mode.next().await;
1932 assert!(result2.is_ok());
1933 assert!(result2.unwrap().is_none());
1934 }
1935
1936 #[tokio::test]
1937 async fn test_cached_mode_empty_headers() {
1938 let provider = MockEthProvider::default();
1939 let eth_api = build_test_eth_api(provider);
1940
1941 let eth_filter =
1942 super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1943 let filter_inner = eth_filter.inner;
1944
1945 let headers: Vec<SealedHeader<alloy_consensus::Header>> = vec![];
1946
1947 let mut cached_mode = CachedMode { filter_inner, headers_iter: headers.into_iter() };
1948
1949 let result = cached_mode.next().await.expect("next should succeed");
1951 assert!(result.is_none());
1952 }
1953
1954 #[tokio::test]
1955 async fn test_log_limit_retry_range_excludes_overflow_block() {
1956 let provider = MockEthProvider::default();
1957
1958 use alloy_consensus::TxLegacy;
1959 use reth_db_api::models::StoredBlockBodyIndices;
1960 use reth_ethereum_primitives::{TransactionSigned, TxType};
1961
1962 let tx_inner = TxLegacy {
1963 chain_id: Some(1),
1964 nonce: 0,
1965 gas_price: 21_000,
1966 gas_limit: 21_000,
1967 to: alloy_primitives::TxKind::Call(alloy_primitives::Address::ZERO),
1968 value: alloy_primitives::U256::ZERO,
1969 input: alloy_primitives::Bytes::new(),
1970 };
1971 let signature = alloy_primitives::Signature::test_signature();
1972 let tx = TransactionSigned::new_unhashed(tx_inner.into(), signature);
1973
1974 let mock_log = alloy_primitives::Log {
1975 address: alloy_primitives::Address::ZERO,
1976 data: alloy_primitives::LogData::new_unchecked(vec![], alloy_primitives::Bytes::new()),
1977 };
1978
1979 let receipt = reth_ethereum_primitives::Receipt {
1980 tx_type: TxType::Legacy,
1981 cumulative_gas_used: 21_000,
1982 logs: vec![mock_log],
1983 success: true,
1984 };
1985
1986 let mut prev_hash = alloy_primitives::B256::default();
1987 for (idx, block_number) in (100u64..=102).enumerate() {
1988 let header = alloy_consensus::Header {
1989 number: block_number,
1990 parent_hash: prev_hash,
1991 logs_bloom: alloy_primitives::Bloom::from([1u8; 256]),
1992 ..Default::default()
1993 };
1994 let hash = header.hash_slow();
1995 prev_hash = hash;
1996
1997 let block = reth_ethereum_primitives::Block {
1998 header,
1999 body: reth_ethereum_primitives::BlockBody {
2000 transactions: vec![tx.clone()],
2001 ..Default::default()
2002 },
2003 };
2004 provider.add_block(hash, block);
2005 provider.add_receipts(block_number, vec![receipt.clone()]);
2006 provider.add_block_body_indices(
2007 block_number,
2008 StoredBlockBodyIndices { first_tx_num: idx as u64, tx_count: 1 },
2009 );
2010 }
2011
2012 let eth_api = build_test_eth_api(provider);
2013 let eth_filter = EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
2014 let err = eth_filter
2015 .inner
2016 .clone()
2017 .get_logs_in_block_range(
2018 Filter::default(),
2019 100,
2020 102,
2021 QueryLimits { max_blocks_per_filter: None, max_logs_per_response: Some(2) },
2022 )
2023 .await
2024 .expect_err("range should exceed max logs");
2025
2026 let EthFilterError::QueryExceedsMaxResults { max_logs, from_block, to_block } = err else {
2027 panic!("unexpected error: {err:?}");
2028 };
2029
2030 assert_eq!(max_logs, 2);
2031 assert_eq!(from_block, 100);
2032 assert_eq!(to_block, 101);
2033 }
2034
2035 #[tokio::test]
2036 async fn test_non_consecutive_headers_after_bloom_filter() {
2037 let provider = MockEthProvider::default();
2038
2039 let mut expected_hashes = vec![];
2041 let mut prev_hash = alloy_primitives::B256::default();
2042
2043 use alloy_consensus::TxLegacy;
2045 use reth_ethereum_primitives::{TransactionSigned, TxType};
2046
2047 let tx_inner = TxLegacy {
2048 chain_id: Some(1),
2049 nonce: 0,
2050 gas_price: 21_000,
2051 gas_limit: 21_000,
2052 to: alloy_primitives::TxKind::Call(alloy_primitives::Address::ZERO),
2053 value: alloy_primitives::U256::ZERO,
2054 input: alloy_primitives::Bytes::new(),
2055 };
2056 let signature = alloy_primitives::Signature::test_signature();
2057 let tx = TransactionSigned::new_unhashed(tx_inner.into(), signature);
2058
2059 for i in 100u64..=103 {
2060 let header = alloy_consensus::Header {
2061 number: i,
2062 parent_hash: prev_hash,
2063 logs_bloom: if i == 100 || i == 102 {
2065 alloy_primitives::Bloom::from([1u8; 256])
2066 } else {
2067 alloy_primitives::Bloom::default()
2068 },
2069 ..Default::default()
2070 };
2071
2072 let hash = header.hash_slow();
2073 expected_hashes.push(hash);
2074 prev_hash = hash;
2075
2076 let transactions = if i == 100 || i == 102 { vec![tx.clone()] } else { vec![] };
2078
2079 let block = reth_ethereum_primitives::Block {
2080 header,
2081 body: reth_ethereum_primitives::BlockBody { transactions, ..Default::default() },
2082 };
2083 provider.add_block(hash, block);
2084 }
2085
2086 let mock_log = alloy_primitives::Log {
2088 address: alloy_primitives::Address::ZERO,
2089 data: alloy_primitives::LogData::new_unchecked(vec![], alloy_primitives::Bytes::new()),
2090 };
2091
2092 let receipt = reth_ethereum_primitives::Receipt {
2093 tx_type: TxType::Legacy,
2094 cumulative_gas_used: 21_000,
2095 logs: vec![mock_log],
2096 success: true,
2097 };
2098
2099 provider.add_receipts(100, vec![receipt.clone()]);
2100 provider.add_receipts(101, vec![]);
2101 provider.add_receipts(102, vec![receipt.clone()]);
2102 provider.add_receipts(103, vec![]);
2103
2104 use reth_db_api::models::StoredBlockBodyIndices;
2106 provider
2107 .add_block_body_indices(100, StoredBlockBodyIndices { first_tx_num: 0, tx_count: 1 });
2108 provider
2109 .add_block_body_indices(101, StoredBlockBodyIndices { first_tx_num: 1, tx_count: 0 });
2110 provider
2111 .add_block_body_indices(102, StoredBlockBodyIndices { first_tx_num: 1, tx_count: 1 });
2112 provider
2113 .add_block_body_indices(103, StoredBlockBodyIndices { first_tx_num: 2, tx_count: 0 });
2114
2115 let eth_api = build_test_eth_api(provider);
2116 let eth_filter = EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
2117
2118 let filter = Filter::default();
2120
2121 let logs = eth_filter
2123 .inner
2124 .clone()
2125 .get_logs_in_block_range(filter, 100, 103, QueryLimits::default())
2126 .await
2127 .expect("should succeed");
2128
2129 assert_eq!(logs.len(), 2);
2131
2132 assert_eq!(logs[0].block_number, Some(100));
2133 assert_eq!(logs[1].block_number, Some(102));
2134
2135 assert_eq!(logs[0].block_hash, Some(expected_hashes[0])); assert_eq!(logs[1].block_hash, Some(expected_hashes[2])); }
2139
2140 #[tokio::test]
2141 async fn test_range_scan_waits_for_blocking_io_permit() {
2142 use reth_rpc_eth_api::helpers::SpawnBlocking;
2143
2144 let provider = MockEthProvider::default();
2145 let header = alloy_consensus::Header::default();
2146 provider.add_header(header.hash_slow(), header);
2147 provider.add_receipts(0, vec![]);
2148 let eth_api = build_test_eth_api(provider);
2149
2150 let guard = eth_api.blocking_io_task_guard().clone();
2152 let permits =
2153 guard.clone().acquire_many_owned(guard.available_permits() as u32).await.unwrap();
2154
2155 let eth_filter = EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
2156 let scan = eth_filter.inner.clone().get_logs_in_block_range(
2157 Filter::default(),
2158 0,
2159 0,
2160 QueryLimits::default(),
2161 );
2162 tokio::pin!(scan);
2163 assert!(
2164 tokio::time::timeout(Duration::from_millis(100), &mut scan).await.is_err(),
2165 "scan must wait for a blocking IO permit"
2166 );
2167
2168 drop(permits);
2169 assert!(scan.await.unwrap().is_empty());
2170 }
2171
2172 #[tokio::test]
2173 async fn test_range_scan_across_windows() {
2174 use alloy_consensus::TxLegacy;
2175 use alloy_primitives::{Address, Bloom, Bytes, Log, LogData, Signature};
2176 use reth_db_api::models::StoredBlockBodyIndices;
2177 use reth_ethereum_primitives::{Block, BlockBody, Receipt, TransactionSigned};
2178
2179 let provider = MockEthProvider::default();
2180 let tx = TransactionSigned::new_unhashed(
2181 TxLegacy {
2182 chain_id: Some(1),
2183 gas_price: 21_000,
2184 gas_limit: 21_000,
2185 ..Default::default()
2186 }
2187 .into(),
2188 Signature::test_signature(),
2189 );
2190 let receipt = Receipt {
2191 tx_type: TxType::Legacy,
2192 cumulative_gas_used: 21_000,
2193 logs: vec![Log {
2194 address: Address::ZERO,
2195 data: LogData::new_unchecked(vec![], Bytes::new()),
2196 }],
2197 success: true,
2198 };
2199
2200 let matching_blocks = [5u64, 1_200, 2_400];
2203 let mut parent_hash = FixedBytes::default();
2204 for number in 0..=2_500u64 {
2205 let matches = matching_blocks.contains(&number);
2206 let header = alloy_consensus::Header {
2207 number,
2208 parent_hash,
2209 logs_bloom: if matches { Bloom::from([1u8; 256]) } else { Bloom::default() },
2210 ..Default::default()
2211 };
2212 parent_hash = header.hash_slow();
2213 let transactions = if matches { vec![tx.clone()] } else { Vec::new() };
2214 provider.add_block(
2215 parent_hash,
2216 Block { header, body: BlockBody { transactions, ..Default::default() } },
2217 );
2218 if matches {
2219 let tx_num = matching_blocks.iter().position(|b| *b == number).unwrap() as u64;
2220 provider.add_receipts(number, vec![receipt.clone()]);
2221 provider.add_block_body_indices(
2222 number,
2223 StoredBlockBodyIndices { first_tx_num: tx_num, tx_count: 1 },
2224 );
2225 } else {
2226 provider.add_receipts(number, vec![]);
2227 }
2228 }
2229
2230 let eth_filter = EthFilter::new(
2231 build_test_eth_api(provider),
2232 EthFilterConfig::default(),
2233 Runtime::test(),
2234 );
2235 let logs = eth_filter
2236 .inner
2237 .clone()
2238 .get_logs_in_block_range(Filter::default(), 0, 2_500, QueryLimits::default())
2239 .await
2240 .unwrap();
2241 let blocks = logs.iter().map(|log| log.block_number.unwrap()).collect::<Vec<_>>();
2242 assert_eq!(blocks, matching_blocks);
2243
2244 let err = eth_filter
2247 .inner
2248 .clone()
2249 .get_logs_in_block_range(
2250 Filter::default(),
2251 0,
2252 2_500,
2253 QueryLimits { max_blocks_per_filter: None, max_logs_per_response: Some(1) },
2254 )
2255 .await
2256 .unwrap_err();
2257 assert!(
2258 matches!(
2259 err,
2260 EthFilterError::QueryExceedsMaxResults {
2261 max_logs: 1,
2262 from_block: 0,
2263 to_block: 1_199
2264 }
2265 ),
2266 "{err:?}"
2267 );
2268 }
2269
2270 #[tokio::test]
2271 async fn test_logs_for_filter_pending_to_block_ends_at_head() {
2272 let provider = MockEthProvider::default();
2273 let header = alloy_consensus::Header { number: 2, ..Default::default() };
2274 let hash = header.hash_slow();
2275 provider.add_header(hash, header);
2276 provider.add_receipts(2, vec![]);
2277 provider.set_pending_block_num_hash(Some(alloy_eips::BlockNumHash::new(3, hash)));
2279
2280 let eth_filter = EthFilter::new(
2281 build_test_eth_api(provider),
2282 EthFilterConfig::default(),
2283 Runtime::test(),
2284 );
2285 for filter in [
2286 Filter::new().from_block(0u64).to_block(BlockNumberOrTag::Pending),
2287 Filter::new().select(BlockNumberOrTag::Pending..),
2288 ] {
2289 let logs = eth_filter
2290 .inner
2291 .clone()
2292 .logs_for_filter(filter, QueryLimits::default())
2293 .await
2294 .unwrap();
2295 assert!(logs.is_empty());
2296 }
2297 }
2298
2299 #[tokio::test]
2300 async fn test_logs_for_filter_over_pruned_receipts() {
2301 let provider = MockEthProvider::default();
2302 let mut parent_hash = FixedBytes::default();
2303 for number in 0..=2u64 {
2304 let header = alloy_consensus::Header {
2305 number,
2306 parent_hash,
2307 logs_bloom: alloy_primitives::Bloom::from([1u8; 256]),
2308 ..Default::default()
2309 };
2310 parent_hash = header.hash_slow();
2311 provider.add_block(
2312 parent_hash,
2313 reth_ethereum_primitives::Block { header, body: Default::default() },
2314 );
2315 if number != 1 {
2317 provider.add_receipts(number, vec![]);
2318 }
2319 }
2320
2321 let eth_filter = EthFilter::new(
2322 build_test_eth_api(provider),
2323 EthFilterConfig::default(),
2324 Runtime::test(),
2325 );
2326
2327 let err = eth_filter
2329 .inner
2330 .clone()
2331 .logs_for_filter(Filter::new().from_block(0u64).to_block(2u64), QueryLimits::default())
2332 .await
2333 .unwrap_err();
2334 assert!(matches!(err, EthFilterError::ReceiptsUnavailable(1)), "{err:?}");
2335
2336 let logs = eth_filter
2338 .inner
2339 .clone()
2340 .logs_for_filter(Filter::new().from_block(2u64).to_block(2u64), QueryLimits::default())
2341 .await
2342 .unwrap();
2343 assert!(logs.is_empty());
2344 }
2345}