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