Skip to main content

reth_storage_overlay/
manager.rs

1//! State trie and execution overlays for in-memory blocks.
2//!
3//! Payload validation needs a view of the state trie as of an in-memory parent block even when that
4//! parent has not been persisted yet. [`OverlayManager`] tracks those in-memory blocks and builds
5//! reusable state trie and execution overlays on demand.
6
7use 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/// Manages state trie and execution overlays for in-memory blocks.
39///
40/// The manager owns the in-memory block graph, changeset cache, and caches keyed by
41/// `(anchor_hash, tip_hash)`.
42#[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    /// Create a new [`OverlayManager`] backed by the given worker pool.
83    #[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    /// Creates an overlay builder for `parent_hash`.
98    ///
99    /// This rebuilds the in-memory chain ending at `parent_hash` from the manager's block graph.
100    /// Prefer [`Self::overlay_builder_for_state`] whenever the caller already holds the chain.
101    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    /// Creates an overlay builder for an already materialized in-memory chain.
106    ///
107    /// The chain tip is used as the parent hash, so no block graph lookup or chain rebuild is
108    /// performed. The builder only reads `state`, which makes it safe to share the same
109    /// [`BlockState`] with other holders.
110    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    /// Gets or computes cached changesets for an inclusive block range.
128    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    /// Evicts cached changesets for blocks below `up_to_block`.
171    pub fn evict_cached_changesets(&self, up_to_block: BlockNumber) {
172        self.changeset_cache.evict(up_to_block);
173    }
174
175    /// Computes the trie updates produced by `block_number`.
176    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    /// Takes the preserved sparse trie if present.
194    pub fn take_sparse_trie(&self) -> Option<PreservedSparseTrie> {
195        self.preserved_sparse_trie.lock().take()
196    }
197
198    /// Stores a preserved sparse trie for later reuse.
199    pub fn store_sparse_trie(&self, trie: PreservedSparseTrie) {
200        *self.preserved_sparse_trie.lock() = Some(trie);
201    }
202
203    /// Clears any preserved sparse trie state.
204    pub fn clear_sparse_trie(&self) {
205        *self.preserved_sparse_trie.lock() = None;
206    }
207
208    /// Waits until the sparse trie lock becomes available.
209    ///
210    /// This acquires and immediately releases the lock, ensuring that any ongoing operations
211    /// complete before returning. Returns the time spent waiting for the lock.
212    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    /// Inserts an executed in-memory block into the state trie overlay manager.
227    #[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        // First add the block to the live graph; duplicate inserts do not need cache work.
243        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        // Snapshot matching parent overlays before spawning so DashMap iteration guards are
260        // dropped.
261        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        // When a new block is inserted we optimistically and asynchronously flatten an execution
285        // overlay for it
286        #[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    /// Optimistically computes an execution overlay from `anchor_hash` to `tip_hash`.
295    ///
296    /// This returns without waiting for computation. No work is scheduled without a worker pool
297    /// or when the tip is already persisted at the anchor. A concurrent persistence or reorg may
298    /// make the requested range unavailable, in which case the background task skips it.
299    #[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    /// Removes blocks from the live block graph and prunes cached overlays that can no longer be
323    /// built from the remaining blocks.
324    #[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        // Remove blocks first, then prune overlays against the remaining block graph.
338        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    /// Returns the flattened overlay from `anchor_hash` to `parent_hash`.
370    #[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    /// Returns execution data for the in-memory chain from `anchor_hash` to `parent_hash`.
416    #[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        // Resolve the block path and any cached parent overlay before locking the child entry.
517        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                    // The entry may have been pruned while the overlay was computing. Only cache
597                    // the result if the map still points at the waiter installed by this task.
598                    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    /// Returns every in-memory block in the chain whose tip is `parent_hash`.
634    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    /// Returns true if `hash` is in the parent chain segment from `anchor_hash` inclusive to
644    /// `parent_hash` inclusive.
645    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(&current_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/// Controls how an overlay computation interacts with the manager cache.
719#[derive(Clone, Copy, Debug)]
720pub(crate) struct OverlayCacheConfig {
721    /// Whether this is a best-effort cache fill that must not wait for an existing computation.
722    pub(crate) precompute: bool,
723    /// Whether to retain the computed overlay in the manager cache.
724    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/// Error returned when a state trie overlay cannot be built from the manager's current block set.
734#[derive(Debug)]
735pub(crate) struct StateTrieOverlayError {
736    /// Requested in-memory tip hash.
737    pub(crate) tip_hash: B256,
738    /// Requested anchor hash.
739    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    /// Returns a ready entry without removing it from the cache.
786    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    /// Removes and returns a ready entry.
794    ///
795    /// Transferring a parent entry lets `Arc::make_mut` extend it in place when no caller retains
796    /// it. Keeping the cache entry would otherwise guarantee a clone.
797    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        // The state/trie frontier advances to block 1. Block 2 may already be persisted at the
1398        // Finish frontier, but its execution state must remain in the overlay until trie
1399        // persistence.
1400        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}