Skip to main content

reth_storage_overlay/
changeset_cache.rs

1//! Trie changeset caching utilities.
2//!
3//! This module provides functionality to compute trie changesets for a given block,
4//! which represent the old trie node values before the block was processed.
5//!
6//! It also provides an efficient in-memory cache for these changesets, which is essential for:
7//! - **Reorg support**: Quickly access changesets to revert blocks during chain reorganizations
8//! - **Memory efficiency**: Explicit eviction releases persisted changesets
9
10use crate::{database_state_frontiers, OverlayManager, OverlayStateProvider};
11use alloy_eips::BlockNumHash;
12use alloy_primitives::{map::B256Map, BlockNumber, B256};
13use parking_lot::RwLock;
14use reth_metrics::{
15    metrics::{Counter, Gauge},
16    Metrics,
17};
18use reth_primitives_traits::{FastInstant as Instant, NodePrimitives};
19use reth_storage_api::{
20    BlockNumReader, ChangeSetReader, DBProvider, PruneCheckpointReader, StageCheckpointReader,
21    StorageChangeSetReader, StorageSettingsCache,
22};
23use reth_storage_errors::provider::{ProviderError, ProviderResult};
24use reth_trie::trie_cursor::{InMemoryTrieCursorFactory, TrieCursor, TrieCursorFactory};
25use reth_trie_common::updates::{StorageTrieUpdatesSorted, TrieUpdatesSorted};
26use reth_trie_db::{DatabaseTrieCursorFactory, TrieTableAdapter};
27use std::{
28    collections::{BTreeMap, HashMap},
29    ops::RangeInclusive,
30    sync::Arc,
31};
32use tracing::{debug, warn};
33
34#[cfg(test)]
35use reth_trie::{changesets::compute_trie_changesets, HashedPostStateSorted, TrieInputSorted};
36#[cfg(test)]
37use reth_trie_db::{DatabaseHashedCursorFactory, DatabaseHashedPostState, DatabaseStateRoot};
38
39/// Computes block trie updates using the changeset cache.
40///
41/// # Algorithm
42///
43/// For block N:
44/// 1. Get cumulative trie reverts from block N+1 to db tip using the cache
45/// 2. Create an overlay cursor factory with these reverts (representing trie state after block N)
46/// 3. Walk through account trie changesets for block N
47/// 4. For each changed path, look up the current value using the overlay cursor
48/// 5. Walk through storage trie changesets for block N
49/// 6. For each changed path, look up the current value using the overlay cursor
50/// 7. Return the collected trie updates
51///
52/// # Arguments
53///
54/// * `provider` - Database provider for accessing changesets and block data
55/// * `block_number` - Block number to compute trie updates for
56///
57/// # Returns
58///
59/// Trie updates representing the state of trie nodes after the block was processed
60///
61/// # Errors
62///
63/// Returns error if:
64/// - Block number exceeds database tip
65/// - Database access fails
66/// - Cache retrieval fails
67pub(crate) fn compute_block_trie_updates<N, Provider>(
68    overlay_manager: &OverlayManager<N>,
69    provider: &Provider,
70    block_number: BlockNumber,
71) -> ProviderResult<TrieUpdatesSorted>
72where
73    N: NodePrimitives,
74    Provider: DBProvider
75        + ChangeSetReader
76        + StorageChangeSetReader
77        + PruneCheckpointReader
78        + StageCheckpointReader
79        + BlockNumReader
80        + StorageSettingsCache,
81{
82    reth_trie_db::with_adapter!(provider, |A| {
83        compute_block_trie_updates_inner::<_, _, A>(overlay_manager, provider, block_number)
84    })
85}
86
87fn compute_block_trie_updates_inner<N, Provider, A>(
88    overlay_manager: &OverlayManager<N>,
89    provider: &Provider,
90    block_number: BlockNumber,
91) -> ProviderResult<TrieUpdatesSorted>
92where
93    N: NodePrimitives,
94    Provider: DBProvider
95        + ChangeSetReader
96        + StorageChangeSetReader
97        + PruneCheckpointReader
98        + StageCheckpointReader
99        + BlockNumReader
100        + StorageSettingsCache,
101    A: TrieTableAdapter,
102{
103    let tx = provider.tx_ref();
104    let cache = overlay_manager.changeset_cache();
105    let (partial_state_trie, finish) = database_state_frontiers(provider)?;
106
107    // Step 1: Get the trie changesets for the target block from cache
108    let changesets = cache.get_or_compute(
109        overlay_manager,
110        provider,
111        block_number,
112        partial_state_trie,
113        finish,
114    )?;
115
116    // Step 2: Get the trie reverts for the state after the target block using the cache
117    let reverts = cache.get_or_compute_range(
118        overlay_manager,
119        provider,
120        (block_number + 1)..=finish.number,
121        partial_state_trie,
122        finish,
123    )?;
124
125    // Step 3: Create an InMemoryTrieCursorFactory with the reverts
126    // This gives us the trie state as it was after the target block was processed
127    let db_cursor_factory = DatabaseTrieCursorFactory::<_, A>::new(tx);
128    let cursor_factory = InMemoryTrieCursorFactory::new(db_cursor_factory, &reverts);
129
130    // Step 4: Collect all account trie nodes that changed in the target block
131    let account_nodes_ref = changesets.account_nodes_ref();
132    let mut account_nodes = Vec::with_capacity(account_nodes_ref.len());
133    let mut account_cursor = cursor_factory.account_trie_cursor()?;
134
135    // Iterate over the account nodes from the changesets
136    for (nibbles, _old_node) in account_nodes_ref {
137        // Look up the current value of this trie node using the overlay cursor
138        let node_value = account_cursor.seek_exact(*nibbles)?.map(|(_, node)| node);
139        account_nodes.push((*nibbles, node_value));
140    }
141
142    // Step 5: Collect all storage trie nodes that changed in the target block
143    let mut storage_tries = B256Map::default();
144
145    // Iterate over the storage tries from the changesets
146    for (hashed_address, storage_changeset) in changesets.storage_tries_ref() {
147        let mut storage_cursor = cursor_factory.storage_trie_cursor(*hashed_address)?;
148        let storage_nodes_ref = storage_changeset.storage_nodes_ref();
149        let mut storage_nodes = Vec::with_capacity(storage_nodes_ref.len());
150
151        // Iterate over the storage nodes for this account
152        for (nibbles, _old_node) in storage_nodes_ref {
153            // Look up the current value of this storage trie node
154            let node_value = storage_cursor.seek_exact(*nibbles)?.map(|(_, node)| node);
155            storage_nodes.push((*nibbles, node_value));
156        }
157
158        storage_tries.insert(
159            *hashed_address,
160            StorageTrieUpdatesSorted { storage_nodes, is_deleted: storage_changeset.is_deleted },
161        );
162    }
163
164    Ok(TrieUpdatesSorted::new(account_nodes, storage_tries))
165}
166
167/// Thread-safe changeset cache.
168///
169/// This type wraps a shared, mutable reference to the cache inner.
170/// The `RwLock` enables concurrent reads while ensuring exclusive access for writes.
171#[derive(Debug, Clone)]
172pub(crate) struct ChangesetCache {
173    inner: Arc<RwLock<ChangesetCacheInner>>,
174}
175
176impl Default for ChangesetCache {
177    fn default() -> Self {
178        Self::new()
179    }
180}
181
182impl ChangesetCache {
183    /// Creates a new cache.
184    ///
185    /// The cache has no capacity limit and relies on explicit eviction
186    /// via the `evict()` method to manage memory usage.
187    pub(crate) fn new() -> Self {
188        Self { inner: Arc::new(RwLock::new(ChangesetCacheInner::new())) }
189    }
190
191    /// Evicts changesets for blocks below the given block number.
192    ///
193    /// This should be called after blocks are persisted to the database to free
194    /// memory for changesets that are no longer needed in the cache.
195    ///
196    /// # Arguments
197    ///
198    /// * `up_to_block` - Evict blocks with number < this value. Blocks with number >= this value
199    ///   are retained.
200    pub(crate) fn evict(&self, up_to_block: BlockNumber) {
201        self.inner.write().evict(up_to_block)
202    }
203
204    /// Gets changesets from cache, or computes them on-the-fly if missing.
205    ///
206    /// This is the primary API for retrieving changesets. It checks the cache first, then falls
207    /// back to computing from database state if missing.
208    ///
209    /// # Arguments
210    ///
211    /// * `block_number` - Block number (for cache insertion and logging)
212    /// * `provider` - Database provider for DB access
213    ///
214    /// # Returns
215    ///
216    /// Changesets for the block, either from cache or computed on-the-fly.
217    pub(crate) fn get_or_compute<N, P>(
218        &self,
219        overlay_manager: &OverlayManager<N>,
220        provider: &P,
221        block_number: BlockNumber,
222        partial_state_trie: BlockNumHash,
223        finish: BlockNumHash,
224    ) -> ProviderResult<Arc<TrieUpdatesSorted>>
225    where
226        N: NodePrimitives,
227        P: DBProvider
228            + ChangeSetReader
229            + StorageChangeSetReader
230            + PruneCheckpointReader
231            + BlockNumReader
232            + StorageSettingsCache,
233    {
234        self.get_or_compute_range(
235            overlay_manager,
236            provider,
237            block_number..=block_number,
238            partial_state_trie,
239            finish,
240        )
241    }
242
243    /// Gets or computes trie reverts for a range of blocks.
244    ///
245    /// If all blocks in the range are cached, this method retrieves and accumulates those
246    /// per-block trie changesets (reverts) in reverse order (newest to oldest), so that older
247    /// values take precedence when there are conflicts.
248    ///
249    /// If any block is missing from cache, this falls back to one aggregate database computation
250    /// for the whole range. The aggregate result restores the trie to the state before the range
251    /// and is inserted into the range cache.
252    ///
253    /// # Arguments
254    ///
255    /// * `provider` - Database provider for DB access and block lookups
256    /// * `range` - Block range to accumulate reverts for (inclusive)
257    ///
258    /// # Returns
259    ///
260    /// Accumulated trie reverts for all blocks in the specified range
261    ///
262    /// # Errors
263    ///
264    /// Returns error if:
265    /// - Any block in the range is beyond the database tip
266    /// - Database access fails
267    /// - Block hash lookup fails
268    /// - Changeset computation fails
269    pub(crate) fn get_or_compute_range<N, P>(
270        &self,
271        overlay_manager: &OverlayManager<N>,
272        provider: &P,
273        range: RangeInclusive<BlockNumber>,
274        partial_state_trie: BlockNumHash,
275        finish: BlockNumHash,
276    ) -> ProviderResult<Arc<TrieUpdatesSorted>>
277    where
278        N: NodePrimitives,
279        P: DBProvider
280            + ChangeSetReader
281            + StorageChangeSetReader
282            + PruneCheckpointReader
283            + BlockNumReader
284            + StorageSettingsCache,
285    {
286        let start_block = *range.start();
287        let end_block = *range.end();
288        let timer = Instant::now();
289
290        if end_block > finish.number {
291            return Err(ProviderError::InsufficientChangesets {
292                requested: end_block,
293                available: 0..=finish.number,
294            });
295        }
296
297        debug!(
298            target: "trie::changeset_cache",
299            start_block,
300            end_block,
301            ?partial_state_trie,
302            ?finish,
303            "Starting get_or_compute_range"
304        );
305
306        if start_block > end_block {
307            debug!(
308                target: "trie::changeset_cache",
309                start_block,
310                end_block,
311                "Empty changeset range requested"
312            );
313            return Ok(Arc::new(TrieUpdatesSorted::default()))
314        }
315
316        let end_block_hash = provider.block_hash(end_block)?.ok_or_else(|| {
317            ProviderError::other(std::io::Error::new(
318                std::io::ErrorKind::NotFound,
319                format!("block hash not found for block number {}", end_block),
320            ))
321        })?;
322        let range_key = ChangesetRangeKey::new(start_block, end_block, end_block_hash);
323
324        if let Some(accumulated_reverts) = self.inner.read().get(&range_key) {
325            let elapsed = timer.elapsed();
326
327            debug!(
328                target: "trie::changeset_cache",
329                ?elapsed,
330                start_block,
331                end_block,
332                ?end_block_hash,
333                num_blocks = end_block.saturating_sub(start_block).saturating_add(1),
334                "Changeset cache HIT for block range"
335            );
336
337            return Ok(accumulated_reverts)
338        }
339
340        let mut cached_reverts =
341            Vec::with_capacity(end_block.saturating_sub(start_block).saturating_add(1) as usize);
342        let mut all_cached = true;
343
344        for block_number in range.rev() {
345            // Get the block hash for this block number
346            let block_hash = if block_number == end_block {
347                end_block_hash
348            } else {
349                provider.block_hash(block_number)?.ok_or_else(|| {
350                    ProviderError::other(std::io::Error::new(
351                        std::io::ErrorKind::NotFound,
352                        format!("block hash not found for block number {}", block_number),
353                    ))
354                })?
355            };
356
357            debug!(
358                target: "trie::changeset_cache",
359                block_number,
360                ?block_hash,
361                "Looked up block hash for block number in range"
362            );
363
364            let block_key = ChangesetRangeKey::single(block_number, block_hash);
365            if let Some(changesets) = self.inner.read().get(&block_key) {
366                cached_reverts.push(changesets);
367            } else {
368                all_cached = false;
369                break
370            }
371        }
372
373        if all_cached {
374            // `merge_slice` gives precedence to earlier items, so pass reverts oldest-to-newest.
375            cached_reverts.reverse();
376            let accumulated_reverts = Arc::new(TrieUpdatesSorted::merge_slice(&cached_reverts));
377            let elapsed = timer.elapsed();
378
379            let num_account_nodes = accumulated_reverts.account_nodes_ref().len();
380            let num_storage_tries = accumulated_reverts.storage_tries_ref().len();
381
382            debug!(
383                target: "trie::changeset_cache",
384                ?elapsed,
385                start_block,
386                end_block,
387                num_blocks = end_block.saturating_sub(start_block).saturating_add(1),
388                num_account_nodes,
389                num_storage_tries,
390                "Finished accumulating cached trie reverts for block range"
391            );
392
393            self.inner.write().insert(range_key, Arc::clone(&accumulated_reverts));
394            return Ok(accumulated_reverts)
395        }
396
397        warn!(
398            target: "trie::changeset_cache",
399            start_block,
400            end_block,
401            "Changeset cache MISS in range, falling back to aggregate DB-based computation"
402        );
403
404        let overlay = overlay_manager
405            .overlay_builder(finish.hash)
406            .with_no_reverts()
407            .build_overlay_at_frontiers(provider, partial_state_trie, finish)?;
408        let state_trie_provider = OverlayStateProvider::new(
409            provider,
410            overlay,
411            provider.cached_storage_settings().is_v2(),
412        );
413
414        let accumulated_reverts = Arc::new(reth_trie_db::compute_range_trie_changesets(
415            provider,
416            &state_trie_provider,
417            start_block..=end_block,
418            finish.number,
419        )?);
420
421        let elapsed = timer.elapsed();
422
423        let num_account_nodes = accumulated_reverts.account_nodes_ref().len();
424        let num_storage_tries = accumulated_reverts.storage_tries_ref().len();
425
426        debug!(
427            target: "trie::changeset_cache",
428            ?elapsed,
429            start_block,
430            end_block,
431            ?end_block_hash,
432            num_blocks = end_block.saturating_sub(start_block).saturating_add(1),
433            num_account_nodes,
434            num_storage_tries,
435            "Finished accumulating trie reverts for block range"
436        );
437
438        self.inner.write().insert(range_key, Arc::clone(&accumulated_reverts));
439
440        Ok(accumulated_reverts)
441    }
442}
443
444/// Cache key for one contiguous range of canonical trie changesets.
445///
446/// The end block hash disambiguates canonical rewrites where the same block numbers later refer to
447/// a different chain. For a single block, `start_block == end_block` and `end_block_hash` is that
448/// block's hash.
449#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
450struct ChangesetRangeKey {
451    start_block: BlockNumber,
452    end_block: BlockNumber,
453    end_block_hash: B256,
454}
455
456impl ChangesetRangeKey {
457    const fn new(start_block: BlockNumber, end_block: BlockNumber, end_block_hash: B256) -> Self {
458        Self { start_block, end_block, end_block_hash }
459    }
460
461    const fn single(block_number: BlockNumber, block_hash: B256) -> Self {
462        Self::new(block_number, block_number, block_hash)
463    }
464}
465
466/// In-memory cache for trie changesets with explicit eviction policy.
467///
468/// Holds changesets for blocks or block ranges that have been validated but not yet persisted.
469/// Keyed by canonical block range. Eviction is controlled
470/// explicitly by the engine API tree handler when persistence completes.
471///
472/// ## Eviction Policy
473///
474/// Unlike traditional caches with automatic eviction, this cache requires explicit
475/// eviction calls. The engine API tree handler calls `evict(block_number)` after
476/// blocks are persisted to the database, ensuring changesets remain available
477/// until their corresponding blocks are safely on disk.
478///
479/// ## Metrics
480///
481/// The cache maintains several metrics for observability:
482/// - `hits`: Number of successful cache lookups
483/// - `misses`: Number of failed cache lookups
484/// - `evictions`: Number of blocks evicted
485/// - `size`: Current number of cached blocks
486#[derive(Debug)]
487struct ChangesetCacheInner {
488    /// Cache entries keyed by inclusive block range plus the range's canonical end hash.
489    entries: HashMap<ChangesetRangeKey, Arc<TrieUpdatesSorted>>,
490
491    /// Range start block to cache keys mapping for eviction.
492    range_starts: BTreeMap<BlockNumber, Vec<ChangesetRangeKey>>,
493
494    /// Metrics for monitoring cache behavior
495    metrics: ChangesetCacheMetrics,
496}
497
498/// Metrics for the changeset cache.
499///
500/// These metrics provide visibility into cache performance and help identify
501/// potential issues like high miss rates.
502#[derive(Metrics, Clone)]
503#[metrics(scope = "trie.changeset_cache")]
504struct ChangesetCacheMetrics {
505    /// Cache hit counter
506    hits: Counter,
507
508    /// Cache miss counter
509    misses: Counter,
510
511    /// Eviction counter
512    evictions: Counter,
513
514    /// Current cache size (number of entries)
515    size: Gauge,
516}
517
518impl Default for ChangesetCacheInner {
519    fn default() -> Self {
520        Self::new()
521    }
522}
523
524impl ChangesetCacheInner {
525    /// Creates a new empty changeset cache.
526    ///
527    /// The cache has no capacity limit and relies on explicit eviction
528    /// via the `evict()` method to manage memory usage.
529    fn new() -> Self {
530        Self { entries: HashMap::new(), range_starts: BTreeMap::new(), metrics: Default::default() }
531    }
532
533    fn get(&self, key: &ChangesetRangeKey) -> Option<Arc<TrieUpdatesSorted>> {
534        match self.entries.get(key) {
535            Some(changesets) => {
536                self.metrics.hits.increment(1);
537                Some(Arc::clone(changesets))
538            }
539            None => {
540                self.metrics.misses.increment(1);
541                None
542            }
543        }
544    }
545
546    fn insert(&mut self, key: ChangesetRangeKey, changesets: Arc<TrieUpdatesSorted>) {
547        debug!(
548            target: "trie::changeset_cache",
549            ?key,
550            cache_size_before = self.entries.len(),
551            "Inserting changeset into cache"
552        );
553
554        let is_new_entry = self.entries.insert(key, changesets).is_none();
555
556        if is_new_entry {
557            self.range_starts.entry(key.start_block).or_default().push(key);
558        }
559
560        // Update size metric
561        self.metrics.size.set(self.entries.len() as f64);
562
563        debug!(
564            target: "trie::changeset_cache",
565            ?key,
566            cache_size_after = self.entries.len(),
567            "Changeset inserted into cache"
568        );
569    }
570
571    fn evict(&mut self, up_to_block: BlockNumber) {
572        debug!(
573            target: "trie::changeset_cache",
574            up_to_block,
575            cache_size_before = self.entries.len(),
576            "Starting cache eviction"
577        );
578
579        // Find all block numbers that should be evicted (< up_to_block)
580        let range_starts_to_evict: Vec<u64> =
581            self.range_starts.range(..up_to_block).map(|(num, _)| *num).collect();
582
583        // Remove entries for each block number below threshold
584        let mut evicted_count = 0;
585
586        for start_block in &range_starts_to_evict {
587            if let Some(keys) = self.range_starts.remove(start_block) {
588                debug!(
589                    target: "trie::changeset_cache",
590                    start_block,
591                    num_ranges = keys.len(),
592                    "Evicting ranges from cache"
593                );
594                for key in keys {
595                    if self.entries.remove(&key).is_some() {
596                        evicted_count += 1;
597                    }
598                }
599            }
600        }
601
602        debug!(
603            target: "trie::changeset_cache",
604            up_to_block,
605            evicted_count,
606            cache_size_after = self.entries.len(),
607            "Finished cache eviction"
608        );
609
610        // Update metrics if we evicted anything
611        if evicted_count > 0 {
612            self.metrics.evictions.increment(evicted_count as u64);
613            self.metrics.size.set(self.entries.len() as f64);
614        }
615    }
616}
617
618#[cfg(test)]
619mod tests {
620    use super::*;
621    use crate::Overlay;
622    use alloy_consensus::Header;
623    use alloy_primitives::{
624        keccak256,
625        map::{B256Map, HashMap},
626        Address, U256,
627    };
628    use reth_db::{
629        models::{AccountBeforeTx, BlockNumberAddress},
630        tables,
631        transaction::DbTxMut,
632    };
633    use reth_primitives_traits::{Account, StorageEntry};
634    use reth_provider::{
635        test_utils::create_test_provider_factory, StaticFileProviderFactory, StaticFileSegment,
636        StaticFileWriter,
637    };
638    use reth_stages_types::{StageCheckpoint, StageId};
639    use reth_storage_api::{StageCheckpointWriter, TrieWriter};
640    use reth_trie::{BranchNodeCompact, Nibbles, StateRoot};
641
642    // Helper function to create empty TrieUpdatesSorted for testing
643    fn create_test_changesets() -> Arc<TrieUpdatesSorted> {
644        Arc::new(TrieUpdatesSorted::new(vec![], B256Map::default()))
645    }
646
647    fn empty_overlay() -> Overlay {
648        Overlay { trie_updates: Arc::default(), hashed_post_state: Arc::default() }
649    }
650
651    fn insert_test_changesets(
652        cache: &mut ChangesetCacheInner,
653        block_hash: B256,
654        block_number: BlockNumber,
655        changesets: Arc<TrieUpdatesSorted>,
656    ) {
657        cache.insert(ChangesetRangeKey::single(block_number, block_hash), changesets);
658    }
659
660    fn get_test_changesets(
661        cache: &ChangesetCacheInner,
662        block_hash: B256,
663        block_number: BlockNumber,
664    ) -> Option<Arc<TrieUpdatesSorted>> {
665        cache.get(&ChangesetRangeKey::single(block_number, block_hash))
666    }
667
668    fn test_account(balance: u64) -> Account {
669        Account { balance: U256::from(balance), ..Default::default() }
670    }
671
672    fn test_storage(slot: u64, value: u64) -> StorageEntry {
673        StorageEntry { key: B256::from(U256::from(slot)), value: U256::from(value) }
674    }
675
676    fn seed_headers(
677        factory: &impl StaticFileProviderFactory<
678            Primitives: reth_primitives_traits::NodePrimitives<BlockHeader = Header>,
679        >,
680        end_block: BlockNumber,
681    ) {
682        let static_file_provider = factory.static_file_provider();
683        let mut header_writer =
684            static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
685        for block_number in 0..=end_block {
686            let header = Header { number: block_number, ..Default::default() };
687            header_writer
688                .append_header(&header, &B256::with_last_byte(block_number as u8))
689                .unwrap();
690        }
691        header_writer.commit().unwrap();
692    }
693
694    fn legacy_compute_range_trie_changesets<Provider>(
695        provider: &Provider,
696        range: RangeInclusive<BlockNumber>,
697    ) -> TrieUpdatesSorted
698    where
699        Provider: DBProvider
700            + ChangeSetReader
701            + StorageChangeSetReader
702            + BlockNumReader
703            + StorageSettingsCache,
704    {
705        let mut accumulated_reverts = TrieUpdatesSorted::default();
706        for block_number in range.rev() {
707            let changesets = legacy_compute_block_trie_changesets(provider, block_number);
708            accumulated_reverts.extend_ref_and_sort(&changesets);
709        }
710        accumulated_reverts
711    }
712
713    fn legacy_compute_block_trie_changesets<Provider>(
714        provider: &Provider,
715        block_number: BlockNumber,
716    ) -> TrieUpdatesSorted
717    where
718        Provider: DBProvider
719            + ChangeSetReader
720            + StorageChangeSetReader
721            + BlockNumReader
722            + StorageSettingsCache,
723    {
724        reth_trie_db::with_adapter!(provider, |A| {
725            legacy_compute_block_trie_changesets_inner::<_, A>(provider, block_number)
726        })
727    }
728
729    fn legacy_compute_block_trie_changesets_inner<Provider, A>(
730        provider: &Provider,
731        block_number: BlockNumber,
732    ) -> TrieUpdatesSorted
733    where
734        Provider: DBProvider
735            + ChangeSetReader
736            + StorageChangeSetReader
737            + BlockNumReader
738            + StorageSettingsCache,
739        A: TrieTableAdapter,
740    {
741        let individual_state_revert =
742            HashedPostStateSorted::from_reverts(provider, block_number..=block_number).unwrap();
743        let cumulative_state_revert =
744            HashedPostStateSorted::from_reverts(provider, (block_number + 1)..).unwrap();
745
746        let mut cumulative_state_revert_prev = cumulative_state_revert.clone();
747        cumulative_state_revert_prev.extend_ref_and_sort(&individual_state_revert);
748
749        type DbStateRoot<'a, TX, A> =
750            StateRoot<DatabaseTrieCursorFactory<&'a TX, A>, DatabaseHashedCursorFactory<&'a TX>>;
751
752        let input_prev = TrieInputSorted::new(
753            Arc::default(),
754            Arc::new(cumulative_state_revert_prev.clone()),
755            cumulative_state_revert_prev.construct_prefix_sets(),
756        );
757        let cumulative_trie_updates_prev =
758            DbStateRoot::<_, A>::overlay_root_from_nodes_with_updates(
759                provider.tx_ref(),
760                input_prev,
761            )
762            .unwrap()
763            .1
764            .into_sorted();
765
766        let input = TrieInputSorted::new(
767            Arc::new(cumulative_trie_updates_prev.clone()),
768            Arc::new(cumulative_state_revert),
769            individual_state_revert.construct_prefix_sets(),
770        );
771        let trie_updates =
772            DbStateRoot::<_, A>::overlay_root_from_nodes_with_updates(provider.tx_ref(), input)
773                .unwrap()
774                .1
775                .into_sorted();
776
777        let db_cursor_factory = DatabaseTrieCursorFactory::<_, A>::new(provider.tx_ref());
778        let overlay_factory =
779            InMemoryTrieCursorFactory::new(db_cursor_factory, &cumulative_trie_updates_prev);
780
781        compute_trie_changesets(&overlay_factory, &trie_updates).unwrap()
782    }
783
784    fn seed_tip_trie_tables<Provider, A>(provider: &Provider)
785    where
786        Provider: DBProvider + TrieWriter,
787        A: TrieTableAdapter,
788    {
789        type DbStateRoot<'a, TX, A> =
790            StateRoot<DatabaseTrieCursorFactory<&'a TX, A>, DatabaseHashedCursorFactory<&'a TX>>;
791
792        let (_, trie_updates) =
793            DbStateRoot::<_, A>::from_tx(provider.tx_ref()).root_with_updates().unwrap();
794        provider.write_trie_updates(trie_updates).unwrap();
795    }
796
797    #[test]
798    fn cached_range_merge_keeps_oldest_revert_values() {
799        let factory = create_test_provider_factory();
800        seed_headers(&factory, 2);
801
802        let provider = factory.provider_rw().unwrap();
803        provider.save_stage_checkpoint(StageId::Finish, StageCheckpoint::new(2)).unwrap();
804
805        let cache = ChangesetCache::new();
806        let path = Nibbles::from_nibbles([0x1, 0x2]);
807        let older_node = BranchNodeCompact::new(0b0001, 0, 0, vec![], None);
808        let newer_node = BranchNodeCompact::new(0b0010, 0, 0, vec![], None);
809
810        {
811            let mut cache = cache.inner.write();
812            insert_test_changesets(
813                &mut cache,
814                B256::with_last_byte(1),
815                1,
816                Arc::new(TrieUpdatesSorted::new(
817                    vec![(path, Some(older_node.clone()))],
818                    B256Map::default(),
819                )),
820            );
821            insert_test_changesets(
822                &mut cache,
823                B256::with_last_byte(2),
824                2,
825                Arc::new(TrieUpdatesSorted::new(
826                    vec![(path, Some(newer_node))],
827                    B256Map::default(),
828                )),
829            );
830        }
831
832        let overlay_manager = OverlayManager::<reth_ethereum_primitives::EthPrimitives>::default();
833        let (partial_state_trie, finish) = database_state_frontiers(&*provider).unwrap();
834        let accumulated = cache
835            .get_or_compute_range(&overlay_manager, &*provider, 1..=2, partial_state_trie, finish)
836            .unwrap();
837        assert_eq!(accumulated.account_nodes_ref(), &[(path, Some(older_node))]);
838    }
839
840    #[test]
841    fn aggregate_range_reverts_to_pre_range_state() {
842        let factory = create_test_provider_factory();
843        seed_headers(&factory, 3);
844
845        let provider = factory.provider_rw().unwrap();
846        let address = Address::with_last_byte(1);
847        let hashed_address = keccak256(address);
848        let slot1 = B256::from(U256::from(1));
849        let slot2 = B256::from(U256::from(2));
850        let account1 = test_account(10);
851        let account2 = test_account(20);
852        let account3 = test_account(30);
853
854        provider.tx_ref().put::<tables::HashedAccounts>(hashed_address, account3).unwrap();
855        provider
856            .tx_ref()
857            .put::<tables::HashedStorages>(
858                hashed_address,
859                StorageEntry { key: keccak256(slot1), value: U256::from(25) },
860            )
861            .unwrap();
862        provider
863            .tx_ref()
864            .put::<tables::HashedStorages>(
865                hashed_address,
866                StorageEntry { key: keccak256(slot2), value: U256::from(20) },
867            )
868            .unwrap();
869
870        provider
871            .tx_ref()
872            .put::<tables::AccountChangeSets>(1, AccountBeforeTx { address, info: None })
873            .unwrap();
874        provider
875            .tx_ref()
876            .put::<tables::AccountChangeSets>(2, AccountBeforeTx { address, info: Some(account1) })
877            .unwrap();
878        provider
879            .tx_ref()
880            .put::<tables::AccountChangeSets>(3, AccountBeforeTx { address, info: Some(account2) })
881            .unwrap();
882
883        provider
884            .tx_ref()
885            .put::<tables::StorageChangeSets>(BlockNumberAddress((1, address)), test_storage(1, 0))
886            .unwrap();
887        provider
888            .tx_ref()
889            .put::<tables::StorageChangeSets>(BlockNumberAddress((1, address)), test_storage(2, 0))
890            .unwrap();
891        provider
892            .tx_ref()
893            .put::<tables::StorageChangeSets>(
894                BlockNumberAddress((2, address)),
895                StorageEntry { key: slot1, value: U256::from(10) },
896            )
897            .unwrap();
898        provider
899            .tx_ref()
900            .put::<tables::StorageChangeSets>(
901                BlockNumberAddress((3, address)),
902                StorageEntry { key: slot1, value: U256::from(15) },
903            )
904            .unwrap();
905
906        provider.save_stage_checkpoint(StageId::Finish, StageCheckpoint::new(3)).unwrap();
907        reth_trie_db::with_adapter!(provider, |A| seed_tip_trie_tables::<_, A>(&*provider));
908
909        let overlay = empty_overlay();
910        let state_trie_provider = OverlayStateProvider::new(
911            &*provider,
912            overlay,
913            provider.cached_storage_settings().is_v2(),
914        );
915        let actual =
916            reth_trie_db::compute_range_trie_changesets(&*provider, &state_trie_provider, 1..=3, 3)
917                .unwrap();
918        let storage_revert = actual
919            .storage_tries_ref()
920            .get(&hashed_address)
921            .expect("created account storage trie should be deleted by range revert");
922        assert!(storage_revert.is_deleted());
923        assert!(storage_revert.storage_nodes_ref().is_empty());
924
925        let cache = ChangesetCache::new();
926        let overlay_manager = OverlayManager::<reth_ethereum_primitives::EthPrimitives>::default();
927        let (partial_state_trie, finish) = database_state_frontiers(&*provider).unwrap();
928        let from_cache_api = cache
929            .get_or_compute_range(&overlay_manager, &*provider, 1..=3, partial_state_trie, finish)
930            .unwrap();
931        assert_eq!(*from_cache_api, actual);
932        assert_eq!(cache.inner.read().entries.len(), 1);
933
934        let block_changesets = cache
935            .get_or_compute(&overlay_manager, &*provider, 2, partial_state_trie, finish)
936            .unwrap();
937        assert_eq!(*block_changesets, legacy_compute_block_trie_changesets(&*provider, 2));
938        assert_eq!(cache.inner.read().entries.len(), 2);
939    }
940
941    #[test]
942    fn aggregate_range_matches_legacy_per_block_merge_with_storage_wipe() {
943        let factory = create_test_provider_factory();
944        seed_headers(&factory, 3);
945
946        let provider = factory.provider_rw().unwrap();
947        let address = Address::with_last_byte(1);
948        let slot1 = B256::from(U256::from(1));
949        let slot2 = B256::from(U256::from(2));
950        let account1 = test_account(10);
951        let account2 = test_account(20);
952
953        provider
954            .tx_ref()
955            .put::<tables::AccountChangeSets>(1, AccountBeforeTx { address, info: None })
956            .unwrap();
957        provider
958            .tx_ref()
959            .put::<tables::AccountChangeSets>(2, AccountBeforeTx { address, info: Some(account1) })
960            .unwrap();
961        provider
962            .tx_ref()
963            .put::<tables::AccountChangeSets>(3, AccountBeforeTx { address, info: Some(account2) })
964            .unwrap();
965
966        provider
967            .tx_ref()
968            .put::<tables::StorageChangeSets>(BlockNumberAddress((1, address)), test_storage(1, 0))
969            .unwrap();
970        provider
971            .tx_ref()
972            .put::<tables::StorageChangeSets>(BlockNumberAddress((1, address)), test_storage(2, 0))
973            .unwrap();
974        provider
975            .tx_ref()
976            .put::<tables::StorageChangeSets>(
977                BlockNumberAddress((2, address)),
978                StorageEntry { key: slot1, value: U256::from(10) },
979            )
980            .unwrap();
981        provider
982            .tx_ref()
983            .put::<tables::StorageChangeSets>(
984                BlockNumberAddress((3, address)),
985                StorageEntry { key: slot1, value: U256::from(15) },
986            )
987            .unwrap();
988        provider
989            .tx_ref()
990            .put::<tables::StorageChangeSets>(
991                BlockNumberAddress((3, address)),
992                StorageEntry { key: slot2, value: U256::from(20) },
993            )
994            .unwrap();
995
996        provider.save_stage_checkpoint(StageId::Finish, StageCheckpoint::new(3)).unwrap();
997        reth_trie_db::with_adapter!(provider, |A| seed_tip_trie_tables::<_, A>(&*provider));
998
999        let expected = legacy_compute_range_trie_changesets(&*provider, 2..=3);
1000        let overlay = empty_overlay();
1001        let state_trie_provider = OverlayStateProvider::new(
1002            &*provider,
1003            overlay,
1004            provider.cached_storage_settings().is_v2(),
1005        );
1006        let actual =
1007            reth_trie_db::compute_range_trie_changesets(&*provider, &state_trie_provider, 2..=3, 3)
1008                .unwrap();
1009        assert_eq!(actual, expected);
1010    }
1011
1012    #[test]
1013    fn test_insert_and_retrieve_single_entry() {
1014        let mut cache = ChangesetCacheInner::new();
1015        let hash = B256::random();
1016        let changesets = create_test_changesets();
1017
1018        insert_test_changesets(&mut cache, hash, 100, Arc::clone(&changesets));
1019
1020        // Should be able to retrieve it
1021        let retrieved = get_test_changesets(&cache, hash, 100);
1022        assert!(retrieved.is_some());
1023        assert_eq!(cache.entries.len(), 1);
1024    }
1025
1026    #[test]
1027    fn test_insert_multiple_entries() {
1028        let mut cache = ChangesetCacheInner::new();
1029
1030        // Insert 10 blocks
1031        let mut hashes = Vec::new();
1032        for i in 0..10 {
1033            let hash = B256::random();
1034            insert_test_changesets(&mut cache, hash, 100 + i, create_test_changesets());
1035            hashes.push((100 + i, hash));
1036        }
1037
1038        // Should be able to retrieve all
1039        assert_eq!(cache.entries.len(), 10);
1040        for (block_number, hash) in hashes {
1041            assert!(get_test_changesets(&cache, hash, block_number).is_some());
1042        }
1043    }
1044
1045    #[test]
1046    fn test_eviction_when_explicitly_called() {
1047        let mut cache = ChangesetCacheInner::new();
1048
1049        // Insert 15 blocks (0-14)
1050        let mut hashes = Vec::new();
1051        for i in 0..15 {
1052            let hash = B256::random();
1053            insert_test_changesets(&mut cache, hash, i, create_test_changesets());
1054            hashes.push((i, hash));
1055        }
1056
1057        // All blocks should be present (no automatic eviction)
1058        assert_eq!(cache.entries.len(), 15);
1059
1060        // Explicitly evict blocks < 4
1061        cache.evict(4);
1062
1063        // Blocks 0-3 should be evicted
1064        assert_eq!(cache.entries.len(), 11); // blocks 4-14 = 11 blocks
1065
1066        // Verify blocks 0-3 are evicted
1067        for i in 0..4 {
1068            assert!(
1069                get_test_changesets(&cache, hashes[i as usize].1, i).is_none(),
1070                "Block {} should be evicted",
1071                i
1072            );
1073        }
1074
1075        // Verify blocks 4-14 are still present
1076        for i in 4..15 {
1077            assert!(
1078                get_test_changesets(&cache, hashes[i as usize].1, i).is_some(),
1079                "Block {} should be present",
1080                i
1081            );
1082        }
1083    }
1084
1085    #[test]
1086    fn test_eviction_with_persistence_watermark() {
1087        let mut cache = ChangesetCacheInner::new();
1088
1089        // Insert blocks 100-165
1090        let mut hashes = HashMap::new();
1091        for i in 100..=165 {
1092            let hash = B256::random();
1093            insert_test_changesets(&mut cache, hash, i, create_test_changesets());
1094            hashes.insert(i, hash);
1095        }
1096
1097        // All blocks should be present (no automatic eviction)
1098        assert_eq!(cache.entries.len(), 66);
1099
1100        // Simulate persistence up to block 164, with 64-block retention window
1101        // Eviction threshold = 164 - 64 = 100
1102        cache.evict(100);
1103
1104        // Blocks 100-165 should remain (66 blocks)
1105        assert_eq!(cache.entries.len(), 66);
1106
1107        // Simulate persistence up to block 165
1108        // Eviction threshold = 165 - 64 = 101
1109        cache.evict(101);
1110
1111        // Blocks 101-165 should remain (65 blocks)
1112        assert_eq!(cache.entries.len(), 65);
1113        assert!(get_test_changesets(&cache, hashes[&100], 100).is_none());
1114        assert!(get_test_changesets(&cache, hashes[&101], 101).is_some());
1115    }
1116
1117    #[test]
1118    fn test_out_of_order_inserts_with_explicit_eviction() {
1119        let mut cache = ChangesetCacheInner::new();
1120
1121        // Insert blocks in random order
1122        let hash_10 = B256::random();
1123        insert_test_changesets(&mut cache, hash_10, 10, create_test_changesets());
1124
1125        let hash_5 = B256::random();
1126        insert_test_changesets(&mut cache, hash_5, 5, create_test_changesets());
1127
1128        let hash_15 = B256::random();
1129        insert_test_changesets(&mut cache, hash_15, 15, create_test_changesets());
1130
1131        let hash_3 = B256::random();
1132        insert_test_changesets(&mut cache, hash_3, 3, create_test_changesets());
1133
1134        // All blocks should be present (no automatic eviction)
1135        assert_eq!(cache.entries.len(), 4);
1136
1137        // Explicitly evict blocks < 5
1138        cache.evict(5);
1139
1140        assert!(get_test_changesets(&cache, hash_3, 3).is_none(), "Block 3 should be evicted");
1141        assert!(get_test_changesets(&cache, hash_5, 5).is_some(), "Block 5 should be present");
1142        assert!(get_test_changesets(&cache, hash_10, 10).is_some(), "Block 10 should be present");
1143        assert!(get_test_changesets(&cache, hash_15, 15).is_some(), "Block 15 should be present");
1144    }
1145
1146    #[test]
1147    fn test_multiple_blocks_same_number() {
1148        let mut cache = ChangesetCacheInner::new();
1149
1150        // Insert multiple blocks with same number (side chains)
1151        let hash_1a = B256::random();
1152        let hash_1b = B256::random();
1153        insert_test_changesets(&mut cache, hash_1a, 100, create_test_changesets());
1154        insert_test_changesets(&mut cache, hash_1b, 100, create_test_changesets());
1155
1156        // Both should be retrievable
1157        assert!(get_test_changesets(&cache, hash_1a, 100).is_some());
1158        assert!(get_test_changesets(&cache, hash_1b, 100).is_some());
1159        assert_eq!(cache.entries.len(), 2);
1160    }
1161
1162    #[test]
1163    fn test_ranges_with_same_numbers_and_different_end_hashes_are_distinct() {
1164        let mut cache = ChangesetCacheInner::new();
1165        let path = Nibbles::from_nibbles_unchecked([0x01]);
1166        let hash_a = B256::with_last_byte(1);
1167        let hash_b = B256::with_last_byte(2);
1168        let key_a = ChangesetRangeKey::new(10, 20, hash_a);
1169        let key_b = ChangesetRangeKey::new(10, 20, hash_b);
1170        let changesets_a = Arc::new(TrieUpdatesSorted::new(
1171            vec![(path, Some(BranchNodeCompact::new(0b0001, 0, 0, vec![], None)))],
1172            B256Map::default(),
1173        ));
1174        let changesets_b = Arc::new(TrieUpdatesSorted::new(
1175            vec![(path, Some(BranchNodeCompact::new(0b0010, 0, 0, vec![], None)))],
1176            B256Map::default(),
1177        ));
1178
1179        cache.insert(key_a, Arc::clone(&changesets_a));
1180        cache.insert(key_b, Arc::clone(&changesets_b));
1181
1182        assert_eq!(cache.entries.len(), 2);
1183        assert_eq!(
1184            cache.get(&key_a).unwrap().account_nodes_ref(),
1185            changesets_a.account_nodes_ref()
1186        );
1187        assert_eq!(
1188            cache.get(&key_b).unwrap().account_nodes_ref(),
1189            changesets_b.account_nodes_ref()
1190        );
1191
1192        cache.evict(11);
1193        assert!(cache.get(&key_a).is_none());
1194        assert!(cache.get(&key_b).is_none());
1195    }
1196
1197    #[test]
1198    fn test_eviction_removes_all_side_chains() {
1199        let mut cache = ChangesetCacheInner::new();
1200
1201        // Insert multiple blocks at same height (side chains)
1202        let hash_10a = B256::random();
1203        let hash_10b = B256::random();
1204        let hash_10c = B256::random();
1205        insert_test_changesets(&mut cache, hash_10a, 10, create_test_changesets());
1206        insert_test_changesets(&mut cache, hash_10b, 10, create_test_changesets());
1207        insert_test_changesets(&mut cache, hash_10c, 10, create_test_changesets());
1208
1209        let hash_20 = B256::random();
1210        insert_test_changesets(&mut cache, hash_20, 20, create_test_changesets());
1211
1212        assert_eq!(cache.entries.len(), 4);
1213
1214        // Evict blocks < 15 - should remove all three side chains at height 10
1215        cache.evict(15);
1216
1217        assert_eq!(cache.entries.len(), 1);
1218        assert!(get_test_changesets(&cache, hash_10a, 10).is_none());
1219        assert!(get_test_changesets(&cache, hash_10b, 10).is_none());
1220        assert!(get_test_changesets(&cache, hash_10c, 10).is_none());
1221        assert!(get_test_changesets(&cache, hash_20, 20).is_some());
1222    }
1223}