Skip to main content

reth_prune/segments/user/
account_history.rs

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
23/// Number of account history tables to prune in one step.
24///
25/// Account History consists of two tables: [`tables::AccountChangeSets`] (either in database or
26/// static files) and [`tables::AccountsHistory`]. We want to prune them to the same block number.
27const 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        // Check where account history indices are stored
77        if provider.cached_storage_settings().storage_v2 {
78            return self.prune_rocksdb(provider, input, range, range_end);
79        }
80
81        // Check where account changesets are stored (MDBX path)
82        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    /// Prunes account history when changesets are stored in static files.
92    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        // The limiter may already be exhausted from a previous segment in the same prune run.
109        // Early exit avoids unnecessary iteration when no budget remains.
110        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        // Deleted account changeset keys (account addresses) with the highest block number deleted
118        // for that key.
119        //
120        // The size of this map is limited by `prune_delete_limit * blocks_since_last_run /
121        // ACCOUNT_HISTORY_TABLES_TO_PRUNE`, and with current default it's usually `3500 * 5
122        // / 2`, so 8750 entries. Each entry is `160 bit + 64 bit`, so the total
123        // size should be up to ~0.25MB + some hashmap overhead. `blocks_since_last_run` is
124        // additionally limited by the `max_reorg_depth`, so no OOM is expected here.
125        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            // The walk itself deletes nothing, so an interrupted block cannot be resumed: giving
134            // up the budget inside block N reports checkpoint N-1 and the next run rereads the
135            // same entries, forever. Stop on block boundaries only, overshooting the budget by at
136            // most the rest of one block.
137            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        // Delete static file jars only when fully processed
150        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        // Deleted account changeset keys (account addresses) with the highest block number deleted
198        // for that key.
199        //
200        // The size of this map is limited by `prune_delete_limit * blocks_since_last_run /
201        // ACCOUNT_HISTORY_TABLES_TO_PRUNE`, and with the current defaults it's usually `3500 * 5 /
202        // 2`, so 8750 entries. Each entry is `160 bit + 64 bit`, so the total size should be up to
203        // ~0.25MB + some hashmap overhead. `blocks_since_last_run` is additionally limited by the
204        // `max_reorg_depth`, so no OOM is expected here.
205        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        // The table walk can stop in the middle of a block, so the interrupted block has to be
220        // pruned again on the next run.
221        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    /// Prunes account history when indices are stored in `RocksDB`.
247    ///
248    /// Reads account changesets from static files and prunes the corresponding
249    /// `RocksDB` history shards.
250    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        // Unlike MDBX path, we don't divide the limit by 2 because RocksDB path only prunes
261        // history shards (no separate changeset table to delete from). The changesets are in
262        // static files which are deleted separately.
263        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        // Walk account changesets from static files using a streaming iterator.
278        // For each changeset, track the highest block number seen for each address
279        // to determine which history shard entries need pruning.
280        let walker = StaticFileAccountChangesetWalker::new(provider, range);
281        for result in walker {
282            let (block_number, changeset) = result?;
283            // Static file changesets are not deleted here, so an interrupted block cannot be
284            // resumed: giving up the budget inside block N reports checkpoint N-1 and the next run
285            // rereads the same entries, forever. Stop on block boundaries only, overshooting the
286            // budget by at most the rest of one block.
287            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        // Prune RocksDB history shards for affected accounts
303        let mut deleted_shards = 0usize;
304        let mut updated_shards = 0usize;
305
306        // Sort by address for better RocksDB cache locality
307        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        // Delete static file jars only when fully processed. During provider.commit(), RocksDB
325        // batch is committed before the MDBX checkpoint. If crash occurs after RocksDB commit
326        // but before MDBX commit, on restart the pruner checkpoint indicates data needs
327        // re-pruning, but the RocksDB shards are already pruned - this is safe because pruning
328        // is idempotent (re-pruning already-pruned shards is a no-op).
329        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                // Skip what we've pruned so far, subtracting one to get last pruned block number
465                // further down
466                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    /// Tests that when a limiter stops mid-block (with multiple changes for the same block),
630    /// the checkpoint is set to `block_number - 1` to avoid dangling index entries.
631    #[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        // Create blocks 0..=10
641        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        // Create specific changesets where block 5 has 4 account changes
649        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        // Build changesets: blocks 0-4 have 1 change each, block 5 has 4 changes, block 6 has 1
658        let changesets: Vec<ChangeSet> = vec![
659            vec![(addr1, account, vec![])], // block 0
660            vec![(addr1, account, vec![])], // block 1
661            vec![(addr1, account, vec![])], // block 2
662            vec![(addr1, account, vec![])], // block 3
663            vec![(addr1, account, vec![])], // block 4
664            // block 5: 4 different account changes (sorted by address for consistency)
665            vec![
666                (addr1, account, vec![]),
667                (addr2, account, vec![]),
668                (addr3, account, vec![]),
669                (addr4, account, vec![]),
670            ],
671            vec![(addr5, account, vec![])], // block 6
672        ];
673
674        db.insert_changesets(changesets.clone(), None).expect("insert changesets");
675        db.insert_history(changesets.clone(), None).expect("insert history");
676
677        // Total changesets: 5 (blocks 0-4) + 4 (block 5) + 1 (block 6) = 10
678        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        // Set limiter to stop after 7 entries (mid-block 5: 5 from blocks 0-4, then 2 of 4 from
686        // block 5). Due to ACCOUNT_HISTORY_TABLES_TO_PRUNE=2, actual limit is 7/2=3
687        // changesets. So we'll process blocks 0, 1, 2 (3 changesets), stopping before block
688        // 3. Actually, let's use a higher limit to reach block 5. With limit=14, we get 7
689        // changeset slots. Blocks 0-4 use 5 slots, leaving 2 for block 5 (which has 4), so
690        // we stop mid-block 5.
691        let deleted_entries_limit = 14; // 14/2 = 7 changeset entries before limit
692        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        // Should report that there's more data
702        assert!(!result.progress.is_finished(), "Expected HasMoreData since we stopped mid-block");
703
704        // Save checkpoint and commit
705        segment
706            .save_checkpoint(&provider, result.checkpoint.unwrap().as_prune_checkpoint(prune_mode))
707            .unwrap();
708        provider.commit().expect("commit");
709
710        // Verify checkpoint is set to block 4 (not 5), since block 5 is incomplete
711        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        // Verify remaining changesets (block 5 and 6 should still have entries)
726        let remaining_changesets = db.table::<tables::AccountChangeSets>().unwrap();
727        // After pruning blocks 0-4, remaining should be block 5 (4 entries) + block 6 (1 entry) = 5
728        // But since we stopped mid-block 5, some of block 5 might be pruned
729        // However, checkpoint is 4, so on re-run we should re-process from block 5
730        assert!(
731            !remaining_changesets.is_empty(),
732            "Should have remaining changesets for blocks 5-6"
733        );
734
735        // Verify no dangling history indices for blocks that weren't fully pruned
736        // The indices for block 5 should still reference blocks <= 5 appropriately
737        let history = db.table::<tables::AccountsHistory>().unwrap();
738        for (key, _blocks) in &history {
739            // All blocks in the history should be > checkpoint block number
740            // OR the shard's highest_block_number should be > checkpoint
741            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        // Run prune again to complete - should finish processing block 5 and 6
749        let input2 = PruneInput {
750            previous_checkpoint: Some(checkpoint),
751            to_block: 10,
752            limiter: PruneLimiter::default().set_deleted_entries_limit(100), // high limit
753        };
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        // Verify final checkpoint
770        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        // Should now be at block 6 (the last block with changesets)
779        assert_eq!(final_checkpoint.block_number, Some(6), "Final checkpoint should be at block 6");
780
781        // All changesets should be pruned
782        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    /// A block holding at least a whole run's budget of changesets must not stall the `RocksDB`
787    /// path: the walk deletes no changesets, so a checkpoint rewound below such a block would make
788    /// every later run reread it and never advance.
789    #[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        // `random_changeset_range` emits exactly 2 account changesets per block (sender +
813        // recipient), so a budget of 2 makes every block "dense".
814        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        // History index lives in RocksDB on the v2 path.
820        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        // Start from a checkpoint in the middle so a rewind can't be masked by block 0.
843        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        // The RocksDB path does not halve the limit, so a budget of 2 is exactly one dense block.
874        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        // With enough budget the remainder of the range completes in one run.
892        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    /// Same guarantee for `prune_static_files`, which shares the changeset walk. `Segment::prune`
898    /// cannot reach it while `storage_v2` short-circuits to the `RocksDB` path, so it is called
899    /// directly.
900    #[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        // Changesets in static files. `StaticFileAccountChangesetWalker` resolves entries through
923        // `ChangeSetReader`, which only reads static files when `storage_v2` is set, so the walk
924        // needs v2 settings even though v2 routes `prune()` to the RocksDB path.
925        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                // Halved internally by ACCOUNT_HISTORY_TABLES_TO_PRUNE, so 4 == one dense block.
942                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}