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    #[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        // Create blocks 0..=10
642        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        // Create specific changesets where block 5 has 4 account changes
650        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        // Build changesets: blocks 0-4 have 1 change each, block 5 has 4 changes, block 6 has 1
659        let changesets: Vec<ChangeSet> = vec![
660            vec![(addr1, account.clone(), vec![])], // block 0
661            vec![(addr1, account.clone(), vec![])], // block 1
662            vec![(addr1, account.clone(), vec![])], // block 2
663            vec![(addr1, account.clone(), vec![])], // block 3
664            vec![(addr1, account.clone(), vec![])], // block 4
665            // block 5: 4 different account changes (sorted by address for consistency)
666            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![])], // block 6
673        ];
674
675        db.insert_changesets(changesets.clone(), None).expect("insert changesets");
676        db.insert_history(changesets.clone(), None).expect("insert history");
677
678        // Total changesets: 5 (blocks 0-4) + 4 (block 5) + 1 (block 6) = 10
679        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        // Set limiter to stop after 7 entries (mid-block 5: 5 from blocks 0-4, then 2 of 4 from
687        // block 5). Due to ACCOUNT_HISTORY_TABLES_TO_PRUNE=2, actual limit is 7/2=3
688        // changesets. So we'll process blocks 0, 1, 2 (3 changesets), stopping before block
689        // 3. Actually, let's use a higher limit to reach block 5. With limit=14, we get 7
690        // changeset slots. Blocks 0-4 use 5 slots, leaving 2 for block 5 (which has 4), so
691        // we stop mid-block 5.
692        let deleted_entries_limit = 14; // 14/2 = 7 changeset entries before limit
693        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        // Should report that there's more data
703        assert!(!result.progress.is_finished(), "Expected HasMoreData since we stopped mid-block");
704
705        // Save checkpoint and commit
706        segment
707            .save_checkpoint(&provider, result.checkpoint.unwrap().as_prune_checkpoint(prune_mode))
708            .unwrap();
709        provider.commit().expect("commit");
710
711        // Verify checkpoint is set to block 4 (not 5), since block 5 is incomplete
712        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        // Verify remaining changesets (block 5 and 6 should still have entries)
727        let remaining_changesets = db.table::<tables::AccountChangeSets>().unwrap();
728        // After pruning blocks 0-4, remaining should be block 5 (4 entries) + block 6 (1 entry) = 5
729        // But since we stopped mid-block 5, some of block 5 might be pruned
730        // However, checkpoint is 4, so on re-run we should re-process from block 5
731        assert!(
732            !remaining_changesets.is_empty(),
733            "Should have remaining changesets for blocks 5-6"
734        );
735
736        // Verify no dangling history indices for blocks that weren't fully pruned
737        // The indices for block 5 should still reference blocks <= 5 appropriately
738        let history = db.table::<tables::AccountsHistory>().unwrap();
739        for (key, _blocks) in &history {
740            // All blocks in the history should be > checkpoint block number
741            // OR the shard's highest_block_number should be > checkpoint
742            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        // Run prune again to complete - should finish processing block 5 and 6
750        let input2 = PruneInput {
751            previous_checkpoint: Some(checkpoint),
752            to_block: 10,
753            limiter: PruneLimiter::default().set_deleted_entries_limit(100), // high limit
754        };
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        // Verify final checkpoint
771        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        // Should now be at block 6 (the last block with changesets)
780        assert_eq!(final_checkpoint.block_number, Some(6), "Final checkpoint should be at block 6");
781
782        // All changesets should be pruned
783        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    /// A block holding at least a whole run's budget of changesets must not stall the `RocksDB`
788    /// path: the walk deletes no changesets, so a checkpoint rewound below such a block would make
789    /// every later run reread it and never advance.
790    #[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        // `random_changeset_range` emits exactly 2 account changesets per block (sender +
814        // recipient), so a budget of 2 makes every block "dense".
815        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        // History index lives in RocksDB on the v2 path.
821        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        // Start from a checkpoint in the middle so a rewind can't be masked by block 0.
844        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        // The RocksDB path does not halve the limit, so a budget of 2 is exactly one dense block.
875        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        // With enough budget the remainder of the range completes in one run.
893        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    /// Same guarantee for `prune_static_files`, which shares the changeset walk. `Segment::prune`
899    /// cannot reach it while `storage_v2` short-circuits to the `RocksDB` path, so it is called
900    /// directly.
901    #[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        // Changesets in static files. `StaticFileAccountChangesetWalker` resolves entries through
924        // `ChangeSetReader`, which only reads static files when `storage_v2` is set, so the walk
925        // needs v2 settings even though v2 routes `prune()` to the RocksDB path.
926        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                // Halved internally by ACCOUNT_HISTORY_TABLES_TO_PRUNE, so 4 == one dense block.
943                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}