1use 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
48type BlockWithSendersResponseSender<B> =
50 oneshot::Sender<ProviderResult<Option<Arc<RecoveredBlock<B>>>>>;
51
52type 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
62type TransactionHashResponseSender<B, R> = oneshot::Sender<Option<CachedTransaction<B, R>>>;
64
65type 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#[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 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 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 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 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 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 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 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 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 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 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 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 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 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#[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#[must_use = "Type does nothing unless spawned"]
319pub(crate) struct EthStateCacheService<Provider>
320where
321 Provider: BlockReader + BalProvider,
322{
323 provider: Provider,
325 full_block_cache: BlockLruCache<Provider::Block>,
327 receipts_cache: ReceiptsLruCache<Provider::Receipt>,
329 bal_cache: BalLruCache,
331 idle_timeout: Duration,
333 eviction_interval: Interval,
335 eviction_pending: bool,
337 action_tx: UnboundedSender<CacheAction<Provider::Block, Provider::Receipt>>,
339 action_rx: UnboundedReceiverStream<CacheAction<Provider::Block, Provider::Receipt>>,
341 action_task_spawner: Runtime,
343 rate_limiter: Arc<Semaphore>,
347 pending_fetches: VecDeque<CacheFetch>,
352 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 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 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 for tx in queued {
449 let _ = tx.send(res.clone());
450 }
451 }
452
453 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 for tx in queued {
468 let _ = tx.send(res.clone());
469 }
470 }
471
472 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 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 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 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 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 let now = LazyCell::new(Instant::now);
559
560 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 cx.waker().wake_by_ref();
572 }
573 }
574
575 loop {
576 let Poll::Ready(action) = this.action_rx.poll_next_unpin(cx) else {
577 this.shrink_queues();
579 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 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 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 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 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 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
712enum 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 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
772struct 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 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#[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#[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
877pub 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#[derive(Clone, Debug)]
904pub(crate) struct CachedRevmBal(Arc<DecodedRevmBal>);
905
906impl CachedRevmBal {
907 #[inline]
909 fn new(bal: DecodedRevmBal) -> Self {
910 Self(Arc::new(bal))
911 }
912
913 #[inline]
916 const fn from_shared(bal: Arc<DecodedRevmBal>) -> Self {
917 Self(bal)
918 }
919
920 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 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 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}