1use 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
39pub(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 let changesets = cache.get_or_compute(
109 overlay_manager,
110 provider,
111 block_number,
112 partial_state_trie,
113 finish,
114 )?;
115
116 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 let db_cursor_factory = DatabaseTrieCursorFactory::<_, A>::new(tx);
128 let cursor_factory = InMemoryTrieCursorFactory::new(db_cursor_factory, &reverts);
129
130 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 for (nibbles, _old_node) in account_nodes_ref {
137 let node_value = account_cursor.seek_exact(*nibbles)?.map(|(_, node)| node);
139 account_nodes.push((*nibbles, node_value));
140 }
141
142 let mut storage_tries = B256Map::default();
144
145 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 for (nibbles, _old_node) in storage_nodes_ref {
153 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#[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 pub(crate) fn new() -> Self {
185 Self { inner: Arc::new(RwLock::new(ChangesetCacheInner::new())) }
186 }
187
188 pub(crate) fn evict(&self, up_to_block: BlockNumber) {
198 self.inner.write().evict(up_to_block)
199 }
200
201 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 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 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 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#[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#[derive(Debug)]
486struct ChangesetCacheInner {
487 entries: HashMap<ChangesetRangeKey, Arc<TrieUpdatesSorted>>,
489
490 range_starts: BTreeMap<BlockNumber, Vec<ChangesetRangeKey>>,
492
493 metrics: ChangesetCacheMetrics,
495}
496
497#[derive(Metrics, Clone)]
502#[metrics(scope = "trie.changeset_cache")]
503struct ChangesetCacheMetrics {
504 hits: Counter,
506
507 misses: Counter,
509
510 evictions: Counter,
512
513 size: Gauge,
515}
516
517impl Default for ChangesetCacheInner {
518 fn default() -> Self {
519 Self::new()
520 }
521}
522
523impl ChangesetCacheInner {
524 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 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 let range_starts_to_evict: Vec<u64> =
580 self.range_starts.range(..up_to_block).map(|(num, _)| *num).collect();
581
582 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 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 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 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 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 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 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 assert_eq!(cache.entries.len(), 15);
1055
1056 cache.evict(4);
1058
1059 assert_eq!(cache.entries.len(), 11); 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 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 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 assert_eq!(cache.entries.len(), 66);
1095
1096 cache.evict(100);
1099
1100 assert_eq!(cache.entries.len(), 66);
1102
1103 cache.evict(101);
1106
1107 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 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 assert_eq!(cache.entries.len(), 4);
1132
1133 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 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 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 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 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}