1use crate::{
2 db_ext::DbTxPruneExt,
3 segments::{
4 user::history::{finalize_history_prune, HistoryPruneResult},
5 PruneInput, Segment,
6 },
7 PrunerError,
8};
9use alloy_primitives::{Address, BlockNumber, B256};
10use reth_db_api::{
11 models::{storage_sharded_key::StorageShardedKey, BlockNumberAddress},
12 tables,
13 transaction::DbTxMut,
14};
15use reth_provider::{DBProvider, EitherWriter, RocksDBProviderFactory, StaticFileProviderFactory};
16use reth_prune_types::{
17 PruneMode, PrunePurpose, PruneSegment, SegmentOutput, SegmentOutputCheckpoint,
18};
19use reth_static_file_types::StaticFileSegment;
20use reth_storage_api::{StorageChangeSetReader, StorageSettingsCache};
21use rustc_hash::FxHashMap;
22use tracing::{instrument, trace};
23
24const STORAGE_HISTORY_TABLES_TO_PRUNE: usize = 2;
29
30#[derive(Debug)]
31pub struct StorageHistory {
32 mode: PruneMode,
33}
34
35impl StorageHistory {
36 pub const fn new(mode: PruneMode) -> Self {
37 Self { mode }
38 }
39}
40
41impl<Provider> Segment<Provider> for StorageHistory
42where
43 Provider: DBProvider<Tx: DbTxMut>
44 + StaticFileProviderFactory
45 + StorageChangeSetReader
46 + StorageSettingsCache
47 + RocksDBProviderFactory,
48{
49 fn segment(&self) -> PruneSegment {
50 PruneSegment::StorageHistory
51 }
52
53 fn mode(&self) -> Option<PruneMode> {
54 Some(self.mode)
55 }
56
57 fn purpose(&self) -> PrunePurpose {
58 PrunePurpose::User
59 }
60
61 #[instrument(
62 name = "StorageHistory::prune",
63 target = "pruner",
64 skip(self, provider),
65 ret(level = "trace")
66 )]
67 fn prune(&self, provider: &Provider, input: PruneInput) -> Result<SegmentOutput, PrunerError> {
68 let range = match input.get_next_block_range() {
69 Some(range) => range,
70 None => {
71 trace!(target: "pruner", "No storage history to prune");
72 return Ok(SegmentOutput::done())
73 }
74 };
75 let range_end = *range.end();
76
77 if provider.cached_storage_settings().storage_v2 {
79 return self.prune_rocksdb(provider, input, range, range_end);
80 }
81
82 if EitherWriter::storage_changesets_destination(provider).is_static_file() {
84 self.prune_static_files(provider, input, range, range_end)
85 } else {
86 self.prune_database(provider, input, range, range_end)
87 }
88 }
89}
90
91impl StorageHistory {
92 fn prune_static_files<Provider>(
94 &self,
95 provider: &Provider,
96 input: PruneInput,
97 range: std::ops::RangeInclusive<BlockNumber>,
98 range_end: BlockNumber,
99 ) -> Result<SegmentOutput, PrunerError>
100 where
101 Provider: DBProvider<Tx: DbTxMut> + StaticFileProviderFactory,
102 {
103 let mut limiter = if let Some(limit) = input.limiter.deleted_entries_limit() {
104 input.limiter.set_deleted_entries_limit(limit / STORAGE_HISTORY_TABLES_TO_PRUNE)
105 } else {
106 input.limiter
107 };
108
109 if limiter.is_limit_reached() {
112 return Ok(SegmentOutput::not_done(
113 limiter.interrupt_reason(),
114 input.previous_checkpoint.map(SegmentOutputCheckpoint::from_prune_checkpoint),
115 ))
116 }
117
118 let mut highest_deleted_storages = FxHashMap::default();
124 let mut last_changeset_pruned_block = None;
125 let mut pruned_changesets = 0;
126 let mut done = true;
127
128 let walker = provider.static_file_provider().walk_storage_changeset_range(range);
129 for result in walker {
130 let (block_address, entry) = result?;
131 let block_number = block_address.block_number();
132 let address = block_address.address();
133 if limiter.is_limit_reached() &&
138 last_changeset_pruned_block.is_some_and(|last| last != block_number)
139 {
140 done = false;
141 break;
142 }
143 highest_deleted_storages.insert((address, entry.key), block_number);
144 last_changeset_pruned_block = Some(block_number);
145 pruned_changesets += 1;
146 limiter.increment_deleted_entries_count();
147 }
148
149 if done && let Some(last_block) = last_changeset_pruned_block {
151 provider
152 .static_file_provider()
153 .delete_segment_below_block(StaticFileSegment::StorageChangeSets, last_block + 1)?;
154 }
155 trace!(target: "pruner", pruned = %pruned_changesets, %done, "Pruned storage history (changesets from static files)");
156
157 let result = HistoryPruneResult {
158 highest_deleted: highest_deleted_storages,
159 last_pruned_block: last_changeset_pruned_block,
160 pruned_count: pruned_changesets,
161 done,
162 };
163 finalize_history_prune::<_, tables::StoragesHistory, (Address, B256), _>(
164 provider,
165 result,
166 range_end,
167 &limiter,
168 |(address, storage_key), block_number| {
169 StorageShardedKey::new(address, storage_key, block_number)
170 },
171 |a, b| a.address == b.address && a.sharded_key.key == b.sharded_key.key,
172 )
173 .map_err(Into::into)
174 }
175
176 fn prune_database<Provider>(
177 &self,
178 provider: &Provider,
179 input: PruneInput,
180 range: std::ops::RangeInclusive<BlockNumber>,
181 range_end: BlockNumber,
182 ) -> Result<SegmentOutput, PrunerError>
183 where
184 Provider: DBProvider<Tx: DbTxMut>,
185 {
186 let mut limiter = if let Some(limit) = input.limiter.deleted_entries_limit() {
187 input.limiter.set_deleted_entries_limit(limit / STORAGE_HISTORY_TABLES_TO_PRUNE)
188 } else {
189 input.limiter
190 };
191
192 if limiter.is_limit_reached() {
193 return Ok(SegmentOutput::not_done(
194 limiter.interrupt_reason(),
195 input.previous_checkpoint.map(SegmentOutputCheckpoint::from_prune_checkpoint),
196 ))
197 }
198
199 let mut last_changeset_pruned_block = None;
208 let mut highest_deleted_storages = FxHashMap::default();
209 let (pruned_changesets, done) =
210 provider.tx_ref().prune_table_with_range::<tables::StorageChangeSets>(
211 BlockNumberAddress::range(range),
212 &mut limiter,
213 |_| false,
214 |(BlockNumberAddress((block_number, address)), entry)| {
215 highest_deleted_storages.insert((address, entry.key), block_number);
216 last_changeset_pruned_block = Some(block_number);
217 },
218 )?;
219 trace!(target: "pruner", deleted = %pruned_changesets, %done, "Pruned storage history (changesets)");
220
221 let last_pruned_block = last_changeset_pruned_block.map(|block_number| {
224 if done {
225 block_number
226 } else {
227 block_number.saturating_sub(1)
228 }
229 });
230
231 let result = HistoryPruneResult {
232 highest_deleted: highest_deleted_storages,
233 last_pruned_block,
234 pruned_count: pruned_changesets,
235 done,
236 };
237 finalize_history_prune::<_, tables::StoragesHistory, (Address, B256), _>(
238 provider,
239 result,
240 range_end,
241 &limiter,
242 |(address, storage_key), block_number| {
243 StorageShardedKey::new(address, storage_key, block_number)
244 },
245 |a, b| a.address == b.address && a.sharded_key.key == b.sharded_key.key,
246 )
247 .map_err(Into::into)
248 }
249
250 fn prune_rocksdb<Provider>(
255 &self,
256 provider: &Provider,
257 input: PruneInput,
258 range: std::ops::RangeInclusive<BlockNumber>,
259 range_end: BlockNumber,
260 ) -> Result<SegmentOutput, PrunerError>
261 where
262 Provider: DBProvider + StaticFileProviderFactory + RocksDBProviderFactory,
263 {
264 let mut limiter = input.limiter;
265
266 if limiter.is_limit_reached() {
267 return Ok(SegmentOutput::not_done(
268 limiter.interrupt_reason(),
269 input.previous_checkpoint.map(SegmentOutputCheckpoint::from_prune_checkpoint),
270 ))
271 }
272
273 let mut highest_deleted_storages: FxHashMap<_, _> = FxHashMap::default();
274 let mut last_changeset_pruned_block = None;
275 let mut changesets_processed = 0usize;
276 let mut done = true;
277
278 let walker = provider.static_file_provider().walk_storage_changeset_range(range);
282 for result in walker {
283 let (block_address, entry) = result?;
284 let block_number = block_address.block_number();
285 let address = block_address.address();
286 if limiter.is_limit_reached() &&
291 last_changeset_pruned_block.is_some_and(|last| last != block_number)
292 {
293 done = false;
294 break;
295 }
296 highest_deleted_storages.insert((address, entry.key), block_number);
297 last_changeset_pruned_block = Some(block_number);
298 changesets_processed += 1;
299 limiter.increment_deleted_entries_count();
300 }
301
302 trace!(target: "pruner", processed = %changesets_processed, %done, "Scanned storage changesets from static files");
303
304 let last_changeset_pruned_block = last_changeset_pruned_block.unwrap_or(range_end);
305
306 let mut deleted_shards = 0usize;
308 let mut updated_shards = 0usize;
309
310 let mut sorted_storages: Vec<_> = highest_deleted_storages.into_iter().collect();
312 sorted_storages.sort_unstable_by_key(|((addr, key), _)| (*addr, *key));
313
314 provider.with_rocksdb_batch(|mut batch| {
315 let targets: Vec<_> = sorted_storages
316 .iter()
317 .map(|((addr, key), highest)| {
318 ((*addr, *key), (*highest).min(last_changeset_pruned_block))
319 })
320 .collect();
321
322 let outcomes = batch.prune_storage_history_batch(&targets)?;
323 deleted_shards = outcomes.deleted;
324 updated_shards = outcomes.updated;
325
326 Ok(((), Some(batch.into_inner())))
327 })?;
328
329 trace!(target: "pruner", deleted = deleted_shards, updated = updated_shards, %done, "Pruned storage history (RocksDB indices)");
330
331 if done {
337 provider.static_file_provider().delete_segment_below_block(
338 StaticFileSegment::StorageChangeSets,
339 last_changeset_pruned_block + 1,
340 )?;
341 }
342
343 let progress = limiter.progress(done);
344
345 Ok(SegmentOutput {
346 progress,
347 pruned: changesets_processed + deleted_shards + updated_shards,
348 checkpoint: Some(SegmentOutputCheckpoint {
349 block_number: Some(last_changeset_pruned_block),
350 tx_number: None,
351 }),
352 })
353 }
354}
355
356#[cfg(test)]
357mod tests {
358 use super::STORAGE_HISTORY_TABLES_TO_PRUNE;
359 use crate::segments::{PruneInput, PruneLimiter, Segment, SegmentOutput, StorageHistory};
360 use alloy_primitives::{BlockNumber, B256};
361 use assert_matches::assert_matches;
362 use reth_db_api::{models::StorageSettings, tables, BlockNumberList};
363 use reth_provider::{DBProvider, DatabaseProviderFactory, PruneCheckpointReader};
364 use reth_prune_types::{
365 PruneCheckpoint, PruneInterruptReason, PruneMode, PruneProgress, PruneSegment,
366 };
367 use reth_stages::test_utils::{StorageKind, TestStageDB};
368 use reth_storage_api::StorageSettingsCache;
369 use reth_testing_utils::generators::{
370 self, random_block_range, random_changeset_range, random_eoa_accounts, BlockRangeParams,
371 };
372 use std::{collections::BTreeMap, ops::AddAssign};
373
374 #[test]
375 fn prune_legacy() {
376 let db = TestStageDB::default();
377 let mut rng = generators::rng();
378
379 let blocks = random_block_range(
380 &mut rng,
381 0..=5000,
382 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
383 );
384 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
385
386 let accounts = random_eoa_accounts(&mut rng, 2).into_iter().collect::<BTreeMap<_, _>>();
387
388 let (changesets, _) = random_changeset_range(
389 &mut rng,
390 blocks.iter(),
391 accounts.into_iter().map(|(addr, acc)| (addr, (acc, Vec::new()))),
392 1..2,
393 1..2,
394 );
395 db.insert_changesets(changesets.clone(), None).expect("insert changesets");
396 db.insert_history(changesets.clone(), None).expect("insert history");
397
398 let storage_occurrences = db.table::<tables::StoragesHistory>().unwrap().into_iter().fold(
399 BTreeMap::<_, usize>::new(),
400 |mut map, (key, _)| {
401 map.entry((key.address, key.sharded_key.key)).or_default().add_assign(1);
402 map
403 },
404 );
405 assert!(storage_occurrences.into_iter().any(|(_, occurrences)| occurrences > 1));
406
407 assert_eq!(
408 db.table::<tables::StorageChangeSets>().unwrap().len(),
409 changesets.iter().flatten().flat_map(|(_, _, entries)| entries).count()
410 );
411
412 let original_shards = db.table::<tables::StoragesHistory>().unwrap();
413
414 let test_prune = |to_block: BlockNumber,
415 run: usize,
416 expected_result: (PruneProgress, usize)| {
417 let prune_mode = PruneMode::Before(to_block);
418 let deleted_entries_limit = 1000;
419 let mut limiter =
420 PruneLimiter::default().set_deleted_entries_limit(deleted_entries_limit);
421 let input = PruneInput {
422 previous_checkpoint: db
423 .factory
424 .provider()
425 .unwrap()
426 .get_prune_checkpoint(PruneSegment::StorageHistory)
427 .unwrap(),
428 to_block,
429 limiter: limiter.clone(),
430 };
431 let segment = StorageHistory::new(prune_mode);
432
433 let provider = db.factory.database_provider_rw().unwrap();
434 provider.set_storage_settings_cache(StorageSettings::v1());
435 let result = segment.prune(&provider, input).unwrap();
436 limiter.increment_deleted_entries_count_by(result.pruned);
437
438 assert_matches!(
439 result,
440 SegmentOutput {progress, pruned, checkpoint: Some(_)}
441 if (progress, pruned) == expected_result
442 );
443
444 segment
445 .save_checkpoint(
446 &provider,
447 result.checkpoint.unwrap().as_prune_checkpoint(prune_mode),
448 )
449 .unwrap();
450 provider.commit().expect("commit");
451
452 let changesets = changesets
453 .iter()
454 .enumerate()
455 .flat_map(|(block_number, changeset)| {
456 changeset.iter().flat_map(move |(address, _, entries)| {
457 entries.iter().map(move |entry| (block_number, address, entry))
458 })
459 })
460 .collect::<Vec<_>>();
461
462 #[expect(clippy::skip_while_next)]
463 let pruned = changesets
464 .iter()
465 .enumerate()
466 .skip_while(|(i, (block_number, _, _))| {
467 *i < deleted_entries_limit / STORAGE_HISTORY_TABLES_TO_PRUNE * run &&
468 *block_number <= to_block as usize
469 })
470 .next()
471 .map(|(i, _)| i)
472 .unwrap_or_default();
473
474 let mut pruned_changesets = changesets.iter().skip(pruned.saturating_sub(1));
477
478 let last_pruned_block_number = pruned_changesets
479 .next()
480 .map(|(block_number, _, _)| {
481 (if result.progress.is_finished() {
482 *block_number
483 } else {
484 block_number.saturating_sub(1)
485 }) as BlockNumber
486 })
487 .unwrap_or(to_block);
488
489 let pruned_changesets = pruned_changesets.fold(
490 BTreeMap::<_, Vec<_>>::new(),
491 |mut acc, (block_number, address, entry)| {
492 acc.entry((block_number, address)).or_default().push(entry);
493 acc
494 },
495 );
496
497 assert_eq!(
498 db.table::<tables::StorageChangeSets>().unwrap().len(),
499 pruned_changesets.values().flatten().count()
500 );
501
502 let actual_shards = db.table::<tables::StoragesHistory>().unwrap();
503
504 let expected_shards = original_shards
505 .iter()
506 .filter(|(key, _)| key.sharded_key.highest_block_number > last_pruned_block_number)
507 .map(|(key, blocks)| {
508 let new_blocks =
509 blocks.iter().skip_while(|block| *block <= last_pruned_block_number);
510 (key.clone(), BlockNumberList::new_pre_sorted(new_blocks))
511 })
512 .collect::<Vec<_>>();
513
514 assert_eq!(actual_shards, expected_shards);
515
516 assert_eq!(
517 db.factory
518 .provider()
519 .unwrap()
520 .get_prune_checkpoint(PruneSegment::StorageHistory)
521 .unwrap(),
522 Some(PruneCheckpoint {
523 block_number: Some(last_pruned_block_number),
524 tx_number: None,
525 prune_mode
526 })
527 );
528 };
529
530 test_prune(
531 998,
532 1,
533 (PruneProgress::HasMoreData(PruneInterruptReason::DeletedEntriesLimitReached), 500),
534 );
535 test_prune(998, 2, (PruneProgress::Finished, 499));
536 test_prune(1200, 3, (PruneProgress::Finished, 202));
537 }
538
539 #[test]
542 #[allow(clippy::clone_on_copy)]
543 fn prune_partial_progress_mid_block() {
544 use alloy_primitives::{Address, U256};
545 use reth_primitives_traits::Account;
546 use reth_testing_utils::generators::ChangeSet;
547
548 let db = TestStageDB::default();
549 let mut rng = generators::rng();
550
551 let blocks = random_block_range(
553 &mut rng,
554 0..=10,
555 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
556 );
557 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
558
559 let addr1 = Address::with_last_byte(1);
561 let addr2 = Address::with_last_byte(2);
562
563 let account = Account { nonce: 1, balance: U256::from(100), ..Default::default() };
564
565 let storage_entry = |key: u8| reth_primitives_traits::StorageEntry {
567 key: B256::with_last_byte(key),
568 value: U256::from(100),
569 };
570
571 let changesets: Vec<ChangeSet> = vec![
574 vec![(addr1, account.clone(), vec![storage_entry(1)])], vec![(addr1, account.clone(), vec![storage_entry(1)])], vec![(addr1, account.clone(), vec![storage_entry(1)])], vec![(addr1, account.clone(), vec![storage_entry(1)])], vec![(addr1, account.clone(), vec![storage_entry(1)])], vec![
582 (addr1, account.clone(), vec![storage_entry(1), storage_entry(2)]),
583 (addr2, account.clone(), vec![storage_entry(1), storage_entry(2)]),
584 ],
585 vec![(addr1, account, vec![storage_entry(3)])], ];
587
588 db.insert_changesets(changesets.clone(), None).expect("insert changesets");
589 db.insert_history(changesets.clone(), None).expect("insert history");
590
591 let total_storage_entries: usize =
593 changesets.iter().flat_map(|c| c.iter()).map(|(_, _, entries)| entries.len()).sum();
594 assert_eq!(db.table::<tables::StorageChangeSets>().unwrap().len(), total_storage_entries);
595
596 let prune_mode = PruneMode::Before(10);
597
598 let deleted_entries_limit = 14; let limiter = PruneLimiter::default().set_deleted_entries_limit(deleted_entries_limit);
603
604 let input = PruneInput { previous_checkpoint: None, to_block: 10, limiter };
605 let segment = StorageHistory::new(prune_mode);
606
607 let provider = db.factory.database_provider_rw().unwrap();
608 provider.set_storage_settings_cache(StorageSettings::v1());
609 let result = segment.prune(&provider, input).unwrap();
610
611 assert!(!result.progress.is_finished(), "Expected HasMoreData since we stopped mid-block");
613
614 segment
616 .save_checkpoint(&provider, result.checkpoint.unwrap().as_prune_checkpoint(prune_mode))
617 .unwrap();
618 provider.commit().expect("commit");
619
620 let checkpoint = db
622 .factory
623 .provider()
624 .unwrap()
625 .get_prune_checkpoint(PruneSegment::StorageHistory)
626 .unwrap()
627 .expect("checkpoint should exist");
628
629 assert_eq!(
630 checkpoint.block_number,
631 Some(4),
632 "Checkpoint should be block 4 (block before incomplete block 5)"
633 );
634
635 let remaining_changesets = db.table::<tables::StorageChangeSets>().unwrap();
637 assert!(
638 !remaining_changesets.is_empty(),
639 "Should have remaining changesets for blocks 5-6"
640 );
641
642 let history = db.table::<tables::StoragesHistory>().unwrap();
644 for (key, _blocks) in &history {
645 assert!(
646 key.sharded_key.highest_block_number > 4,
647 "Found stale history shard with highest_block_number {} <= checkpoint 4",
648 key.sharded_key.highest_block_number
649 );
650 }
651
652 let input2 = PruneInput {
654 previous_checkpoint: Some(checkpoint),
655 to_block: 10,
656 limiter: PruneLimiter::default().set_deleted_entries_limit(100), };
658
659 let provider2 = db.factory.database_provider_rw().unwrap();
660 provider2.set_storage_settings_cache(StorageSettings::v1());
661 let result2 = segment.prune(&provider2, input2).unwrap();
662
663 assert!(result2.progress.is_finished(), "Second run should complete");
664
665 segment
666 .save_checkpoint(
667 &provider2,
668 result2.checkpoint.unwrap().as_prune_checkpoint(prune_mode),
669 )
670 .unwrap();
671 provider2.commit().expect("commit");
672
673 let final_checkpoint = db
675 .factory
676 .provider()
677 .unwrap()
678 .get_prune_checkpoint(PruneSegment::StorageHistory)
679 .unwrap()
680 .expect("checkpoint should exist");
681
682 assert_eq!(final_checkpoint.block_number, Some(6), "Final checkpoint should be at block 6");
684
685 let final_changesets = db.table::<tables::StorageChangeSets>().unwrap();
687 assert!(final_changesets.is_empty(), "All changesets up to block 10 should be pruned");
688 }
689
690 #[test]
691 fn prune_rocksdb() {
692 use reth_db_api::models::storage_sharded_key::StorageShardedKey;
693 use reth_provider::RocksDBProviderFactory;
694 use reth_storage_api::StorageSettings;
695
696 let db = TestStageDB::default();
697 let mut rng = generators::rng();
698
699 let blocks = random_block_range(
700 &mut rng,
701 0..=100,
702 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
703 );
704 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
705
706 let accounts = random_eoa_accounts(&mut rng, 2).into_iter().collect::<BTreeMap<_, _>>();
707
708 let (changesets, _) = random_changeset_range(
709 &mut rng,
710 blocks.iter(),
711 accounts.into_iter().map(|(addr, acc)| (addr, (acc, Vec::new()))),
712 1..2,
713 1..2,
714 );
715
716 db.insert_changesets_to_static_files(changesets.clone(), None)
717 .expect("insert changesets to static files");
718
719 let mut storage_indices: BTreeMap<(alloy_primitives::Address, B256), Vec<u64>> =
720 BTreeMap::new();
721 for (block, changeset) in changesets.iter().enumerate() {
722 for (address, _, storage_entries) in changeset {
723 for entry in storage_entries {
724 storage_indices.entry((*address, entry.key)).or_default().push(block as u64);
725 }
726 }
727 }
728
729 {
730 let rocksdb = db.factory.rocksdb_provider();
731 let mut batch = rocksdb.batch();
732 for ((address, storage_key), block_numbers) in &storage_indices {
733 let shard = BlockNumberList::new_pre_sorted(block_numbers.clone());
734 batch
735 .put::<tables::StoragesHistory>(
736 StorageShardedKey::last(*address, *storage_key),
737 &shard,
738 )
739 .expect("insert storage history shard");
740 }
741 batch.commit().expect("commit rocksdb batch");
742 }
743
744 {
745 let rocksdb = db.factory.rocksdb_provider();
746 for (address, storage_key) in storage_indices.keys() {
747 let shards = rocksdb.storage_history_shards(*address, *storage_key).unwrap();
748 assert!(!shards.is_empty(), "RocksDB should contain storage history before prune");
749 }
750 }
751
752 let to_block = 50u64;
753 let prune_mode = PruneMode::Before(to_block);
754 let input =
755 PruneInput { previous_checkpoint: None, to_block, limiter: PruneLimiter::default() };
756 let segment = StorageHistory::new(prune_mode);
757
758 let provider = db.factory.database_provider_rw().unwrap();
759 provider.set_storage_settings_cache(StorageSettings::v2());
760 let result = segment.prune(&provider, input).unwrap();
761 provider.commit().expect("commit");
762
763 assert_matches!(
764 result,
765 SegmentOutput { progress: PruneProgress::Finished, checkpoint: Some(_), .. }
766 );
767
768 {
769 let rocksdb = db.factory.rocksdb_provider();
770 for ((address, storage_key), block_numbers) in &storage_indices {
771 let shards = rocksdb.storage_history_shards(*address, *storage_key).unwrap();
772
773 let remaining_blocks: Vec<u64> =
774 block_numbers.iter().copied().filter(|&b| b > to_block).collect();
775
776 if remaining_blocks.is_empty() {
777 assert!(
778 shards.is_empty(),
779 "Shard for {:?}/{:?} should be deleted when all blocks pruned",
780 address,
781 storage_key
782 );
783 } else {
784 assert!(!shards.is_empty(), "Shard should exist with remaining blocks");
785 let actual_blocks: Vec<u64> =
786 shards.iter().flat_map(|(_, list)| list.iter()).collect();
787 assert_eq!(
788 actual_blocks, remaining_blocks,
789 "RocksDB shard should only contain blocks > {}",
790 to_block
791 );
792 }
793 }
794 }
795 }
796
797 #[test]
801 #[allow(clippy::clone_on_copy)]
802 fn dense_block_advances_rocksdb_checkpoint() {
803 use alloy_primitives::U256;
804 use reth_db_api::models::storage_sharded_key::StorageShardedKey;
805 use reth_primitives_traits::StorageEntry;
806 use reth_provider::RocksDBProviderFactory;
807 use reth_storage_api::StorageSettings;
808
809 let db = TestStageDB::default();
810 let mut rng = generators::rng();
811
812 let blocks = random_block_range(
813 &mut rng,
814 0..=20,
815 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
816 );
817 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
818
819 const ENTRIES_PER_BLOCK: usize = 2;
821 let (address, account) = random_eoa_accounts(&mut rng, 1).into_iter().next().unwrap();
822 let keys = [B256::with_last_byte(1), B256::with_last_byte(2)];
823 let changesets = (0..=20)
824 .map(|_| {
825 vec![(
826 address,
827 account.clone(),
828 keys.iter()
829 .map(|key| StorageEntry { key: *key, value: U256::from(1) })
830 .collect(),
831 )]
832 })
833 .collect::<Vec<_>>();
834 db.insert_changesets_to_static_files(changesets, None)
835 .expect("insert changesets to static files");
836
837 {
838 let rocksdb = db.factory.rocksdb_provider();
839 let mut batch = rocksdb.batch();
840 for key in keys {
841 batch
842 .put::<tables::StoragesHistory>(
843 StorageShardedKey::last(address, key),
844 &BlockNumberList::new_pre_sorted(0..=20),
845 )
846 .expect("insert storage history shard");
847 }
848 batch.commit().expect("commit rocksdb batch");
849 }
850
851 let to_block = 15u64;
852 let prune_mode = PruneMode::Before(to_block);
853 let segment = StorageHistory::new(prune_mode);
854
855 let mut checkpoint = PruneCheckpoint { block_number: Some(4), tx_number: None, prune_mode };
857
858 let run_prune = |checkpoint: PruneCheckpoint, limit: usize| {
859 let input = PruneInput {
860 previous_checkpoint: Some(checkpoint),
861 to_block,
862 limiter: PruneLimiter::default().set_deleted_entries_limit(limit),
863 };
864
865 let provider = db.factory.database_provider_rw().unwrap();
866 provider.set_storage_settings_cache(StorageSettings::v2());
867 let result = segment.prune(&provider, input).unwrap();
868 segment
869 .save_checkpoint(
870 &provider,
871 result.checkpoint.unwrap().as_prune_checkpoint(prune_mode),
872 )
873 .unwrap();
874 provider.commit().expect("commit");
875
876 let checkpoint = db
877 .factory
878 .provider()
879 .unwrap()
880 .get_prune_checkpoint(PruneSegment::StorageHistory)
881 .unwrap()
882 .unwrap();
883 (result, checkpoint)
884 };
885
886 for _ in 0..3 {
888 let previous = checkpoint.block_number;
889 let (result, next) = run_prune(checkpoint, ENTRIES_PER_BLOCK);
890 checkpoint = next;
891
892 assert!(
893 !result.progress.is_finished(),
894 "the range is longer than one run's budget allows"
895 );
896 assert!(
897 checkpoint.block_number > previous,
898 "checkpoint must advance past the dense block, got {:?} after {previous:?}",
899 checkpoint.block_number
900 );
901 }
902 assert_eq!(checkpoint.block_number, Some(7), "one dense block cleared per run");
903
904 let (result, checkpoint) = run_prune(checkpoint, 1000);
906 assert!(result.progress.is_finished());
907 assert_eq!(checkpoint.block_number, Some(to_block));
908 }
909}