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::BlockNumber;
10use reth_db_api::{models::ShardedKey, tables, transaction::DbTxMut};
11use reth_provider::{
12 changeset_walker::StaticFileAccountChangesetWalker, DBProvider, EitherWriter,
13 RocksDBProviderFactory, StaticFileProviderFactory,
14};
15use reth_prune_types::{
16 PruneMode, PrunePurpose, PruneSegment, SegmentOutput, SegmentOutputCheckpoint,
17};
18use reth_static_file_types::StaticFileSegment;
19use reth_storage_api::{ChangeSetReader, StorageSettingsCache};
20use rustc_hash::FxHashMap;
21use tracing::{instrument, trace};
22
23const ACCOUNT_HISTORY_TABLES_TO_PRUNE: usize = 2;
28
29#[derive(Debug)]
30pub struct AccountHistory {
31 mode: PruneMode,
32}
33
34impl AccountHistory {
35 pub const fn new(mode: PruneMode) -> Self {
36 Self { mode }
37 }
38}
39
40impl<Provider> Segment<Provider> for AccountHistory
41where
42 Provider: DBProvider<Tx: DbTxMut>
43 + StaticFileProviderFactory
44 + StorageSettingsCache
45 + ChangeSetReader
46 + RocksDBProviderFactory,
47{
48 fn segment(&self) -> PruneSegment {
49 PruneSegment::AccountHistory
50 }
51
52 fn mode(&self) -> Option<PruneMode> {
53 Some(self.mode)
54 }
55
56 fn purpose(&self) -> PrunePurpose {
57 PrunePurpose::User
58 }
59
60 #[instrument(
61 name = "AccountHistory::prune",
62 target = "pruner",
63 skip(self, provider),
64 ret(level = "trace")
65 )]
66 fn prune(&self, provider: &Provider, input: PruneInput) -> Result<SegmentOutput, PrunerError> {
67 let range = match input.get_next_block_range() {
68 Some(range) => range,
69 None => {
70 trace!(target: "pruner", "No account history to prune");
71 return Ok(SegmentOutput::done())
72 }
73 };
74 let range_end = *range.end();
75
76 if provider.cached_storage_settings().storage_v2 {
78 return self.prune_rocksdb(provider, input, range, range_end);
79 }
80
81 if EitherWriter::account_changesets_destination(provider).is_static_file() {
83 self.prune_static_files(provider, input, range, range_end)
84 } else {
85 self.prune_database(provider, input, range, range_end)
86 }
87 }
88}
89
90impl AccountHistory {
91 fn prune_static_files<Provider>(
93 &self,
94 provider: &Provider,
95 input: PruneInput,
96 range: std::ops::RangeInclusive<BlockNumber>,
97 range_end: BlockNumber,
98 ) -> Result<SegmentOutput, PrunerError>
99 where
100 Provider: DBProvider<Tx: DbTxMut> + StaticFileProviderFactory + ChangeSetReader,
101 {
102 let mut limiter = if let Some(limit) = input.limiter.deleted_entries_limit() {
103 input.limiter.set_deleted_entries_limit(limit / ACCOUNT_HISTORY_TABLES_TO_PRUNE)
104 } else {
105 input.limiter
106 };
107
108 if limiter.is_limit_reached() {
111 return Ok(SegmentOutput::not_done(
112 limiter.interrupt_reason(),
113 input.previous_checkpoint.map(SegmentOutputCheckpoint::from_prune_checkpoint),
114 ))
115 }
116
117 let mut highest_deleted_accounts = FxHashMap::default();
126 let mut last_changeset_pruned_block = None;
127 let mut pruned_changesets = 0;
128 let mut done = true;
129
130 let walker = StaticFileAccountChangesetWalker::new(provider, range);
131 for result in walker {
132 let (block_number, changeset) = result?;
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_accounts.insert(changeset.address, 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::AccountChangeSets, last_block + 1)?;
154 }
155 trace!(target: "pruner", pruned = %pruned_changesets, %done, "Pruned account history (changesets from static files)");
156
157 let result = HistoryPruneResult {
158 highest_deleted: highest_deleted_accounts,
159 last_pruned_block: last_changeset_pruned_block,
160 pruned_count: pruned_changesets,
161 done,
162 };
163 finalize_history_prune::<_, tables::AccountsHistory, _, _>(
164 provider,
165 result,
166 range_end,
167 &limiter,
168 ShardedKey::new,
169 |a, b| a.key == b.key,
170 )
171 .map_err(Into::into)
172 }
173
174 fn prune_database<Provider>(
175 &self,
176 provider: &Provider,
177 input: PruneInput,
178 range: std::ops::RangeInclusive<BlockNumber>,
179 range_end: BlockNumber,
180 ) -> Result<SegmentOutput, PrunerError>
181 where
182 Provider: DBProvider<Tx: DbTxMut>,
183 {
184 let mut limiter = if let Some(limit) = input.limiter.deleted_entries_limit() {
185 input.limiter.set_deleted_entries_limit(limit / ACCOUNT_HISTORY_TABLES_TO_PRUNE)
186 } else {
187 input.limiter
188 };
189
190 if limiter.is_limit_reached() {
191 return Ok(SegmentOutput::not_done(
192 limiter.interrupt_reason(),
193 input.previous_checkpoint.map(SegmentOutputCheckpoint::from_prune_checkpoint),
194 ))
195 }
196
197 let mut last_changeset_pruned_block = None;
206 let mut highest_deleted_accounts = FxHashMap::default();
207 let (pruned_changesets, done) =
208 provider.tx_ref().prune_table_with_range::<tables::AccountChangeSets>(
209 range,
210 &mut limiter,
211 |_| false,
212 |(block_number, account)| {
213 highest_deleted_accounts.insert(account.address, block_number);
214 last_changeset_pruned_block = Some(block_number);
215 },
216 )?;
217 trace!(target: "pruner", pruned = %pruned_changesets, %done, "Pruned account history (changesets from database)");
218
219 let last_pruned_block = last_changeset_pruned_block.map(|block_number| {
222 if done {
223 block_number
224 } else {
225 block_number.saturating_sub(1)
226 }
227 });
228
229 let result = HistoryPruneResult {
230 highest_deleted: highest_deleted_accounts,
231 last_pruned_block,
232 pruned_count: pruned_changesets,
233 done,
234 };
235 finalize_history_prune::<_, tables::AccountsHistory, _, _>(
236 provider,
237 result,
238 range_end,
239 &limiter,
240 ShardedKey::new,
241 |a, b| a.key == b.key,
242 )
243 .map_err(Into::into)
244 }
245
246 fn prune_rocksdb<Provider>(
251 &self,
252 provider: &Provider,
253 input: PruneInput,
254 range: std::ops::RangeInclusive<BlockNumber>,
255 range_end: BlockNumber,
256 ) -> Result<SegmentOutput, PrunerError>
257 where
258 Provider: DBProvider + StaticFileProviderFactory + ChangeSetReader + RocksDBProviderFactory,
259 {
260 let mut limiter = input.limiter;
264
265 if limiter.is_limit_reached() {
266 return Ok(SegmentOutput::not_done(
267 limiter.interrupt_reason(),
268 input.previous_checkpoint.map(SegmentOutputCheckpoint::from_prune_checkpoint),
269 ))
270 }
271
272 let mut highest_deleted_accounts = FxHashMap::default();
273 let mut last_changeset_pruned_block = None;
274 let mut changesets_processed = 0usize;
275 let mut done = true;
276
277 let walker = StaticFileAccountChangesetWalker::new(provider, range);
281 for result in walker {
282 let (block_number, changeset) = result?;
283 if limiter.is_limit_reached() &&
288 last_changeset_pruned_block.is_some_and(|last| last != block_number)
289 {
290 done = false;
291 break;
292 }
293 highest_deleted_accounts.insert(changeset.address, block_number);
294 last_changeset_pruned_block = Some(block_number);
295 changesets_processed += 1;
296 limiter.increment_deleted_entries_count();
297 }
298 trace!(target: "pruner", processed = %changesets_processed, %done, "Scanned account changesets from static files");
299
300 let last_changeset_pruned_block = last_changeset_pruned_block.unwrap_or(range_end);
301
302 let mut deleted_shards = 0usize;
304 let mut updated_shards = 0usize;
305
306 let mut sorted_accounts: Vec<_> = highest_deleted_accounts.into_iter().collect();
308 sorted_accounts.sort_unstable_by_key(|(addr, _)| *addr);
309
310 provider.with_rocksdb_batch(|mut batch| {
311 let targets: Vec<_> = sorted_accounts
312 .iter()
313 .map(|(addr, highest)| (*addr, (*highest).min(last_changeset_pruned_block)))
314 .collect();
315
316 let outcomes = batch.prune_account_history_batch(&targets)?;
317 deleted_shards = outcomes.deleted;
318 updated_shards = outcomes.updated;
319
320 Ok(((), Some(batch.into_inner())))
321 })?;
322 trace!(target: "pruner", deleted = deleted_shards, updated = updated_shards, %done, "Pruned account history (RocksDB indices)");
323
324 if done {
330 provider.static_file_provider().delete_segment_below_block(
331 StaticFileSegment::AccountChangeSets,
332 last_changeset_pruned_block + 1,
333 )?;
334 }
335
336 let progress = limiter.progress(done);
337
338 Ok(SegmentOutput {
339 progress,
340 pruned: changesets_processed + deleted_shards + updated_shards,
341 checkpoint: Some(SegmentOutputCheckpoint {
342 block_number: Some(last_changeset_pruned_block),
343 tx_number: None,
344 }),
345 })
346 }
347}
348
349#[cfg(test)]
350mod tests {
351 use super::ACCOUNT_HISTORY_TABLES_TO_PRUNE;
352 use crate::segments::{AccountHistory, PruneInput, PruneLimiter, Segment, SegmentOutput};
353 use alloy_primitives::{BlockNumber, B256};
354 use assert_matches::assert_matches;
355 use reth_db_api::{models::StorageSettings, tables, BlockNumberList};
356 use reth_provider::{DBProvider, DatabaseProviderFactory, PruneCheckpointReader};
357 use reth_prune_types::{
358 PruneCheckpoint, PruneInterruptReason, PruneMode, PruneProgress, PruneSegment,
359 };
360 use reth_stages::test_utils::{StorageKind, TestStageDB};
361 use reth_storage_api::StorageSettingsCache;
362 use reth_testing_utils::generators::{
363 self, random_block_range, random_changeset_range, random_eoa_accounts, BlockRangeParams,
364 };
365 use std::{collections::BTreeMap, ops::AddAssign};
366
367 #[test]
368 fn prune_legacy() {
369 let db = TestStageDB::default();
370 let mut rng = generators::rng();
371
372 let blocks = random_block_range(
373 &mut rng,
374 0..=5000,
375 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
376 );
377 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
378
379 let accounts = random_eoa_accounts(&mut rng, 2).into_iter().collect::<BTreeMap<_, _>>();
380
381 let (changesets, _) = random_changeset_range(
382 &mut rng,
383 blocks.iter(),
384 accounts.into_iter().map(|(addr, acc)| (addr, (acc, Vec::new()))),
385 0..0,
386 0..0,
387 );
388 db.insert_changesets(changesets.clone(), None).expect("insert changesets");
389 db.insert_history(changesets.clone(), None).expect("insert history");
390
391 let account_occurrences = db.table::<tables::AccountsHistory>().unwrap().into_iter().fold(
392 BTreeMap::<_, usize>::new(),
393 |mut map, (key, _)| {
394 map.entry(key.key).or_default().add_assign(1);
395 map
396 },
397 );
398 assert!(account_occurrences.into_iter().any(|(_, occurrences)| occurrences > 1));
399
400 assert_eq!(
401 db.table::<tables::AccountChangeSets>().unwrap().len(),
402 changesets.iter().flatten().count()
403 );
404
405 let original_shards = db.table::<tables::AccountsHistory>().unwrap();
406
407 let test_prune =
408 |to_block: BlockNumber, run: usize, expected_result: (PruneProgress, usize)| {
409 let prune_mode = PruneMode::Before(to_block);
410 let deleted_entries_limit = 2000;
411 let mut limiter =
412 PruneLimiter::default().set_deleted_entries_limit(deleted_entries_limit);
413 let input = PruneInput {
414 previous_checkpoint: db
415 .factory
416 .provider()
417 .unwrap()
418 .get_prune_checkpoint(PruneSegment::AccountHistory)
419 .unwrap(),
420 to_block,
421 limiter: limiter.clone(),
422 };
423 let segment = AccountHistory::new(prune_mode);
424
425 let provider = db.factory.database_provider_rw().unwrap();
426 provider.set_storage_settings_cache(StorageSettings::v1());
427 let result = segment.prune(&provider, input).unwrap();
428 limiter.increment_deleted_entries_count_by(result.pruned);
429
430 assert_matches!(
431 result,
432 SegmentOutput {progress, pruned, checkpoint: Some(_)}
433 if (progress, pruned) == expected_result
434 );
435
436 segment
437 .save_checkpoint(
438 &provider,
439 result.checkpoint.unwrap().as_prune_checkpoint(prune_mode),
440 )
441 .unwrap();
442 provider.commit().expect("commit");
443
444 let changesets = changesets
445 .iter()
446 .enumerate()
447 .flat_map(|(block_number, changeset)| {
448 changeset.iter().map(move |change| (block_number, change))
449 })
450 .collect::<Vec<_>>();
451
452 #[expect(clippy::skip_while_next)]
453 let pruned = changesets
454 .iter()
455 .enumerate()
456 .skip_while(|(i, (block_number, _))| {
457 *i < deleted_entries_limit / ACCOUNT_HISTORY_TABLES_TO_PRUNE * run &&
458 *block_number <= to_block as usize
459 })
460 .next()
461 .map(|(i, _)| i)
462 .unwrap_or_default();
463
464 let mut pruned_changesets = changesets.iter().skip(pruned.saturating_sub(1));
467
468 let last_pruned_block_number = pruned_changesets
469 .next()
470 .map(|(block_number, _)| if result.progress.is_finished() {
471 *block_number
472 } else {
473 block_number.saturating_sub(1)
474 } as BlockNumber)
475 .unwrap_or(to_block);
476
477 let pruned_changesets = pruned_changesets.fold(
478 BTreeMap::<_, Vec<_>>::new(),
479 |mut acc, (block_number, change)| {
480 acc.entry(block_number).or_default().push(change);
481 acc
482 },
483 );
484
485 assert_eq!(
486 db.table::<tables::AccountChangeSets>().unwrap().len(),
487 pruned_changesets.values().flatten().count()
488 );
489
490 let actual_shards = db.table::<tables::AccountsHistory>().unwrap();
491
492 let expected_shards = original_shards
493 .iter()
494 .filter(|(key, _)| key.highest_block_number > last_pruned_block_number)
495 .map(|(key, blocks)| {
496 let new_blocks =
497 blocks.iter().skip_while(|block| *block <= last_pruned_block_number);
498 (key.clone(), BlockNumberList::new_pre_sorted(new_blocks))
499 })
500 .collect::<Vec<_>>();
501
502 assert_eq!(actual_shards, expected_shards);
503
504 assert_eq!(
505 db.factory
506 .provider()
507 .unwrap()
508 .get_prune_checkpoint(PruneSegment::AccountHistory)
509 .unwrap(),
510 Some(PruneCheckpoint {
511 block_number: Some(last_pruned_block_number),
512 tx_number: None,
513 prune_mode
514 })
515 );
516 };
517
518 test_prune(
519 998,
520 1,
521 (PruneProgress::HasMoreData(PruneInterruptReason::DeletedEntriesLimitReached), 1000),
522 );
523 test_prune(998, 2, (PruneProgress::Finished, 998));
524 test_prune(1400, 3, (PruneProgress::Finished, 804));
525 }
526
527 #[test]
528 fn prune_rocksdb_path() {
529 use reth_db_api::models::ShardedKey;
530 use reth_provider::{RocksDBProviderFactory, StaticFileProviderFactory};
531
532 let db = TestStageDB::default();
533 let mut rng = generators::rng();
534
535 let blocks = random_block_range(
536 &mut rng,
537 0..=100,
538 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
539 );
540 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
541
542 let accounts = random_eoa_accounts(&mut rng, 2).into_iter().collect::<BTreeMap<_, _>>();
543
544 let (changesets, _) = random_changeset_range(
545 &mut rng,
546 blocks.iter(),
547 accounts.into_iter().map(|(addr, acc)| (addr, (acc, Vec::new()))),
548 0..0,
549 0..0,
550 );
551
552 db.insert_changesets_to_static_files(changesets.clone(), None)
553 .expect("insert changesets to static files");
554
555 let mut account_blocks: BTreeMap<_, Vec<u64>> = BTreeMap::new();
556 for (block, changeset) in changesets.iter().enumerate() {
557 for (address, _, _) in changeset {
558 account_blocks.entry(*address).or_default().push(block as u64);
559 }
560 }
561
562 let rocksdb = db.factory.rocksdb_provider();
563 let mut batch = rocksdb.batch();
564 for (address, block_numbers) in &account_blocks {
565 let shard = BlockNumberList::new_pre_sorted(block_numbers.iter().copied());
566 batch
567 .put::<tables::AccountsHistory>(ShardedKey::new(*address, u64::MAX), &shard)
568 .unwrap();
569 }
570 batch.commit().unwrap();
571
572 for (address, expected_blocks) in &account_blocks {
573 let shards = rocksdb.account_history_shards(*address).unwrap();
574 assert_eq!(shards.len(), 1);
575 assert_eq!(shards[0].1.iter().collect::<Vec<_>>(), *expected_blocks);
576 }
577
578 let to_block: BlockNumber = 50;
579 let prune_mode = PruneMode::Before(to_block);
580 let input =
581 PruneInput { previous_checkpoint: None, to_block, limiter: PruneLimiter::default() };
582 let segment = AccountHistory::new(prune_mode);
583
584 db.factory.set_storage_settings_cache(StorageSettings::v2());
585
586 let provider = db.factory.database_provider_rw().unwrap();
587 let result = segment.prune(&provider, input).unwrap();
588 provider.commit().expect("commit");
589
590 assert_matches!(
591 result,
592 SegmentOutput { progress: PruneProgress::Finished, pruned, checkpoint: Some(_) }
593 if pruned > 0
594 );
595
596 for (address, original_blocks) in &account_blocks {
597 let shards = rocksdb.account_history_shards(*address).unwrap();
598
599 let expected_blocks: Vec<u64> =
600 original_blocks.iter().copied().filter(|b| *b > to_block).collect();
601
602 if expected_blocks.is_empty() {
603 assert!(
604 shards.is_empty(),
605 "Expected no shards for address {address:?} after pruning"
606 );
607 } else {
608 assert_eq!(shards.len(), 1, "Expected 1 shard for address {address:?}");
609 assert_eq!(
610 shards[0].1.iter().collect::<Vec<_>>(),
611 expected_blocks,
612 "Shard blocks mismatch for address {address:?}"
613 );
614 }
615 }
616
617 let static_file_provider = db.factory.static_file_provider();
618 let highest_block = static_file_provider.get_highest_static_file_block(
619 reth_static_file_types::StaticFileSegment::AccountChangeSets,
620 );
621 if let Some(block) = highest_block {
622 assert!(
623 block > to_block,
624 "Static files should only contain blocks above to_block ({to_block}), got {block}"
625 );
626 }
627 }
628
629 #[test]
632 #[allow(clippy::clone_on_copy)]
633 fn prune_partial_progress_mid_block() {
634 use alloy_primitives::{Address, U256};
635 use reth_primitives_traits::Account;
636 use reth_testing_utils::generators::ChangeSet;
637
638 let db = TestStageDB::default();
639 let mut rng = generators::rng();
640
641 let blocks = random_block_range(
643 &mut rng,
644 0..=10,
645 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
646 );
647 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
648
649 let addr1 = Address::with_last_byte(1);
651 let addr2 = Address::with_last_byte(2);
652 let addr3 = Address::with_last_byte(3);
653 let addr4 = Address::with_last_byte(4);
654 let addr5 = Address::with_last_byte(5);
655
656 let account = Account { nonce: 1, balance: U256::from(100), ..Default::default() };
657
658 let changesets: Vec<ChangeSet> = vec![
660 vec![(addr1, account.clone(), vec![])], vec![(addr1, account.clone(), vec![])], vec![(addr1, account.clone(), vec![])], vec![(addr1, account.clone(), vec![])], vec![(addr1, account.clone(), vec![])], vec![
667 (addr1, account.clone(), vec![]),
668 (addr2, account.clone(), vec![]),
669 (addr3, account.clone(), vec![]),
670 (addr4, account.clone(), vec![]),
671 ],
672 vec![(addr5, account, vec![])], ];
674
675 db.insert_changesets(changesets.clone(), None).expect("insert changesets");
676 db.insert_history(changesets.clone(), None).expect("insert history");
677
678 assert_eq!(
680 db.table::<tables::AccountChangeSets>().unwrap().len(),
681 changesets.iter().flatten().count()
682 );
683
684 let prune_mode = PruneMode::Before(10);
685
686 let deleted_entries_limit = 14; let limiter = PruneLimiter::default().set_deleted_entries_limit(deleted_entries_limit);
694
695 let input = PruneInput { previous_checkpoint: None, to_block: 10, limiter };
696 let segment = AccountHistory::new(prune_mode);
697
698 let provider = db.factory.database_provider_rw().unwrap();
699 provider.set_storage_settings_cache(StorageSettings::v1());
700 let result = segment.prune(&provider, input).unwrap();
701
702 assert!(!result.progress.is_finished(), "Expected HasMoreData since we stopped mid-block");
704
705 segment
707 .save_checkpoint(&provider, result.checkpoint.unwrap().as_prune_checkpoint(prune_mode))
708 .unwrap();
709 provider.commit().expect("commit");
710
711 let checkpoint = db
713 .factory
714 .provider()
715 .unwrap()
716 .get_prune_checkpoint(PruneSegment::AccountHistory)
717 .unwrap()
718 .expect("checkpoint should exist");
719
720 assert_eq!(
721 checkpoint.block_number,
722 Some(4),
723 "Checkpoint should be block 4 (block before incomplete block 5)"
724 );
725
726 let remaining_changesets = db.table::<tables::AccountChangeSets>().unwrap();
728 assert!(
732 !remaining_changesets.is_empty(),
733 "Should have remaining changesets for blocks 5-6"
734 );
735
736 let history = db.table::<tables::AccountsHistory>().unwrap();
739 for (key, _blocks) in &history {
740 assert!(
743 key.highest_block_number > 4,
744 "Found stale history shard with highest_block_number {} <= checkpoint 4",
745 key.highest_block_number
746 );
747 }
748
749 let input2 = PruneInput {
751 previous_checkpoint: Some(checkpoint),
752 to_block: 10,
753 limiter: PruneLimiter::default().set_deleted_entries_limit(100), };
755
756 let provider2 = db.factory.database_provider_rw().unwrap();
757 provider2.set_storage_settings_cache(StorageSettings::v1());
758 let result2 = segment.prune(&provider2, input2).unwrap();
759
760 assert!(result2.progress.is_finished(), "Second run should complete");
761
762 segment
763 .save_checkpoint(
764 &provider2,
765 result2.checkpoint.unwrap().as_prune_checkpoint(prune_mode),
766 )
767 .unwrap();
768 provider2.commit().expect("commit");
769
770 let final_checkpoint = db
772 .factory
773 .provider()
774 .unwrap()
775 .get_prune_checkpoint(PruneSegment::AccountHistory)
776 .unwrap()
777 .expect("checkpoint should exist");
778
779 assert_eq!(final_checkpoint.block_number, Some(6), "Final checkpoint should be at block 6");
781
782 let final_changesets = db.table::<tables::AccountChangeSets>().unwrap();
784 assert!(final_changesets.is_empty(), "All changesets up to block 10 should be pruned");
785 }
786
787 #[test]
791 fn dense_block_advances_rocksdb_checkpoint() {
792 use reth_db_api::models::ShardedKey;
793 use reth_provider::RocksDBProviderFactory;
794
795 let db = TestStageDB::default();
796 let mut rng = generators::rng();
797
798 let blocks = random_block_range(
799 &mut rng,
800 0..=20,
801 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
802 );
803 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
804
805 let accounts = random_eoa_accounts(&mut rng, 2).into_iter().collect::<BTreeMap<_, _>>();
806 let (changesets, _) = random_changeset_range(
807 &mut rng,
808 blocks.iter(),
809 accounts.into_iter().map(|(addr, acc)| (addr, (acc, Vec::new()))),
810 0..0,
811 0..0,
812 );
813 assert!(changesets.iter().all(|changeset| changeset.len() == 2));
816
817 db.insert_changesets_to_static_files(changesets.clone(), None)
818 .expect("insert changesets to static files");
819
820 let mut account_blocks: BTreeMap<_, Vec<u64>> = BTreeMap::new();
822 for (block, changeset) in changesets.iter().enumerate() {
823 for (address, _, _) in changeset {
824 account_blocks.entry(*address).or_default().push(block as u64);
825 }
826 }
827 let rocksdb = db.factory.rocksdb_provider();
828 let mut batch = rocksdb.batch();
829 for (address, block_numbers) in &account_blocks {
830 let shard = BlockNumberList::new_pre_sorted(block_numbers.iter().copied());
831 batch
832 .put::<tables::AccountsHistory>(ShardedKey::new(*address, u64::MAX), &shard)
833 .unwrap();
834 }
835 batch.commit().unwrap();
836
837 db.factory.set_storage_settings_cache(StorageSettings::v2());
838
839 let to_block: BlockNumber = 15;
840 let prune_mode = PruneMode::Before(to_block);
841 let segment = AccountHistory::new(prune_mode);
842
843 let mut checkpoint = PruneCheckpoint { block_number: Some(4), tx_number: None, prune_mode };
845
846 let run_prune = |checkpoint: PruneCheckpoint, limit: usize| {
847 let input = PruneInput {
848 previous_checkpoint: Some(checkpoint),
849 to_block,
850 limiter: PruneLimiter::default().set_deleted_entries_limit(limit),
851 };
852
853 let provider = db.factory.database_provider_rw().unwrap();
854 provider.set_storage_settings_cache(StorageSettings::v2());
855 let result = segment.prune(&provider, input).unwrap();
856 segment
857 .save_checkpoint(
858 &provider,
859 result.checkpoint.unwrap().as_prune_checkpoint(prune_mode),
860 )
861 .unwrap();
862 provider.commit().expect("commit");
863
864 let checkpoint = db
865 .factory
866 .provider()
867 .unwrap()
868 .get_prune_checkpoint(PruneSegment::AccountHistory)
869 .unwrap()
870 .unwrap();
871 (result, checkpoint)
872 };
873
874 for _ in 0..3 {
876 let previous = checkpoint.block_number;
877 let (result, next) = run_prune(checkpoint, 2);
878 checkpoint = next;
879
880 assert!(
881 !result.progress.is_finished(),
882 "the range is longer than one run's budget allows"
883 );
884 assert!(
885 checkpoint.block_number > previous,
886 "checkpoint must advance past the dense block, got {:?} after {previous:?}",
887 checkpoint.block_number
888 );
889 }
890 assert_eq!(checkpoint.block_number, Some(7), "one dense block cleared per run");
891
892 let (result, checkpoint) = run_prune(checkpoint, 1000);
894 assert!(result.progress.is_finished());
895 assert_eq!(checkpoint.block_number, Some(to_block));
896 }
897
898 #[test]
902 fn dense_block_advances_static_file_checkpoint() {
903 let db = TestStageDB::default();
904 let mut rng = generators::rng();
905
906 let blocks = random_block_range(
907 &mut rng,
908 0..=20,
909 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
910 );
911 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
912
913 let accounts = random_eoa_accounts(&mut rng, 2).into_iter().collect::<BTreeMap<_, _>>();
914 let (changesets, _) = random_changeset_range(
915 &mut rng,
916 blocks.iter(),
917 accounts.into_iter().map(|(addr, acc)| (addr, (acc, Vec::new()))),
918 0..0,
919 0..0,
920 );
921 assert!(changesets.iter().all(|changeset| changeset.len() == 2));
922
923 db.insert_changesets_to_static_files(changesets, None)
927 .expect("insert changesets to static files");
928 db.factory.set_storage_settings_cache(StorageSettings::v2());
929 assert!(db.table::<tables::AccountChangeSets>().unwrap().is_empty());
930
931 let to_block: BlockNumber = 15;
932 let prune_mode = PruneMode::Before(to_block);
933 let segment = AccountHistory::new(prune_mode);
934
935 let mut checkpoint = PruneCheckpoint { block_number: Some(4), tx_number: None, prune_mode };
936
937 for _ in 0..3 {
938 let previous = checkpoint.block_number;
939 let input = PruneInput {
940 previous_checkpoint: Some(checkpoint),
941 to_block,
942 limiter: PruneLimiter::default()
944 .set_deleted_entries_limit(2 * ACCOUNT_HISTORY_TABLES_TO_PRUNE),
945 };
946 let range = input.get_next_block_range().unwrap();
947 let range_end = *range.end();
948
949 let provider = db.factory.database_provider_rw().unwrap();
950 provider.set_storage_settings_cache(StorageSettings::v2());
951 let result = segment.prune_static_files(&provider, input, range, range_end).unwrap();
952 segment
953 .save_checkpoint(
954 &provider,
955 result.checkpoint.unwrap().as_prune_checkpoint(prune_mode),
956 )
957 .unwrap();
958 provider.commit().expect("commit");
959
960 checkpoint = db
961 .factory
962 .provider()
963 .unwrap()
964 .get_prune_checkpoint(PruneSegment::AccountHistory)
965 .unwrap()
966 .unwrap();
967
968 assert!(!result.progress.is_finished());
969 assert!(
970 checkpoint.block_number > previous,
971 "checkpoint must advance past the dense block, got {:?} after {previous:?}",
972 checkpoint.block_number
973 );
974 }
975 assert_eq!(checkpoint.block_number, Some(7), "one dense block cleared per run");
976 }
977}