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 fn prune_partial_progress_mid_block() {
543 use alloy_primitives::{Address, U256};
544 use reth_primitives_traits::Account;
545 use reth_testing_utils::generators::ChangeSet;
546
547 let db = TestStageDB::default();
548 let mut rng = generators::rng();
549
550 let blocks = random_block_range(
552 &mut rng,
553 0..=10,
554 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
555 );
556 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
557
558 let addr1 = Address::with_last_byte(1);
560 let addr2 = Address::with_last_byte(2);
561
562 let account = Account { nonce: 1, balance: U256::from(100), bytecode_hash: None };
563
564 let storage_entry = |key: u8| reth_primitives_traits::StorageEntry {
566 key: B256::with_last_byte(key),
567 value: U256::from(100),
568 };
569
570 let changesets: Vec<ChangeSet> = vec![
573 vec![(addr1, account, vec![storage_entry(1)])], vec![(addr1, account, vec![storage_entry(1)])], vec![(addr1, account, vec![storage_entry(1)])], vec![(addr1, account, vec![storage_entry(1)])], vec![(addr1, account, vec![storage_entry(1)])], vec![
581 (addr1, account, vec![storage_entry(1), storage_entry(2)]),
582 (addr2, account, vec![storage_entry(1), storage_entry(2)]),
583 ],
584 vec![(addr1, account, vec![storage_entry(3)])], ];
586
587 db.insert_changesets(changesets.clone(), None).expect("insert changesets");
588 db.insert_history(changesets.clone(), None).expect("insert history");
589
590 let total_storage_entries: usize =
592 changesets.iter().flat_map(|c| c.iter()).map(|(_, _, entries)| entries.len()).sum();
593 assert_eq!(db.table::<tables::StorageChangeSets>().unwrap().len(), total_storage_entries);
594
595 let prune_mode = PruneMode::Before(10);
596
597 let deleted_entries_limit = 14; let limiter = PruneLimiter::default().set_deleted_entries_limit(deleted_entries_limit);
602
603 let input = PruneInput { previous_checkpoint: None, to_block: 10, limiter };
604 let segment = StorageHistory::new(prune_mode);
605
606 let provider = db.factory.database_provider_rw().unwrap();
607 provider.set_storage_settings_cache(StorageSettings::v1());
608 let result = segment.prune(&provider, input).unwrap();
609
610 assert!(!result.progress.is_finished(), "Expected HasMoreData since we stopped mid-block");
612
613 segment
615 .save_checkpoint(&provider, result.checkpoint.unwrap().as_prune_checkpoint(prune_mode))
616 .unwrap();
617 provider.commit().expect("commit");
618
619 let checkpoint = db
621 .factory
622 .provider()
623 .unwrap()
624 .get_prune_checkpoint(PruneSegment::StorageHistory)
625 .unwrap()
626 .expect("checkpoint should exist");
627
628 assert_eq!(
629 checkpoint.block_number,
630 Some(4),
631 "Checkpoint should be block 4 (block before incomplete block 5)"
632 );
633
634 let remaining_changesets = db.table::<tables::StorageChangeSets>().unwrap();
636 assert!(
637 !remaining_changesets.is_empty(),
638 "Should have remaining changesets for blocks 5-6"
639 );
640
641 let history = db.table::<tables::StoragesHistory>().unwrap();
643 for (key, _blocks) in &history {
644 assert!(
645 key.sharded_key.highest_block_number > 4,
646 "Found stale history shard with highest_block_number {} <= checkpoint 4",
647 key.sharded_key.highest_block_number
648 );
649 }
650
651 let input2 = PruneInput {
653 previous_checkpoint: Some(checkpoint),
654 to_block: 10,
655 limiter: PruneLimiter::default().set_deleted_entries_limit(100), };
657
658 let provider2 = db.factory.database_provider_rw().unwrap();
659 provider2.set_storage_settings_cache(StorageSettings::v1());
660 let result2 = segment.prune(&provider2, input2).unwrap();
661
662 assert!(result2.progress.is_finished(), "Second run should complete");
663
664 segment
665 .save_checkpoint(
666 &provider2,
667 result2.checkpoint.unwrap().as_prune_checkpoint(prune_mode),
668 )
669 .unwrap();
670 provider2.commit().expect("commit");
671
672 let final_checkpoint = db
674 .factory
675 .provider()
676 .unwrap()
677 .get_prune_checkpoint(PruneSegment::StorageHistory)
678 .unwrap()
679 .expect("checkpoint should exist");
680
681 assert_eq!(final_checkpoint.block_number, Some(6), "Final checkpoint should be at block 6");
683
684 let final_changesets = db.table::<tables::StorageChangeSets>().unwrap();
686 assert!(final_changesets.is_empty(), "All changesets up to block 10 should be pruned");
687 }
688
689 #[test]
690 fn prune_rocksdb() {
691 use reth_db_api::models::storage_sharded_key::StorageShardedKey;
692 use reth_provider::RocksDBProviderFactory;
693 use reth_storage_api::StorageSettings;
694
695 let db = TestStageDB::default();
696 let mut rng = generators::rng();
697
698 let blocks = random_block_range(
699 &mut rng,
700 0..=100,
701 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
702 );
703 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
704
705 let accounts = random_eoa_accounts(&mut rng, 2).into_iter().collect::<BTreeMap<_, _>>();
706
707 let (changesets, _) = random_changeset_range(
708 &mut rng,
709 blocks.iter(),
710 accounts.into_iter().map(|(addr, acc)| (addr, (acc, Vec::new()))),
711 1..2,
712 1..2,
713 );
714
715 db.insert_changesets_to_static_files(changesets.clone(), None)
716 .expect("insert changesets to static files");
717
718 let mut storage_indices: BTreeMap<(alloy_primitives::Address, B256), Vec<u64>> =
719 BTreeMap::new();
720 for (block, changeset) in changesets.iter().enumerate() {
721 for (address, _, storage_entries) in changeset {
722 for entry in storage_entries {
723 storage_indices.entry((*address, entry.key)).or_default().push(block as u64);
724 }
725 }
726 }
727
728 {
729 let rocksdb = db.factory.rocksdb_provider();
730 let mut batch = rocksdb.batch();
731 for ((address, storage_key), block_numbers) in &storage_indices {
732 let shard = BlockNumberList::new_pre_sorted(block_numbers.clone());
733 batch
734 .put::<tables::StoragesHistory>(
735 StorageShardedKey::last(*address, *storage_key),
736 &shard,
737 )
738 .expect("insert storage history shard");
739 }
740 batch.commit().expect("commit rocksdb batch");
741 }
742
743 {
744 let rocksdb = db.factory.rocksdb_provider();
745 for (address, storage_key) in storage_indices.keys() {
746 let shards = rocksdb.storage_history_shards(*address, *storage_key).unwrap();
747 assert!(!shards.is_empty(), "RocksDB should contain storage history before prune");
748 }
749 }
750
751 let to_block = 50u64;
752 let prune_mode = PruneMode::Before(to_block);
753 let input =
754 PruneInput { previous_checkpoint: None, to_block, limiter: PruneLimiter::default() };
755 let segment = StorageHistory::new(prune_mode);
756
757 let provider = db.factory.database_provider_rw().unwrap();
758 provider.set_storage_settings_cache(StorageSettings::v2());
759 let result = segment.prune(&provider, input).unwrap();
760 provider.commit().expect("commit");
761
762 assert_matches!(
763 result,
764 SegmentOutput { progress: PruneProgress::Finished, checkpoint: Some(_), .. }
765 );
766
767 {
768 let rocksdb = db.factory.rocksdb_provider();
769 for ((address, storage_key), block_numbers) in &storage_indices {
770 let shards = rocksdb.storage_history_shards(*address, *storage_key).unwrap();
771
772 let remaining_blocks: Vec<u64> =
773 block_numbers.iter().copied().filter(|&b| b > to_block).collect();
774
775 if remaining_blocks.is_empty() {
776 assert!(
777 shards.is_empty(),
778 "Shard for {:?}/{:?} should be deleted when all blocks pruned",
779 address,
780 storage_key
781 );
782 } else {
783 assert!(!shards.is_empty(), "Shard should exist with remaining blocks");
784 let actual_blocks: Vec<u64> =
785 shards.iter().flat_map(|(_, list)| list.iter()).collect();
786 assert_eq!(
787 actual_blocks, remaining_blocks,
788 "RocksDB shard should only contain blocks > {}",
789 to_block
790 );
791 }
792 }
793 }
794 }
795
796 #[test]
800 fn dense_block_advances_rocksdb_checkpoint() {
801 use alloy_primitives::U256;
802 use reth_db_api::models::storage_sharded_key::StorageShardedKey;
803 use reth_primitives_traits::StorageEntry;
804 use reth_provider::RocksDBProviderFactory;
805 use reth_storage_api::StorageSettings;
806
807 let db = TestStageDB::default();
808 let mut rng = generators::rng();
809
810 let blocks = random_block_range(
811 &mut rng,
812 0..=20,
813 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
814 );
815 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
816
817 const ENTRIES_PER_BLOCK: usize = 2;
819 let (address, account) = random_eoa_accounts(&mut rng, 1).into_iter().next().unwrap();
820 let keys = [B256::with_last_byte(1), B256::with_last_byte(2)];
821 let changesets = (0..=20)
822 .map(|_| {
823 vec![(
824 address,
825 account,
826 keys.iter()
827 .map(|key| StorageEntry { key: *key, value: U256::from(1) })
828 .collect(),
829 )]
830 })
831 .collect::<Vec<_>>();
832 db.insert_changesets_to_static_files(changesets, None)
833 .expect("insert changesets to static files");
834
835 {
836 let rocksdb = db.factory.rocksdb_provider();
837 let mut batch = rocksdb.batch();
838 for key in keys {
839 batch
840 .put::<tables::StoragesHistory>(
841 StorageShardedKey::last(address, key),
842 &BlockNumberList::new_pre_sorted(0..=20),
843 )
844 .expect("insert storage history shard");
845 }
846 batch.commit().expect("commit rocksdb batch");
847 }
848
849 let to_block = 15u64;
850 let prune_mode = PruneMode::Before(to_block);
851 let segment = StorageHistory::new(prune_mode);
852
853 let mut checkpoint = PruneCheckpoint { block_number: Some(4), tx_number: None, prune_mode };
855
856 let run_prune = |checkpoint: PruneCheckpoint, limit: usize| {
857 let input = PruneInput {
858 previous_checkpoint: Some(checkpoint),
859 to_block,
860 limiter: PruneLimiter::default().set_deleted_entries_limit(limit),
861 };
862
863 let provider = db.factory.database_provider_rw().unwrap();
864 provider.set_storage_settings_cache(StorageSettings::v2());
865 let result = segment.prune(&provider, input).unwrap();
866 segment
867 .save_checkpoint(
868 &provider,
869 result.checkpoint.unwrap().as_prune_checkpoint(prune_mode),
870 )
871 .unwrap();
872 provider.commit().expect("commit");
873
874 let checkpoint = db
875 .factory
876 .provider()
877 .unwrap()
878 .get_prune_checkpoint(PruneSegment::StorageHistory)
879 .unwrap()
880 .unwrap();
881 (result, checkpoint)
882 };
883
884 for _ in 0..3 {
886 let previous = checkpoint.block_number;
887 let (result, next) = run_prune(checkpoint, ENTRIES_PER_BLOCK);
888 checkpoint = next;
889
890 assert!(
891 !result.progress.is_finished(),
892 "the range is longer than one run's budget allows"
893 );
894 assert!(
895 checkpoint.block_number > previous,
896 "checkpoint must advance past the dense block, got {:?} after {previous:?}",
897 checkpoint.block_number
898 );
899 }
900 assert_eq!(checkpoint.block_number, Some(7), "one dense block cleared per run");
901
902 let (result, checkpoint) = run_prune(checkpoint, 1000);
904 assert!(result.progress.is_finished());
905 assert_eq!(checkpoint.block_number, Some(to_block));
906 }
907}