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 PARALLEL_PROCESSING_THRESHOLD: usize = 1000;
86
87const DEFAULT_PARALLEL_CONCURRENCY: usize = 4;
89
90pub struct EthFilter<Eth: EthApiTypes> {
94 inner: Arc<EthFilterInner<Eth>>,
96}
97
98impl<Eth> Clone for EthFilter<Eth>
99where
100 Eth: EthApiTypes,
101{
102 fn clone(&self) -> Self {
103 Self { inner: self.inner.clone() }
104 }
105}
106
107impl<Eth> EthFilter<Eth>
108where
109 Eth: EthApiTypes + 'static,
110{
111 pub fn new(eth_api: Eth, config: EthFilterConfig, task_spawner: Runtime) -> Self {
139 let EthFilterConfig { max_blocks_per_filter, max_logs_per_response, stale_filter_ttl } =
140 config;
141 let inner = EthFilterInner {
142 eth_api,
143 active_filters: ActiveFilters::new(),
144 id_provider: Arc::new(EthSubscriptionIdProvider::default()),
145 max_headers_range: MAX_HEADERS_RANGE,
146 task_spawner,
147 stale_filter_ttl,
148 query_limits: QueryLimits { max_blocks_per_filter, max_logs_per_response },
149 };
150
151 let eth_filter = Self { inner: Arc::new(inner) };
152
153 let this = eth_filter.clone();
154 eth_filter.inner.task_spawner.spawn_critical_task(
155 "eth-filters_stale-filters-clean",
156 async move {
157 this.watch_and_clear_stale_filters().await;
158 },
159 );
160
161 eth_filter
162 }
163
164 pub fn active_filters(&self) -> &ActiveFilters<RpcTransaction<Eth::NetworkTypes>> {
166 &self.inner.active_filters
167 }
168
169 async fn watch_and_clear_stale_filters(&self) {
172 let mut interval = tokio::time::interval_at(
173 tokio::time::Instant::now() + self.inner.stale_filter_ttl,
174 self.inner.stale_filter_ttl,
175 );
176 interval.set_missed_tick_behavior(MissedTickBehavior::Delay);
177 loop {
178 interval.tick().await;
179 self.clear_stale_filters(Instant::now()).await;
180 }
181 }
182
183 pub async fn clear_stale_filters(&self, now: Instant) {
186 trace!(target: "rpc::eth", "clear stale filters");
187 let mut filters = self.active_filters().inner.lock().await;
188 filters.retain(|id, filter| {
189 let is_valid = (now - filter.last_poll_timestamp) < self.inner.stale_filter_ttl;
190
191 if !is_valid {
192 trace!(target: "rpc::eth", "evict filter with id: {:?}", id);
193 }
194
195 is_valid
196 });
197 filters.shrink_to_fit();
198 }
199}
200
201impl<Eth> EthFilter<Eth>
202where
203 Eth: FullEthApiTypes<Provider: BlockReader + BlockIdReader>
204 + RpcNodeCoreExt
205 + LoadReceipt
206 + EthBlocks
207 + 'static,
208{
209 fn provider(&self) -> &Eth::Provider {
211 self.inner.eth_api.provider()
212 }
213
214 fn pool(&self) -> &Eth::Pool {
216 self.inner.eth_api.pool()
217 }
218
219 pub async fn filter_changes(
221 &self,
222 id: FilterId,
223 ) -> Result<
224 FilterChanges<RpcTransaction<Eth::NetworkTypes>, RpcLog<Eth::NetworkTypes>>,
225 EthFilterError,
226 > {
227 let info = self.provider().chain_info()?;
228 let best_number = info.best_number;
229
230 let (start_block, kind) = {
233 let mut filters = self.inner.active_filters.inner.lock().await;
234 let filter = filters.get_mut(&id).ok_or(EthFilterError::FilterNotFound(id))?;
235
236 if filter.block > best_number {
237 return Ok(FilterChanges::Empty)
239 }
240
241 let mut block = best_number + 1;
245 std::mem::swap(&mut filter.block, &mut block);
246 filter.last_poll_timestamp = Instant::now();
247
248 (block, filter.kind.clone())
249 };
250
251 match kind {
252 FilterKind::PendingTransaction(filter) => Ok(match filter.drain().await {
253 FilterChanges::Empty => FilterChanges::Empty,
254 FilterChanges::Hashes(hashes) => FilterChanges::Hashes(hashes),
255 FilterChanges::Transactions(transactions) => {
256 FilterChanges::Transactions(transactions)
257 }
258 FilterChanges::Logs(_) => unreachable!("pending transaction filter returned logs"),
259 }),
260 FilterKind::Block => {
261 let end_block = best_number + 1;
264 let block_hashes =
265 self.provider().canonical_hashes_range(start_block, end_block).map_err(
266 |_| EthApiError::HeaderRangeNotFound(start_block.into(), end_block.into()),
267 )?;
268 Ok(FilterChanges::Hashes(block_hashes))
269 }
270 FilterKind::Log(filter) => {
271 let (from_block_number, to_block_number) = match filter.block_option {
272 FilterBlockOption::Range { from_block, to_block } => {
273 let from = from_block
274 .map(|num| self.provider().convert_block_number(num))
275 .transpose()?
276 .flatten();
277 let to = to_block
278 .map(|num| self.provider().convert_block_number(num))
279 .transpose()?
280 .flatten();
281 logs_utils::get_filter_block_range(from, to, start_block, info)?
282 }
283 FilterBlockOption::AtBlockHash(block_hash) => {
284 let block_number = self
288 .provider()
289 .block_number(block_hash)?
290 .ok_or(ProviderError::HeaderNotFound(block_hash.into()))?;
291 (block_number, block_number)
292 }
293 };
294 let logs = self
295 .inner
296 .clone()
297 .get_logs_in_block_range(
298 *filter,
299 from_block_number,
300 to_block_number,
301 self.inner.query_limits,
302 )
303 .await?;
304 Ok(FilterChanges::Logs(logs))
305 }
306 }
307 }
308
309 pub async fn filter_logs(
315 &self,
316 id: FilterId,
317 ) -> Result<Vec<RpcLog<Eth::NetworkTypes>>, EthFilterError> {
318 let filter = {
319 let mut filters = self.inner.active_filters.inner.lock().await;
320 let filter =
321 filters.get_mut(&id).ok_or_else(|| EthFilterError::FilterNotFound(id.clone()))?;
322 if let FilterKind::Log(ref inner_filter) = filter.kind {
323 filter.last_poll_timestamp = Instant::now();
324 *inner_filter.clone()
325 } else {
326 return Err(EthFilterError::FilterNotFound(id))
328 }
329 };
330
331 self.logs_for_filter(filter, self.inner.query_limits).await
332 }
333
334 async fn logs_for_filter(
336 &self,
337 filter: Filter,
338 limits: QueryLimits,
339 ) -> Result<Vec<RpcLog<Eth::NetworkTypes>>, EthFilterError> {
340 self.inner.clone().logs_for_filter(filter, limits).await
341 }
342}
343
344#[async_trait]
345impl<Eth> EthFilterApiServer<RpcTransaction<Eth::NetworkTypes>, RpcLog<Eth::NetworkTypes>>
346 for EthFilter<Eth>
347where
348 Eth: FullEthApiTypes + RpcNodeCoreExt + LoadReceipt + EthBlocks + 'static,
349{
350 async fn new_filter(&self, filter: Filter) -> RpcResult<FilterId> {
352 trace!(target: "rpc::eth", "Serving eth_newFilter");
353 self.inner
354 .install_filter(FilterKind::<RpcTransaction<Eth::NetworkTypes>>::Log(Box::new(filter)))
355 .await
356 }
357
358 async fn new_block_filter(&self) -> RpcResult<FilterId> {
360 trace!(target: "rpc::eth", "Serving eth_newBlockFilter");
361 self.inner.install_filter(FilterKind::<RpcTransaction<Eth::NetworkTypes>>::Block).await
362 }
363
364 async fn new_pending_transaction_filter(
366 &self,
367 kind: Option<PendingTransactionFilterKind>,
368 ) -> RpcResult<FilterId> {
369 trace!(target: "rpc::eth", "Serving eth_newPendingTransactionFilter");
370
371 let transaction_kind = match kind.unwrap_or_default() {
372 PendingTransactionFilterKind::Hashes => {
373 let receiver = self.pool().pending_transactions_listener();
374 let pending_txs_receiver = PendingTransactionsReceiver::new(receiver);
375 FilterKind::PendingTransaction(PendingTransactionKind::Hashes(pending_txs_receiver))
376 }
377 PendingTransactionFilterKind::Full => {
378 let stream = self.pool().new_pending_pool_transactions_listener();
379 let full_txs_receiver = FullTransactionsReceiver::new(
380 stream,
381 dyn_clone::clone(self.inner.eth_api.converter()),
382 );
383 FilterKind::PendingTransaction(PendingTransactionKind::FullTransaction(Arc::new(
384 full_txs_receiver,
385 )))
386 }
387 };
388
389 self.inner.install_filter(transaction_kind).await
391 }
392
393 async fn filter_changes(
395 &self,
396 id: FilterId,
397 ) -> RpcResult<FilterChanges<RpcTransaction<Eth::NetworkTypes>, RpcLog<Eth::NetworkTypes>>>
398 {
399 trace!(target: "rpc::eth", "Serving eth_getFilterChanges");
400 Ok(Self::filter_changes(self, id).await?)
401 }
402
403 async fn filter_logs(&self, id: FilterId) -> RpcResult<Vec<RpcLog<Eth::NetworkTypes>>> {
409 trace!(target: "rpc::eth", "Serving eth_getFilterLogs");
410 Ok(Self::filter_logs(self, id).await?)
411 }
412
413 async fn uninstall_filter(&self, id: FilterId) -> RpcResult<bool> {
415 trace!(target: "rpc::eth", "Serving eth_uninstallFilter");
416 let mut filters = self.inner.active_filters.inner.lock().await;
417 if filters.remove(&id).is_some() {
418 trace!(target: "rpc::eth::filter", ?id, "uninstalled filter");
419 Ok(true)
420 } else {
421 Ok(false)
422 }
423 }
424
425 async fn logs(&self, filter: Filter) -> RpcResult<Vec<RpcLog<Eth::NetworkTypes>>> {
429 trace!(target: "rpc::eth", "Serving eth_getLogs");
430 Ok(self.logs_for_filter(filter, self.inner.query_limits).await?)
431 }
432}
433
434impl<Eth> std::fmt::Debug for EthFilter<Eth>
435where
436 Eth: EthApiTypes,
437{
438 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
439 f.debug_struct("EthFilter").finish_non_exhaustive()
440 }
441}
442
443#[derive(Debug)]
445struct EthFilterInner<Eth: EthApiTypes> {
446 eth_api: Eth,
448 active_filters: ActiveFilters<RpcTransaction<Eth::NetworkTypes>>,
450 id_provider: Arc<dyn IdProvider>,
452 query_limits: QueryLimits,
454 max_headers_range: u64,
456 task_spawner: Runtime,
458 stale_filter_ttl: Duration,
460}
461
462impl<Eth> EthFilterInner<Eth>
463where
464 Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
465 + EthApiTypes<NetworkTypes: reth_rpc_eth_api::types::RpcTypes>
466 + LoadReceipt
467 + EthBlocks
468 + 'static,
469{
470 fn provider(&self) -> &Eth::Provider {
472 self.eth_api.provider()
473 }
474
475 fn eth_cache(&self) -> &EthStateCache<Eth::Primitives> {
477 self.eth_api.cache()
478 }
479
480 async fn logs_for_filter(
482 self: Arc<Self>,
483 filter: Filter,
484 limits: QueryLimits,
485 ) -> Result<Vec<RpcLog<Eth::NetworkTypes>>, EthFilterError> {
486 match filter.block_option {
487 FilterBlockOption::AtBlockHash(block_hash) => {
488 let Some((receipts, maybe_block)) =
490 self.eth_cache().get_receipts_and_maybe_block(block_hash).await?
491 else {
492 return Err(ProviderError::HeaderNotFound(block_hash.into()).into())
493 };
494
495 let header = if let Some(block) = &maybe_block {
496 block.clone_sealed_header()
497 } else {
498 let header = self
499 .provider()
500 .header_by_hash_or_number(block_hash.into())?
501 .ok_or_else(|| ProviderError::HeaderNotFound(block_hash.into()))?;
502 SealedHeader::new(header, block_hash)
503 };
504
505 let earliest_block = self.provider().earliest_block_number()?;
507 if header.number() < earliest_block {
508 return Err(EthApiError::PrunedHistoryUnavailable {
509 requested: header.number(),
510 earliest_available: earliest_block,
511 }
512 .into());
513 }
514
515 let mut all_logs = Vec::new();
516 append_matching_block_logs(
517 &mut all_logs,
518 self.eth_api.converter(),
519 maybe_block
520 .map(ProviderOrBlock::Block)
521 .unwrap_or_else(|| ProviderOrBlock::Provider(self.provider())),
522 &filter,
523 &header,
524 &receipts,
525 false,
526 )?;
527 Ok(all_logs)
528 }
529 FilterBlockOption::Range { from_block, to_block } => {
530 if from_block.is_some_and(|b| b.is_pending()) {
532 let to_block = to_block.unwrap_or(BlockNumberOrTag::Pending);
533 if !(to_block.is_pending() || to_block.is_number()) {
534 return Ok(Vec::new());
536 }
537 if let Ok(Some(pending_block)) = self.eth_api.local_pending_block().await {
539 if let BlockNumberOrTag::Number(to_block) = to_block &&
540 to_block < pending_block.block.number()
541 {
542 return Ok(Vec::new());
544 }
545
546 let info = self.provider().chain_info()?;
547 if pending_block.block.number() > info.best_number {
548 let mut all_logs = Vec::new();
550 let header = pending_block.block.clone_sealed_header();
551 append_matching_block_logs(
552 &mut all_logs,
553 self.eth_api.converter(),
554 ProviderOrBlock::<Eth::Provider>::Block(pending_block.block),
555 &filter,
556 &header,
557 &pending_block.receipts,
558 false, )?;
560 return Ok(all_logs)
561 }
562 }
563 }
564
565 let info = self.provider().chain_info()?;
566 let start_block = info.best_number;
567 let from = from_block
568 .map(|num| self.provider().convert_block_number(num))
569 .transpose()?
570 .flatten();
571 let to = to_block
572 .map(|num| self.provider().convert_block_number(num))
573 .transpose()?
574 .flatten();
575
576 if let Some(t) = to &&
578 t > info.best_number
579 {
580 return Err(EthFilterError::BlockRangeExceedsHead {
581 requested: t,
582 head: info.best_number,
583 });
584 }
585
586 if let Some(f) = from &&
587 f > info.best_number
588 {
589 return Ok(Vec::new());
591 }
592
593 let (from_block_number, to_block_number) =
594 logs_utils::get_filter_block_range(from, to, start_block, info)?;
595
596 let earliest_block = self.provider().earliest_block_number()?;
598 if from_block_number < earliest_block {
599 return Err(EthApiError::PrunedHistoryUnavailable {
600 requested: from_block_number,
601 earliest_available: earliest_block,
602 }
603 .into());
604 }
605
606 self.get_logs_in_block_range(filter, from_block_number, to_block_number, limits)
607 .await
608 }
609 }
610 }
611
612 async fn install_filter(
614 &self,
615 kind: FilterKind<RpcTransaction<Eth::NetworkTypes>>,
616 ) -> RpcResult<FilterId> {
617 let last_poll_block_number = self.provider().best_block_number().to_rpc_result()?;
618 let subscription_id = self.id_provider.next_id();
619
620 let id = match subscription_id {
621 jsonrpsee_types::SubscriptionId::Num(n) => FilterId::Num(n),
622 jsonrpsee_types::SubscriptionId::Str(s) => FilterId::Str(s.into_owned()),
623 };
624 let mut filters = self.active_filters.inner.lock().await;
625 filters.insert(
626 id.clone(),
627 ActiveFilter {
628 block: last_poll_block_number,
629 last_poll_timestamp: Instant::now(),
630 kind,
631 },
632 );
633 Ok(id)
634 }
635
636 async fn get_logs_in_block_range(
642 self: Arc<Self>,
643 filter: Filter,
644 from_block: u64,
645 to_block: u64,
646 limits: QueryLimits,
647 ) -> Result<Vec<RpcLog<Eth::NetworkTypes>>, EthFilterError> {
648 trace!(target: "rpc::eth::filter", from=from_block, to=to_block, ?filter, "finding logs in range");
649
650 if to_block < from_block {
652 return Err(EthFilterError::InvalidBlockRangeParams)
653 }
654
655 if let Some(max_blocks_per_filter) =
656 limits.max_blocks_per_filter.filter(|limit| to_block - from_block > *limit)
657 {
658 return Err(EthFilterError::QueryExceedsMaxBlocks(max_blocks_per_filter))
659 }
660
661 let (tx, rx) = oneshot::channel();
662 let this = self.clone();
663 self.task_spawner.spawn_blocking_task(async move {
664 let res =
665 this.get_logs_in_block_range_inner(&filter, from_block, to_block, limits).await;
666 let _ = tx.send(res);
667 });
668
669 rx.await.map_err(|_| EthFilterError::InternalError)?
670 }
671
672 async fn get_logs_in_block_range_inner(
681 self: Arc<Self>,
682 filter: &Filter,
683 from_block: u64,
684 to_block: u64,
685 limits: QueryLimits,
686 ) -> Result<Vec<RpcLog<Eth::NetworkTypes>>, EthFilterError> {
687 let mut all_logs = Vec::new();
688 let mut matching_headers = Vec::new();
689
690 let chain_tip = self.provider().best_block_number()?;
692
693 for (from, to) in
695 BlockRangeInclusiveIter::new(from_block..=to_block, self.max_headers_range)
696 {
697 let headers = self.provider().headers_range(from..=to)?;
698
699 let mut headers_iter = headers.into_iter().peekable();
700
701 while let Some(header) = headers_iter.next() {
702 if !filter.matches_bloom(header.logs_bloom()) {
703 continue
704 }
705
706 let current_number = header.number();
707
708 let block_hash = match headers_iter.peek() {
709 Some(next_header) if next_header.number() == current_number + 1 => {
710 next_header.parent_hash()
712 }
713 _ => {
714 header.hash_slow()
716 }
717 };
718
719 matching_headers.push(SealedHeader::new(header, block_hash));
720 }
721 }
722
723 let mut range_mode = RangeMode::new(
725 self.clone(),
726 matching_headers,
727 from_block,
728 to_block,
729 self.max_headers_range,
730 chain_tip,
731 );
732
733 while let Some(ReceiptBlockResult { receipts, recovered_block, header }) =
735 range_mode.next().await?
736 {
737 let num_hash = header.num_hash();
738 append_matching_block_logs(
739 &mut all_logs,
740 self.eth_api.converter(),
741 recovered_block
742 .map(ProviderOrBlock::Block)
743 .unwrap_or_else(|| ProviderOrBlock::Provider(self.provider())),
744 filter,
745 &header,
746 &receipts,
747 false,
748 )?;
749
750 let is_multi_block_range = from_block != to_block;
753 if let Some(max_logs_per_response) = limits.max_logs_per_response &&
754 is_multi_block_range &&
755 all_logs.len() > max_logs_per_response
756 {
757 let retry_to_block =
758 if num_hash.number == from_block { from_block } else { num_hash.number - 1 };
759
760 debug!(
761 target: "rpc::eth::filter",
762 logs_found = all_logs.len(),
763 max_logs_per_response,
764 from_block,
765 to_block = retry_to_block,
766 "Query exceeded max logs per response limit"
767 );
768 return Err(EthFilterError::QueryExceedsMaxResults {
769 max_logs: max_logs_per_response,
770 from_block,
771 to_block: retry_to_block,
772 });
773 }
774 }
775
776 Ok(all_logs)
777 }
778}
779
780#[derive(Debug, Clone, Default)]
782pub struct ActiveFilters<T> {
783 inner: Arc<Mutex<HashMap<FilterId, ActiveFilter<T>>>>,
784}
785
786impl<T> ActiveFilters<T> {
787 pub fn new() -> Self {
789 Self { inner: Arc::new(Mutex::new(HashMap::default())) }
790 }
791
792 pub async fn contains(&self, id: &FilterId) -> bool {
794 self.inner.lock().await.contains_key(id)
795 }
796
797 pub async fn len(&self) -> usize {
799 self.inner.lock().await.len()
800 }
801
802 pub async fn is_empty(&self) -> bool {
804 self.inner.lock().await.is_empty()
805 }
806
807 pub async fn ids(&self) -> Vec<FilterId> {
809 self.inner.lock().await.keys().cloned().collect()
810 }
811}
812
813#[derive(Debug)]
815struct ActiveFilter<T> {
816 block: u64,
818 last_poll_timestamp: Instant,
820 kind: FilterKind<T>,
822}
823
824#[derive(Debug, Clone)]
826struct PendingTransactionsReceiver {
827 txs_receiver: Arc<Mutex<Receiver<TxHash>>>,
828}
829
830impl PendingTransactionsReceiver {
831 fn new(receiver: Receiver<TxHash>) -> Self {
832 Self { txs_receiver: Arc::new(Mutex::new(receiver)) }
833 }
834
835 async fn drain<T>(&self) -> FilterChanges<T> {
837 let mut pending_txs = Vec::new();
838 let mut prepared_stream = self.txs_receiver.lock().await;
839
840 while let Ok(tx_hash) = prepared_stream.try_recv() {
841 pending_txs.push(tx_hash);
842 }
843
844 FilterChanges::Hashes(pending_txs)
846 }
847}
848
849#[derive(Debug, Clone)]
851struct FullTransactionsReceiver<T: PoolTransaction, TxCompat> {
852 txs_stream: Arc<Mutex<NewSubpoolTransactionStream<T>>>,
853 converter: TxCompat,
854}
855
856impl<T, TxCompat> FullTransactionsReceiver<T, TxCompat>
857where
858 T: PoolTransaction + 'static,
859 TxCompat: RpcConvert<Primitives: NodePrimitives<SignedTx = T::Consensus>>,
860{
861 fn new(stream: NewSubpoolTransactionStream<T>, converter: TxCompat) -> Self {
863 Self { txs_stream: Arc::new(Mutex::new(stream)), converter }
864 }
865
866 async fn drain(&self) -> FilterChanges<RpcTransaction<TxCompat::Network>> {
868 let mut pending_txs = Vec::new();
869 let mut prepared_stream = self.txs_stream.lock().await;
870
871 while let Ok(tx) = prepared_stream.try_recv() {
872 match self.converter.fill_pending(tx.transaction.to_consensus()) {
873 Ok(tx) => pending_txs.push(tx),
874 Err(err) => {
875 error!(target: "rpc",
876 %err,
877 "Failed to fill txn with block context"
878 );
879 }
880 }
881 }
882 FilterChanges::Transactions(pending_txs)
883 }
884}
885
886#[async_trait]
888trait FullTransactionsFilter<T>: fmt::Debug + Send + Sync + Unpin + 'static {
889 async fn drain(&self) -> FilterChanges<T>;
890}
891
892#[async_trait]
893impl<T, TxCompat> FullTransactionsFilter<RpcTransaction<TxCompat::Network>>
894 for FullTransactionsReceiver<T, TxCompat>
895where
896 T: PoolTransaction + 'static,
897 TxCompat: RpcConvert<Primitives: NodePrimitives<SignedTx = T::Consensus>> + 'static,
898{
899 async fn drain(&self) -> FilterChanges<RpcTransaction<TxCompat::Network>> {
900 Self::drain(self).await
901 }
902}
903
904#[derive(Debug, Clone)]
910enum PendingTransactionKind<T> {
911 Hashes(PendingTransactionsReceiver),
912 FullTransaction(Arc<dyn FullTransactionsFilter<T>>),
913}
914
915impl<T: 'static> PendingTransactionKind<T> {
916 async fn drain(&self) -> FilterChanges<T> {
917 match self {
918 Self::Hashes(receiver) => receiver.drain().await,
919 Self::FullTransaction(receiver) => receiver.drain().await,
920 }
921 }
922}
923
924#[derive(Clone, Debug)]
925enum FilterKind<T> {
926 Log(Box<Filter>),
927 Block,
928 PendingTransaction(PendingTransactionKind<T>),
929}
930
931#[derive(Debug)]
933struct BlockRangeInclusiveIter {
934 iter: StepBy<RangeInclusive<u64>>,
935 step: u64,
936 end: u64,
937}
938
939impl BlockRangeInclusiveIter {
940 fn new(range: RangeInclusive<u64>, step: u64) -> Self {
941 Self { end: *range.end(), iter: range.step_by(step as usize + 1), step }
942 }
943}
944
945impl Iterator for BlockRangeInclusiveIter {
946 type Item = (u64, u64);
947
948 fn next(&mut self) -> Option<Self::Item> {
949 let start = self.iter.next()?;
950 let end = (start + self.step).min(self.end);
951 if start > end {
952 return None
953 }
954 Some((start, end))
955 }
956}
957
958#[derive(Debug, thiserror::Error)]
960pub enum EthFilterError {
961 #[error("filter not found")]
963 FilterNotFound(FilterId),
964 #[error("invalid block range params")]
966 InvalidBlockRangeParams,
967 #[error("block range extends beyond current head block: requested {requested}, head {head}")]
969 BlockRangeExceedsHead {
970 requested: u64,
972 head: u64,
974 },
975 #[error("query exceeds max block range {0}")]
977 QueryExceedsMaxBlocks(u64),
978 #[error("query exceeds max results {max_logs}, retry with the range {from_block}-{to_block}")]
980 QueryExceedsMaxResults {
981 max_logs: usize,
983 from_block: u64,
985 to_block: u64,
987 },
988 #[error(transparent)]
990 EthAPIError(#[from] EthApiError),
991 #[error("internal filter error")]
993 InternalError,
994}
995
996impl From<EthFilterError> for jsonrpsee::types::error::ErrorObject<'static> {
997 fn from(err: EthFilterError) -> Self {
998 match err {
999 EthFilterError::FilterNotFound(_) => rpc_error_with_code(
1000 jsonrpsee::types::error::INVALID_PARAMS_CODE,
1001 "filter not found",
1002 ),
1003 err @ EthFilterError::InternalError => {
1004 rpc_error_with_code(jsonrpsee::types::error::INTERNAL_ERROR_CODE, err.to_string())
1005 }
1006 EthFilterError::EthAPIError(err) => err.into(),
1007 err @ (EthFilterError::InvalidBlockRangeParams |
1008 EthFilterError::QueryExceedsMaxBlocks(_) |
1009 EthFilterError::QueryExceedsMaxResults { .. } |
1010 EthFilterError::BlockRangeExceedsHead { .. }) => {
1011 rpc_error_with_code(jsonrpsee::types::error::INVALID_PARAMS_CODE, err.to_string())
1012 }
1013 }
1014 }
1015}
1016
1017impl From<ProviderError> for EthFilterError {
1018 fn from(err: ProviderError) -> Self {
1019 Self::EthAPIError(err.into())
1020 }
1021}
1022
1023impl From<logs_utils::FilterBlockRangeError> for EthFilterError {
1024 fn from(err: logs_utils::FilterBlockRangeError) -> Self {
1025 match err {
1026 logs_utils::FilterBlockRangeError::InvalidBlockRange => Self::InvalidBlockRangeParams,
1027 logs_utils::FilterBlockRangeError::BlockRangeExceedsHead { requested, head } => {
1028 Self::BlockRangeExceedsHead { requested, head }
1029 }
1030 }
1031 }
1032}
1033
1034struct ReceiptBlockResult<P>
1037where
1038 P: ReceiptProvider + BlockReader,
1039{
1040 receipts: Arc<Vec<ProviderReceipt<P>>>,
1042 recovered_block: Option<Arc<reth_primitives_traits::RecoveredBlock<ProviderBlock<P>>>>,
1044 header: SealedHeader<<P as HeaderProvider>::Header>,
1046}
1047
1048enum RangeMode<
1050 Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1051 + EthApiTypes
1052 + LoadReceipt
1053 + EthBlocks
1054 + 'static,
1055> {
1056 Cached(CachedMode<Eth>),
1058 Range(RangeBlockMode<Eth>),
1060}
1061
1062impl<
1063 Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1064 + EthApiTypes
1065 + LoadReceipt
1066 + EthBlocks
1067 + 'static,
1068 > RangeMode<Eth>
1069{
1070 fn new(
1072 filter_inner: Arc<EthFilterInner<Eth>>,
1073 sealed_headers: Vec<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>,
1074 from_block: u64,
1075 to_block: u64,
1076 max_headers_range: u64,
1077 chain_tip: u64,
1078 ) -> Self {
1079 let block_count = to_block - from_block + 1;
1080 let distance_from_tip = chain_tip.saturating_sub(to_block);
1081
1082 let use_cached_mode =
1084 Self::should_use_cached_mode(&sealed_headers, block_count, distance_from_tip);
1085
1086 if use_cached_mode && !sealed_headers.is_empty() {
1087 Self::Cached(CachedMode { filter_inner, headers_iter: sealed_headers.into_iter() })
1088 } else {
1089 Self::Range(RangeBlockMode {
1090 filter_inner,
1091 iter: sealed_headers.into_iter().peekable(),
1092 next: VecDeque::new(),
1093 max_range: max_headers_range as usize,
1094 pending_tasks: FuturesOrdered::new(),
1095 })
1096 }
1097 }
1098
1099 const fn should_use_cached_mode(
1101 headers: &[SealedHeader<<Eth::Provider as HeaderProvider>::Header>],
1102 block_count: u64,
1103 distance_from_tip: u64,
1104 ) -> bool {
1105 let bloom_matches = headers.len();
1107
1108 let adjusted_threshold = Self::calculate_adjusted_threshold(block_count, bloom_matches);
1110
1111 block_count <= adjusted_threshold && distance_from_tip <= adjusted_threshold
1112 }
1113
1114 const fn calculate_adjusted_threshold(block_count: u64, bloom_matches: usize) -> u64 {
1116 if block_count <= BLOOM_ADJUSTMENT_MIN_BLOCKS {
1118 return CACHED_MODE_BLOCK_THRESHOLD;
1119 }
1120
1121 match bloom_matches {
1122 n if n > HIGH_BLOOM_MATCH_THRESHOLD => CACHED_MODE_BLOCK_THRESHOLD / 2,
1123 n if n > MODERATE_BLOOM_MATCH_THRESHOLD => (CACHED_MODE_BLOCK_THRESHOLD * 3) / 4,
1124 _ => CACHED_MODE_BLOCK_THRESHOLD,
1125 }
1126 }
1127
1128 async fn next(&mut self) -> Result<Option<ReceiptBlockResult<Eth::Provider>>, EthFilterError> {
1130 match self {
1131 Self::Cached(cached) => cached.next().await,
1132 Self::Range(range) => range.next().await,
1133 }
1134 }
1135}
1136
1137struct CachedMode<
1139 Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1140 + EthApiTypes
1141 + LoadReceipt
1142 + EthBlocks
1143 + 'static,
1144> {
1145 filter_inner: Arc<EthFilterInner<Eth>>,
1146 headers_iter: std::vec::IntoIter<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>,
1147}
1148
1149impl<
1150 Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1151 + EthApiTypes
1152 + LoadReceipt
1153 + EthBlocks
1154 + 'static,
1155 > CachedMode<Eth>
1156{
1157 async fn next(&mut self) -> Result<Option<ReceiptBlockResult<Eth::Provider>>, EthFilterError> {
1158 for header in self.headers_iter.by_ref() {
1159 if let Some((receipts, maybe_block)) =
1161 self.filter_inner.eth_cache().get_receipts_and_maybe_block(header.hash()).await?
1162 {
1163 return Ok(Some(ReceiptBlockResult {
1164 receipts,
1165 recovered_block: maybe_block,
1166 header,
1167 }));
1168 }
1169 }
1170
1171 Ok(None) }
1173}
1174
1175type ReceiptFetchFuture<P> =
1177 Pin<Box<dyn Future<Output = Result<Vec<ReceiptBlockResult<P>>, EthFilterError>> + Send>>;
1178
1179struct RangeBlockMode<
1181 Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1182 + EthApiTypes
1183 + LoadReceipt
1184 + EthBlocks
1185 + 'static,
1186> {
1187 filter_inner: Arc<EthFilterInner<Eth>>,
1188 iter: Peekable<std::vec::IntoIter<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>>,
1189 next: VecDeque<ReceiptBlockResult<Eth::Provider>>,
1190 max_range: usize,
1191 pending_tasks: FuturesOrdered<ReceiptFetchFuture<Eth::Provider>>,
1193}
1194
1195impl<
1196 Eth: RpcNodeCoreExt<Provider: BlockIdReader, Pool: TransactionPool>
1197 + EthApiTypes
1198 + LoadReceipt
1199 + EthBlocks
1200 + 'static,
1201 > RangeBlockMode<Eth>
1202{
1203 async fn next(&mut self) -> Result<Option<ReceiptBlockResult<Eth::Provider>>, EthFilterError> {
1204 loop {
1205 if let Some(result) = self.next.pop_front() {
1207 return Ok(Some(result));
1208 }
1209
1210 if let Some(task_result) = self.pending_tasks.next().await {
1212 self.next.extend(task_result?);
1213 continue;
1214 }
1215
1216 let Some(next_header) = self.iter.next() else {
1218 return Ok(None);
1220 };
1221
1222 let mut range_headers = Vec::with_capacity(self.max_range);
1223 range_headers.push(next_header);
1224
1225 while range_headers.len() < self.max_range {
1227 let Some(peeked) = self.iter.peek() else { break };
1228 let Some(last_header) = range_headers.last() else { break };
1229
1230 let expected_next = last_header.number() + 1;
1231 if peeked.number() != expected_next {
1232 trace!(
1233 target: "rpc::eth::filter",
1234 last_block = last_header.number(),
1235 next_block = peeked.number(),
1236 expected = expected_next,
1237 range_size = range_headers.len(),
1238 "Non-consecutive block detected, stopping range collection"
1239 );
1240 break; }
1242
1243 let Some(next_header) = self.iter.next() else { break };
1244 range_headers.push(next_header);
1245 }
1246
1247 let remaining_headers = self.iter.len() + range_headers.len();
1249 if remaining_headers >= PARALLEL_PROCESSING_THRESHOLD {
1250 self.spawn_parallel_tasks(range_headers);
1251 } else {
1253 if let Some(result) = self.process_small_range(range_headers).await? {
1255 return Ok(Some(result));
1256 }
1257 }
1259 }
1260 }
1261
1262 async fn process_small_range(
1266 &mut self,
1267 range_headers: Vec<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>,
1268 ) -> Result<Option<ReceiptBlockResult<Eth::Provider>>, EthFilterError> {
1269 for header in range_headers {
1271 let (maybe_block, maybe_receipts) = self
1273 .filter_inner
1274 .eth_cache()
1275 .maybe_cached_block_and_receipts(header.hash())
1276 .await?;
1277
1278 let receipts = match maybe_receipts {
1279 Some(receipts) => receipts,
1280 None => {
1281 match self.filter_inner.provider().receipts_by_block(header.hash().into())? {
1283 Some(receipts) => Arc::new(receipts),
1284 None => continue, }
1286 }
1287 };
1288
1289 if !receipts.is_empty() {
1290 self.next.push_back(ReceiptBlockResult {
1291 receipts,
1292 recovered_block: maybe_block,
1293 header,
1294 });
1295 }
1296 }
1297
1298 Ok(self.next.pop_front())
1299 }
1300
1301 fn spawn_parallel_tasks(
1306 &mut self,
1307 range_headers: Vec<SealedHeader<<Eth::Provider as HeaderProvider>::Header>>,
1308 ) {
1309 let chunk_size = std::cmp::max(range_headers.len() / DEFAULT_PARALLEL_CONCURRENCY, 1);
1311 let header_chunks = range_headers
1312 .into_iter()
1313 .chunks(chunk_size)
1314 .into_iter()
1315 .map(|chunk| chunk.collect::<Vec<_>>())
1316 .collect::<Vec<_>>();
1317
1318 for chunk_headers in header_chunks {
1320 let filter_inner = self.filter_inner.clone();
1321 let chunk_task = Box::pin(async move {
1322 let chunk_task = tokio::task::spawn_blocking(move || {
1323 let mut chunk_results = Vec::with_capacity(chunk_headers.len());
1324
1325 for header in chunk_headers {
1326 let receipts = match filter_inner
1329 .provider()
1330 .receipts_by_block(header.hash().into())?
1331 {
1332 Some(receipts) => Arc::new(receipts),
1333 None => continue, };
1335
1336 if !receipts.is_empty() {
1337 chunk_results.push(ReceiptBlockResult {
1338 receipts,
1339 recovered_block: None,
1340 header,
1341 });
1342 }
1343 }
1344
1345 Ok(chunk_results)
1346 });
1347
1348 match chunk_task.await {
1350 Ok(Ok(chunk_results)) => Ok(chunk_results),
1351 Ok(Err(e)) => Err(e),
1352 Err(join_err) => {
1353 trace!(target: "rpc::eth::filter", error = ?join_err, "Task join error");
1354 Err(EthFilterError::InternalError)
1355 }
1356 }
1357 });
1358
1359 self.pending_tasks.push_back(chunk_task);
1360 }
1361 }
1362}
1363
1364#[cfg(test)]
1365mod tests {
1366 use super::*;
1367 use crate::{eth::EthApi, EthApiBuilder};
1368 use alloy_network::Ethereum;
1369 use alloy_primitives::FixedBytes;
1370 use rand::Rng;
1371 use reth_chainspec::{ChainSpec, ChainSpecProvider};
1372 use reth_ethereum_primitives::TxType;
1373 use reth_evm_ethereum::EthEvmConfig;
1374 use reth_network_api::noop::NoopNetwork;
1375 use reth_provider::test_utils::MockEthProvider;
1376 use reth_rpc_convert::RpcConverter;
1377 use reth_rpc_eth_api::node::RpcNodeCoreAdapter;
1378 use reth_rpc_eth_types::receipt::EthReceiptConverter;
1379 use reth_tasks::Runtime;
1380 use reth_testing_utils::generators;
1381 use reth_transaction_pool::test_utils::{testing_pool, TestPool};
1382 use std::{collections::VecDeque, sync::Arc};
1383
1384 #[test]
1385 fn test_block_range_iter() {
1386 let mut rng = generators::rng();
1387
1388 let start = rng.random::<u32>() as u64;
1389 let end = start.saturating_add(rng.random::<u32>() as u64);
1390 let step = rng.random::<u16>() as u64;
1391 let range = start..=end;
1392 let mut iter = BlockRangeInclusiveIter::new(range.clone(), step);
1393 let (from, mut end) = iter.next().unwrap();
1394 assert_eq!(from, start);
1395 assert_eq!(end, (from + step).min(*range.end()));
1396
1397 for (next_from, next_end) in iter {
1398 assert_eq!(next_from, end + 1);
1400 end = next_end;
1401 }
1402
1403 assert_eq!(end, *range.end());
1404 }
1405
1406 #[expect(clippy::type_complexity)]
1408 fn build_test_eth_api(
1409 provider: MockEthProvider,
1410 ) -> EthApi<
1411 RpcNodeCoreAdapter<MockEthProvider, TestPool, NoopNetwork, EthEvmConfig>,
1412 RpcConverter<Ethereum, EthEvmConfig, EthReceiptConverter<ChainSpec>>,
1413 > {
1414 EthApiBuilder::new(
1415 provider.clone(),
1416 testing_pool(),
1417 NoopNetwork::default(),
1418 EthEvmConfig::new(provider.chain_spec()),
1419 )
1420 .build()
1421 }
1422
1423 #[tokio::test]
1424 async fn test_range_block_mode_empty_range() {
1425 let provider = MockEthProvider::default();
1426 let eth_api = build_test_eth_api(provider);
1427
1428 let eth_filter =
1429 super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1430 let filter_inner = eth_filter.inner;
1431
1432 let headers = vec![];
1433 let max_range = 100;
1434
1435 let mut range_mode = RangeBlockMode {
1436 filter_inner,
1437 iter: headers.into_iter().peekable(),
1438 next: VecDeque::new(),
1439 max_range,
1440 pending_tasks: FuturesOrdered::new(),
1441 };
1442
1443 let result = range_mode.next().await;
1444 assert!(result.is_ok());
1445 assert!(result.unwrap().is_none());
1446 }
1447
1448 #[tokio::test]
1449 async fn test_range_block_mode_queued_results_priority() {
1450 let provider = MockEthProvider::default();
1451 let eth_api = build_test_eth_api(provider);
1452
1453 let eth_filter =
1454 super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1455 let filter_inner = eth_filter.inner;
1456
1457 let headers = vec![
1458 SealedHeader::new(
1459 alloy_consensus::Header { number: 100, ..Default::default() },
1460 FixedBytes::random(),
1461 ),
1462 SealedHeader::new(
1463 alloy_consensus::Header { number: 101, ..Default::default() },
1464 FixedBytes::random(),
1465 ),
1466 ];
1467
1468 let expected_block_hash_1 = FixedBytes::from([1u8; 32]);
1470 let expected_block_hash_2 = FixedBytes::from([2u8; 32]);
1471
1472 let mock_receipt_1 = reth_ethereum_primitives::Receipt {
1474 tx_type: TxType::Legacy,
1475 cumulative_gas_used: 100_000,
1476 logs: vec![],
1477 success: true,
1478 };
1479 let mock_receipt_2 = reth_ethereum_primitives::Receipt {
1480 tx_type: TxType::Eip1559,
1481 cumulative_gas_used: 200_000,
1482 logs: vec![],
1483 success: true,
1484 };
1485 let mock_receipt_3 = reth_ethereum_primitives::Receipt {
1486 tx_type: TxType::Eip2930,
1487 cumulative_gas_used: 150_000,
1488 logs: vec![],
1489 success: false, };
1491
1492 let mock_result_1 = ReceiptBlockResult {
1493 receipts: Arc::new(vec![mock_receipt_1.clone(), mock_receipt_2.clone()]),
1494 recovered_block: None,
1495 header: SealedHeader::new(
1496 alloy_consensus::Header { number: 42, ..Default::default() },
1497 expected_block_hash_1,
1498 ),
1499 };
1500
1501 let mock_result_2 = ReceiptBlockResult {
1502 receipts: Arc::new(vec![mock_receipt_3.clone()]),
1503 recovered_block: None,
1504 header: SealedHeader::new(
1505 alloy_consensus::Header { number: 43, ..Default::default() },
1506 expected_block_hash_2,
1507 ),
1508 };
1509
1510 let mut range_mode = RangeBlockMode {
1511 filter_inner,
1512 iter: headers.into_iter().peekable(),
1513 next: VecDeque::from([mock_result_1, mock_result_2]), max_range: 100,
1515 pending_tasks: FuturesOrdered::new(),
1516 };
1517
1518 let result1 = range_mode.next().await;
1520 assert!(result1.is_ok());
1521 let receipt_result1 = result1.unwrap().unwrap();
1522 assert_eq!(receipt_result1.header.hash(), expected_block_hash_1);
1523 assert_eq!(receipt_result1.header.number, 42);
1524
1525 assert_eq!(receipt_result1.receipts.len(), 2);
1527 assert_eq!(receipt_result1.receipts[0].tx_type, mock_receipt_1.tx_type);
1528 assert_eq!(
1529 receipt_result1.receipts[0].cumulative_gas_used,
1530 mock_receipt_1.cumulative_gas_used
1531 );
1532 assert_eq!(receipt_result1.receipts[0].success, mock_receipt_1.success);
1533 assert_eq!(receipt_result1.receipts[1].tx_type, mock_receipt_2.tx_type);
1534 assert_eq!(
1535 receipt_result1.receipts[1].cumulative_gas_used,
1536 mock_receipt_2.cumulative_gas_used
1537 );
1538 assert_eq!(receipt_result1.receipts[1].success, mock_receipt_2.success);
1539
1540 let result2 = range_mode.next().await;
1542 assert!(result2.is_ok());
1543 let receipt_result2 = result2.unwrap().unwrap();
1544 assert_eq!(receipt_result2.header.hash(), expected_block_hash_2);
1545 assert_eq!(receipt_result2.header.number, 43);
1546
1547 assert_eq!(receipt_result2.receipts.len(), 1);
1549 assert_eq!(receipt_result2.receipts[0].tx_type, mock_receipt_3.tx_type);
1550 assert_eq!(
1551 receipt_result2.receipts[0].cumulative_gas_used,
1552 mock_receipt_3.cumulative_gas_used
1553 );
1554 assert_eq!(receipt_result2.receipts[0].success, mock_receipt_3.success);
1555
1556 assert!(range_mode.next.is_empty());
1558
1559 let result3 = range_mode.next().await;
1560 assert!(result3.is_ok());
1561 }
1562
1563 #[tokio::test]
1564 async fn test_range_block_mode_single_block_no_receipts() {
1565 let provider = MockEthProvider::default();
1566 let eth_api = build_test_eth_api(provider);
1567
1568 let eth_filter =
1569 super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1570 let filter_inner = eth_filter.inner;
1571
1572 let headers = vec![SealedHeader::new(
1573 alloy_consensus::Header { number: 100, ..Default::default() },
1574 FixedBytes::random(),
1575 )];
1576
1577 let mut range_mode = RangeBlockMode {
1578 filter_inner,
1579 iter: headers.into_iter().peekable(),
1580 next: VecDeque::new(),
1581 max_range: 100,
1582 pending_tasks: FuturesOrdered::new(),
1583 };
1584
1585 let result = range_mode.next().await;
1586 assert!(result.is_ok());
1587 }
1588
1589 #[tokio::test]
1590 async fn test_range_block_mode_provider_receipts() {
1591 let provider = MockEthProvider::default();
1592
1593 let header_1 = alloy_consensus::Header { number: 100, ..Default::default() };
1594 let header_2 = alloy_consensus::Header { number: 101, ..Default::default() };
1595 let header_3 = alloy_consensus::Header { number: 102, ..Default::default() };
1596
1597 let block_hash_1 = FixedBytes::random();
1598 let block_hash_2 = FixedBytes::random();
1599 let block_hash_3 = FixedBytes::random();
1600
1601 provider.add_header(block_hash_1, header_1.clone());
1602 provider.add_header(block_hash_2, header_2.clone());
1603 provider.add_header(block_hash_3, header_3.clone());
1604
1605 let mock_log = alloy_primitives::Log {
1607 address: alloy_primitives::Address::ZERO,
1608 data: alloy_primitives::LogData::new_unchecked(vec![], alloy_primitives::Bytes::new()),
1609 };
1610
1611 let receipt_100_1 = reth_ethereum_primitives::Receipt {
1612 tx_type: TxType::Legacy,
1613 cumulative_gas_used: 21_000,
1614 logs: vec![mock_log.clone()],
1615 success: true,
1616 };
1617 let receipt_100_2 = reth_ethereum_primitives::Receipt {
1618 tx_type: TxType::Eip1559,
1619 cumulative_gas_used: 42_000,
1620 logs: vec![mock_log.clone()],
1621 success: true,
1622 };
1623 let receipt_101_1 = reth_ethereum_primitives::Receipt {
1624 tx_type: TxType::Eip2930,
1625 cumulative_gas_used: 30_000,
1626 logs: vec![mock_log.clone()],
1627 success: false,
1628 };
1629
1630 provider.add_receipts(100, vec![receipt_100_1.clone(), receipt_100_2.clone()]);
1631 provider.add_receipts(101, vec![receipt_101_1.clone()]);
1632
1633 let eth_api = build_test_eth_api(provider);
1634
1635 let eth_filter =
1636 super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1637 let filter_inner = eth_filter.inner;
1638
1639 let headers = vec![
1640 SealedHeader::new(header_1, block_hash_1),
1641 SealedHeader::new(header_2, block_hash_2),
1642 SealedHeader::new(header_3, block_hash_3),
1643 ];
1644
1645 let mut range_mode = RangeBlockMode {
1646 filter_inner,
1647 iter: headers.into_iter().peekable(),
1648 next: VecDeque::new(),
1649 max_range: 3, pending_tasks: FuturesOrdered::new(),
1651 };
1652
1653 let result = range_mode.next().await;
1655 assert!(result.is_ok());
1656 let receipt_result = result.unwrap().unwrap();
1657
1658 assert_eq!(receipt_result.header.hash(), block_hash_1);
1659 assert_eq!(receipt_result.header.number, 100);
1660 assert_eq!(receipt_result.receipts.len(), 2);
1661
1662 assert_eq!(receipt_result.receipts[0].tx_type, receipt_100_1.tx_type);
1664 assert_eq!(
1665 receipt_result.receipts[0].cumulative_gas_used,
1666 receipt_100_1.cumulative_gas_used
1667 );
1668 assert_eq!(receipt_result.receipts[0].success, receipt_100_1.success);
1669
1670 assert_eq!(receipt_result.receipts[1].tx_type, receipt_100_2.tx_type);
1671 assert_eq!(
1672 receipt_result.receipts[1].cumulative_gas_used,
1673 receipt_100_2.cumulative_gas_used
1674 );
1675 assert_eq!(receipt_result.receipts[1].success, receipt_100_2.success);
1676
1677 let result2 = range_mode.next().await;
1679 assert!(result2.is_ok());
1680 let receipt_result2 = result2.unwrap().unwrap();
1681
1682 assert_eq!(receipt_result2.header.hash(), block_hash_2);
1683 assert_eq!(receipt_result2.header.number, 101);
1684 assert_eq!(receipt_result2.receipts.len(), 1);
1685
1686 assert_eq!(receipt_result2.receipts[0].tx_type, receipt_101_1.tx_type);
1688 assert_eq!(
1689 receipt_result2.receipts[0].cumulative_gas_used,
1690 receipt_101_1.cumulative_gas_used
1691 );
1692 assert_eq!(receipt_result2.receipts[0].success, receipt_101_1.success);
1693
1694 let result3 = range_mode.next().await;
1696 assert!(result3.is_ok());
1697 assert!(result3.unwrap().is_none());
1698 }
1699
1700 #[tokio::test]
1701 async fn test_range_block_mode_iterator_exhaustion() {
1702 let provider = MockEthProvider::default();
1703
1704 let header_100 = alloy_consensus::Header { number: 100, ..Default::default() };
1705 let header_101 = alloy_consensus::Header { number: 101, ..Default::default() };
1706
1707 let block_hash_100 = FixedBytes::random();
1708 let block_hash_101 = FixedBytes::random();
1709
1710 provider.add_header(block_hash_100, header_100.clone());
1712 provider.add_header(block_hash_101, header_101.clone());
1713
1714 let mock_receipt = reth_ethereum_primitives::Receipt {
1716 tx_type: TxType::Legacy,
1717 cumulative_gas_used: 21_000,
1718 logs: vec![],
1719 success: true,
1720 };
1721 provider.add_receipts(100, vec![mock_receipt.clone()]);
1722 provider.add_receipts(101, vec![mock_receipt.clone()]);
1723
1724 let eth_api = build_test_eth_api(provider);
1725
1726 let eth_filter =
1727 super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1728 let filter_inner = eth_filter.inner;
1729
1730 let headers = vec![
1731 SealedHeader::new(header_100, block_hash_100),
1732 SealedHeader::new(header_101, block_hash_101),
1733 ];
1734
1735 let mut range_mode = RangeBlockMode {
1736 filter_inner,
1737 iter: headers.into_iter().peekable(),
1738 next: VecDeque::new(),
1739 max_range: 1,
1740 pending_tasks: FuturesOrdered::new(),
1741 };
1742
1743 let result1 = range_mode.next().await;
1744 assert!(result1.is_ok());
1745 assert!(result1.unwrap().is_some()); assert!(range_mode.iter.peek().is_some()); let result2 = range_mode.next().await;
1750 assert!(result2.is_ok());
1751 assert!(result2.unwrap().is_some()); assert!(range_mode.iter.peek().is_none());
1755
1756 let result3 = range_mode.next().await;
1758 assert!(result3.is_ok());
1759 assert!(result3.unwrap().is_none());
1760 }
1761
1762 #[tokio::test]
1763 async fn test_cached_mode_with_mock_receipts() {
1764 let test_hash = FixedBytes::from([42u8; 32]);
1766 let test_block_number = 100u64;
1767 let test_header = SealedHeader::new(
1768 alloy_consensus::Header {
1769 number: test_block_number,
1770 gas_used: 50_000,
1771 ..Default::default()
1772 },
1773 test_hash,
1774 );
1775
1776 let mock_log = alloy_primitives::Log {
1778 address: alloy_primitives::Address::ZERO,
1779 data: alloy_primitives::LogData::new_unchecked(vec![], alloy_primitives::Bytes::new()),
1780 };
1781
1782 let mock_receipt = reth_ethereum_primitives::Receipt {
1783 tx_type: TxType::Legacy,
1784 cumulative_gas_used: 21_000,
1785 logs: vec![mock_log],
1786 success: true,
1787 };
1788
1789 let provider = MockEthProvider::default();
1790 provider.add_header(test_hash, test_header.header().clone());
1791 provider.add_receipts(test_block_number, vec![mock_receipt.clone()]);
1792
1793 let eth_api = build_test_eth_api(provider);
1794 let eth_filter =
1795 super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1796 let filter_inner = eth_filter.inner;
1797
1798 let headers = vec![test_header.clone()];
1799
1800 let mut cached_mode = CachedMode { filter_inner, headers_iter: headers.into_iter() };
1801
1802 let result = cached_mode.next().await.expect("next should succeed");
1804 let receipt_block_result = result.expect("should have receipt result");
1805 assert_eq!(receipt_block_result.header.hash(), test_hash);
1806 assert_eq!(receipt_block_result.header.number, test_block_number);
1807 assert_eq!(receipt_block_result.receipts.len(), 1);
1808 assert_eq!(receipt_block_result.receipts[0].tx_type, mock_receipt.tx_type);
1809 assert_eq!(
1810 receipt_block_result.receipts[0].cumulative_gas_used,
1811 mock_receipt.cumulative_gas_used
1812 );
1813 assert_eq!(receipt_block_result.receipts[0].success, mock_receipt.success);
1814
1815 let result2 = cached_mode.next().await;
1817 assert!(result2.is_ok());
1818 assert!(result2.unwrap().is_none());
1819 }
1820
1821 #[tokio::test]
1822 async fn test_cached_mode_empty_headers() {
1823 let provider = MockEthProvider::default();
1824 let eth_api = build_test_eth_api(provider);
1825
1826 let eth_filter =
1827 super::EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1828 let filter_inner = eth_filter.inner;
1829
1830 let headers: Vec<SealedHeader<alloy_consensus::Header>> = vec![];
1831
1832 let mut cached_mode = CachedMode { filter_inner, headers_iter: headers.into_iter() };
1833
1834 let result = cached_mode.next().await.expect("next should succeed");
1836 assert!(result.is_none());
1837 }
1838
1839 #[tokio::test]
1840 async fn test_log_limit_retry_range_excludes_overflow_block() {
1841 let provider = MockEthProvider::default();
1842
1843 use alloy_consensus::TxLegacy;
1844 use reth_db_api::models::StoredBlockBodyIndices;
1845 use reth_ethereum_primitives::{TransactionSigned, TxType};
1846
1847 let tx_inner = TxLegacy {
1848 chain_id: Some(1),
1849 nonce: 0,
1850 gas_price: 21_000,
1851 gas_limit: 21_000,
1852 to: alloy_primitives::TxKind::Call(alloy_primitives::Address::ZERO),
1853 value: alloy_primitives::U256::ZERO,
1854 input: alloy_primitives::Bytes::new(),
1855 };
1856 let signature = alloy_primitives::Signature::test_signature();
1857 let tx = TransactionSigned::new_unhashed(tx_inner.into(), signature);
1858
1859 let mock_log = alloy_primitives::Log {
1860 address: alloy_primitives::Address::ZERO,
1861 data: alloy_primitives::LogData::new_unchecked(vec![], alloy_primitives::Bytes::new()),
1862 };
1863
1864 let receipt = reth_ethereum_primitives::Receipt {
1865 tx_type: TxType::Legacy,
1866 cumulative_gas_used: 21_000,
1867 logs: vec![mock_log],
1868 success: true,
1869 };
1870
1871 let mut prev_hash = alloy_primitives::B256::default();
1872 for (idx, block_number) in (100u64..=102).enumerate() {
1873 let header = alloy_consensus::Header {
1874 number: block_number,
1875 parent_hash: prev_hash,
1876 logs_bloom: alloy_primitives::Bloom::from([1u8; 256]),
1877 ..Default::default()
1878 };
1879 let hash = header.hash_slow();
1880 prev_hash = hash;
1881
1882 let block = reth_ethereum_primitives::Block {
1883 header,
1884 body: reth_ethereum_primitives::BlockBody {
1885 transactions: vec![tx.clone()],
1886 ..Default::default()
1887 },
1888 };
1889 provider.add_block(hash, block);
1890 provider.add_receipts(block_number, vec![receipt.clone()]);
1891 provider.add_block_body_indices(
1892 block_number,
1893 StoredBlockBodyIndices { first_tx_num: idx as u64, tx_count: 1 },
1894 );
1895 }
1896
1897 let eth_api = build_test_eth_api(provider);
1898 let eth_filter = EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
1899 let err = eth_filter
1900 .inner
1901 .clone()
1902 .get_logs_in_block_range(
1903 Filter::default(),
1904 100,
1905 102,
1906 QueryLimits { max_blocks_per_filter: None, max_logs_per_response: Some(2) },
1907 )
1908 .await
1909 .expect_err("range should exceed max logs");
1910
1911 let EthFilterError::QueryExceedsMaxResults { max_logs, from_block, to_block } = err else {
1912 panic!("unexpected error: {err:?}");
1913 };
1914
1915 assert_eq!(max_logs, 2);
1916 assert_eq!(from_block, 100);
1917 assert_eq!(to_block, 101);
1918 }
1919
1920 #[tokio::test]
1921 async fn test_non_consecutive_headers_after_bloom_filter() {
1922 let provider = MockEthProvider::default();
1923
1924 let mut expected_hashes = vec![];
1926 let mut prev_hash = alloy_primitives::B256::default();
1927
1928 use alloy_consensus::TxLegacy;
1930 use reth_ethereum_primitives::{TransactionSigned, TxType};
1931
1932 let tx_inner = TxLegacy {
1933 chain_id: Some(1),
1934 nonce: 0,
1935 gas_price: 21_000,
1936 gas_limit: 21_000,
1937 to: alloy_primitives::TxKind::Call(alloy_primitives::Address::ZERO),
1938 value: alloy_primitives::U256::ZERO,
1939 input: alloy_primitives::Bytes::new(),
1940 };
1941 let signature = alloy_primitives::Signature::test_signature();
1942 let tx = TransactionSigned::new_unhashed(tx_inner.into(), signature);
1943
1944 for i in 100u64..=103 {
1945 let header = alloy_consensus::Header {
1946 number: i,
1947 parent_hash: prev_hash,
1948 logs_bloom: if i == 100 || i == 102 {
1950 alloy_primitives::Bloom::from([1u8; 256])
1951 } else {
1952 alloy_primitives::Bloom::default()
1953 },
1954 ..Default::default()
1955 };
1956
1957 let hash = header.hash_slow();
1958 expected_hashes.push(hash);
1959 prev_hash = hash;
1960
1961 let transactions = if i == 100 || i == 102 { vec![tx.clone()] } else { vec![] };
1963
1964 let block = reth_ethereum_primitives::Block {
1965 header,
1966 body: reth_ethereum_primitives::BlockBody { transactions, ..Default::default() },
1967 };
1968 provider.add_block(hash, block);
1969 }
1970
1971 let mock_log = alloy_primitives::Log {
1973 address: alloy_primitives::Address::ZERO,
1974 data: alloy_primitives::LogData::new_unchecked(vec![], alloy_primitives::Bytes::new()),
1975 };
1976
1977 let receipt = reth_ethereum_primitives::Receipt {
1978 tx_type: TxType::Legacy,
1979 cumulative_gas_used: 21_000,
1980 logs: vec![mock_log],
1981 success: true,
1982 };
1983
1984 provider.add_receipts(100, vec![receipt.clone()]);
1985 provider.add_receipts(101, vec![]);
1986 provider.add_receipts(102, vec![receipt.clone()]);
1987 provider.add_receipts(103, vec![]);
1988
1989 use reth_db_api::models::StoredBlockBodyIndices;
1991 provider
1992 .add_block_body_indices(100, StoredBlockBodyIndices { first_tx_num: 0, tx_count: 1 });
1993 provider
1994 .add_block_body_indices(101, StoredBlockBodyIndices { first_tx_num: 1, tx_count: 0 });
1995 provider
1996 .add_block_body_indices(102, StoredBlockBodyIndices { first_tx_num: 1, tx_count: 1 });
1997 provider
1998 .add_block_body_indices(103, StoredBlockBodyIndices { first_tx_num: 2, tx_count: 0 });
1999
2000 let eth_api = build_test_eth_api(provider);
2001 let eth_filter = EthFilter::new(eth_api, EthFilterConfig::default(), Runtime::test());
2002
2003 let filter = Filter::default();
2005
2006 let logs = eth_filter
2008 .inner
2009 .clone()
2010 .get_logs_in_block_range(filter, 100, 103, QueryLimits::default())
2011 .await
2012 .expect("should succeed");
2013
2014 assert_eq!(logs.len(), 2);
2016
2017 assert_eq!(logs[0].block_number, Some(100));
2018 assert_eq!(logs[1].block_number, Some(102));
2019
2020 assert_eq!(logs[0].block_hash, Some(expected_hashes[0])); assert_eq!(logs[1].block_hash, Some(expected_hashes[2])); }
2024}