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