Skip to main content

reth_rpc_eth_types/cache/
mod.rs

1//! Async caching support for eth RPC
2
3use super::{EthStateCacheConfig, MultiConsumerLruCache};
4use crate::block::CachedTransaction;
5use alloy_consensus::transaction::TxHashRef;
6use alloy_eip7928::bal::DecodedBal;
7use alloy_eips::BlockHashOrNumber;
8use alloy_primitives::{Address, Bytes, TxHash, B256};
9use futures::{Stream, StreamExt};
10use reth_chain_state::CanonStateNotification;
11use reth_errors::{ProviderError, ProviderResult};
12use reth_execution_types::{Chain, DecodedRevmBal};
13use reth_primitives_traits::{Block, BlockBody, InMemorySize, NodePrimitives, RecoveredBlock};
14use reth_revm::{
15    bytecode::Bytecode,
16    primitives::{StorageKey, StorageValue},
17    state::bal::{
18        AccountBal as RevmAccountBal, AccountInfoBal as RevmAccountInfoBal, Bal as RevmBal,
19        BalWrites as RevmBalWrites, StorageBal as RevmStorageBal,
20    },
21};
22use reth_storage_api::{BalProvider, BlockReader, TransactionVariant};
23use reth_tasks::Runtime;
24use schnellru::{ByLength, LruMap};
25use std::{
26    cell::LazyCell,
27    collections::VecDeque,
28    future::Future,
29    pin::Pin,
30    sync::Arc,
31    task::{Context, Poll},
32    time::Duration,
33};
34use tokio::{
35    sync::{
36        mpsc::{unbounded_channel, UnboundedSender},
37        oneshot, OwnedSemaphorePermit, Semaphore,
38    },
39    time::{Instant, Interval, MissedTickBehavior},
40};
41use tokio_stream::wrappers::UnboundedReceiverStream;
42
43pub mod config;
44pub mod db;
45pub mod metrics;
46pub mod multi_consumer;
47
48/// The type that can send the response to a requested [`RecoveredBlock`]
49type BlockWithSendersResponseSender<B> =
50    oneshot::Sender<ProviderResult<Option<Arc<RecoveredBlock<B>>>>>;
51
52/// The type that can send the response to the requested receipts of a block.
53type ReceiptsResponseSender<R> = oneshot::Sender<ProviderResult<Option<Arc<Vec<R>>>>>;
54
55type CachedBlockResponseSender<B> = oneshot::Sender<Option<Arc<RecoveredBlock<B>>>>;
56
57type CachedBalResponseSender = oneshot::Sender<Option<CachedRevmBal>>;
58
59type CachedBlockAndReceiptsResponseSender<B, R> =
60    oneshot::Sender<(Option<Arc<RecoveredBlock<B>>>, Option<Arc<Vec<R>>>)>;
61
62/// The type that can send the response for a transaction hash lookup
63type TransactionHashResponseSender<B, R> = oneshot::Sender<Option<CachedTransaction<B, R>>>;
64
65/// The type that can send the response to a requested revm BAL.
66type BalResponseSender = oneshot::Sender<ProviderResult<Option<CachedRevmBal>>>;
67
68type BlockLruCache<B> =
69    MultiConsumerLruCache<B256, Arc<RecoveredBlock<B>>, BlockWithSendersResponseSender<B>>;
70
71type ReceiptsLruCache<R> = MultiConsumerLruCache<B256, Arc<Vec<R>>, ReceiptsResponseSender<R>>;
72
73type BalLruCache = MultiConsumerLruCache<B256, CachedRevmBal, BalResponseSender>;
74
75/// Provides async access to cached eth data
76///
77/// This is the frontend for the async caching service which manages cached data on a different
78/// task.
79#[derive(Debug)]
80pub struct EthStateCache<N: NodePrimitives> {
81    to_service: UnboundedSender<CacheAction<N::Block, N::Receipt>>,
82}
83
84impl<N: NodePrimitives> Clone for EthStateCache<N> {
85    fn clone(&self) -> Self {
86        Self { to_service: self.to_service.clone() }
87    }
88}
89
90impl<N: NodePrimitives> EthStateCache<N> {
91    /// Creates and returns both [`EthStateCache`] frontend and the memory bound service.
92    fn create<Provider>(
93        provider: Provider,
94        action_task_spawner: Runtime,
95        config: EthStateCacheConfig,
96    ) -> (Self, EthStateCacheService<Provider>)
97    where
98        Provider: BlockReader<Block = N::Block, Receipt = N::Receipt> + BalProvider,
99    {
100        let EthStateCacheConfig {
101            max_blocks,
102            max_receipts,
103            max_bals,
104            max_blocks_bytes,
105            max_receipts_bytes,
106            max_bals_bytes,
107            idle_timeout,
108            cache_computed_bals: _,
109            prewarm_bals: _,
110            max_concurrent_db_requests,
111            max_cached_tx_hashes,
112        } = config;
113        let eviction_interval = {
114            let _guard = action_task_spawner.handle().enter();
115            // Clamp short timeouts to avoid busy polling. Expiration is best effort between ticks.
116            let period = (idle_timeout / 4).clamp(Duration::from_secs(1), Duration::from_secs(60));
117            let mut interval = tokio::time::interval_at(Instant::now() + period, period);
118            interval.set_missed_tick_behavior(MissedTickBehavior::Skip);
119            interval
120        };
121        let (to_service, rx) = unbounded_channel();
122
123        let service = EthStateCacheService {
124            provider,
125            full_block_cache: BlockLruCache::new(max_blocks, max_blocks_bytes, "blocks"),
126            receipts_cache: ReceiptsLruCache::new(max_receipts, max_receipts_bytes, "receipts"),
127            bal_cache: BalLruCache::new(max_bals, max_bals_bytes, "bals"),
128            idle_timeout,
129            eviction_interval,
130            eviction_pending: false,
131            action_tx: to_service.clone(),
132            action_rx: UnboundedReceiverStream::new(rx),
133            action_task_spawner,
134            rate_limiter: Arc::new(Semaphore::new(max_concurrent_db_requests)),
135            pending_fetches: VecDeque::new(),
136            tx_hash_index: LruMap::new(ByLength::new(max_cached_tx_hashes)),
137        };
138        let cache = Self { to_service };
139        (cache, service)
140    }
141
142    /// Creates a new async LRU backed cache service task and spawns it to a new task via the given
143    /// spawner.
144    ///
145    /// Entry counts and estimated payload bytes bound the cache. Idle entries expire periodically.
146    pub fn spawn_with<Provider>(
147        provider: Provider,
148        config: EthStateCacheConfig,
149        executor: Runtime,
150    ) -> Self
151    where
152        Provider: BlockReader<Block = N::Block, Receipt = N::Receipt>
153            + BalProvider
154            + Clone
155            + Unpin
156            + 'static,
157    {
158        let (this, service) = Self::create(provider, executor.clone(), config);
159        executor.spawn_critical_task("eth state cache", service);
160        this
161    }
162
163    /// Requests the  [`RecoveredBlock`] for the block hash
164    ///
165    /// Returns `None` if the block does not exist.
166    pub async fn get_recovered_block(
167        &self,
168        block_hash: B256,
169    ) -> ProviderResult<Option<Arc<RecoveredBlock<N::Block>>>> {
170        let (response_tx, rx) = oneshot::channel();
171        let _ = self.to_service.send(CacheAction::GetBlockWithSenders { block_hash, response_tx });
172        rx.await.map_err(|_| CacheServiceUnavailable)?
173    }
174
175    /// Requests the block for the given block hash if it is cached.
176    pub async fn get_maybe_block(
177        &self,
178        block_hash: B256,
179    ) -> ProviderResult<Option<Arc<RecoveredBlock<N::Block>>>> {
180        let (response_tx, rx) = oneshot::channel();
181        let _ = self.to_service.send(CacheAction::GetCachedBlock { block_hash, response_tx });
182        rx.await.map_err(|_| CacheServiceUnavailable.into())
183    }
184
185    /// Requests the receipts for the block hash
186    ///
187    /// Returns `None` if the block was not found.
188    pub async fn get_receipts(
189        &self,
190        block_hash: B256,
191    ) -> ProviderResult<Option<Arc<Vec<N::Receipt>>>> {
192        let (response_tx, rx) = oneshot::channel();
193        let _ = self.to_service.send(CacheAction::GetReceipts { block_hash, response_tx });
194        rx.await.map_err(|_| CacheServiceUnavailable)?
195    }
196
197    /// Fetches both receipts and block for the given block hash.
198    pub async fn get_block_and_receipts(
199        &self,
200        block_hash: B256,
201    ) -> ProviderResult<Option<(Arc<RecoveredBlock<N::Block>>, Arc<Vec<N::Receipt>>)>> {
202        let block = self.get_recovered_block(block_hash);
203        let receipts = self.get_receipts(block_hash);
204
205        let (block, receipts) = futures::try_join!(block, receipts)?;
206
207        Ok(block.zip(receipts))
208    }
209
210    /// Retrieves the block and its BAL if the BAL is already cached.
211    ///
212    /// The block is fetched if it is not cached. Returns `None` if the block does not exist.
213    pub async fn get_recovered_block_and_maybe_bal(
214        &self,
215        block_hash: B256,
216    ) -> ProviderResult<
217        Option<(Arc<RecoveredBlock<N::Block>>, Option<Arc<DecodedBal<Arc<RevmBal>>>>)>,
218    > {
219        let (response_tx, rx) = oneshot::channel();
220        let _ = self.to_service.send(CacheAction::GetCachedBal { block_hash, response_tx });
221
222        let block = self.get_recovered_block(block_hash);
223        let (block, bal) = futures::join!(block, rx);
224
225        let bal = bal.map_err(|_| CacheServiceUnavailable)?.map(|cached| cached.0);
226        Ok(block?.map(|block| (block, bal)))
227    }
228
229    /// Retrieves receipts and blocks from cache if block is in the cache, otherwise only receipts.
230    pub async fn get_receipts_and_maybe_block(
231        &self,
232        block_hash: B256,
233    ) -> ProviderResult<Option<(Arc<Vec<N::Receipt>>, Option<Arc<RecoveredBlock<N::Block>>>)>> {
234        let (response_tx, rx) = oneshot::channel();
235        let _ = self.to_service.send(CacheAction::GetCachedBlock { block_hash, response_tx });
236
237        let receipts = self.get_receipts(block_hash);
238
239        let (receipts, block) = futures::join!(receipts, rx);
240
241        let block = block.map_err(|_| CacheServiceUnavailable)?;
242        Ok(receipts?.map(|r| (r, block)))
243    }
244
245    /// Retrieves both block and receipts from cache if available.
246    pub async fn maybe_cached_block_and_receipts(
247        &self,
248        block_hash: B256,
249    ) -> ProviderResult<(Option<Arc<RecoveredBlock<N::Block>>>, Option<Arc<Vec<N::Receipt>>>)> {
250        let (response_tx, rx) = oneshot::channel();
251        let _ = self
252            .to_service
253            .send(CacheAction::GetCachedBlockAndReceipts { block_hash, response_tx });
254        rx.await.map_err(|_| CacheServiceUnavailable.into())
255    }
256
257    /// Looks up a transaction by its hash in the cache index.
258    ///
259    /// Returns the cached block, transaction index, and optionally receipts if the transaction
260    /// is in a cached block.
261    pub async fn get_transaction_by_hash(
262        &self,
263        tx_hash: TxHash,
264    ) -> Option<CachedTransaction<N::Block, N::Receipt>> {
265        let (response_tx, rx) = oneshot::channel();
266        let _ = self.to_service.send(CacheAction::GetTransactionByHash { tx_hash, response_tx });
267        rx.await.ok()?
268    }
269
270    /// Requests the revm BAL for the block hash.
271    ///
272    /// Returns `None` if the BAL does not exist.
273    pub async fn get_bal(
274        &self,
275        block_hash: B256,
276    ) -> ProviderResult<Option<Arc<DecodedBal<Arc<RevmBal>>>>> {
277        let (response_tx, rx) = oneshot::channel();
278        let _ = self.to_service.send(CacheAction::GetBal { block_hash, response_tx });
279        rx.await
280            .map_err(|_| CacheServiceUnavailable)?
281            .map(|maybe_bal| maybe_bal.map(|cached| cached.0))
282    }
283
284    /// Inserts a decoded revm BAL into the cache.
285    pub fn insert_bal(&self, block_hash: B256, bal: DecodedBal<Arc<RevmBal>>) {
286        let _ = self
287            .to_service
288            .send(CacheAction::InsertBal { block_hash, bal: CachedRevmBal::new(bal) });
289    }
290}
291/// Thrown when the cache service task dropped.
292#[derive(Debug, thiserror::Error)]
293#[error("cache service task stopped")]
294pub struct CacheServiceUnavailable;
295
296impl From<CacheServiceUnavailable> for ProviderError {
297    fn from(err: CacheServiceUnavailable) -> Self {
298        Self::other(err)
299    }
300}
301
302/// A task that manages caches for data required by the `eth` rpc implementation.
303///
304/// It provides a caching layer on top of the given
305/// [`StateProvider`](reth_storage_api::StateProvider) and keeps data fetched via the provider in
306/// memory in an LRU cache. If the requested data is missing in the cache it is fetched and inserted
307/// into the cache afterwards. While fetching data from disk is sync, this service is async since
308/// requests and data is shared via channels.
309///
310/// This type is an endless future that listens for incoming messages from the user facing
311/// [`EthStateCache`] via a channel. If the requested data is not cached then it spawns a new task
312/// that does the IO and sends the result back to it. This way the caching service only
313/// handles messages and does LRU lookups and never blocking IO.
314///
315/// Caution: The channel for the data is _unbounded_ it is assumed that this is mainly used by the
316/// `reth_rpc::EthApi` which is typically invoked by the RPC server, which already uses
317/// permits to limit concurrent requests.
318#[must_use = "Type does nothing unless spawned"]
319pub(crate) struct EthStateCacheService<Provider>
320where
321    Provider: BlockReader + BalProvider,
322{
323    /// The type used to lookup data from disk
324    provider: Provider,
325    /// The LRU cache for full blocks grouped by their block hash.
326    full_block_cache: BlockLruCache<Provider::Block>,
327    /// The LRU cache for block receipts grouped by the block hash.
328    receipts_cache: ReceiptsLruCache<Provider::Receipt>,
329    /// The LRU cache for revm BALs grouped by the block hash.
330    bal_cache: BalLruCache,
331    /// Maximum time without a cache hit. Zero disables idle eviction.
332    idle_timeout: Duration,
333    /// Reclaims idle entries even when there are no requests or new blocks.
334    eviction_interval: Interval,
335    /// Whether an idle sweep has more entries to remove in a subsequent poll.
336    eviction_pending: bool,
337    /// Sender half of the action channel.
338    action_tx: UnboundedSender<CacheAction<Provider::Block, Provider::Receipt>>,
339    /// Receiver half of the action channel.
340    action_rx: UnboundedReceiverStream<CacheAction<Provider::Block, Provider::Receipt>>,
341    /// The type that's used to spawn tasks that do the actual work
342    action_task_spawner: Runtime,
343    /// Rate limiter for spawned fetch tasks.
344    ///
345    /// This restricts the max concurrent fetch tasks at the same time.
346    rate_limiter: Arc<Semaphore>,
347    /// Cache misses waiting for a database request slot.
348    ///
349    /// A request is only moved to the blocking pool after it acquires a permit. This prevents
350    /// cache misses that are waiting for a slot from occupying Tokio blocking threads.
351    pending_fetches: VecDeque<CacheFetch>,
352    /// LRU index mapping transaction hashes to their block hash and index within the block.
353    tx_hash_index: LruMap<TxHash, (B256, usize), ByLength>,
354}
355
356impl<Provider> EthStateCacheService<Provider>
357where
358    Provider: BlockReader + BalProvider + Clone + Unpin + 'static,
359{
360    fn queue_fetch(&mut self, fetch: CacheFetch) {
361        if self.pending_fetches.is_empty() &&
362            let Ok(permit) = self.rate_limiter.clone().try_acquire_owned()
363        {
364            self.spawn_fetch(fetch, permit);
365            return
366        }
367
368        self.pending_fetches.push_back(fetch);
369        self.spawn_pending_fetches();
370    }
371
372    fn spawn_pending_fetches(&mut self) {
373        while let Some(fetch) = self.pending_fetches.pop_front() {
374            let Ok(permit) = self.rate_limiter.clone().try_acquire_owned() else {
375                self.pending_fetches.push_front(fetch);
376                return
377            };
378
379            self.spawn_fetch(fetch, permit);
380        }
381    }
382
383    fn spawn_fetch(&self, fetch: CacheFetch, permit: OwnedSemaphorePermit) {
384        let provider = self.provider.clone();
385        let action_tx = self.action_tx.clone();
386
387        match fetch {
388            CacheFetch::Block(block_hash) => {
389                let mut action_sender = ActionSender::new(CacheKind::Block, block_hash, action_tx);
390                self.action_task_spawner.spawn_blocking_task(async move {
391                    let block_sender = provider
392                        .sealed_block_with_senders(
393                            BlockHashOrNumber::Hash(block_hash),
394                            TransactionVariant::WithHash,
395                        )
396                        .map(|maybe_block| maybe_block.map(Arc::new));
397                    drop(permit);
398                    action_sender.send_block(block_sender);
399                });
400            }
401            CacheFetch::Receipts(block_hash) => {
402                let mut action_sender =
403                    ActionSender::new(CacheKind::Receipt, block_hash, action_tx);
404                self.action_task_spawner.spawn_blocking_task(async move {
405                    let res = provider
406                        .receipts_by_block(block_hash.into())
407                        .map(|maybe_receipts| maybe_receipts.map(Arc::new));
408                    drop(permit);
409                    action_sender.send_receipts(res);
410                });
411            }
412            CacheFetch::Bal(block_hash) => {
413                let mut action_sender = ActionSender::new(CacheKind::Bal, block_hash, action_tx);
414                self.action_task_spawner.spawn_blocking_task(async move {
415                    let res = provider.get_bal_by_hash(block_hash).and_then(|maybe_bal| {
416                        maybe_bal.map(CachedRevmBal::try_from_raw).transpose()
417                    });
418                    drop(permit);
419                    action_sender.send_bal(res);
420                });
421            }
422        }
423    }
424
425    /// Indexes all transactions in a block by transaction hash.
426    fn index_block_transactions(&mut self, block: &RecoveredBlock<Provider::Block>) {
427        let block_hash = block.hash();
428        for (tx_idx, tx) in block.body().transactions().iter().enumerate() {
429            self.tx_hash_index.insert(*tx.tx_hash(), (block_hash, tx_idx));
430        }
431    }
432
433    /// Removes transaction index entries for a reorged block.
434    fn remove_block_transactions(&mut self, block: &RecoveredBlock<Provider::Block>) {
435        for tx in block.body().transactions() {
436            self.tx_hash_index.remove(tx.tx_hash());
437        }
438    }
439
440    fn on_new_block(
441        &mut self,
442        block_hash: B256,
443        res: ProviderResult<Option<Arc<RecoveredBlock<Provider::Block>>>>,
444        now: Instant,
445    ) {
446        if let Some(queued) = self.full_block_cache.remove(&block_hash) {
447            // send the response to queued senders
448            for tx in queued {
449                let _ = tx.send(res.clone());
450            }
451        }
452
453        // cache good block
454        if let Ok(Some(block)) = res {
455            self.full_block_cache.insert_at(block_hash, block, now);
456        }
457    }
458
459    fn on_new_receipts(
460        &mut self,
461        block_hash: B256,
462        res: ProviderResult<Option<Arc<Vec<Provider::Receipt>>>>,
463        now: Instant,
464    ) {
465        if let Some(queued) = self.receipts_cache.remove(&block_hash) {
466            // send the response to queued senders
467            for tx in queued {
468                let _ = tx.send(res.clone());
469            }
470        }
471
472        // cache good receipts
473        if let Ok(Some(receipts)) = res {
474            self.receipts_cache.insert_at(block_hash, receipts, now);
475        }
476    }
477
478    fn on_new_bal(
479        &mut self,
480        block_hash: B256,
481        res: ProviderResult<Option<CachedRevmBal>>,
482        now: Instant,
483    ) {
484        // A local replay may finish before an older provider lookup returns.
485        if self.bal_cache.contains_key(&block_hash) {
486            return
487        }
488
489        if let Some(queued) = self.bal_cache.remove(&block_hash) {
490            for tx in queued {
491                let _ = tx.send(res.clone());
492            }
493        }
494
495        if let Ok(Some(bal)) = res {
496            self.bal_cache.insert_at(block_hash, bal, now);
497        }
498    }
499
500    fn on_reorg_block(
501        &mut self,
502        block_hash: B256,
503        res: ProviderResult<Option<Arc<RecoveredBlock<Provider::Block>>>>,
504    ) {
505        if let Some(queued) = self.full_block_cache.remove(&block_hash) {
506            // send the response to queued senders
507            for tx in queued {
508                let _ = tx.send(res.clone());
509            }
510        }
511    }
512
513    fn on_reorg_receipts(
514        &mut self,
515        block_hash: B256,
516        res: ProviderResult<Option<Arc<Vec<Provider::Receipt>>>>,
517    ) {
518        if let Some(queued) = self.receipts_cache.remove(&block_hash) {
519            // send the response to queued senders
520            for tx in queued {
521                let _ = tx.send(res.clone());
522            }
523        }
524    }
525
526    fn on_reorg_bal(&mut self, block_hash: B256, res: ProviderResult<Option<CachedRevmBal>>) {
527        if let Some(queued) = self.bal_cache.remove(&block_hash) {
528            for tx in queued {
529                let _ = tx.send(res.clone());
530            }
531        }
532    }
533
534    /// Shrinks the queues but leaves some space for the next requests
535    fn shrink_queues(&mut self) {
536        let min_capacity = 2;
537        self.full_block_cache.shrink_to(min_capacity);
538        self.receipts_cache.shrink_to(min_capacity);
539        self.bal_cache.shrink_to(min_capacity);
540    }
541
542    fn update_cached_metrics(&mut self) {
543        self.full_block_cache.update_cached_metrics();
544        self.receipts_cache.update_cached_metrics();
545        self.bal_cache.update_cached_metrics();
546    }
547}
548
549impl<Provider> Future for EthStateCacheService<Provider>
550where
551    Provider: BlockReader + BalProvider + Clone + Unpin + 'static,
552{
553    type Output = ();
554
555    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
556        let this = self.get_mut();
557        // All inserts and hits in this poll share a clock read, preserving timestamp/LRU order.
558        let now = LazyCell::new(Instant::now);
559
560        // Poll until pending to register the next timer wakeup before draining requests.
561        while !this.idle_timeout.is_zero() && this.eviction_interval.poll_tick(cx).is_ready() {
562            this.eviction_pending = true;
563        }
564        if this.eviction_pending {
565            let now = *now;
566            this.eviction_pending = this.full_block_cache.evict_expired(now, this.idle_timeout);
567            this.eviction_pending |= this.receipts_cache.evict_expired(now, this.idle_timeout);
568            this.eviction_pending |= this.bal_cache.evict_expired(now, this.idle_timeout);
569            if this.eviction_pending {
570                // Process requests before continuing the sweep in another poll.
571                cx.waker().wake_by_ref();
572            }
573        }
574
575        loop {
576            let Poll::Ready(action) = this.action_rx.poll_next_unpin(cx) else {
577                // shrink queues if we don't have any work to do
578                this.shrink_queues();
579                // gauges only need to be accurate once the batch of messages is drained, and only
580                // caches that changed during the batch republish theirs
581                this.update_cached_metrics();
582                return Poll::Pending;
583            };
584
585            match action {
586                None => {
587                    unreachable!("can't close")
588                }
589                Some(action) => {
590                    let now = *now;
591                    match action {
592                        CacheAction::GetCachedBlock { block_hash, response_tx } => {
593                            let _ = response_tx
594                                .send(this.full_block_cache.get_at(&block_hash, now).cloned());
595                        }
596                        CacheAction::GetCachedBal { block_hash, response_tx } => {
597                            let _ =
598                                response_tx.send(this.bal_cache.get_at(&block_hash, now).cloned());
599                        }
600                        CacheAction::GetCachedBlockAndReceipts { block_hash, response_tx } => {
601                            let block = this.full_block_cache.get_at(&block_hash, now).cloned();
602                            let receipts = this.receipts_cache.get_at(&block_hash, now).cloned();
603                            let _ = response_tx.send((block, receipts));
604                        }
605                        CacheAction::GetBlockWithSenders { block_hash, response_tx } => {
606                            if let Some(block) =
607                                this.full_block_cache.get_at(&block_hash, now).cloned()
608                            {
609                                let _ = response_tx.send(Ok(Some(block)));
610                                continue
611                            }
612
613                            // block is not in the cache, request it if this is the first consumer
614                            if this.full_block_cache.queue(block_hash, response_tx) {
615                                this.queue_fetch(CacheFetch::Block(block_hash));
616                            }
617                        }
618                        CacheAction::GetReceipts { block_hash, response_tx } => {
619                            // check if block is cached
620                            if let Some(receipts) =
621                                this.receipts_cache.get_at(&block_hash, now).cloned()
622                            {
623                                let _ = response_tx.send(Ok(Some(receipts)));
624                                continue
625                            }
626
627                            // block is not in the cache, request it if this is the first consumer
628                            if this.receipts_cache.queue(block_hash, response_tx) {
629                                this.queue_fetch(CacheFetch::Receipts(block_hash));
630                            }
631                        }
632                        CacheAction::GetBal { block_hash, response_tx } => {
633                            if let Some(bal) = this.bal_cache.get_at(&block_hash, now).cloned() {
634                                let _ = response_tx.send(Ok(Some(bal)));
635                                continue
636                            }
637
638                            if this.bal_cache.queue(block_hash, response_tx) {
639                                this.queue_fetch(CacheFetch::Bal(block_hash));
640                            }
641                        }
642                        CacheAction::ReceiptsResult { block_hash, res } => {
643                            this.on_new_receipts(block_hash, res, now);
644                        }
645                        CacheAction::BalResult { block_hash, res } => {
646                            this.on_new_bal(block_hash, res, now);
647                        }
648                        CacheAction::InsertBal { block_hash, bal } => {
649                            this.on_new_bal(block_hash, Ok(Some(bal)), now);
650                        }
651                        CacheAction::BlockWithSendersResult { block_hash, res } => {
652                            this.on_new_block(block_hash, res, now);
653                        }
654                        CacheAction::CacheNewCanonicalChain { chain_change } => {
655                            for block in chain_change.blocks {
656                                // Index transactions before caching the block
657                                this.index_block_transactions(&block);
658                                this.on_new_block(block.hash(), Ok(Some(block)), now);
659                            }
660
661                            for block_receipts in chain_change.receipts {
662                                this.on_new_receipts(
663                                    block_receipts.block_hash,
664                                    Ok(Some(block_receipts.receipts)),
665                                    now,
666                                );
667                            }
668
669                            for (block_hash, bal) in chain_change.bals {
670                                this.on_new_bal(
671                                    block_hash,
672                                    Ok(Some(CachedRevmBal::from_shared(bal))),
673                                    now,
674                                );
675                            }
676                        }
677                        CacheAction::RemoveReorgedChain { chain_change } => {
678                            for block in chain_change.blocks {
679                                let block_hash = block.hash();
680                                // Remove transaction index entries for reorged blocks
681                                this.remove_block_transactions(&block);
682                                this.on_reorg_block(block_hash, Ok(Some(block)));
683                                this.on_reorg_bal(block_hash, Ok(None));
684                            }
685
686                            for block_receipts in chain_change.receipts {
687                                this.on_reorg_receipts(
688                                    block_receipts.block_hash,
689                                    Ok(Some(block_receipts.receipts)),
690                                );
691                            }
692                        }
693                        CacheAction::GetTransactionByHash { tx_hash, response_tx } => {
694                            let result =
695                                this.tx_hash_index.get(&tx_hash).and_then(|(block_hash, idx)| {
696                                    let block =
697                                        this.full_block_cache.get_at(block_hash, now).cloned()?;
698                                    let receipts =
699                                        this.receipts_cache.get_at(block_hash, now).cloned();
700                                    Some(CachedTransaction::new(block, *idx, receipts))
701                                });
702                            let _ = response_tx.send(result);
703                        }
704                    };
705                    this.spawn_pending_fetches();
706                }
707            }
708        }
709    }
710}
711
712/// All message variants sent through the channel
713enum CacheAction<B: Block, R> {
714    GetBlockWithSenders {
715        block_hash: B256,
716        response_tx: BlockWithSendersResponseSender<B>,
717    },
718    GetReceipts {
719        block_hash: B256,
720        response_tx: ReceiptsResponseSender<R>,
721    },
722    GetBal {
723        block_hash: B256,
724        response_tx: BalResponseSender,
725    },
726    GetCachedBlock {
727        block_hash: B256,
728        response_tx: CachedBlockResponseSender<B>,
729    },
730    GetCachedBal {
731        block_hash: B256,
732        response_tx: CachedBalResponseSender,
733    },
734    GetCachedBlockAndReceipts {
735        block_hash: B256,
736        response_tx: CachedBlockAndReceiptsResponseSender<B, R>,
737    },
738    BlockWithSendersResult {
739        block_hash: B256,
740        res: ProviderResult<Option<Arc<RecoveredBlock<B>>>>,
741    },
742    ReceiptsResult {
743        block_hash: B256,
744        res: ProviderResult<Option<Arc<Vec<R>>>>,
745    },
746    BalResult {
747        block_hash: B256,
748        res: ProviderResult<Option<CachedRevmBal>>,
749    },
750    InsertBal {
751        block_hash: B256,
752        bal: CachedRevmBal,
753    },
754    CacheNewCanonicalChain {
755        chain_change: ChainChange<B, R>,
756    },
757    RemoveReorgedChain {
758        chain_change: ChainChange<B, R>,
759    },
760    /// Look up a transaction's cached data by its hash
761    GetTransactionByHash {
762        tx_hash: TxHash,
763        response_tx: TransactionHashResponseSender<B, R>,
764    },
765}
766
767struct BlockReceipts<R> {
768    block_hash: B256,
769    receipts: Arc<Vec<R>>,
770}
771
772/// A change of the canonical chain
773struct ChainChange<B: Block, R> {
774    blocks: Vec<Arc<RecoveredBlock<B>>>,
775    receipts: Vec<BlockReceipts<R>>,
776    bals: Vec<(B256, Arc<DecodedRevmBal>)>,
777}
778
779impl<B: Block, R: Clone> ChainChange<B, R> {
780    fn new<N>(chain: Arc<Chain<N>>) -> Self
781    where
782        N: NodePrimitives<Block = B, Receipt = R>,
783    {
784        let (blocks, receipts): (Vec<_>, Vec<_>) = chain
785            .blocks_and_receipts()
786            .map(|(block, receipts)| {
787                let block_receipts = BlockReceipts {
788                    block_hash: block.hash(),
789                    receipts: Arc::new(receipts.clone()),
790                };
791                (Arc::clone(block), block_receipts)
792            })
793            .unzip();
794        // Blocks without an attached BAL are absent here; their BAL is fetched from the store on
795        // demand.
796        let bals =
797            chain.blocks_and_bals().map(|(block, bal)| (block.hash(), Arc::clone(bal))).collect();
798        Self { blocks, receipts, bals }
799    }
800}
801
802/// Identifier for the caches.
803#[derive(Copy, Clone, Debug)]
804enum CacheKind {
805    Block,
806    Receipt,
807    Bal,
808}
809
810#[derive(Copy, Clone, Debug)]
811enum CacheFetch {
812    Block(B256),
813    Receipts(B256),
814    Bal(B256),
815}
816
817/// Drop aware sender struct that ensures a response is always emitted even if the db task panics
818/// before a result could be sent.
819///
820/// This type wraps a sender and in case the sender is still present on drop emit an error response.
821#[derive(Debug)]
822struct ActionSender<B: Block, R: Send + Sync> {
823    kind: CacheKind,
824    blockhash: B256,
825    tx: Option<UnboundedSender<CacheAction<B, R>>>,
826}
827
828impl<R: Send + Sync, B: Block> ActionSender<B, R> {
829    const fn new(kind: CacheKind, blockhash: B256, tx: UnboundedSender<CacheAction<B, R>>) -> Self {
830        Self { kind, blockhash, tx: Some(tx) }
831    }
832
833    fn send_block(&mut self, block_sender: Result<Option<Arc<RecoveredBlock<B>>>, ProviderError>) {
834        if let Some(tx) = self.tx.take() {
835            let _ = tx.send(CacheAction::BlockWithSendersResult {
836                block_hash: self.blockhash,
837                res: block_sender,
838            });
839        }
840    }
841
842    fn send_receipts(&mut self, receipts: Result<Option<Arc<Vec<R>>>, ProviderError>) {
843        if let Some(tx) = self.tx.take() {
844            let _ =
845                tx.send(CacheAction::ReceiptsResult { block_hash: self.blockhash, res: receipts });
846        }
847    }
848
849    fn send_bal(&mut self, bal: Result<Option<CachedRevmBal>, ProviderError>) {
850        if let Some(tx) = self.tx.take() {
851            let _ = tx.send(CacheAction::BalResult { block_hash: self.blockhash, res: bal });
852        }
853    }
854}
855impl<R: Send + Sync, B: Block> Drop for ActionSender<B, R> {
856    fn drop(&mut self) {
857        if let Some(tx) = self.tx.take() {
858            let msg = match self.kind {
859                CacheKind::Block => CacheAction::BlockWithSendersResult {
860                    block_hash: self.blockhash,
861                    res: Err(CacheServiceUnavailable.into()),
862                },
863                CacheKind::Receipt => CacheAction::ReceiptsResult {
864                    block_hash: self.blockhash,
865                    res: Err(CacheServiceUnavailable.into()),
866                },
867                CacheKind::Bal => CacheAction::BalResult {
868                    block_hash: self.blockhash,
869                    res: Err(CacheServiceUnavailable.into()),
870                },
871            };
872            let _ = tx.send(msg);
873        }
874    }
875}
876
877/// Awaits for new chain events and directly inserts them into the cache so they're available
878/// immediately before they need to be fetched from disk.
879///
880/// Reorged blocks are removed from the cache.
881pub async fn cache_new_blocks_task<St, N: NodePrimitives>(
882    eth_state_cache: EthStateCache<N>,
883    mut events: St,
884) where
885    St: Stream<Item = CanonStateNotification<N>> + Unpin + 'static,
886{
887    while let Some(event) = events.next().await {
888        if let Some(reverted) = event.reverted() {
889            let chain_change = ChainChange::new(reverted);
890
891            let _ =
892                eth_state_cache.to_service.send(CacheAction::RemoveReorgedChain { chain_change });
893        }
894
895        let chain_change = ChainChange::new(event.committed());
896
897        let _ =
898            eth_state_cache.to_service.send(CacheAction::CacheNewCanonicalChain { chain_change });
899    }
900}
901
902/// Cached decoded revm BAL.
903#[derive(Clone, Debug)]
904pub(crate) struct CachedRevmBal(Arc<DecodedRevmBal>);
905
906impl CachedRevmBal {
907    /// Creates a cached revm BAL from an owned decoded BAL.
908    #[inline]
909    fn new(bal: DecodedRevmBal) -> Self {
910        Self(Arc::new(bal))
911    }
912
913    /// Creates a cached revm BAL from a BAL that is already shared, for example one prepared
914    /// during block validation.
915    #[inline]
916    const fn from_shared(bal: Arc<DecodedRevmBal>) -> Self {
917        Self(bal)
918    }
919
920    /// Decodes raw BAL bytes into the representation used by revm.
921    fn try_from_raw(raw: Bytes) -> ProviderResult<Self> {
922        DecodedBal::from_rlp_bytes(raw)
923            .map_err(Into::into)
924            .and_then(|decoded| {
925                decoded.try_map(|bal| {
926                    RevmBal::try_from(Vec::from(bal)).map(Arc::new).map_err(ProviderError::other)
927                })
928            })
929            .map(Self::new)
930    }
931}
932
933impl InMemorySize for CachedRevmBal {
934    fn size(&self) -> usize {
935        core::mem::size_of::<Self>() + decoded_revm_bal_size(&self.0)
936    }
937}
938
939fn decoded_revm_bal_size(bal: &DecodedRevmBal) -> usize {
940    core::mem::size_of::<DecodedRevmBal>() + bal.as_raw().len() + revm_bal_size(bal.as_bal())
941}
942
943fn revm_bal_size(bal: &Arc<RevmBal>) -> usize {
944    core::mem::size_of::<RevmBal>() +
945        bal.accounts.capacity() * core::mem::size_of::<(Address, RevmAccountBal)>() +
946        bal.accounts.values().map(revm_account_bal_heap_size).sum::<usize>()
947}
948
949fn revm_account_bal_heap_size(account: &RevmAccountBal) -> usize {
950    revm_account_info_bal_heap_size(&account.account_info) +
951        revm_storage_bal_heap_size(&account.storage)
952}
953
954fn revm_account_info_bal_heap_size(account_info: &RevmAccountInfoBal) -> usize {
955    revm_bal_writes_heap_size(&account_info.nonce, |_| 0) +
956        revm_bal_writes_heap_size(&account_info.balance, |_| 0) +
957        revm_bal_writes_heap_size(&account_info.code, revm_code_write_heap_size)
958}
959
960fn revm_storage_bal_heap_size(storage: &RevmStorageBal) -> usize {
961    storage.storage.len() * core::mem::size_of::<(StorageKey, RevmBalWrites<StorageValue>)>() +
962        storage
963            .storage
964            .values()
965            .map(|writes| revm_bal_writes_heap_size(writes, |_| 0))
966            .sum::<usize>()
967}
968
969fn revm_bal_writes_heap_size<T, F>(writes: &RevmBalWrites<T>, mut item_heap_size: F) -> usize
970where
971    T: PartialEq + Clone,
972    F: FnMut(&T) -> usize,
973{
974    writes.writes.capacity() * core::mem::size_of::<(u64, T)>() +
975        writes.writes.iter().map(|(_, item)| item_heap_size(item)).sum::<usize>()
976}
977
978fn revm_code_write_heap_size((_, bytecode): &(B256, Bytecode)) -> usize {
979    bytecode.bytes_ref().len()
980}
981
982#[cfg(test)]
983mod tests {
984    use super::*;
985    use alloy_consensus::{transaction::TransactionMeta, Header};
986    use alloy_eip7928::BlockAccessIndex;
987    use alloy_eips::{BlockHashOrNumber, NumHash};
988    use alloy_primitives::{Address, BlockHash, BlockNumber, Bytes, Signature, TxHash, TxNumber};
989    use core::ops::{RangeBounds, RangeInclusive};
990    use reth_db_models::StoredBlockBodyIndices;
991    use reth_ethereum_primitives::{
992        Block, BlockBody, EthPrimitives, Receipt, Transaction, TransactionSigned,
993    };
994    use reth_execution_types::{ExecutionOutcome, RecoveredBlockAndExecutionOutput};
995    use reth_primitives_traits::{RecoveredBlock, SealedHeader};
996    use reth_storage_api::{
997        noop::NoopProvider, BalProvider, BalStore, BalStoreHandle, BlockBodyIndicesProvider,
998        BlockHashReader, BlockNumReader, BlockReader, BlockSource, HeaderProvider, ReceiptProvider,
999        TransactionVariant, TransactionsProvider,
1000    };
1001    use std::{
1002        sync::{
1003            atomic::{AtomicUsize, Ordering},
1004            mpsc::{self, Receiver, SyncSender},
1005            Mutex,
1006        },
1007        task::{Wake, Waker},
1008        thread,
1009        time::Duration,
1010    };
1011
1012    #[derive(Default)]
1013    struct WakeCounter(AtomicUsize);
1014
1015    impl Wake for WakeCounter {
1016        fn wake(self: Arc<Self>) {
1017            self.0.fetch_add(1, Ordering::Relaxed);
1018        }
1019    }
1020
1021    fn test_service() -> EthStateCacheService<NoopProvider> {
1022        let (_cache, service) = EthStateCache::<EthPrimitives>::create(
1023            NoopProvider::default(),
1024            Runtime::test(),
1025            EthStateCacheConfig {
1026                max_blocks: 4,
1027                max_receipts: 4,
1028                max_bals: 4,
1029                cache_computed_bals: false,
1030                prewarm_bals: None,
1031                max_concurrent_db_requests: 1,
1032                max_cached_tx_hashes: 16,
1033                ..Default::default()
1034            },
1035        );
1036        service
1037    }
1038
1039    fn test_decoded_revm_bal() -> DecodedBal<Arc<RevmBal>> {
1040        DecodedBal::new(Arc::new(RevmBal::default()), Bytes::from_static(&[0xc0]))
1041    }
1042
1043    fn test_block() -> RecoveredBlock<Block> {
1044        RecoveredBlock::new_unhashed(
1045            Block {
1046                header: Header { number: 1, ..Default::default() },
1047                body: BlockBody {
1048                    transactions: vec![TransactionSigned::new_unhashed(
1049                        Transaction::Legacy(Default::default()),
1050                        Signature::test_signature(),
1051                    )],
1052                    ..Default::default()
1053                },
1054            },
1055            vec![Address::ZERO],
1056        )
1057    }
1058
1059    #[tokio::test(start_paused = true)]
1060    async fn idle_timer_continues_waking_after_its_first_tick() {
1061        let (cache, mut service) = EthStateCache::<EthPrimitives>::create(
1062            NoopProvider::default(),
1063            Runtime::test(),
1064            EthStateCacheConfig { idle_timeout: Duration::from_secs(10), ..Default::default() },
1065        );
1066        tokio::time::advance(Duration::from_secs(3)).await;
1067        let block = Arc::new(test_block());
1068        let hash = block.hash();
1069        let receipts = Arc::new(vec![]);
1070        let bal = CachedRevmBal::new(test_decoded_revm_bal());
1071        let retained_block = Arc::downgrade(&block);
1072        let retained_receipts = Arc::downgrade(&receipts);
1073        let retained_bal = Arc::downgrade(&bal.0);
1074        service.on_new_block(hash, Ok(Some(block)), Instant::now());
1075        service.on_new_receipts(hash, Ok(Some(receipts)), Instant::now());
1076        service.on_new_bal(hash, Ok(Some(bal)), Instant::now());
1077        let (response_tx, response_rx) = oneshot::channel();
1078        service.full_block_cache.queue(hash, response_tx);
1079        let service_task = tokio::spawn(service);
1080        assert!(cache.get_maybe_block(hash).await.unwrap().is_some());
1081
1082        tokio::time::advance(Duration::from_secs(7)).await;
1083        tokio::task::yield_now().await;
1084        assert!(retained_block.upgrade().is_some());
1085        assert!(retained_receipts.upgrade().is_some());
1086        assert!(retained_bal.upgrade().is_some());
1087        // No more cache actions: only a later timer tick can release these payloads.
1088        tokio::time::advance(Duration::from_secs(11)).await;
1089        tokio::task::yield_now().await;
1090        let expired = retained_block.upgrade().is_none() &&
1091            retained_receipts.upgrade().is_none() &&
1092            retained_bal.upgrade().is_none();
1093
1094        // Expiration must leave pending consumers available for the next result.
1095        cache
1096            .to_service
1097            .send(CacheAction::BlockWithSendersResult {
1098                block_hash: hash,
1099                res: Ok(Some(Arc::new(test_block()))),
1100            })
1101            .unwrap();
1102        let response = response_rx.await.unwrap();
1103        service_task.abort();
1104        assert!(expired, "idle service did not register its next timer wakeup");
1105        assert!(response.unwrap().is_some());
1106    }
1107
1108    #[tokio::test(start_paused = true)]
1109    async fn idle_eviction_yields_between_batches_and_serves_requests() {
1110        let (cache, mut service) = EthStateCache::<EthPrimitives>::create(
1111            NoopProvider::default(),
1112            Runtime::test(),
1113            EthStateCacheConfig { idle_timeout: Duration::from_secs(10), ..Default::default() },
1114        );
1115        let mut retained = Vec::new();
1116        for key in 0..12 {
1117            let hash = B256::with_last_byte(key);
1118            let block = Arc::new(test_block());
1119            let receipts = Arc::new(vec![]);
1120            let bal = CachedRevmBal::new(test_decoded_revm_bal());
1121            retained.push((
1122                Arc::downgrade(&block),
1123                Arc::downgrade(&receipts),
1124                Arc::downgrade(&bal.0),
1125            ));
1126            service.on_new_block(hash, Ok(Some(block)), Instant::now());
1127            service.on_new_receipts(hash, Ok(Some(receipts)), Instant::now());
1128            service.on_new_bal(hash, Ok(Some(bal)), Instant::now());
1129        }
1130        let (response_tx, mut response_rx) = oneshot::channel();
1131        cache
1132            .to_service
1133            .send(CacheAction::GetCachedBlock { block_hash: B256::repeat_byte(0xff), response_tx })
1134            .unwrap();
1135        tokio::time::advance(Duration::from_secs(10)).await;
1136
1137        let wakes = Arc::new(WakeCounter::default());
1138        let waker = Waker::from(wakes.clone());
1139        let mut cx = Context::from_waker(&waker);
1140        for (remaining, expected_wakes) in [(7, 1), (2, 2), (0, 2)] {
1141            assert!(Pin::new(&mut service).poll(&mut cx).is_pending());
1142            assert_eq!(wakes.0.load(Ordering::Relaxed), expected_wakes);
1143            assert_eq!(
1144                retained.iter().filter(|(block, _, _)| block.strong_count() > 0).count(),
1145                remaining
1146            );
1147            assert_eq!(
1148                retained.iter().filter(|(_, receipts, _)| receipts.strong_count() > 0).count(),
1149                remaining
1150            );
1151            assert_eq!(
1152                retained.iter().filter(|(_, _, bal)| bal.strong_count() > 0).count(),
1153                remaining
1154            );
1155            if remaining == 7 {
1156                assert!(response_rx.try_recv().unwrap().is_none());
1157            }
1158        }
1159    }
1160
1161    #[tokio::test(start_paused = true)]
1162    async fn disabled_idle_timeout_does_not_schedule_cleanup() {
1163        let (_cache, mut service) = EthStateCache::<EthPrimitives>::create(
1164            NoopProvider::default(),
1165            Runtime::test(),
1166            EthStateCacheConfig { idle_timeout: Duration::ZERO, ..Default::default() },
1167        );
1168        let wakes = Arc::new(WakeCounter::default());
1169        let waker = Waker::from(wakes.clone());
1170        let mut cx = Context::from_waker(&waker);
1171        assert!(Pin::new(&mut service).poll(&mut cx).is_pending());
1172
1173        tokio::time::advance(Duration::from_secs(10)).await;
1174        assert_eq!(wakes.0.load(Ordering::Relaxed), 0);
1175    }
1176
1177    #[test]
1178    fn oversized_results_reach_queued_consumers() {
1179        let block = Arc::new(test_block());
1180        let receipts = Arc::new(vec![Receipt::default()]);
1181        let bal = CachedRevmBal::new(test_decoded_revm_bal());
1182        let hash = block.hash();
1183        let (_cache, mut service) = EthStateCache::<EthPrimitives>::create(
1184            NoopProvider::default(),
1185            Runtime::test(),
1186            EthStateCacheConfig {
1187                max_blocks_bytes: block.size() - 1,
1188                max_receipts_bytes: receipts.size() - 1,
1189                max_bals_bytes: bal.size() - 1,
1190                ..Default::default()
1191            },
1192        );
1193        let (block_tx, mut block_rx) = oneshot::channel();
1194        let (receipts_tx, mut receipts_rx) = oneshot::channel();
1195        let (bal_tx, mut bal_rx) = oneshot::channel();
1196        service.full_block_cache.queue(hash, block_tx);
1197        service.receipts_cache.queue(hash, receipts_tx);
1198        service.bal_cache.queue(hash, bal_tx);
1199
1200        service.on_new_block(hash, Ok(Some(block.clone())), Instant::now());
1201        service.on_new_receipts(hash, Ok(Some(receipts.clone())), Instant::now());
1202        service.on_new_bal(hash, Ok(Some(bal.clone())), Instant::now());
1203
1204        assert!(Arc::ptr_eq(&block_rx.try_recv().unwrap().unwrap().unwrap(), &block));
1205        assert!(Arc::ptr_eq(&receipts_rx.try_recv().unwrap().unwrap().unwrap(), &receipts));
1206        assert!(Arc::ptr_eq(&bal_rx.try_recv().unwrap().unwrap().unwrap().0, &bal.0));
1207        assert!(service.full_block_cache.get(&hash).is_none());
1208        assert!(service.receipts_cache.get(&hash).is_none());
1209        assert!(service.bal_cache.get(&hash).is_none());
1210    }
1211
1212    #[tokio::test(start_paused = true)]
1213    async fn cached_transaction_lookup_renews_idle_timeout() {
1214        let (cache, mut service) = EthStateCache::<EthPrimitives>::create(
1215            NoopProvider::default(),
1216            Runtime::test(),
1217            EthStateCacheConfig { idle_timeout: Duration::from_secs(10), ..Default::default() },
1218        );
1219        let block = Arc::new(test_block());
1220        let hash = block.hash();
1221        let tx_hash = *block.body().transactions().next().unwrap().tx_hash();
1222        service.index_block_transactions(&block);
1223        service.on_new_block(hash, Ok(Some(block)), Instant::now());
1224        service.on_new_receipts(hash, Ok(Some(Arc::new(vec![Receipt::default()]))), Instant::now());
1225
1226        for delay in [6, 6, 11] {
1227            tokio::time::advance(Duration::from_secs(delay)).await;
1228            let request = cache.get_transaction_by_hash(tx_hash);
1229            futures::pin_mut!(request);
1230            assert!(futures::poll!(&mut request).is_pending());
1231            assert!(futures::poll!(&mut service).is_pending());
1232            let result = request.await;
1233            if delay < 10 {
1234                let result = result.expect("transaction is still cached after a renewed hit");
1235                assert_eq!(result.block.hash(), hash);
1236                assert!(result.receipts.is_some());
1237            } else {
1238                assert!(result.is_none());
1239            }
1240        }
1241    }
1242
1243    #[test]
1244    fn reorg_removes_tx_hash_index_entries_unconditionally() {
1245        let mut service = test_service();
1246        let block = test_block();
1247        let tx_hash = *block.body().transactions().next().expect("test transaction").tx_hash();
1248
1249        service.tx_hash_index.insert(tx_hash, (B256::repeat_byte(0x33), 0));
1250
1251        service.remove_block_transactions(&block);
1252
1253        assert!(service.tx_hash_index.get(&tx_hash).is_none());
1254    }
1255
1256    #[test]
1257    fn reorg_evicts_cached_bal() {
1258        let mut service = test_service();
1259        let block_hash = B256::repeat_byte(0x44);
1260
1261        assert!(service.bal_cache.insert(block_hash, CachedRevmBal::new(test_decoded_revm_bal())));
1262        assert!(service.bal_cache.get(&block_hash).is_some());
1263
1264        service.on_reorg_bal(block_hash, Ok(None));
1265
1266        assert!(service.bal_cache.get(&block_hash).is_none());
1267    }
1268
1269    #[test]
1270    fn late_provider_miss_preserves_generated_bal() {
1271        let mut service = test_service();
1272        let hash = B256::repeat_byte(0x69);
1273        let (tx, mut rx) = oneshot::channel();
1274        assert!(service.bal_cache.queue(hash, tx));
1275        service.on_new_bal(
1276            hash,
1277            Ok(Some(CachedRevmBal::new(test_decoded_revm_bal()))),
1278            Instant::now(),
1279        );
1280        assert!(rx.try_recv().unwrap().unwrap().is_some());
1281        service.on_new_bal(hash, Ok(None), Instant::now());
1282        assert!(service.bal_cache.contains_key(&hash));
1283    }
1284
1285    #[tokio::test]
1286    async fn cache_misses_wait_in_service_queue_for_db_slot() {
1287        let (started_tx, started_rx) = mpsc::sync_channel(1);
1288        let (release_tx, release_rx) = mpsc::sync_channel(1);
1289        let (_cache, mut service) = EthStateCache::<EthPrimitives>::create(
1290            TestBalProvider::new_blocking(started_tx, release_rx),
1291            Runtime::test(),
1292            EthStateCacheConfig { max_concurrent_db_requests: 1, ..Default::default() },
1293        );
1294        let permit = service.rate_limiter.clone().try_acquire_owned().unwrap();
1295
1296        service.queue_fetch(CacheFetch::Bal(B256::repeat_byte(0x01)));
1297        service.queue_fetch(CacheFetch::Bal(B256::repeat_byte(0x02)));
1298
1299        assert_eq!(service.pending_fetches.len(), 2);
1300
1301        drop(permit);
1302        service.spawn_pending_fetches();
1303
1304        started_rx.recv_timeout(Duration::from_secs(1)).expect("first fetch started");
1305        assert_eq!(service.pending_fetches.len(), 1);
1306        release_tx.send(()).unwrap();
1307    }
1308
1309    #[test]
1310    fn pending_cache_miss_does_not_block_unrelated_blocking_work() {
1311        let (started_tx, started_rx) = mpsc::sync_channel(1);
1312        let (release_tx, release_rx) = mpsc::sync_channel(1);
1313        let (unrelated_tx, unrelated_rx) = mpsc::sync_channel(1);
1314
1315        let test = thread::spawn(move || {
1316            let runtime = tokio::runtime::Builder::new_multi_thread()
1317                .worker_threads(1)
1318                .max_blocking_threads(1)
1319                .enable_all()
1320                .build()
1321                .unwrap();
1322
1323            runtime.block_on(async move {
1324                let provider = TestBalProvider::new_blocking(started_tx, release_rx);
1325                let (_cache, mut service) = EthStateCache::<EthPrimitives>::create(
1326                    provider,
1327                    Runtime::test(),
1328                    EthStateCacheConfig {
1329                        max_blocks: 0,
1330                        max_receipts: 0,
1331                        max_bals: 0,
1332                        cache_computed_bals: false,
1333                        prewarm_bals: None,
1334                        max_concurrent_db_requests: 1,
1335                        max_cached_tx_hashes: 0,
1336                        ..Default::default()
1337                    },
1338                );
1339
1340                service.queue_fetch(CacheFetch::Bal(B256::repeat_byte(0x01)));
1341                started_rx.recv_timeout(Duration::from_secs(1)).expect("first fetch started");
1342
1343                service.queue_fetch(CacheFetch::Bal(B256::repeat_byte(0x02)));
1344                assert_eq!(service.pending_fetches.len(), 1);
1345
1346                tokio::task::spawn_blocking(move || unrelated_tx.send(()).unwrap());
1347                release_tx.send(()).unwrap();
1348
1349                unrelated_rx
1350                    .recv_timeout(Duration::from_secs(1))
1351                    .expect("unrelated blocking work was not starved");
1352            });
1353        });
1354
1355        test.join().unwrap();
1356    }
1357
1358    #[test]
1359    fn reorg_forwards_bal_to_queued_requests() {
1360        let mut service = test_service();
1361        let block_hash = B256::repeat_byte(0x55);
1362        let (response_tx, mut response_rx) = oneshot::channel();
1363        let bal = CachedRevmBal::new(test_decoded_revm_bal());
1364
1365        assert!(service.bal_cache.queue(block_hash, response_tx));
1366
1367        service.on_reorg_bal(block_hash, Ok(Some(bal)));
1368
1369        let bal = response_rx.try_recv().expect("queued BAL response").expect("BAL result");
1370
1371        assert!(bal.is_some());
1372    }
1373
1374    #[test]
1375    fn cached_revm_bal_size_accounts_for_nested_allocations() {
1376        let mut account = RevmAccountBal::default();
1377        account.account_info.nonce.writes.push((BlockAccessIndex::new(1), 1));
1378        account
1379            .account_info
1380            .balance
1381            .writes
1382            .push((BlockAccessIndex::new(2), StorageValue::from(1u64)));
1383        account.account_info.code.writes.push((
1384            BlockAccessIndex::new(3),
1385            (B256::repeat_byte(0xaa), Bytecode::new_raw(Bytes::from_static(&[0x60, 0x00]))),
1386        ));
1387        account.storage.storage.insert(
1388            StorageKey::from(1u64),
1389            RevmBalWrites::new(vec![(BlockAccessIndex::new(4), StorageValue::from(2u64))]),
1390        );
1391
1392        let mut bal = RevmBal::default();
1393        bal.accounts.insert(Address::ZERO, account);
1394
1395        let raw = Bytes::from_static(&[0xc0, 0x01, 0x02]);
1396        let previous_estimate = core::mem::size_of::<CachedRevmBal>() +
1397            core::mem::size_of::<DecodedBal<Arc<RevmBal>>>() +
1398            raw.len() +
1399            core::mem::size_of::<RevmBal>();
1400        assert!(CachedRevmBal::new(DecodedBal::new(Arc::new(bal), raw)).size() > previous_estimate);
1401    }
1402
1403    #[tokio::test]
1404    async fn get_bal_only_caches_results_within_budget() {
1405        for (max_bals_bytes, expected_fetches) in [(usize::MAX, 1), (1, 2), (0, 2)] {
1406            let fetches = Arc::new(AtomicUsize::default());
1407            let cache = EthStateCache::<EthPrimitives>::spawn_with(
1408                TestBalProvider::new(fetches.clone()),
1409                EthStateCacheConfig {
1410                    max_blocks: 0,
1411                    max_receipts: 0,
1412                    max_bals: 4,
1413                    max_bals_bytes,
1414                    cache_computed_bals: false,
1415                    prewarm_bals: None,
1416                    max_concurrent_db_requests: 1,
1417                    max_cached_tx_hashes: 0,
1418                    ..Default::default()
1419                },
1420                Runtime::test(),
1421            );
1422            let block_hash = B256::repeat_byte(0x66);
1423
1424            assert!(cache.get_bal(block_hash).await.unwrap().is_some());
1425            assert!(cache.get_bal(block_hash).await.unwrap().is_some());
1426
1427            assert_eq!(fetches.load(Ordering::SeqCst), expected_fetches);
1428        }
1429    }
1430
1431    #[tokio::test]
1432    async fn get_recovered_block_and_maybe_bal_does_not_fetch_bal() {
1433        let bal_fetches = Arc::new(AtomicUsize::default());
1434        let block = test_block();
1435        let block_hash = block.hash();
1436        let provider = TestBalProvider::new(bal_fetches.clone()).with_block(block);
1437        let cache = EthStateCache::<EthPrimitives>::spawn_with(
1438            provider,
1439            EthStateCacheConfig {
1440                max_blocks: 4,
1441                max_receipts: 0,
1442                max_bals: 4,
1443                cache_computed_bals: false,
1444                prewarm_bals: None,
1445                max_concurrent_db_requests: 1,
1446                max_cached_tx_hashes: 0,
1447                ..Default::default()
1448            },
1449            Runtime::test(),
1450        );
1451
1452        let (returned_block, bal) = cache
1453            .get_recovered_block_and_maybe_bal(block_hash)
1454            .await
1455            .unwrap()
1456            .expect("block exists");
1457        assert_eq!(returned_block.hash(), block_hash);
1458        assert!(bal.is_none());
1459        assert_eq!(bal_fetches.load(Ordering::SeqCst), 0);
1460
1461        assert!(cache.get_bal(block_hash).await.unwrap().is_some());
1462
1463        let (_, bal) = cache
1464            .get_recovered_block_and_maybe_bal(block_hash)
1465            .await
1466            .unwrap()
1467            .expect("block exists");
1468        assert!(bal.is_some());
1469        assert_eq!(bal_fetches.load(Ordering::SeqCst), 1);
1470    }
1471
1472    #[tokio::test]
1473    async fn insert_bal_populates_cache_without_provider_fetch() {
1474        let fetches = Arc::new(AtomicUsize::default());
1475        let provider = TestBalProvider::new(fetches.clone());
1476        let cache = EthStateCache::<EthPrimitives>::spawn_with(
1477            provider,
1478            EthStateCacheConfig { max_bals: 4, ..Default::default() },
1479            Runtime::test(),
1480        );
1481        let block_hash = B256::repeat_byte(0x68);
1482        cache.insert_bal(block_hash, test_decoded_revm_bal());
1483        assert!(cache.get_bal(block_hash).await.unwrap().is_some());
1484        assert_eq!(fetches.load(Ordering::SeqCst), 0);
1485    }
1486
1487    #[tokio::test]
1488    async fn canonical_chain_notification_caches_prepared_bal() {
1489        let fetches = Arc::new(AtomicUsize::default());
1490        let provider = TestBalProvider::new(fetches.clone());
1491        let cache = EthStateCache::<EthPrimitives>::spawn_with(
1492            provider,
1493            EthStateCacheConfig { max_blocks: 4, max_bals: 4, ..Default::default() },
1494            Runtime::test(),
1495        );
1496
1497        let block = test_block();
1498        let block_hash = block.hash();
1499        let block_number = block.number;
1500        let mut chain: Chain<EthPrimitives> = Chain::new(
1501            [block],
1502            ExecutionOutcome::new(Default::default(), vec![vec![]], block_number, vec![]),
1503            Default::default(),
1504        );
1505        chain.insert_bal(block_number, Arc::new(test_decoded_revm_bal()));
1506
1507        cache_new_blocks_task(
1508            cache.clone(),
1509            tokio_stream::iter([CanonStateNotification::Commit { new: Arc::new(chain) }]),
1510        )
1511        .await;
1512
1513        assert!(cache.get_bal(block_hash).await.unwrap().is_some());
1514        assert_eq!(fetches.load(Ordering::SeqCst), 0);
1515    }
1516
1517    #[tokio::test]
1518    async fn concurrent_get_bal_requests_share_fetch() {
1519        let fetches = Arc::new(AtomicUsize::default());
1520        let provider = TestBalProvider::new(fetches.clone());
1521        let cache = EthStateCache::<EthPrimitives>::spawn_with(
1522            provider,
1523            EthStateCacheConfig {
1524                max_blocks: 0,
1525                max_receipts: 0,
1526                max_bals: 4,
1527                cache_computed_bals: false,
1528                prewarm_bals: None,
1529                max_concurrent_db_requests: 1,
1530                max_cached_tx_hashes: 0,
1531                ..Default::default()
1532            },
1533            Runtime::test(),
1534        );
1535        let block_hash = B256::repeat_byte(0x77);
1536
1537        let (first, second) = tokio::join!(cache.get_bal(block_hash), cache.get_bal(block_hash));
1538
1539        assert!(first.unwrap().is_some());
1540        assert!(second.unwrap().is_some());
1541        assert_eq!(fetches.load(Ordering::SeqCst), 1);
1542    }
1543
1544    #[derive(Clone, Debug, Default)]
1545    struct TestBalProvider {
1546        bal_store: BalStoreHandle,
1547        block: Option<RecoveredBlock<Block>>,
1548    }
1549
1550    impl TestBalProvider {
1551        fn new(fetches: Arc<AtomicUsize>) -> Self {
1552            Self {
1553                bal_store: BalStoreHandle::new(TestBalStore { fetches, blocking_fetch: None }),
1554                block: None,
1555            }
1556        }
1557
1558        fn with_block(mut self, block: RecoveredBlock<Block>) -> Self {
1559            self.block = Some(block);
1560            self
1561        }
1562
1563        fn new_blocking(started: SyncSender<()>, release: Receiver<()>) -> Self {
1564            Self {
1565                bal_store: BalStoreHandle::new(TestBalStore {
1566                    fetches: Arc::new(AtomicUsize::default()),
1567                    blocking_fetch: Some(BlockingFetch { started, release: Mutex::new(release) }),
1568                }),
1569                block: None,
1570            }
1571        }
1572    }
1573
1574    impl BalProvider for TestBalProvider {
1575        fn bal_store(&self) -> &BalStoreHandle {
1576            &self.bal_store
1577        }
1578    }
1579
1580    #[derive(Debug)]
1581    struct TestBalStore {
1582        fetches: Arc<AtomicUsize>,
1583        blocking_fetch: Option<BlockingFetch>,
1584    }
1585
1586    #[derive(Debug)]
1587    struct BlockingFetch {
1588        started: SyncSender<()>,
1589        release: Mutex<Receiver<()>>,
1590    }
1591
1592    impl BalStore for TestBalStore {
1593        fn insert(&self, _num_hash: NumHash, _bal: reth_storage_api::RawBal) -> ProviderResult<()> {
1594            Ok(())
1595        }
1596
1597        fn prune(&self, _tip: BlockNumber) -> ProviderResult<usize> {
1598            Ok(0)
1599        }
1600
1601        fn get_by_hashes(&self, block_hashes: &[BlockHash]) -> ProviderResult<Vec<Option<Bytes>>> {
1602            self.fetches.fetch_add(1, Ordering::SeqCst);
1603            if let Some(blocking_fetch) = &self.blocking_fetch {
1604                blocking_fetch.started.send(()).unwrap();
1605                blocking_fetch.release.lock().unwrap().recv().unwrap();
1606            }
1607            Ok(block_hashes.iter().map(|_| Some(Bytes::from_static(&[0xc0]))).collect())
1608        }
1609    }
1610
1611    impl BlockHashReader for TestBalProvider {
1612        fn block_hash(&self, _number: BlockNumber) -> ProviderResult<Option<B256>> {
1613            Ok(None)
1614        }
1615
1616        fn canonical_hashes_range(
1617            &self,
1618            _start: BlockNumber,
1619            _end: BlockNumber,
1620        ) -> ProviderResult<Vec<B256>> {
1621            Ok(Vec::new())
1622        }
1623    }
1624
1625    impl BlockNumReader for TestBalProvider {
1626        fn chain_info(&self) -> ProviderResult<reth_chainspec::ChainInfo> {
1627            Ok(reth_chainspec::ChainInfo::default())
1628        }
1629
1630        fn best_block_number(&self) -> ProviderResult<BlockNumber> {
1631            Ok(0)
1632        }
1633
1634        fn last_block_number(&self) -> ProviderResult<BlockNumber> {
1635            Ok(0)
1636        }
1637
1638        fn block_number(&self, _hash: B256) -> ProviderResult<Option<BlockNumber>> {
1639            Ok(None)
1640        }
1641    }
1642
1643    impl HeaderProvider for TestBalProvider {
1644        type Header = Header;
1645
1646        fn header(&self, _block_hash: BlockHash) -> ProviderResult<Option<Self::Header>> {
1647            Ok(None)
1648        }
1649
1650        fn header_by_number(&self, _num: u64) -> ProviderResult<Option<Self::Header>> {
1651            Ok(None)
1652        }
1653
1654        fn headers_range(
1655            &self,
1656            _range: impl RangeBounds<BlockNumber>,
1657        ) -> ProviderResult<Vec<Self::Header>> {
1658            Ok(Vec::new())
1659        }
1660
1661        fn sealed_header(
1662            &self,
1663            _number: BlockNumber,
1664        ) -> ProviderResult<Option<SealedHeader<Self::Header>>> {
1665            Ok(None)
1666        }
1667
1668        fn sealed_headers_while(
1669            &self,
1670            _range: impl RangeBounds<BlockNumber>,
1671            _predicate: impl FnMut(&SealedHeader<Self::Header>) -> bool,
1672        ) -> ProviderResult<Vec<SealedHeader<Self::Header>>> {
1673            Ok(Vec::new())
1674        }
1675    }
1676
1677    impl BlockBodyIndicesProvider for TestBalProvider {
1678        fn block_body_indices(&self, _num: u64) -> ProviderResult<Option<StoredBlockBodyIndices>> {
1679            Ok(None)
1680        }
1681
1682        fn block_body_indices_range(
1683            &self,
1684            _range: RangeInclusive<BlockNumber>,
1685        ) -> ProviderResult<Vec<StoredBlockBodyIndices>> {
1686            Ok(Vec::new())
1687        }
1688    }
1689
1690    impl TransactionsProvider for TestBalProvider {
1691        type Transaction = TransactionSigned;
1692
1693        fn transaction_id(&self, _tx_hash: TxHash) -> ProviderResult<Option<TxNumber>> {
1694            Ok(None)
1695        }
1696
1697        fn transaction_by_id(&self, _id: TxNumber) -> ProviderResult<Option<Self::Transaction>> {
1698            Ok(None)
1699        }
1700
1701        fn transaction_by_id_unhashed(
1702            &self,
1703            _id: TxNumber,
1704        ) -> ProviderResult<Option<Self::Transaction>> {
1705            Ok(None)
1706        }
1707
1708        fn transaction_by_hash(&self, _hash: TxHash) -> ProviderResult<Option<Self::Transaction>> {
1709            Ok(None)
1710        }
1711
1712        fn transaction_by_hash_with_meta(
1713            &self,
1714            _hash: TxHash,
1715        ) -> ProviderResult<Option<(Self::Transaction, TransactionMeta)>> {
1716            Ok(None)
1717        }
1718
1719        fn transactions_by_block(
1720            &self,
1721            _block: BlockHashOrNumber,
1722        ) -> ProviderResult<Option<Vec<Self::Transaction>>> {
1723            Ok(None)
1724        }
1725
1726        fn transactions_by_block_range(
1727            &self,
1728            _range: impl RangeBounds<BlockNumber>,
1729        ) -> ProviderResult<Vec<Vec<Self::Transaction>>> {
1730            Ok(Vec::new())
1731        }
1732
1733        fn transactions_by_tx_range(
1734            &self,
1735            _range: impl RangeBounds<TxNumber>,
1736        ) -> ProviderResult<Vec<Self::Transaction>> {
1737            Ok(Vec::new())
1738        }
1739
1740        fn senders_by_tx_range(
1741            &self,
1742            _range: impl RangeBounds<TxNumber>,
1743        ) -> ProviderResult<Vec<Address>> {
1744            Ok(Vec::new())
1745        }
1746
1747        fn transaction_sender(&self, _id: TxNumber) -> ProviderResult<Option<Address>> {
1748            Ok(None)
1749        }
1750    }
1751
1752    impl ReceiptProvider for TestBalProvider {
1753        type Receipt = Receipt;
1754
1755        fn receipt(&self, _id: TxNumber) -> ProviderResult<Option<Self::Receipt>> {
1756            Ok(None)
1757        }
1758
1759        fn receipt_by_hash(&self, _hash: TxHash) -> ProviderResult<Option<Self::Receipt>> {
1760            Ok(None)
1761        }
1762
1763        fn receipts_by_block(
1764            &self,
1765            _block: BlockHashOrNumber,
1766        ) -> ProviderResult<Option<Vec<Self::Receipt>>> {
1767            Ok(None)
1768        }
1769
1770        fn receipts_by_tx_range(
1771            &self,
1772            _range: impl RangeBounds<TxNumber>,
1773        ) -> ProviderResult<Vec<Self::Receipt>> {
1774            Ok(Vec::new())
1775        }
1776
1777        fn receipts_by_block_range(
1778            &self,
1779            _block_range: RangeInclusive<BlockNumber>,
1780        ) -> ProviderResult<Vec<Vec<Self::Receipt>>> {
1781            Ok(Vec::new())
1782        }
1783    }
1784
1785    impl BlockReader for TestBalProvider {
1786        type Block = Block;
1787
1788        fn find_block_by_hash(
1789            &self,
1790            _hash: B256,
1791            _source: BlockSource,
1792        ) -> ProviderResult<Option<Self::Block>> {
1793            Ok(None)
1794        }
1795
1796        fn block(&self, _id: BlockHashOrNumber) -> ProviderResult<Option<Self::Block>> {
1797            Ok(None)
1798        }
1799
1800        fn pending_block(&self) -> ProviderResult<Option<Arc<RecoveredBlock<Self::Block>>>> {
1801            Ok(None)
1802        }
1803
1804        fn pending_block_and_receipts(
1805            &self,
1806        ) -> ProviderResult<Option<RecoveredBlockAndExecutionOutput<Self::Block, Self::Receipt>>>
1807        {
1808            Ok(None)
1809        }
1810
1811        fn recovered_block(
1812            &self,
1813            _id: BlockHashOrNumber,
1814            _transaction_kind: TransactionVariant,
1815        ) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
1816            Ok(None)
1817        }
1818
1819        fn sealed_block_with_senders(
1820            &self,
1821            _id: BlockHashOrNumber,
1822            _transaction_kind: TransactionVariant,
1823        ) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
1824            Ok(self.block.clone())
1825        }
1826
1827        fn block_range(
1828            &self,
1829            _range: RangeInclusive<BlockNumber>,
1830        ) -> ProviderResult<Vec<Self::Block>> {
1831            Ok(Vec::new())
1832        }
1833
1834        fn block_with_senders_range(
1835            &self,
1836            _range: RangeInclusive<BlockNumber>,
1837        ) -> ProviderResult<Vec<RecoveredBlock<Self::Block>>> {
1838            Ok(Vec::new())
1839        }
1840
1841        fn recovered_block_range(
1842            &self,
1843            _range: RangeInclusive<BlockNumber>,
1844        ) -> ProviderResult<Vec<RecoveredBlock<Self::Block>>> {
1845            Ok(Vec::new())
1846        }
1847
1848        fn block_by_transaction_id(&self, _id: TxNumber) -> ProviderResult<Option<BlockNumber>> {
1849            Ok(None)
1850        }
1851    }
1852}