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 fn prune_partial_progress_mid_block() {
633 use alloy_primitives::{Address, U256};
634 use reth_primitives_traits::Account;
635 use reth_testing_utils::generators::ChangeSet;
636
637 let db = TestStageDB::default();
638 let mut rng = generators::rng();
639
640 let blocks = random_block_range(
642 &mut rng,
643 0..=10,
644 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
645 );
646 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
647
648 let addr1 = Address::with_last_byte(1);
650 let addr2 = Address::with_last_byte(2);
651 let addr3 = Address::with_last_byte(3);
652 let addr4 = Address::with_last_byte(4);
653 let addr5 = Address::with_last_byte(5);
654
655 let account = Account { nonce: 1, balance: U256::from(100), bytecode_hash: None };
656
657 let changesets: Vec<ChangeSet> = vec![
659 vec![(addr1, account, vec![])], vec![(addr1, account, vec![])], vec![(addr1, account, vec![])], vec![(addr1, account, vec![])], vec![(addr1, account, vec![])], vec![
666 (addr1, account, vec![]),
667 (addr2, account, vec![]),
668 (addr3, account, vec![]),
669 (addr4, account, vec![]),
670 ],
671 vec![(addr5, account, vec![])], ];
673
674 db.insert_changesets(changesets.clone(), None).expect("insert changesets");
675 db.insert_history(changesets.clone(), None).expect("insert history");
676
677 assert_eq!(
679 db.table::<tables::AccountChangeSets>().unwrap().len(),
680 changesets.iter().flatten().count()
681 );
682
683 let prune_mode = PruneMode::Before(10);
684
685 let deleted_entries_limit = 14; let limiter = PruneLimiter::default().set_deleted_entries_limit(deleted_entries_limit);
693
694 let input = PruneInput { previous_checkpoint: None, to_block: 10, limiter };
695 let segment = AccountHistory::new(prune_mode);
696
697 let provider = db.factory.database_provider_rw().unwrap();
698 provider.set_storage_settings_cache(StorageSettings::v1());
699 let result = segment.prune(&provider, input).unwrap();
700
701 assert!(!result.progress.is_finished(), "Expected HasMoreData since we stopped mid-block");
703
704 segment
706 .save_checkpoint(&provider, result.checkpoint.unwrap().as_prune_checkpoint(prune_mode))
707 .unwrap();
708 provider.commit().expect("commit");
709
710 let checkpoint = db
712 .factory
713 .provider()
714 .unwrap()
715 .get_prune_checkpoint(PruneSegment::AccountHistory)
716 .unwrap()
717 .expect("checkpoint should exist");
718
719 assert_eq!(
720 checkpoint.block_number,
721 Some(4),
722 "Checkpoint should be block 4 (block before incomplete block 5)"
723 );
724
725 let remaining_changesets = db.table::<tables::AccountChangeSets>().unwrap();
727 assert!(
731 !remaining_changesets.is_empty(),
732 "Should have remaining changesets for blocks 5-6"
733 );
734
735 let history = db.table::<tables::AccountsHistory>().unwrap();
738 for (key, _blocks) in &history {
739 assert!(
742 key.highest_block_number > 4,
743 "Found stale history shard with highest_block_number {} <= checkpoint 4",
744 key.highest_block_number
745 );
746 }
747
748 let input2 = PruneInput {
750 previous_checkpoint: Some(checkpoint),
751 to_block: 10,
752 limiter: PruneLimiter::default().set_deleted_entries_limit(100), };
754
755 let provider2 = db.factory.database_provider_rw().unwrap();
756 provider2.set_storage_settings_cache(StorageSettings::v1());
757 let result2 = segment.prune(&provider2, input2).unwrap();
758
759 assert!(result2.progress.is_finished(), "Second run should complete");
760
761 segment
762 .save_checkpoint(
763 &provider2,
764 result2.checkpoint.unwrap().as_prune_checkpoint(prune_mode),
765 )
766 .unwrap();
767 provider2.commit().expect("commit");
768
769 let final_checkpoint = db
771 .factory
772 .provider()
773 .unwrap()
774 .get_prune_checkpoint(PruneSegment::AccountHistory)
775 .unwrap()
776 .expect("checkpoint should exist");
777
778 assert_eq!(final_checkpoint.block_number, Some(6), "Final checkpoint should be at block 6");
780
781 let final_changesets = db.table::<tables::AccountChangeSets>().unwrap();
783 assert!(final_changesets.is_empty(), "All changesets up to block 10 should be pruned");
784 }
785
786 #[test]
790 fn dense_block_advances_rocksdb_checkpoint() {
791 use reth_db_api::models::ShardedKey;
792 use reth_provider::RocksDBProviderFactory;
793
794 let db = TestStageDB::default();
795 let mut rng = generators::rng();
796
797 let blocks = random_block_range(
798 &mut rng,
799 0..=20,
800 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
801 );
802 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
803
804 let accounts = random_eoa_accounts(&mut rng, 2).into_iter().collect::<BTreeMap<_, _>>();
805 let (changesets, _) = random_changeset_range(
806 &mut rng,
807 blocks.iter(),
808 accounts.into_iter().map(|(addr, acc)| (addr, (acc, Vec::new()))),
809 0..0,
810 0..0,
811 );
812 assert!(changesets.iter().all(|changeset| changeset.len() == 2));
815
816 db.insert_changesets_to_static_files(changesets.clone(), None)
817 .expect("insert changesets to static files");
818
819 let mut account_blocks: BTreeMap<_, Vec<u64>> = BTreeMap::new();
821 for (block, changeset) in changesets.iter().enumerate() {
822 for (address, _, _) in changeset {
823 account_blocks.entry(*address).or_default().push(block as u64);
824 }
825 }
826 let rocksdb = db.factory.rocksdb_provider();
827 let mut batch = rocksdb.batch();
828 for (address, block_numbers) in &account_blocks {
829 let shard = BlockNumberList::new_pre_sorted(block_numbers.iter().copied());
830 batch
831 .put::<tables::AccountsHistory>(ShardedKey::new(*address, u64::MAX), &shard)
832 .unwrap();
833 }
834 batch.commit().unwrap();
835
836 db.factory.set_storage_settings_cache(StorageSettings::v2());
837
838 let to_block: BlockNumber = 15;
839 let prune_mode = PruneMode::Before(to_block);
840 let segment = AccountHistory::new(prune_mode);
841
842 let mut checkpoint = PruneCheckpoint { block_number: Some(4), tx_number: None, prune_mode };
844
845 let run_prune = |checkpoint: PruneCheckpoint, limit: usize| {
846 let input = PruneInput {
847 previous_checkpoint: Some(checkpoint),
848 to_block,
849 limiter: PruneLimiter::default().set_deleted_entries_limit(limit),
850 };
851
852 let provider = db.factory.database_provider_rw().unwrap();
853 provider.set_storage_settings_cache(StorageSettings::v2());
854 let result = segment.prune(&provider, input).unwrap();
855 segment
856 .save_checkpoint(
857 &provider,
858 result.checkpoint.unwrap().as_prune_checkpoint(prune_mode),
859 )
860 .unwrap();
861 provider.commit().expect("commit");
862
863 let checkpoint = db
864 .factory
865 .provider()
866 .unwrap()
867 .get_prune_checkpoint(PruneSegment::AccountHistory)
868 .unwrap()
869 .unwrap();
870 (result, checkpoint)
871 };
872
873 for _ in 0..3 {
875 let previous = checkpoint.block_number;
876 let (result, next) = run_prune(checkpoint, 2);
877 checkpoint = next;
878
879 assert!(
880 !result.progress.is_finished(),
881 "the range is longer than one run's budget allows"
882 );
883 assert!(
884 checkpoint.block_number > previous,
885 "checkpoint must advance past the dense block, got {:?} after {previous:?}",
886 checkpoint.block_number
887 );
888 }
889 assert_eq!(checkpoint.block_number, Some(7), "one dense block cleared per run");
890
891 let (result, checkpoint) = run_prune(checkpoint, 1000);
893 assert!(result.progress.is_finished());
894 assert_eq!(checkpoint.block_number, Some(to_block));
895 }
896
897 #[test]
901 fn dense_block_advances_static_file_checkpoint() {
902 let db = TestStageDB::default();
903 let mut rng = generators::rng();
904
905 let blocks = random_block_range(
906 &mut rng,
907 0..=20,
908 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
909 );
910 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
911
912 let accounts = random_eoa_accounts(&mut rng, 2).into_iter().collect::<BTreeMap<_, _>>();
913 let (changesets, _) = random_changeset_range(
914 &mut rng,
915 blocks.iter(),
916 accounts.into_iter().map(|(addr, acc)| (addr, (acc, Vec::new()))),
917 0..0,
918 0..0,
919 );
920 assert!(changesets.iter().all(|changeset| changeset.len() == 2));
921
922 db.insert_changesets_to_static_files(changesets, None)
926 .expect("insert changesets to static files");
927 db.factory.set_storage_settings_cache(StorageSettings::v2());
928 assert!(db.table::<tables::AccountChangeSets>().unwrap().is_empty());
929
930 let to_block: BlockNumber = 15;
931 let prune_mode = PruneMode::Before(to_block);
932 let segment = AccountHistory::new(prune_mode);
933
934 let mut checkpoint = PruneCheckpoint { block_number: Some(4), tx_number: None, prune_mode };
935
936 for _ in 0..3 {
937 let previous = checkpoint.block_number;
938 let input = PruneInput {
939 previous_checkpoint: Some(checkpoint),
940 to_block,
941 limiter: PruneLimiter::default()
943 .set_deleted_entries_limit(2 * ACCOUNT_HISTORY_TABLES_TO_PRUNE),
944 };
945 let range = input.get_next_block_range().unwrap();
946 let range_end = *range.end();
947
948 let provider = db.factory.database_provider_rw().unwrap();
949 provider.set_storage_settings_cache(StorageSettings::v2());
950 let result = segment.prune_static_files(&provider, input, range, range_end).unwrap();
951 segment
952 .save_checkpoint(
953 &provider,
954 result.checkpoint.unwrap().as_prune_checkpoint(prune_mode),
955 )
956 .unwrap();
957 provider.commit().expect("commit");
958
959 checkpoint = db
960 .factory
961 .provider()
962 .unwrap()
963 .get_prune_checkpoint(PruneSegment::AccountHistory)
964 .unwrap()
965 .unwrap();
966
967 assert!(!result.progress.is_finished());
968 assert!(
969 checkpoint.block_number > previous,
970 "checkpoint must advance past the dense block, got {:?} after {previous:?}",
971 checkpoint.block_number
972 );
973 }
974 assert_eq!(checkpoint.block_number, Some(7), "one dense block cleared per run");
975 }
976}