1use crate::{
8 changeset_cache::compute_block_trie_updates,
9 database_state_frontiers,
10 manager_metrics::{ExecutionOverlayMetrics, OverlayCacheMetrics, StateTrieOverlayMetrics},
11 ChangesetCache, ExecutionOverlay, OverlayBuilder,
12};
13use alloy_eips::BlockNumHash;
14use alloy_primitives::{BlockNumber, B256};
15use parking_lot::Mutex;
16use reth_chain_state::{BlockState, ExecutedBlock, PreservedSparseTrie};
17use reth_errors::ProviderResult;
18use reth_ethereum_primitives::EthPrimitives;
19use reth_primitives_traits::{
20 dashmap::{mapref::entry::Entry, DashMap},
21 AlloyBlockHeader, FastInstant, NodePrimitives,
22};
23use reth_storage_api::{
24 BlockNumReader, ChangeSetReader, DBProvider, PruneCheckpointReader, StageCheckpointReader,
25 StorageChangeSetReader, StorageSettingsCache,
26};
27#[cfg(feature = "rayon")]
28use reth_tasks::WorkerPool;
29use reth_trie::{updates::TrieUpdatesSorted, HashedPostStateSorted, TrieInputSorted};
30use std::{
31 fmt,
32 ops::RangeInclusive,
33 sync::{Arc, OnceLock},
34 time::Instant,
35};
36use tracing::{debug, trace};
37
38#[derive(Clone)]
43pub struct OverlayManager<N: NodePrimitives = EthPrimitives> {
44 blocks: Arc<DashMap<B256, ExecutedBlock<N>>>,
45 state_trie_overlays: OverlayCache<TrieInputSorted>,
46 execution_overlays: OverlayCache<ExecutionOverlay>,
47 changeset_cache: ChangesetCache,
48 preserved_sparse_trie: Arc<Mutex<Option<PreservedSparseTrie>>>,
49 #[cfg(feature = "rayon")]
50 worker_pool: Option<Arc<WorkerPool>>,
51 metrics: StateTrieOverlayMetrics,
52 execution_metrics: ExecutionOverlayMetrics,
53}
54
55impl<N: NodePrimitives> Default for OverlayManager<N> {
56 fn default() -> Self {
57 Self {
58 blocks: Default::default(),
59 state_trie_overlays: Default::default(),
60 execution_overlays: Default::default(),
61 changeset_cache: Default::default(),
62 preserved_sparse_trie: Default::default(),
63 #[cfg(feature = "rayon")]
64 worker_pool: None,
65 metrics: Default::default(),
66 execution_metrics: Default::default(),
67 }
68 }
69}
70
71impl<N: NodePrimitives> std::fmt::Debug for OverlayManager<N> {
72 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
73 f.debug_struct("OverlayManager")
74 .field("blocks", &self.blocks.len())
75 .field("state_trie_overlays", &self.state_trie_overlays.len())
76 .field("execution_overlays", &self.execution_overlays.len())
77 .finish()
78 }
79}
80
81impl<N: NodePrimitives> OverlayManager<N> {
82 #[cfg(feature = "rayon")]
84 pub fn new(worker_pool: Arc<WorkerPool>) -> Self {
85 Self {
86 blocks: Default::default(),
87 state_trie_overlays: Default::default(),
88 execution_overlays: Default::default(),
89 changeset_cache: Default::default(),
90 preserved_sparse_trie: Default::default(),
91 worker_pool: Some(worker_pool),
92 metrics: Default::default(),
93 execution_metrics: Default::default(),
94 }
95 }
96
97 pub fn overlay_builder(&self, parent_hash: B256) -> OverlayBuilder<N> {
102 OverlayBuilder::new(parent_hash, self.block_state(parent_hash).map(Arc::new), self.clone())
103 }
104
105 pub fn overlay_builder_for_state(&self, state: Arc<BlockState<N>>) -> OverlayBuilder<N> {
111 OverlayBuilder::new(state.hash(), Some(state), self.clone())
112 }
113
114 pub(crate) fn block_state(&self, parent_hash: B256) -> Option<BlockState<N>> {
115 let mut blocks = self.parent_chain(parent_hash).collect::<Vec<_>>();
116 blocks.pop().map(|oldest| {
117 blocks.into_iter().rev().fold(BlockState::new(oldest), |parent, block| {
118 BlockState::with_parent(block, Some(Arc::new(parent)))
119 })
120 })
121 }
122
123 pub(crate) const fn changeset_cache(&self) -> &ChangesetCache {
124 &self.changeset_cache
125 }
126
127 pub fn get_or_compute_cached_changesets_range<P>(
129 &self,
130 provider: &P,
131 range: RangeInclusive<BlockNumber>,
132 ) -> ProviderResult<Arc<TrieUpdatesSorted>>
133 where
134 P: DBProvider
135 + ChangeSetReader
136 + StorageChangeSetReader
137 + StageCheckpointReader
138 + PruneCheckpointReader
139 + BlockNumReader
140 + StorageSettingsCache,
141 {
142 let (partial_state_trie, finish) = database_state_frontiers(provider)?;
143 self.get_or_compute_cached_changesets_range_at_frontiers(
144 provider,
145 range,
146 partial_state_trie,
147 finish,
148 )
149 }
150
151 pub(crate) fn get_or_compute_cached_changesets_range_at_frontiers<P>(
152 &self,
153 provider: &P,
154 range: RangeInclusive<BlockNumber>,
155 partial_state_trie: BlockNumHash,
156 finish: BlockNumHash,
157 ) -> ProviderResult<Arc<TrieUpdatesSorted>>
158 where
159 P: DBProvider
160 + ChangeSetReader
161 + StorageChangeSetReader
162 + StageCheckpointReader
163 + PruneCheckpointReader
164 + BlockNumReader
165 + StorageSettingsCache,
166 {
167 self.changeset_cache.get_or_compute_range(self, provider, range, partial_state_trie, finish)
168 }
169
170 pub fn evict_cached_changesets(&self, up_to_block: BlockNumber) {
172 self.changeset_cache.evict(up_to_block);
173 }
174
175 pub fn compute_block_trie_updates<P>(
177 &self,
178 provider: &P,
179 block_number: BlockNumber,
180 ) -> ProviderResult<TrieUpdatesSorted>
181 where
182 P: DBProvider
183 + ChangeSetReader
184 + StorageChangeSetReader
185 + PruneCheckpointReader
186 + StageCheckpointReader
187 + BlockNumReader
188 + StorageSettingsCache,
189 {
190 compute_block_trie_updates(self, provider, block_number)
191 }
192
193 pub fn take_sparse_trie(&self) -> Option<PreservedSparseTrie> {
195 self.preserved_sparse_trie.lock().take()
196 }
197
198 pub fn store_sparse_trie(&self, trie: PreservedSparseTrie) {
200 *self.preserved_sparse_trie.lock() = Some(trie);
201 }
202
203 pub fn clear_sparse_trie(&self) {
205 *self.preserved_sparse_trie.lock() = None;
206 }
207
208 pub fn wait_for_sparse_trie_availability(&self) -> std::time::Duration {
213 let start = FastInstant::now();
214 let _guard = self.preserved_sparse_trie.lock();
215 let elapsed = start.elapsed();
216 if elapsed.as_millis() > 5 {
217 debug!(
218 target: "storage::overlay::manager",
219 blocked_for=?elapsed,
220 "Waited for preserved sparse trie to become available"
221 );
222 }
223 elapsed
224 }
225
226 #[tracing::instrument(
228 level = "trace",
229 target = "storage::overlay::manager",
230 skip_all,
231 fields(
232 block_hash = %block.recovered_block().hash(),
233 parent_hash = %block.recovered_block().parent_hash(),
234 duplicate = false,
235 )
236 )]
237 pub fn insert_block(&self, block: ExecutedBlock<N>) {
238 let hash = block.recovered_block().hash();
239 let parent_hash = block.recovered_block().parent_hash();
240 let span = tracing::Span::current();
241
242 match self.blocks.entry(hash) {
244 Entry::Occupied(_) => {
245 span.record("duplicate", true);
246 debug!(
247 target: "storage::overlay::manager",
248 %hash,
249 %parent_hash,
250 "state trie overlay block already inserted"
251 );
252 return
253 }
254 Entry::Vacant(entry) => {
255 entry.insert(block);
256 }
257 }
258
259 let cached_parent_overlays = self
262 .execution_overlays
263 .entries
264 .iter()
265 .filter_map(|entry| {
266 let key = *entry.key();
267 (key.tip_hash == parent_hash).then_some(key.anchor_hash)
268 })
269 .collect::<Vec<_>>();
270
271 debug!(
272 target: "storage::overlay::manager",
273 %hash,
274 %parent_hash,
275 "inserted block into state trie overlay manager"
276 );
277 if cached_parent_overlays.is_empty() {
278 return
279 }
280
281 #[cfg(not(feature = "rayon"))]
282 let _ = cached_parent_overlays;
283
284 #[cfg(feature = "rayon")]
287 {
288 for anchor_hash in cached_parent_overlays {
289 self.precompute_execution_overlay(hash, anchor_hash);
290 }
291 }
292 }
293
294 #[cfg(feature = "rayon")]
300 pub fn precompute_execution_overlay(&self, tip_hash: B256, anchor_hash: B256) {
301 if tip_hash == anchor_hash {
302 return
303 }
304 let Some(worker_pool) = &self.worker_pool else { return };
305 let manager = self.clone();
306 let parent_span = tracing::Span::current();
307 worker_pool.spawn(move || {
308 let _span = tracing::trace_span!(
309 target: "storage::overlay::manager",
310 parent: parent_span,
311 "precompute_execution_overlay",
312 %tip_hash,
313 %anchor_hash,
314 )
315 .entered();
316 if let Err(err) = manager.precompute_execution_overlay_for_parent(tip_hash, anchor_hash) {
317 debug!(target: "storage::overlay::manager", %err, "Skipping execution overlay precompute");
318 }
319 });
320 }
321
322 #[tracing::instrument(
325 level = "trace",
326 target = "storage::overlay::manager",
327 skip_all,
328 fields(
329 block_count = tracing::field::Empty,
330 removed_blocks = tracing::field::Empty,
331 pruned_overlays = tracing::field::Empty,
332 )
333 )]
334 pub fn remove_blocks(&self, hashes: impl IntoIterator<Item = B256>) {
335 let span = tracing::Span::current();
336
337 let mut block_count = 0usize;
339 let mut removed_blocks = 0usize;
340 let mut pruned_overlays = 0usize;
341 for hash in hashes {
342 block_count += 1;
343 removed_blocks += self.blocks.remove(&hash).is_some() as usize;
344 }
345 span.record("block_count", block_count);
346 span.record("removed_blocks", removed_blocks);
347
348 if removed_blocks > 0 {
349 let overlays_before = self.state_trie_overlays.len() + self.execution_overlays.len();
350 self.state_trie_overlays.retain(|key, _| {
351 self.contains_hash(key.tip_hash, key.anchor_hash, key.anchor_hash)
352 });
353 self.execution_overlays.retain(|key, _| {
354 self.contains_hash(key.tip_hash, key.anchor_hash, key.anchor_hash)
355 });
356 pruned_overlays = overlays_before
357 .saturating_sub(self.state_trie_overlays.len() + self.execution_overlays.len());
358 span.record("pruned_overlays", pruned_overlays);
359 }
360 debug!(
361 target: "storage::overlay::manager",
362 block_count,
363 removed_blocks,
364 pruned_overlays,
365 "removed blocks from state trie overlay manager"
366 );
367 }
368
369 #[tracing::instrument(
371 level = "trace",
372 target = "storage::overlay::manager",
373 skip_all,
374 fields(tip_hash = %parent_state.hash(), anchor_hash = %anchor_hash)
375 )]
376 pub(crate) fn overlay_for_parent(
377 &self,
378 parent_state: &BlockState<N>,
379 anchor_hash: B256,
380 cache_config: OverlayCacheConfig,
381 ) -> Result<(Arc<TrieUpdatesSorted>, Arc<HashedPostStateSorted>), StateTrieOverlayError> {
382 let parent_hash = parent_state.hash();
383 if parent_hash == anchor_hash {
384 return Ok((
385 Arc::new(TrieUpdatesSorted::default()),
386 Arc::new(HashedPostStateSorted::default()),
387 ))
388 }
389 debug!(
390 target: "storage::overlay::manager",
391 tip_hash = %parent_hash,
392 %anchor_hash,
393 "loading state trie overlay for parent"
394 );
395 let input = self
396 .get_or_compute_overlay(
397 &self.state_trie_overlays,
398 &self.metrics,
399 anchor_hash,
400 parent_state,
401 cache_config,
402 |input, span| {
403 self.compute_state_trie_overlay(
404 input,
405 anchor_hash,
406 span,
407 cache_config.write_to_cache,
408 )
409 },
410 )?
411 .expect("required overlay lookup cannot skip an in-progress computation");
412 Ok((Arc::clone(&input.nodes), Arc::clone(&input.state)))
413 }
414
415 #[tracing::instrument(
417 level = "trace",
418 target = "storage::overlay::manager",
419 skip_all,
420 fields(tip_hash = %parent_state.hash(), anchor_hash = %anchor_hash)
421 )]
422 pub(crate) fn execution_overlay_for_block_state(
423 &self,
424 parent_state: &BlockState<N>,
425 anchor_hash: B256,
426 cache_config: OverlayCacheConfig,
427 ) -> Result<Arc<ExecutionOverlay>, StateTrieOverlayError> {
428 Ok(self
429 .execution_overlay_for_parent_inner(parent_state, anchor_hash, cache_config)?
430 .expect("required overlay lookup cannot skip an in-progress computation"))
431 }
432
433 #[cfg(feature = "rayon")]
434 fn precompute_execution_overlay_for_parent(
435 &self,
436 parent_hash: B256,
437 anchor_hash: B256,
438 ) -> Result<(), StateTrieOverlayError> {
439 let parent_state = self
440 .block_state(parent_hash)
441 .ok_or(StateTrieOverlayError { tip_hash: parent_hash, anchor_hash })?;
442 self.execution_overlay_for_parent_inner(
443 &parent_state,
444 anchor_hash,
445 OverlayCacheConfig { precompute: true, write_to_cache: true },
446 )
447 .map(drop)
448 }
449
450 fn execution_overlay_for_parent_inner(
451 &self,
452 parent_state: &BlockState<N>,
453 anchor_hash: B256,
454 cache_config: OverlayCacheConfig,
455 ) -> Result<Option<Arc<ExecutionOverlay>>, StateTrieOverlayError> {
456 let parent_hash = parent_state.hash();
457 if parent_hash == anchor_hash {
458 return Ok(Some(Arc::new(ExecutionOverlay::default())))
459 }
460
461 self.get_or_compute_overlay(
462 &self.execution_overlays,
463 &self.execution_metrics,
464 anchor_hash,
465 parent_state,
466 cache_config,
467 |input, span| {
468 self.compute_execution_overlay(
469 input,
470 anchor_hash,
471 span,
472 cache_config.write_to_cache,
473 )
474 },
475 )
476 }
477
478 #[tracing::instrument(
479 level = "trace",
480 target = "storage::overlay::manager",
481 skip_all,
482 fields(
483 tip_hash = %parent_state.hash(),
484 anchor_hash = %anchor_hash,
485 cache_reused = tracing::field::Empty,
486 block_count = tracing::field::Empty,
487 parent_overlay_reused = tracing::field::Empty,
488 )
489 )]
490 fn get_or_compute_overlay<T, M>(
491 &self,
492 cache: &OverlayCache<T>,
493 metrics: &M,
494 anchor_hash: B256,
495 parent_state: &BlockState<N>,
496 cache_config: OverlayCacheConfig,
497 compute: impl FnOnce(ComputeOverlayInput<N, T>, tracing::Span) -> T,
498 ) -> Result<Option<Arc<T>>, StateTrieOverlayError>
499 where
500 M: OverlayCacheMetrics,
501 {
502 let tip_hash = parent_state.hash();
503 let key = OverlayCacheKey { anchor_hash, tip_hash };
504 let span = tracing::Span::current();
505 if let Some(entry) = cache.entries.get(&key).map(|entry| entry.value().clone()) {
506 metrics.record_cache_reuse();
507 span.record("cache_reused", true);
508 return match entry {
509 OverlayCacheEntry::Ready(input) => Ok(Some(input)),
510 OverlayCacheEntry::Computing(_) if cache_config.precompute => Ok(None),
511 OverlayCacheEntry::Computing(waiter) => Ok(Some(waiter.wait(metrics))),
512 }
513 }
514 span.record("cache_reused", false);
515
516 let mut blocks = Self::blocks_from_parent_state(parent_state, anchor_hash)?;
518 span.record("block_count", blocks.len());
519
520 if !cache_config.write_to_cache {
521 let parent_input = blocks.first().and_then(|block| {
522 let parent_hash = block.recovered_block().parent_hash();
523 (parent_hash != anchor_hash)
524 .then(|| cache.ready(&OverlayCacheKey { anchor_hash, tip_hash: parent_hash }))
525 .flatten()
526 });
527 span.record("parent_overlay_reused", parent_input.is_some());
528 let compute_input = match parent_input {
529 Some(parent_input) => {
530 ComputeOverlayInput::ExtendCached { block: blocks.swap_remove(0), parent_input }
531 }
532 None => ComputeOverlayInput::MergeBlocks(blocks),
533 };
534 return Ok(Some(Arc::new(compute(compute_input, span))))
535 }
536
537 enum CacheAction<T> {
538 Ready(Arc<T>),
539 Wait(Arc<OverlayWaiter<T>>),
540 Compute(Arc<OverlayWaiter<T>>),
541 }
542
543 let parent_hash = parent_state.block_ref().recovered_block().parent_hash();
544 cache.retain(|sibling_key, entry| {
545 sibling_key.tip_hash == tip_hash ||
546 !matches!(entry, OverlayCacheEntry::Ready(_)) ||
547 self.blocks
548 .get(&sibling_key.tip_hash)
549 .is_none_or(|block| block.recovered_block().parent_hash() != parent_hash)
550 });
551
552 let action = match cache.entries.entry(key) {
553 Entry::Occupied(entry) => {
554 let entry = entry.get().clone();
555 metrics.record_cache_reuse();
556 span.record("cache_reused", true);
557 match entry {
558 OverlayCacheEntry::Ready(input) => CacheAction::Ready(input),
559 OverlayCacheEntry::Computing(_) if cache_config.precompute => return Ok(None),
560 OverlayCacheEntry::Computing(waiter) => CacheAction::Wait(waiter),
561 }
562 }
563 Entry::Vacant(entry) => {
564 metrics.record_cache_fill();
565 let waiter = Arc::new(OverlayWaiter::new());
566 entry.insert(OverlayCacheEntry::Computing(Arc::clone(&waiter)));
567 CacheAction::Compute(waiter)
568 }
569 };
570
571 match action {
572 CacheAction::Ready(input) => Ok(Some(input)),
573 CacheAction::Wait(waiter) => Ok(Some(waiter.wait(metrics))),
574 CacheAction::Compute(waiter) => {
575 let parent_input = blocks.first().and_then(|block| {
576 let parent_hash = block.recovered_block().parent_hash();
577 (parent_hash != anchor_hash)
578 .then(|| {
579 cache
580 .take_ready(&OverlayCacheKey { anchor_hash, tip_hash: parent_hash })
581 })
582 .flatten()
583 });
584 span.record("parent_overlay_reused", parent_input.is_some());
585 let compute_input = match parent_input {
586 Some(parent_input) => ComputeOverlayInput::ExtendCached {
587 block: blocks.swap_remove(0),
588 parent_input,
589 },
590 None => ComputeOverlayInput::MergeBlocks(blocks),
591 };
592 let input = Arc::new(compute(compute_input, span));
593 waiter.finish(Arc::clone(&input));
594
595 if let Entry::Occupied(mut entry) = cache.entries.entry(key) {
596 let should_publish = match entry.get() {
599 OverlayCacheEntry::Computing(existing) => Arc::ptr_eq(existing, &waiter),
600 OverlayCacheEntry::Ready(_) => false,
601 };
602 if should_publish {
603 entry.insert(OverlayCacheEntry::Ready(Arc::clone(&input)));
604 }
605 }
606
607 Ok(Some(input))
608 }
609 }
610 }
611
612 fn blocks_from_parent_state(
613 parent_state: &BlockState<N>,
614 anchor_hash: B256,
615 ) -> Result<Vec<ExecutedBlock<N>>, StateTrieOverlayError> {
616 let tip_hash = parent_state.hash();
617 let mut hash = tip_hash;
618 let mut blocks = Vec::new();
619 for state in parent_state.chain() {
620 let block = state.block();
621 if block.recovered_block().hash() != hash {
622 return Err(StateTrieOverlayError { tip_hash, anchor_hash })
623 }
624 hash = block.recovered_block().parent_hash();
625 blocks.push(block);
626 if hash == anchor_hash {
627 return Ok(blocks)
628 }
629 }
630 Err(StateTrieOverlayError { tip_hash, anchor_hash })
631 }
632
633 fn parent_chain(&self, parent_hash: B256) -> impl Iterator<Item = ExecutedBlock<N>> + '_ {
635 let mut hash = parent_hash;
636 std::iter::from_fn(move || {
637 let block = self.blocks.get(&hash)?;
638 hash = block.recovered_block().parent_hash();
639 Some(block.clone())
640 })
641 }
642
643 fn contains_hash(&self, parent_hash: B256, anchor_hash: B256, hash: B256) -> bool {
646 let mut current_hash = parent_hash;
647
648 loop {
649 if current_hash == hash {
650 return true
651 }
652 if current_hash == anchor_hash {
653 return false
654 }
655
656 let Some(block) = self.blocks.get(¤t_hash) else { return false };
657 current_hash = block.recovered_block().parent_hash();
658 }
659 }
660
661 fn compute_state_trie_overlay(
662 &self,
663 compute_input: ComputeOverlayInput<N, TrieInputSorted>,
664 anchor_hash: B256,
665 _span: tracing::Span,
666 write_to_cache: bool,
667 ) -> TrieInputSorted {
668 if !write_to_cache {
669 return compute_overlay(compute_input, anchor_hash, &self.metrics)
670 }
671
672 #[cfg(feature = "rayon")]
673 {
674 if let Some(worker_pool) = &self.worker_pool {
675 let compute_span = _span;
676 let metrics = self.metrics.clone();
677 return worker_pool.spawn_and_wait(move || {
678 let _guard = compute_span.enter();
679 compute_overlay(compute_input, anchor_hash, &metrics)
680 })
681 }
682 }
683
684 compute_overlay(compute_input, anchor_hash, &self.metrics)
685 }
686
687 fn compute_execution_overlay(
688 &self,
689 compute_input: ComputeOverlayInput<N, ExecutionOverlay>,
690 anchor_hash: B256,
691 _span: tracing::Span,
692 write_to_cache: bool,
693 ) -> ExecutionOverlay {
694 if !write_to_cache {
695 return compute_execution_overlay_inner(
696 compute_input,
697 anchor_hash,
698 &self.execution_metrics,
699 )
700 }
701
702 #[cfg(feature = "rayon")]
703 {
704 if let Some(worker_pool) = &self.worker_pool {
705 let compute_span = _span;
706 let metrics = self.execution_metrics.clone();
707 return worker_pool.spawn_and_wait(move || {
708 let _guard = compute_span.enter();
709 compute_execution_overlay_inner(compute_input, anchor_hash, &metrics)
710 })
711 }
712 }
713
714 compute_execution_overlay_inner(compute_input, anchor_hash, &self.execution_metrics)
715 }
716}
717
718#[derive(Clone, Copy, Debug)]
720pub(crate) struct OverlayCacheConfig {
721 pub(crate) precompute: bool,
723 pub(crate) write_to_cache: bool,
725}
726
727impl Default for OverlayCacheConfig {
728 fn default() -> Self {
729 Self { precompute: false, write_to_cache: true }
730 }
731}
732
733#[derive(Debug)]
735pub(crate) struct StateTrieOverlayError {
736 pub(crate) tip_hash: B256,
738 pub(crate) anchor_hash: B256,
740}
741
742impl fmt::Display for StateTrieOverlayError {
743 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
744 write!(
745 f,
746 "state trie overlay for tip {} cannot be anchored to {} with current blocks",
747 self.tip_hash, self.anchor_hash
748 )
749 }
750}
751
752impl std::error::Error for StateTrieOverlayError {}
753
754#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
755struct OverlayCacheKey {
756 anchor_hash: B256,
757 tip_hash: B256,
758}
759
760struct OverlayCache<T> {
761 entries: Arc<DashMap<OverlayCacheKey, OverlayCacheEntry<T>>>,
762}
763
764impl<T> Default for OverlayCache<T> {
765 fn default() -> Self {
766 Self { entries: Default::default() }
767 }
768}
769
770impl<T> Clone for OverlayCache<T> {
771 fn clone(&self) -> Self {
772 Self { entries: Arc::clone(&self.entries) }
773 }
774}
775
776impl<T> OverlayCache<T> {
777 fn len(&self) -> usize {
778 self.entries.len()
779 }
780
781 fn retain(&self, mut keep: impl FnMut(&OverlayCacheKey, &OverlayCacheEntry<T>) -> bool) {
782 self.entries.retain(|key, entry| keep(key, entry))
783 }
784
785 fn ready(&self, key: &OverlayCacheKey) -> Option<Arc<T>> {
787 self.entries.get(key).and_then(|entry| match entry.value() {
788 OverlayCacheEntry::Ready(input) => Some(Arc::clone(input)),
789 OverlayCacheEntry::Computing(_) => None,
790 })
791 }
792
793 fn take_ready(&self, key: &OverlayCacheKey) -> Option<Arc<T>> {
798 let (_, entry) =
799 self.entries.remove_if(key, |_, entry| matches!(entry, OverlayCacheEntry::Ready(_)))?;
800 let OverlayCacheEntry::Ready(input) = entry else { unreachable!() };
801 Some(input)
802 }
803}
804
805enum OverlayCacheEntry<T> {
806 Ready(Arc<T>),
807 Computing(Arc<OverlayWaiter<T>>),
808}
809
810impl<T> Clone for OverlayCacheEntry<T> {
811 fn clone(&self) -> Self {
812 match self {
813 Self::Ready(input) => Self::Ready(Arc::clone(input)),
814 Self::Computing(waiter) => Self::Computing(Arc::clone(waiter)),
815 }
816 }
817}
818
819struct OverlayWaiter<T> {
820 input: OnceLock<Arc<T>>,
821}
822
823impl<T> OverlayWaiter<T> {
824 const fn new() -> Self {
825 Self { input: OnceLock::new() }
826 }
827
828 fn wait(&self, metrics: &impl OverlayCacheMetrics) -> Arc<T> {
829 let start = Instant::now();
830 let input = self.input.wait();
831 metrics.record_wait_duration(start.elapsed());
832 Arc::clone(input)
833 }
834
835 fn finish(&self, computed: Arc<T>) {
836 let _ = self.input.set(computed);
837 }
838}
839
840enum ComputeOverlayInput<N: NodePrimitives, T> {
841 ExtendCached { block: ExecutedBlock<N>, parent_input: Arc<T> },
842 MergeBlocks(Vec<ExecutedBlock<N>>),
843}
844
845#[tracing::instrument(
846 level = "trace",
847 target = "storage::overlay::manager",
848 skip_all,
849 fields(
850 anchor_hash = %anchor_hash,
851 block_count = tracing::field::Empty,
852 parent_overlay = tracing::field::Empty,
853 elapsed_us = tracing::field::Empty,
854 )
855)]
856fn compute_overlay<N: NodePrimitives>(
857 input: ComputeOverlayInput<N, TrieInputSorted>,
858 anchor_hash: B256,
859 metrics: &StateTrieOverlayMetrics,
860) -> TrieInputSorted {
861 let started_at = Instant::now();
862 let block_count = match &input {
863 ComputeOverlayInput::ExtendCached { .. } => 1,
864 ComputeOverlayInput::MergeBlocks(blocks) => blocks.len(),
865 };
866 let parent_overlay = matches!(&input, ComputeOverlayInput::ExtendCached { .. });
867 tracing::Span::current().record("block_count", block_count);
868 tracing::Span::current().record("parent_overlay", parent_overlay);
869
870 let overlay = match input {
871 ComputeOverlayInput::ExtendCached { block, parent_input } => {
872 trace!(
873 target: "storage::overlay::manager",
874 %anchor_hash,
875 head = %block.recovered_block().hash(),
876 "extending cached parent state trie overlay"
877 );
878
879 let mut parent_input = parent_input;
880 extend_overlay(
881 Arc::make_mut(&mut parent_input),
882 block.hashed_state_ref(),
883 block.trie_updates_ref(),
884 );
885 Arc::try_unwrap(parent_input).expect("Arc::make_mut leaves the child overlay unique")
886 }
887 ComputeOverlayInput::MergeBlocks(blocks) => merge_blocks(blocks),
888 };
889
890 let elapsed = started_at.elapsed();
891 metrics.overlay_computation_duration_seconds.record(elapsed.as_secs_f64());
892 tracing::Span::current().record("elapsed_us", elapsed.as_micros() as u64);
893 debug!(
894 target: "storage::overlay::manager",
895 %anchor_hash,
896 block_count,
897 parent_overlay,
898 ?elapsed,
899 "computed state trie overlay"
900 );
901
902 overlay
903}
904
905fn merge_blocks<N: NodePrimitives>(blocks: Vec<ExecutedBlock<N>>) -> TrieInputSorted {
906 let hashed_states = blocks.iter().map(ExecutedBlock::hashed_state).collect::<Vec<_>>();
907
908 #[cfg(feature = "rayon")]
909 let (nodes, state) = rayon::join(
910 || TrieUpdatesSorted::merge_batch(blocks.iter().map(ExecutedBlock::trie_updates)),
911 || HashedPostStateSorted::merge_batch(hashed_states.iter().cloned()),
912 );
913
914 #[cfg(not(feature = "rayon"))]
915 let (nodes, state) = (
916 TrieUpdatesSorted::merge_batch(blocks.iter().map(ExecutedBlock::trie_updates)),
917 HashedPostStateSorted::merge_batch(hashed_states.iter().cloned()),
918 );
919
920 TrieInputSorted::new(nodes, state, Default::default())
921}
922
923fn extend_overlay(
924 overlay: &mut TrieInputSorted,
925 hashed_state: &HashedPostStateSorted,
926 trie_updates: &TrieUpdatesSorted,
927) {
928 #[cfg(feature = "rayon")]
929 {
930 rayon::join(
931 || {
932 if !hashed_state.is_empty() {
933 Arc::make_mut(&mut overlay.state).extend_ref_and_sort(hashed_state);
934 }
935 },
936 || {
937 if !trie_updates.is_empty() {
938 Arc::make_mut(&mut overlay.nodes).extend_ref_and_sort(trie_updates);
939 }
940 },
941 );
942 }
943
944 #[cfg(not(feature = "rayon"))]
945 {
946 if !hashed_state.is_empty() {
947 Arc::make_mut(&mut overlay.state).extend_ref_and_sort(hashed_state);
948 }
949 if !trie_updates.is_empty() {
950 Arc::make_mut(&mut overlay.nodes).extend_ref_and_sort(trie_updates);
951 }
952 }
953}
954
955fn compute_execution_overlay_inner<N: NodePrimitives>(
956 input: ComputeOverlayInput<N, ExecutionOverlay>,
957 anchor_hash: B256,
958 metrics: &ExecutionOverlayMetrics,
959) -> ExecutionOverlay {
960 let started_at = Instant::now();
961 let block_count = match &input {
962 ComputeOverlayInput::ExtendCached { .. } => 1,
963 ComputeOverlayInput::MergeBlocks(blocks) => blocks.len(),
964 };
965 let parent_overlay = matches!(&input, ComputeOverlayInput::ExtendCached { .. });
966 tracing::Span::current().record("block_count", block_count);
967 tracing::Span::current().record("parent_overlay", parent_overlay);
968
969 let overlay = match input {
970 ComputeOverlayInput::ExtendCached { block, parent_input } => {
971 let mut parent_input = parent_input;
972 Arc::make_mut(&mut parent_input).extend_block(&block);
973 Arc::try_unwrap(parent_input).expect("Arc::make_mut leaves the child overlay unique")
974 }
975 ComputeOverlayInput::MergeBlocks(blocks) => {
976 let mut overlay = ExecutionOverlay::default();
977 for block in blocks.iter().rev() {
978 overlay.extend_block(block);
979 }
980 overlay
981 }
982 };
983
984 let elapsed = started_at.elapsed();
985 metrics.overlay_computation_duration_seconds.record(elapsed.as_secs_f64());
986 tracing::Span::current().record("elapsed_us", elapsed.as_micros() as u64);
987 debug!(
988 target: "storage::overlay::manager",
989 %anchor_hash,
990 block_count,
991 parent_overlay,
992 ?elapsed,
993 "computed execution overlay"
994 );
995
996 overlay
997}
998
999#[cfg(test)]
1000mod tests {
1001 use super::*;
1002 use alloy_primitives::{map::HashMap, Address, U256};
1003 use reth_chain_state::{test_utils::TestBlockBuilder, ExecutedBlock, SparseTrie};
1004 use reth_ethereum_primitives::EthPrimitives;
1005 use reth_primitives_traits::Account;
1006 #[cfg(feature = "rayon")]
1007 use reth_tasks::WorkerPool;
1008 use reth_trie::{updates::TrieUpdatesSorted, HashedPostState, HashedStorage};
1009 use revm::{
1010 bytecode::Bytecode,
1011 database::BundleState,
1012 state::{AccountId, AccountInfo},
1013 };
1014 use std::{
1015 sync::{mpsc, Arc},
1016 thread,
1017 time::Duration,
1018 };
1019
1020 fn with_unique_state(
1021 block: &ExecutedBlock<EthPrimitives>,
1022 id: u8,
1023 ) -> ExecutedBlock<EthPrimitives> {
1024 let hashed_address = B256::with_last_byte(id);
1025 let hashed_slot = B256::with_last_byte(id.saturating_add(32));
1026 let hashed_state = HashedPostState::default()
1027 .with_accounts([(hashed_address, Some(Account::default()))])
1028 .with_storages([(
1029 hashed_address,
1030 HashedStorage::from_iter([(hashed_slot, U256::from(id))]),
1031 )])
1032 .into_sorted();
1033 let address = Address::with_last_byte(id);
1034 let slot = U256::from(id);
1035 let code_hash = B256::with_last_byte(id.saturating_add(64));
1036 let state = BundleState::builder(block.block_number()..=block.block_number())
1037 .state_present_account_info(
1038 address,
1039 AccountInfo {
1040 nonce: id as u64,
1041 balance: U256::from(id),
1042 account_id: AccountId::new(id as usize),
1043 ..Default::default()
1044 },
1045 )
1046 .state_storage(address, HashMap::from_iter([(slot, (U256::ZERO, U256::from(id)))]))
1047 .contract(code_hash, Bytecode::new_raw(vec![id].into()))
1048 .build();
1049 let mut execution_output = (*block.execution_output).clone();
1050 execution_output.state = state;
1051
1052 ExecutedBlock::new(
1053 Arc::clone(&block.recovered_block),
1054 Arc::new(execution_output),
1055 Arc::new(hashed_state),
1056 Arc::new(TrieUpdatesSorted::default()),
1057 )
1058 }
1059
1060 fn test_blocks() -> Vec<ExecutedBlock<EthPrimitives>> {
1061 TestBlockBuilder::eth()
1062 .get_executed_blocks(1..4)
1063 .enumerate()
1064 .map(|(index, block)| with_unique_state(&block, index as u8 + 1))
1065 .collect()
1066 }
1067
1068 impl OverlayManager {
1069 fn execution_overlay_for_parent(
1070 &self,
1071 parent_hash: B256,
1072 anchor_hash: B256,
1073 ) -> Result<Arc<ExecutionOverlay>, StateTrieOverlayError> {
1074 if parent_hash == anchor_hash {
1075 return Ok(Arc::new(ExecutionOverlay::default()))
1076 }
1077 let parent_state = self
1078 .block_state(parent_hash)
1079 .ok_or(StateTrieOverlayError { tip_hash: parent_hash, anchor_hash })?;
1080 self.execution_overlay_for_block_state(
1081 &parent_state,
1082 anchor_hash,
1083 OverlayCacheConfig::default(),
1084 )
1085 }
1086 }
1087
1088 fn overlay_for_parent(
1089 manager: &OverlayManager,
1090 parent_hash: B256,
1091 anchor_hash: B256,
1092 ) -> Result<(Arc<TrieUpdatesSorted>, Arc<HashedPostStateSorted>), StateTrieOverlayError> {
1093 let parent_state = manager
1094 .block_state(parent_hash)
1095 .ok_or(StateTrieOverlayError { tip_hash: parent_hash, anchor_hash })?;
1096 manager.overlay_for_parent(&parent_state, anchor_hash, OverlayCacheConfig::default())
1097 }
1098
1099 #[test]
1100 fn errors_for_unknown_parent() {
1101 let manager = OverlayManager::<EthPrimitives>::default();
1102 let parent = B256::random();
1103 let anchor = B256::random();
1104
1105 let err = overlay_for_parent(&manager, parent, anchor).unwrap_err();
1106
1107 assert_eq!(err.tip_hash, parent);
1108 assert_eq!(err.anchor_hash, anchor);
1109 }
1110
1111 #[test]
1112 fn builds_managed_overlay_for_inserted_blocks() {
1113 let manager = OverlayManager::default();
1114 let blocks = test_blocks();
1115 for block in &blocks {
1116 manager.insert_block(block.clone());
1117 }
1118
1119 let anchor_hash = blocks[0].recovered_block().parent_hash();
1120
1121 let (_, state) =
1122 overlay_for_parent(&manager, blocks[2].recovered_block().hash(), anchor_hash).unwrap();
1123 assert_eq!(state.accounts.len(), 3);
1124
1125 let short_anchor = blocks[1].recovered_block().hash();
1126 let (_, short) =
1127 overlay_for_parent(&manager, blocks[2].recovered_block().hash(), short_anchor).unwrap();
1128 assert_eq!(short.accounts.len(), 1);
1129 let (_, cached_short) =
1130 overlay_for_parent(&manager, blocks[2].recovered_block().hash(), short_anchor).unwrap();
1131 assert!(Arc::ptr_eq(&short, &cached_short));
1132 }
1133
1134 #[test]
1135 fn builds_execution_overlay_for_inserted_blocks() {
1136 let manager = OverlayManager::default();
1137 let blocks = test_blocks();
1138 for block in &blocks {
1139 manager.insert_block(block.clone());
1140 }
1141
1142 let anchor_hash = blocks[0].recovered_block().parent_hash();
1143 let overlay = manager
1144 .execution_overlay_for_parent(blocks[2].recovered_block().hash(), anchor_hash)
1145 .unwrap();
1146
1147 for id in 1..=3 {
1148 let address = Address::with_last_byte(id);
1149 let code_hash = B256::with_last_byte(id + 64);
1150 assert_eq!(overlay.accounts()[&address].as_ref().unwrap().nonce, id as u64);
1151 assert_eq!(overlay.accounts()[&address].as_ref().unwrap().account_id, None);
1152 assert_eq!(overlay.storage()[&address][&U256::from(id)], U256::from(id));
1153 assert_eq!(overlay.code_hashes()[&code_hash], Bytecode::new_raw(vec![id].into()));
1154 }
1155 assert_eq!(
1156 overlay.block_hashes(),
1157 blocks[..=2].iter().map(|block| block.recovered_block().num_hash()).collect::<Vec<_>>(),
1158 );
1159
1160 let cached = manager
1161 .execution_overlay_for_parent(blocks[2].recovered_block().hash(), anchor_hash)
1162 .unwrap();
1163 assert!(Arc::ptr_eq(&overlay, &cached));
1164
1165 let short_anchor = blocks[1].recovered_block().hash();
1166 let short = manager
1167 .execution_overlay_for_parent(blocks[2].recovered_block().hash(), short_anchor)
1168 .unwrap();
1169 assert_eq!(short.accounts().len(), 1);
1170 }
1171
1172 #[test]
1173 fn computing_sibling_evicts_ready_cached_overlays() {
1174 let manager = OverlayManager::default();
1175 let mut builder = TestBlockBuilder::eth();
1176 let anchor_hash = B256::random();
1177 let parent = builder.get_executed_block_with_number(1, anchor_hash);
1178 let parent_hash = parent.recovered_block().hash();
1179 let first = builder.get_executed_block_with_number(2, parent_hash);
1180 let sibling = builder.get_executed_block_with_number(2, parent_hash);
1181 let sibling_hash = sibling.recovered_block().hash();
1182 let first_key = OverlayCacheKey { anchor_hash, tip_hash: first.recovered_block().hash() };
1183
1184 manager.insert_block(parent);
1185 manager.insert_block(first);
1186 manager.insert_block(sibling);
1187 manager
1188 .state_trie_overlays
1189 .entries
1190 .insert(first_key, OverlayCacheEntry::Ready(Arc::new(TrieInputSorted::default())));
1191 manager
1192 .execution_overlays
1193 .entries
1194 .insert(first_key, OverlayCacheEntry::Ready(Arc::new(ExecutionOverlay::default())));
1195
1196 overlay_for_parent(&manager, sibling_hash, anchor_hash).unwrap();
1197 manager.execution_overlay_for_parent(sibling_hash, anchor_hash).unwrap();
1198
1199 assert!(!manager.state_trie_overlays.entries.contains_key(&first_key));
1200 assert!(!manager.execution_overlays.entries.contains_key(&first_key));
1201 }
1202
1203 #[test]
1204 fn execution_overlay_for_parent_at_anchor_is_empty() {
1205 let manager = OverlayManager::<EthPrimitives>::default();
1206 let anchor_hash = B256::with_last_byte(1);
1207
1208 let overlay = manager.execution_overlay_for_parent(anchor_hash, anchor_hash).unwrap();
1209
1210 assert!(overlay.accounts().is_empty());
1211 assert!(overlay.storage().is_empty());
1212 assert!(overlay.code_hashes().is_empty());
1213 assert!(overlay.block_hashes().is_empty());
1214 }
1215
1216 #[test]
1217 fn promotes_ready_parent_overlays_to_the_child() {
1218 let manager = OverlayManager::default();
1219 let blocks = test_blocks();
1220 for block in &blocks {
1221 manager.insert_block(block.clone());
1222 }
1223
1224 let anchor_hash = blocks[0].recovered_block().parent_hash();
1225 let parent_hash = blocks[1].recovered_block().hash();
1226 let child_hash = blocks[2].recovered_block().hash();
1227 let parent_key = OverlayCacheKey { anchor_hash, tip_hash: parent_hash };
1228 let child_key = OverlayCacheKey { anchor_hash, tip_hash: child_hash };
1229
1230 overlay_for_parent(&manager, parent_hash, anchor_hash).unwrap();
1231 manager.execution_overlay_for_parent(parent_hash, anchor_hash).unwrap();
1232
1233 overlay_for_parent(&manager, child_hash, anchor_hash).unwrap();
1234 manager.execution_overlay_for_parent(child_hash, anchor_hash).unwrap();
1235
1236 assert!(!manager.state_trie_overlays.entries.contains_key(&parent_key));
1237 assert!(manager.state_trie_overlays.entries.contains_key(&child_key));
1238 assert!(!manager.execution_overlays.entries.contains_key(&parent_key));
1239 assert!(manager.execution_overlays.entries.contains_key(&child_key));
1240 }
1241
1242 #[test]
1243 fn promotes_parent_overlays_held_by_callers() {
1244 let manager = OverlayManager::default();
1245 let blocks = test_blocks();
1246 for block in &blocks {
1247 manager.insert_block(block.clone());
1248 }
1249
1250 let anchor_hash = blocks[0].recovered_block().parent_hash();
1251 let parent_hash = blocks[1].recovered_block().hash();
1252 let child_hash = blocks[2].recovered_block().hash();
1253 let parent_key = OverlayCacheKey { anchor_hash, tip_hash: parent_hash };
1254
1255 overlay_for_parent(&manager, parent_hash, anchor_hash).unwrap();
1256 let state_parent = manager
1257 .state_trie_overlays
1258 .entries
1259 .get(&parent_key)
1260 .and_then(|entry| match entry.value() {
1261 OverlayCacheEntry::Ready(input) => Some(Arc::clone(input)),
1262 OverlayCacheEntry::Computing(_) => None,
1263 })
1264 .unwrap();
1265 let execution_parent =
1266 manager.execution_overlay_for_parent(parent_hash, anchor_hash).unwrap();
1267
1268 let (_, child_state) = overlay_for_parent(&manager, child_hash, anchor_hash).unwrap();
1269 let child_execution =
1270 manager.execution_overlay_for_parent(child_hash, anchor_hash).unwrap();
1271
1272 assert!(!manager.state_trie_overlays.entries.contains_key(&parent_key));
1273 assert!(!manager.execution_overlays.entries.contains_key(&parent_key));
1274 assert_eq!(state_parent.state.accounts.len(), 2);
1275 assert_eq!(execution_parent.accounts().len(), 2);
1276 assert_eq!(child_state.accounts.len(), 3);
1277 assert_eq!(child_execution.accounts().len(), 3);
1278 assert!(child_execution
1279 .accounts()
1280 .values()
1281 .flatten()
1282 .all(|account| account.account_id.is_none()));
1283 }
1284
1285 #[test]
1286 fn does_not_cache_or_take_parent_overlays_for_unmanaged_blocks() {
1287 let manager = OverlayManager::default();
1288 let blocks = test_blocks();
1289 for block in &blocks[..2] {
1290 manager.insert_block(block.clone());
1291 }
1292
1293 let anchor_hash = blocks[0].recovered_block().parent_hash();
1294 let parent_hash = blocks[1].recovered_block().hash();
1295 let child_hash = blocks[2].recovered_block().hash();
1296 let parent_key = OverlayCacheKey { anchor_hash, tip_hash: parent_hash };
1297 let child_key = OverlayCacheKey { anchor_hash, tip_hash: child_hash };
1298 let parent_state = manager.block_state(parent_hash).unwrap();
1299 let child_state = BlockState::with_parent(blocks[2].clone(), Some(Arc::new(parent_state)));
1300 let cache_config = OverlayCacheConfig { precompute: false, write_to_cache: false };
1301
1302 overlay_for_parent(&manager, parent_hash, anchor_hash).unwrap();
1303 manager.execution_overlay_for_parent(parent_hash, anchor_hash).unwrap();
1304
1305 let (_, state) =
1306 manager.overlay_for_parent(&child_state, anchor_hash, cache_config).unwrap();
1307 let execution = manager
1308 .execution_overlay_for_block_state(&child_state, anchor_hash, cache_config)
1309 .unwrap();
1310
1311 assert_eq!(state.accounts.len(), 3);
1312 assert_eq!(execution.accounts().len(), 3);
1313 assert!(manager.state_trie_overlays.entries.contains_key(&parent_key));
1314 assert!(!manager.state_trie_overlays.entries.contains_key(&child_key));
1315 assert!(manager.execution_overlays.entries.contains_key(&parent_key));
1316 assert!(!manager.execution_overlays.entries.contains_key(&child_key));
1317 }
1318
1319 #[cfg(feature = "rayon")]
1320 #[test]
1321 fn uncached_overlays_do_not_use_worker_pool() {
1322 let worker_pool = Arc::new(WorkerPool::new(1, "uncached-overlay-test"));
1323 let manager = OverlayManager::new(Arc::clone(&worker_pool));
1324 let block = test_blocks().remove(0);
1325 let anchor_hash = block.recovered_block().parent_hash();
1326 let parent_state = BlockState::new(block);
1327 let cache_config = OverlayCacheConfig { precompute: false, write_to_cache: false };
1328
1329 let (started_tx, started_rx) = mpsc::channel();
1330 let (release_tx, release_rx) = mpsc::channel();
1331 worker_pool.spawn(move || {
1332 started_tx.send(()).unwrap();
1333 release_rx.recv().unwrap();
1334 });
1335 started_rx.recv().unwrap();
1336
1337 let (completed_tx, completed_rx) = mpsc::channel();
1338 let task = thread::spawn(move || {
1339 let execution =
1340 manager.execution_overlay_for_block_state(&parent_state, anchor_hash, cache_config);
1341 let state = manager.overlay_for_parent(&parent_state, anchor_hash, cache_config);
1342 completed_tx.send((execution, state)).unwrap();
1343 });
1344
1345 let completed = completed_rx.recv_timeout(Duration::from_millis(100));
1346 release_tx.send(()).unwrap();
1347 task.join().unwrap();
1348
1349 assert!(completed.is_ok(), "uncached overlay used the worker pool");
1350 let (execution, state) = completed.unwrap();
1351 assert!(execution.is_ok());
1352 assert!(state.is_ok());
1353 }
1354
1355 #[cfg(feature = "rayon")]
1356 #[test]
1357 fn precomputes_execution_overlay_for_cached_parent() {
1358 let manager = OverlayManager::new(Arc::new(WorkerPool::new(1, "execution-overlay-test")));
1359 let blocks = test_blocks();
1360 let anchor_hash = blocks[0].recovered_block().parent_hash();
1361
1362 manager.insert_block(blocks[0].clone());
1363 manager
1364 .execution_overlay_for_parent(blocks[0].recovered_block().hash(), anchor_hash)
1365 .unwrap();
1366
1367 manager.insert_block(blocks[1].clone());
1368 let key = OverlayCacheKey { anchor_hash, tip_hash: blocks[1].recovered_block().hash() };
1369 let deadline = std::time::Instant::now() + Duration::from_secs(1);
1370 while !manager
1371 .execution_overlays
1372 .entries
1373 .get(&key)
1374 .is_some_and(|entry| matches!(entry.value(), OverlayCacheEntry::Ready(_)))
1375 {
1376 assert!(std::time::Instant::now() < deadline, "execution overlay was not precomputed");
1377 thread::sleep(Duration::from_millis(10));
1378 }
1379 assert!(!manager.execution_overlays.entries.contains_key(&OverlayCacheKey {
1380 anchor_hash,
1381 tip_hash: blocks[0].recovered_block().hash(),
1382 }));
1383 }
1384
1385 #[cfg(feature = "rayon")]
1386 #[test]
1387 fn precomputes_execution_overlay_after_anchor_advances() {
1388 let manager = OverlayManager::new(Arc::new(WorkerPool::new(1, "execution-overlay-test")));
1389 let blocks = test_blocks();
1390 for block in &blocks {
1391 manager.insert_block(block.clone());
1392 }
1393 let tip_hash = blocks[2].recovered_block().hash();
1394 let old_anchor = blocks[0].recovered_block().parent_hash();
1395 manager.execution_overlay_for_parent(tip_hash, old_anchor).unwrap();
1396
1397 let new_anchor = blocks[0].recovered_block().hash();
1401 manager.remove_blocks([new_anchor]);
1402 assert!(!manager
1403 .execution_overlays
1404 .entries
1405 .contains_key(&OverlayCacheKey { anchor_hash: old_anchor, tip_hash }));
1406
1407 manager.precompute_execution_overlay(tip_hash, new_anchor);
1408
1409 let key = OverlayCacheKey { anchor_hash: new_anchor, tip_hash };
1410 let deadline = std::time::Instant::now() + Duration::from_secs(1);
1411 let overlay = loop {
1412 if let Some(overlay) = manager.execution_overlays.ready(&key) {
1413 break overlay
1414 }
1415 assert!(std::time::Instant::now() < deadline, "new anchor overlay was not precomputed");
1416 thread::sleep(Duration::from_millis(10));
1417 };
1418 assert_eq!(
1419 overlay.block_hashes(),
1420 blocks[1..].iter().map(|block| block.recovered_block().num_hash()).collect::<Vec<_>>(),
1421 );
1422 assert!(!overlay.accounts().contains_key(&Address::with_last_byte(1)));
1423 for id in 2..=3 {
1424 let address = Address::with_last_byte(id);
1425 assert_eq!(overlay.accounts()[&address].as_ref().unwrap().nonce, id as u64);
1426 assert_eq!(overlay.storage()[&address][&U256::from(id)], U256::from(id));
1427 }
1428 let cached = manager.execution_overlay_for_parent(tip_hash, new_anchor).unwrap();
1429 assert!(Arc::ptr_eq(&overlay, &cached));
1430 }
1431
1432 #[cfg(feature = "rayon")]
1433 #[test]
1434 fn execution_overlay_precompute_does_not_wait_for_pending_entry() {
1435 let worker_pool = Arc::new(WorkerPool::new(1, "execution-overlay-pending-test"));
1436 let manager = OverlayManager::new(Arc::clone(&worker_pool));
1437 let block = test_blocks().remove(0);
1438 let anchor_hash = block.recovered_block().parent_hash();
1439 let tip_hash = block.recovered_block().hash();
1440 manager.insert_block(block);
1441
1442 let waiter = Arc::new(OverlayWaiter::new());
1443 manager.execution_overlays.entries.insert(
1444 OverlayCacheKey { anchor_hash, tip_hash },
1445 OverlayCacheEntry::Computing(Arc::clone(&waiter)),
1446 );
1447
1448 let (tx, rx) = mpsc::channel();
1449 worker_pool.spawn(move || {
1450 manager.precompute_execution_overlay_for_parent(tip_hash, anchor_hash).unwrap();
1451 tx.send(()).unwrap();
1452 });
1453
1454 let completed = rx.recv_timeout(Duration::from_millis(100));
1455 waiter.finish(Arc::new(ExecutionOverlay::default()));
1456 assert!(completed.is_ok(), "execution overlay precompute waited for pending entry");
1457 }
1458
1459 #[test]
1460 fn contains_hash_detects_hashes_from_anchor_to_parent() {
1461 let manager = OverlayManager::default();
1462 let blocks = test_blocks();
1463 for block in &blocks {
1464 manager.insert_block(block.clone());
1465 }
1466
1467 let anchor_hash = blocks[0].recovered_block().parent_hash();
1468 let parent_hash = blocks[2].recovered_block().hash();
1469
1470 assert!(manager.contains_hash(parent_hash, anchor_hash, anchor_hash));
1471 for block in &blocks {
1472 assert!(manager.contains_hash(
1473 parent_hash,
1474 anchor_hash,
1475 block.recovered_block().hash()
1476 ));
1477 }
1478 assert!(!manager.contains_hash(parent_hash, anchor_hash, B256::random()));
1479 }
1480
1481 #[test]
1482 fn contains_hash_rejects_hash_before_anchor() {
1483 let manager = OverlayManager::default();
1484 let blocks = test_blocks();
1485 for block in &blocks {
1486 manager.insert_block(block.clone());
1487 }
1488
1489 let parent_hash = blocks[2].recovered_block().hash();
1490 let anchor_hash = blocks[1].recovered_block().hash();
1491 let before_anchor_hash = blocks[0].recovered_block().hash();
1492
1493 assert!(manager.contains_hash(parent_hash, anchor_hash, parent_hash));
1494 assert!(manager.contains_hash(parent_hash, anchor_hash, anchor_hash));
1495 assert!(!manager.contains_hash(parent_hash, anchor_hash, before_anchor_hash));
1496 }
1497
1498 #[test]
1499 fn contains_hash_rejects_unknown_anchor() {
1500 let manager = OverlayManager::default();
1501 let blocks = test_blocks();
1502 for block in &blocks {
1503 manager.insert_block(block.clone());
1504 }
1505
1506 let parent_hash = blocks[2].recovered_block().hash();
1507 let anchor_hash = B256::random();
1508
1509 assert!(!manager.contains_hash(parent_hash, anchor_hash, anchor_hash));
1510 }
1511
1512 #[test]
1513 fn taking_sparse_trie_removes_it() {
1514 let manager = OverlayManager::<EthPrimitives>::default();
1515 let block_hash = B256::with_last_byte(1);
1516 let other_block_hash = B256::with_last_byte(2);
1517 let anchor_hash = B256::with_last_byte(3);
1518
1519 manager.store_sparse_trie(PreservedSparseTrie::anchored(
1520 SparseTrie::default(),
1521 block_hash,
1522 anchor_hash,
1523 ));
1524
1525 let preserved = manager.take_sparse_trie().expect("preserved trie should be available");
1526 assert_eq!(preserved.block_hash(), block_hash);
1527 assert_eq!(preserved.anchor_hash(), anchor_hash);
1528 assert!(preserved.into_trie_for(other_block_hash).unwrap().is_none());
1529 assert!(manager.take_sparse_trie().is_none());
1530 }
1531
1532 #[test]
1533 fn required_lookup_waits_for_in_progress_overlay() {
1534 let manager = OverlayManager::<EthPrimitives>::default();
1535 let block = test_blocks().remove(0);
1536 let parent_state = BlockState::new(block);
1537 let key = OverlayCacheKey {
1538 anchor_hash: parent_state.block_ref().recovered_block().parent_hash(),
1539 tip_hash: parent_state.hash(),
1540 };
1541 let waiter = Arc::new(OverlayWaiter::new());
1542 manager
1543 .state_trie_overlays
1544 .entries
1545 .insert(key, OverlayCacheEntry::Computing(Arc::clone(&waiter)));
1546
1547 let (tx, rx) = mpsc::channel();
1548 thread::spawn(move || {
1549 let res = manager
1550 .overlay_for_parent(&parent_state, key.anchor_hash, OverlayCacheConfig::default())
1551 .map(|(_, state)| state);
1552 tx.send(res).unwrap();
1553 });
1554
1555 assert!(matches!(
1556 rx.recv_timeout(Duration::from_millis(50)),
1557 Err(mpsc::RecvTimeoutError::Timeout)
1558 ));
1559
1560 waiter.finish(Arc::new(TrieInputSorted::default()));
1561
1562 let state = rx.recv_timeout(Duration::from_secs(1)).unwrap().unwrap();
1563 assert!(state.is_empty());
1564 }
1565
1566 #[test]
1567 fn prunes_cached_overlays_after_removing_blocks() {
1568 let manager = OverlayManager::default();
1569 let blocks = test_blocks();
1570 for block in &blocks {
1571 manager.insert_block(block.clone());
1572 }
1573
1574 let original_anchor = blocks[0].recovered_block().parent_hash();
1575 overlay_for_parent(&manager, blocks[2].recovered_block().hash(), original_anchor).unwrap();
1576 manager
1577 .execution_overlay_for_parent(blocks[2].recovered_block().hash(), original_anchor)
1578 .unwrap();
1579
1580 manager.remove_blocks([
1581 blocks[0].recovered_block().hash(),
1582 blocks[1].recovered_block().hash(),
1583 ]);
1584
1585 let anchor_hash = blocks[1].recovered_block().hash();
1586 assert!(overlay_for_parent(&manager, blocks[2].recovered_block().hash(), original_anchor)
1587 .is_err());
1588 assert!(manager
1589 .execution_overlay_for_parent(blocks[2].recovered_block().hash(), original_anchor)
1590 .is_err());
1591
1592 let (_, state) =
1593 overlay_for_parent(&manager, blocks[2].recovered_block().hash(), anchor_hash).unwrap();
1594 assert_eq!(state.accounts.len(), 1);
1595 let execution = manager
1596 .execution_overlay_for_parent(blocks[2].recovered_block().hash(), anchor_hash)
1597 .unwrap();
1598 assert_eq!(execution.accounts().len(), 1);
1599 }
1600}