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(
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#[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 pub(crate) fn new() -> Self {
188 Self { inner: Arc::new(RwLock::new(ChangesetCacheInner::new())) }
189 }
190
191 pub(crate) fn evict(&self, up_to_block: BlockNumber) {
201 self.inner.write().evict(up_to_block)
202 }
203
204 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 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 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 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#[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#[derive(Debug)]
487struct ChangesetCacheInner {
488 entries: HashMap<ChangesetRangeKey, Arc<TrieUpdatesSorted>>,
490
491 range_starts: BTreeMap<BlockNumber, Vec<ChangesetRangeKey>>,
493
494 metrics: ChangesetCacheMetrics,
496}
497
498#[derive(Metrics, Clone)]
503#[metrics(scope = "trie.changeset_cache")]
504struct ChangesetCacheMetrics {
505 hits: Counter,
507
508 misses: Counter,
510
511 evictions: Counter,
513
514 size: Gauge,
516}
517
518impl Default for ChangesetCacheInner {
519 fn default() -> Self {
520 Self::new()
521 }
522}
523
524impl ChangesetCacheInner {
525 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 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 let range_starts_to_evict: Vec<u64> =
581 self.range_starts.range(..up_to_block).map(|(num, _)| *num).collect();
582
583 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 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 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 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 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 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 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 assert_eq!(cache.entries.len(), 15);
1059
1060 cache.evict(4);
1062
1063 assert_eq!(cache.entries.len(), 11); 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 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 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 assert_eq!(cache.entries.len(), 66);
1099
1100 cache.evict(100);
1103
1104 assert_eq!(cache.entries.len(), 66);
1106
1107 cache.evict(101);
1110
1111 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 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 assert_eq!(cache.entries.len(), 4);
1136
1137 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 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 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 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 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}