Skip to main content

reth_prune/segments/user/
storage_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::{Address, BlockNumber, B256};
10use reth_db_api::{
11    models::{storage_sharded_key::StorageShardedKey, BlockNumberAddress},
12    tables,
13    transaction::DbTxMut,
14};
15use reth_provider::{DBProvider, EitherWriter, RocksDBProviderFactory, StaticFileProviderFactory};
16use reth_prune_types::{
17    PruneMode, PrunePurpose, PruneSegment, SegmentOutput, SegmentOutputCheckpoint,
18};
19use reth_static_file_types::StaticFileSegment;
20use reth_storage_api::{StorageChangeSetReader, StorageSettingsCache};
21use rustc_hash::FxHashMap;
22use tracing::{instrument, trace};
23
24/// Number of storage history tables to prune in one step.
25///
26/// Storage History consists of two tables: [`tables::StorageChangeSets`] and
27/// [`tables::StoragesHistory`]. We want to prune them to the same block number.
28const STORAGE_HISTORY_TABLES_TO_PRUNE: usize = 2;
29
30#[derive(Debug)]
31pub struct StorageHistory {
32    mode: PruneMode,
33}
34
35impl StorageHistory {
36    pub const fn new(mode: PruneMode) -> Self {
37        Self { mode }
38    }
39}
40
41impl<Provider> Segment<Provider> for StorageHistory
42where
43    Provider: DBProvider<Tx: DbTxMut>
44        + StaticFileProviderFactory
45        + StorageChangeSetReader
46        + StorageSettingsCache
47        + RocksDBProviderFactory,
48{
49    fn segment(&self) -> PruneSegment {
50        PruneSegment::StorageHistory
51    }
52
53    fn mode(&self) -> Option<PruneMode> {
54        Some(self.mode)
55    }
56
57    fn purpose(&self) -> PrunePurpose {
58        PrunePurpose::User
59    }
60
61    #[instrument(
62        name = "StorageHistory::prune",
63        target = "pruner",
64        skip(self, provider),
65        ret(level = "trace")
66    )]
67    fn prune(&self, provider: &Provider, input: PruneInput) -> Result<SegmentOutput, PrunerError> {
68        let range = match input.get_next_block_range() {
69            Some(range) => range,
70            None => {
71                trace!(target: "pruner", "No storage history to prune");
72                return Ok(SegmentOutput::done())
73            }
74        };
75        let range_end = *range.end();
76
77        // Check where storage history indices are stored
78        if provider.cached_storage_settings().storage_v2 {
79            return self.prune_rocksdb(provider, input, range, range_end);
80        }
81
82        // Check where storage changesets are stored (MDBX path)
83        if EitherWriter::storage_changesets_destination(provider).is_static_file() {
84            self.prune_static_files(provider, input, range, range_end)
85        } else {
86            self.prune_database(provider, input, range, range_end)
87        }
88    }
89}
90
91impl StorageHistory {
92    /// Prunes storage history when changesets are stored in static files.
93    fn prune_static_files<Provider>(
94        &self,
95        provider: &Provider,
96        input: PruneInput,
97        range: std::ops::RangeInclusive<BlockNumber>,
98        range_end: BlockNumber,
99    ) -> Result<SegmentOutput, PrunerError>
100    where
101        Provider: DBProvider<Tx: DbTxMut> + StaticFileProviderFactory,
102    {
103        let mut limiter = if let Some(limit) = input.limiter.deleted_entries_limit() {
104            input.limiter.set_deleted_entries_limit(limit / STORAGE_HISTORY_TABLES_TO_PRUNE)
105        } else {
106            input.limiter
107        };
108
109        // The limiter may already be exhausted from a previous segment in the same prune run.
110        // Early exit avoids unnecessary iteration when no budget remains.
111        if limiter.is_limit_reached() {
112            return Ok(SegmentOutput::not_done(
113                limiter.interrupt_reason(),
114                input.previous_checkpoint.map(SegmentOutputCheckpoint::from_prune_checkpoint),
115            ))
116        }
117
118        // The size of this map is limited by `prune_delete_limit * blocks_since_last_run /
119        // STORAGE_HISTORY_TABLES_TO_PRUNE`, and with current defaults it's usually `3500 * 5
120        // / 2`, so 8750 entries. Each entry is `160 bit + 256 bit + 64 bit`, so the total
121        // size should be up to ~0.5MB + some hashmap overhead. `blocks_since_last_run` is
122        // additionally limited by the `max_reorg_depth`, so no OOM is expected here.
123        let mut highest_deleted_storages = FxHashMap::default();
124        let mut last_changeset_pruned_block = None;
125        let mut pruned_changesets = 0;
126        let mut done = true;
127
128        let walker = provider.static_file_provider().walk_storage_changeset_range(range);
129        for result in walker {
130            let (block_address, entry) = result?;
131            let block_number = block_address.block_number();
132            let address = block_address.address();
133            // 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_storages.insert((address, entry.key), block_number);
144            last_changeset_pruned_block = Some(block_number);
145            pruned_changesets += 1;
146            limiter.increment_deleted_entries_count();
147        }
148
149        // 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::StorageChangeSets, last_block + 1)?;
154        }
155        trace!(target: "pruner", pruned = %pruned_changesets, %done, "Pruned storage history (changesets from static files)");
156
157        let result = HistoryPruneResult {
158            highest_deleted: highest_deleted_storages,
159            last_pruned_block: last_changeset_pruned_block,
160            pruned_count: pruned_changesets,
161            done,
162        };
163        finalize_history_prune::<_, tables::StoragesHistory, (Address, B256), _>(
164            provider,
165            result,
166            range_end,
167            &limiter,
168            |(address, storage_key), block_number| {
169                StorageShardedKey::new(address, storage_key, block_number)
170            },
171            |a, b| a.address == b.address && a.sharded_key.key == b.sharded_key.key,
172        )
173        .map_err(Into::into)
174    }
175
176    fn prune_database<Provider>(
177        &self,
178        provider: &Provider,
179        input: PruneInput,
180        range: std::ops::RangeInclusive<BlockNumber>,
181        range_end: BlockNumber,
182    ) -> Result<SegmentOutput, PrunerError>
183    where
184        Provider: DBProvider<Tx: DbTxMut>,
185    {
186        let mut limiter = if let Some(limit) = input.limiter.deleted_entries_limit() {
187            input.limiter.set_deleted_entries_limit(limit / STORAGE_HISTORY_TABLES_TO_PRUNE)
188        } else {
189            input.limiter
190        };
191
192        if limiter.is_limit_reached() {
193            return Ok(SegmentOutput::not_done(
194                limiter.interrupt_reason(),
195                input.previous_checkpoint.map(SegmentOutputCheckpoint::from_prune_checkpoint),
196            ))
197        }
198
199        // Deleted storage changeset keys (account addresses and storage slots) with the highest
200        // block number deleted for that key.
201        //
202        // The size of this map is limited by `prune_delete_limit * blocks_since_last_run /
203        // STORAGE_HISTORY_TABLES_TO_PRUNE`, and with current defaults it's usually `3500 * 5
204        // / 2`, so 8750 entries. Each entry is `160 bit + 256 bit + 64 bit`, so the total
205        // size should be up to ~0.5MB + some hashmap overhead. `blocks_since_last_run` is
206        // additionally limited by the `max_reorg_depth`, so no OOM is expected here.
207        let mut last_changeset_pruned_block = None;
208        let mut highest_deleted_storages = FxHashMap::default();
209        let (pruned_changesets, done) =
210            provider.tx_ref().prune_table_with_range::<tables::StorageChangeSets>(
211                BlockNumberAddress::range(range),
212                &mut limiter,
213                |_| false,
214                |(BlockNumberAddress((block_number, address)), entry)| {
215                    highest_deleted_storages.insert((address, entry.key), block_number);
216                    last_changeset_pruned_block = Some(block_number);
217                },
218            )?;
219        trace!(target: "pruner", deleted = %pruned_changesets, %done, "Pruned storage history (changesets)");
220
221        // The table walk can stop in the middle of a block, so the interrupted block has to be
222        // pruned again on the next run.
223        let last_pruned_block = last_changeset_pruned_block.map(|block_number| {
224            if done {
225                block_number
226            } else {
227                block_number.saturating_sub(1)
228            }
229        });
230
231        let result = HistoryPruneResult {
232            highest_deleted: highest_deleted_storages,
233            last_pruned_block,
234            pruned_count: pruned_changesets,
235            done,
236        };
237        finalize_history_prune::<_, tables::StoragesHistory, (Address, B256), _>(
238            provider,
239            result,
240            range_end,
241            &limiter,
242            |(address, storage_key), block_number| {
243                StorageShardedKey::new(address, storage_key, block_number)
244            },
245            |a, b| a.address == b.address && a.sharded_key.key == b.sharded_key.key,
246        )
247        .map_err(Into::into)
248    }
249
250    /// Prunes storage history when indices are stored in `RocksDB`.
251    ///
252    /// Reads storage changesets from static files and prunes the corresponding
253    /// `RocksDB` history shards.
254    fn prune_rocksdb<Provider>(
255        &self,
256        provider: &Provider,
257        input: PruneInput,
258        range: std::ops::RangeInclusive<BlockNumber>,
259        range_end: BlockNumber,
260    ) -> Result<SegmentOutput, PrunerError>
261    where
262        Provider: DBProvider + StaticFileProviderFactory + RocksDBProviderFactory,
263    {
264        let mut limiter = input.limiter;
265
266        if limiter.is_limit_reached() {
267            return Ok(SegmentOutput::not_done(
268                limiter.interrupt_reason(),
269                input.previous_checkpoint.map(SegmentOutputCheckpoint::from_prune_checkpoint),
270            ))
271        }
272
273        let mut highest_deleted_storages: FxHashMap<_, _> = FxHashMap::default();
274        let mut last_changeset_pruned_block = None;
275        let mut changesets_processed = 0usize;
276        let mut done = true;
277
278        // Walk storage changesets from static files using a streaming iterator.
279        // For each changeset, track the highest block number seen for each (address, storage_key)
280        // pair to determine which history shard entries need pruning.
281        let walker = provider.static_file_provider().walk_storage_changeset_range(range);
282        for result in walker {
283            let (block_address, entry) = result?;
284            let block_number = block_address.block_number();
285            let address = block_address.address();
286            // Static file changesets are not deleted here, so an interrupted block cannot be
287            // resumed: giving up the budget inside block N reports checkpoint N-1 and the next run
288            // rereads the same entries, forever. Stop on block boundaries only, overshooting the
289            // budget by at most the rest of one block.
290            if limiter.is_limit_reached() &&
291                last_changeset_pruned_block.is_some_and(|last| last != block_number)
292            {
293                done = false;
294                break;
295            }
296            highest_deleted_storages.insert((address, entry.key), block_number);
297            last_changeset_pruned_block = Some(block_number);
298            changesets_processed += 1;
299            limiter.increment_deleted_entries_count();
300        }
301
302        trace!(target: "pruner", processed = %changesets_processed, %done, "Scanned storage changesets from static files");
303
304        let last_changeset_pruned_block = last_changeset_pruned_block.unwrap_or(range_end);
305
306        // Prune RocksDB history shards for affected storage slots
307        let mut deleted_shards = 0usize;
308        let mut updated_shards = 0usize;
309
310        // Sort by (address, storage_key) for better RocksDB cache locality
311        let mut sorted_storages: Vec<_> = highest_deleted_storages.into_iter().collect();
312        sorted_storages.sort_unstable_by_key(|((addr, key), _)| (*addr, *key));
313
314        provider.with_rocksdb_batch(|mut batch| {
315            let targets: Vec<_> = sorted_storages
316                .iter()
317                .map(|((addr, key), highest)| {
318                    ((*addr, *key), (*highest).min(last_changeset_pruned_block))
319                })
320                .collect();
321
322            let outcomes = batch.prune_storage_history_batch(&targets)?;
323            deleted_shards = outcomes.deleted;
324            updated_shards = outcomes.updated;
325
326            Ok(((), Some(batch.into_inner())))
327        })?;
328
329        trace!(target: "pruner", deleted = deleted_shards, updated = updated_shards, %done, "Pruned storage history (RocksDB indices)");
330
331        // Delete static file jars only when fully processed. During provider.commit(), RocksDB
332        // batch is committed before the MDBX checkpoint. If crash occurs after RocksDB commit
333        // but before MDBX commit, on restart the pruner checkpoint indicates data needs
334        // re-pruning, but the RocksDB shards are already pruned - this is safe because pruning
335        // is idempotent (re-pruning already-pruned shards is a no-op).
336        if done {
337            provider.static_file_provider().delete_segment_below_block(
338                StaticFileSegment::StorageChangeSets,
339                last_changeset_pruned_block + 1,
340            )?;
341        }
342
343        let progress = limiter.progress(done);
344
345        Ok(SegmentOutput {
346            progress,
347            pruned: changesets_processed + deleted_shards + updated_shards,
348            checkpoint: Some(SegmentOutputCheckpoint {
349                block_number: Some(last_changeset_pruned_block),
350                tx_number: None,
351            }),
352        })
353    }
354}
355
356#[cfg(test)]
357mod tests {
358    use super::STORAGE_HISTORY_TABLES_TO_PRUNE;
359    use crate::segments::{PruneInput, PruneLimiter, Segment, SegmentOutput, StorageHistory};
360    use alloy_primitives::{BlockNumber, B256};
361    use assert_matches::assert_matches;
362    use reth_db_api::{models::StorageSettings, tables, BlockNumberList};
363    use reth_provider::{DBProvider, DatabaseProviderFactory, PruneCheckpointReader};
364    use reth_prune_types::{
365        PruneCheckpoint, PruneInterruptReason, PruneMode, PruneProgress, PruneSegment,
366    };
367    use reth_stages::test_utils::{StorageKind, TestStageDB};
368    use reth_storage_api::StorageSettingsCache;
369    use reth_testing_utils::generators::{
370        self, random_block_range, random_changeset_range, random_eoa_accounts, BlockRangeParams,
371    };
372    use std::{collections::BTreeMap, ops::AddAssign};
373
374    #[test]
375    fn prune_legacy() {
376        let db = TestStageDB::default();
377        let mut rng = generators::rng();
378
379        let blocks = random_block_range(
380            &mut rng,
381            0..=5000,
382            BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
383        );
384        db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
385
386        let accounts = random_eoa_accounts(&mut rng, 2).into_iter().collect::<BTreeMap<_, _>>();
387
388        let (changesets, _) = random_changeset_range(
389            &mut rng,
390            blocks.iter(),
391            accounts.into_iter().map(|(addr, acc)| (addr, (acc, Vec::new()))),
392            1..2,
393            1..2,
394        );
395        db.insert_changesets(changesets.clone(), None).expect("insert changesets");
396        db.insert_history(changesets.clone(), None).expect("insert history");
397
398        let storage_occurrences = db.table::<tables::StoragesHistory>().unwrap().into_iter().fold(
399            BTreeMap::<_, usize>::new(),
400            |mut map, (key, _)| {
401                map.entry((key.address, key.sharded_key.key)).or_default().add_assign(1);
402                map
403            },
404        );
405        assert!(storage_occurrences.into_iter().any(|(_, occurrences)| occurrences > 1));
406
407        assert_eq!(
408            db.table::<tables::StorageChangeSets>().unwrap().len(),
409            changesets.iter().flatten().flat_map(|(_, _, entries)| entries).count()
410        );
411
412        let original_shards = db.table::<tables::StoragesHistory>().unwrap();
413
414        let test_prune = |to_block: BlockNumber,
415                          run: usize,
416                          expected_result: (PruneProgress, usize)| {
417            let prune_mode = PruneMode::Before(to_block);
418            let deleted_entries_limit = 1000;
419            let mut limiter =
420                PruneLimiter::default().set_deleted_entries_limit(deleted_entries_limit);
421            let input = PruneInput {
422                previous_checkpoint: db
423                    .factory
424                    .provider()
425                    .unwrap()
426                    .get_prune_checkpoint(PruneSegment::StorageHistory)
427                    .unwrap(),
428                to_block,
429                limiter: limiter.clone(),
430            };
431            let segment = StorageHistory::new(prune_mode);
432
433            let provider = db.factory.database_provider_rw().unwrap();
434            provider.set_storage_settings_cache(StorageSettings::v1());
435            let result = segment.prune(&provider, input).unwrap();
436            limiter.increment_deleted_entries_count_by(result.pruned);
437
438            assert_matches!(
439                result,
440                SegmentOutput {progress, pruned, checkpoint: Some(_)}
441                    if (progress, pruned) == expected_result
442            );
443
444            segment
445                .save_checkpoint(
446                    &provider,
447                    result.checkpoint.unwrap().as_prune_checkpoint(prune_mode),
448                )
449                .unwrap();
450            provider.commit().expect("commit");
451
452            let changesets = changesets
453                .iter()
454                .enumerate()
455                .flat_map(|(block_number, changeset)| {
456                    changeset.iter().flat_map(move |(address, _, entries)| {
457                        entries.iter().map(move |entry| (block_number, address, entry))
458                    })
459                })
460                .collect::<Vec<_>>();
461
462            #[expect(clippy::skip_while_next)]
463            let pruned = changesets
464                .iter()
465                .enumerate()
466                .skip_while(|(i, (block_number, _, _))| {
467                    *i < deleted_entries_limit / STORAGE_HISTORY_TABLES_TO_PRUNE * run &&
468                        *block_number <= to_block as usize
469                })
470                .next()
471                .map(|(i, _)| i)
472                .unwrap_or_default();
473
474            // Skip what we've pruned so far, subtracting one to get last pruned block number
475            // further down
476            let mut pruned_changesets = changesets.iter().skip(pruned.saturating_sub(1));
477
478            let last_pruned_block_number = pruned_changesets
479                .next()
480                .map(|(block_number, _, _)| {
481                    (if result.progress.is_finished() {
482                        *block_number
483                    } else {
484                        block_number.saturating_sub(1)
485                    }) as BlockNumber
486                })
487                .unwrap_or(to_block);
488
489            let pruned_changesets = pruned_changesets.fold(
490                BTreeMap::<_, Vec<_>>::new(),
491                |mut acc, (block_number, address, entry)| {
492                    acc.entry((block_number, address)).or_default().push(entry);
493                    acc
494                },
495            );
496
497            assert_eq!(
498                db.table::<tables::StorageChangeSets>().unwrap().len(),
499                pruned_changesets.values().flatten().count()
500            );
501
502            let actual_shards = db.table::<tables::StoragesHistory>().unwrap();
503
504            let expected_shards = original_shards
505                .iter()
506                .filter(|(key, _)| key.sharded_key.highest_block_number > last_pruned_block_number)
507                .map(|(key, blocks)| {
508                    let new_blocks =
509                        blocks.iter().skip_while(|block| *block <= last_pruned_block_number);
510                    (key.clone(), BlockNumberList::new_pre_sorted(new_blocks))
511                })
512                .collect::<Vec<_>>();
513
514            assert_eq!(actual_shards, expected_shards);
515
516            assert_eq!(
517                db.factory
518                    .provider()
519                    .unwrap()
520                    .get_prune_checkpoint(PruneSegment::StorageHistory)
521                    .unwrap(),
522                Some(PruneCheckpoint {
523                    block_number: Some(last_pruned_block_number),
524                    tx_number: None,
525                    prune_mode
526                })
527            );
528        };
529
530        test_prune(
531            998,
532            1,
533            (PruneProgress::HasMoreData(PruneInterruptReason::DeletedEntriesLimitReached), 500),
534        );
535        test_prune(998, 2, (PruneProgress::Finished, 499));
536        test_prune(1200, 3, (PruneProgress::Finished, 202));
537    }
538
539    /// Tests that when a limiter stops mid-block (with multiple storage changes for the same
540    /// block), the checkpoint is set to `block_number - 1` to avoid dangling index entries.
541    #[test]
542    fn prune_partial_progress_mid_block() {
543        use alloy_primitives::{Address, U256};
544        use reth_primitives_traits::Account;
545        use reth_testing_utils::generators::ChangeSet;
546
547        let db = TestStageDB::default();
548        let mut rng = generators::rng();
549
550        // Create blocks 0..=10
551        let blocks = random_block_range(
552            &mut rng,
553            0..=10,
554            BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
555        );
556        db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
557
558        // Create specific changesets where block 5 has 4 storage changes
559        let addr1 = Address::with_last_byte(1);
560        let addr2 = Address::with_last_byte(2);
561
562        let account = Account { nonce: 1, balance: U256::from(100), bytecode_hash: None };
563
564        // Create storage entries
565        let storage_entry = |key: u8| reth_primitives_traits::StorageEntry {
566            key: B256::with_last_byte(key),
567            value: U256::from(100),
568        };
569
570        // Build changesets: blocks 0-4 have 1 storage change each, block 5 has 4 changes, block 6
571        // has 1. Entries within each account must be sorted by key.
572        let changesets: Vec<ChangeSet> = vec![
573            vec![(addr1, account, vec![storage_entry(1)])], // block 0
574            vec![(addr1, account, vec![storage_entry(1)])], // block 1
575            vec![(addr1, account, vec![storage_entry(1)])], // block 2
576            vec![(addr1, account, vec![storage_entry(1)])], // block 3
577            vec![(addr1, account, vec![storage_entry(1)])], // block 4
578            // block 5: 4 different storage changes (2 addresses, each with 2 storage slots)
579            // Sorted by address, then by storage key within each address
580            vec![
581                (addr1, account, vec![storage_entry(1), storage_entry(2)]),
582                (addr2, account, vec![storage_entry(1), storage_entry(2)]),
583            ],
584            vec![(addr1, account, vec![storage_entry(3)])], // block 6
585        ];
586
587        db.insert_changesets(changesets.clone(), None).expect("insert changesets");
588        db.insert_history(changesets.clone(), None).expect("insert history");
589
590        // Total storage changesets
591        let total_storage_entries: usize =
592            changesets.iter().flat_map(|c| c.iter()).map(|(_, _, entries)| entries.len()).sum();
593        assert_eq!(db.table::<tables::StorageChangeSets>().unwrap().len(), total_storage_entries);
594
595        let prune_mode = PruneMode::Before(10);
596
597        // Set limiter to stop mid-block 5
598        // With STORAGE_HISTORY_TABLES_TO_PRUNE=2, limit=14 gives us 7 storage entries before limit
599        // Blocks 0-4 use 5 slots, leaving 2 for block 5 (which has 4), so we stop mid-block 5
600        let deleted_entries_limit = 14; // 14/2 = 7 storage entries before limit
601        let limiter = PruneLimiter::default().set_deleted_entries_limit(deleted_entries_limit);
602
603        let input = PruneInput { previous_checkpoint: None, to_block: 10, limiter };
604        let segment = StorageHistory::new(prune_mode);
605
606        let provider = db.factory.database_provider_rw().unwrap();
607        provider.set_storage_settings_cache(StorageSettings::v1());
608        let result = segment.prune(&provider, input).unwrap();
609
610        // Should report that there's more data
611        assert!(!result.progress.is_finished(), "Expected HasMoreData since we stopped mid-block");
612
613        // Save checkpoint and commit
614        segment
615            .save_checkpoint(&provider, result.checkpoint.unwrap().as_prune_checkpoint(prune_mode))
616            .unwrap();
617        provider.commit().expect("commit");
618
619        // Verify checkpoint is set to block 4 (not 5), since block 5 is incomplete
620        let checkpoint = db
621            .factory
622            .provider()
623            .unwrap()
624            .get_prune_checkpoint(PruneSegment::StorageHistory)
625            .unwrap()
626            .expect("checkpoint should exist");
627
628        assert_eq!(
629            checkpoint.block_number,
630            Some(4),
631            "Checkpoint should be block 4 (block before incomplete block 5)"
632        );
633
634        // Verify remaining changesets
635        let remaining_changesets = db.table::<tables::StorageChangeSets>().unwrap();
636        assert!(
637            !remaining_changesets.is_empty(),
638            "Should have remaining changesets for blocks 5-6"
639        );
640
641        // Verify no dangling history indices for blocks that weren't fully pruned
642        let history = db.table::<tables::StoragesHistory>().unwrap();
643        for (key, _blocks) in &history {
644            assert!(
645                key.sharded_key.highest_block_number > 4,
646                "Found stale history shard with highest_block_number {} <= checkpoint 4",
647                key.sharded_key.highest_block_number
648            );
649        }
650
651        // Run prune again to complete - should finish processing block 5 and 6
652        let input2 = PruneInput {
653            previous_checkpoint: Some(checkpoint),
654            to_block: 10,
655            limiter: PruneLimiter::default().set_deleted_entries_limit(100), // high limit
656        };
657
658        let provider2 = db.factory.database_provider_rw().unwrap();
659        provider2.set_storage_settings_cache(StorageSettings::v1());
660        let result2 = segment.prune(&provider2, input2).unwrap();
661
662        assert!(result2.progress.is_finished(), "Second run should complete");
663
664        segment
665            .save_checkpoint(
666                &provider2,
667                result2.checkpoint.unwrap().as_prune_checkpoint(prune_mode),
668            )
669            .unwrap();
670        provider2.commit().expect("commit");
671
672        // Verify final checkpoint
673        let final_checkpoint = db
674            .factory
675            .provider()
676            .unwrap()
677            .get_prune_checkpoint(PruneSegment::StorageHistory)
678            .unwrap()
679            .expect("checkpoint should exist");
680
681        // Should now be at block 6 (the last block with changesets)
682        assert_eq!(final_checkpoint.block_number, Some(6), "Final checkpoint should be at block 6");
683
684        // All changesets should be pruned
685        let final_changesets = db.table::<tables::StorageChangeSets>().unwrap();
686        assert!(final_changesets.is_empty(), "All changesets up to block 10 should be pruned");
687    }
688
689    #[test]
690    fn prune_rocksdb() {
691        use reth_db_api::models::storage_sharded_key::StorageShardedKey;
692        use reth_provider::RocksDBProviderFactory;
693        use reth_storage_api::StorageSettings;
694
695        let db = TestStageDB::default();
696        let mut rng = generators::rng();
697
698        let blocks = random_block_range(
699            &mut rng,
700            0..=100,
701            BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
702        );
703        db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
704
705        let accounts = random_eoa_accounts(&mut rng, 2).into_iter().collect::<BTreeMap<_, _>>();
706
707        let (changesets, _) = random_changeset_range(
708            &mut rng,
709            blocks.iter(),
710            accounts.into_iter().map(|(addr, acc)| (addr, (acc, Vec::new()))),
711            1..2,
712            1..2,
713        );
714
715        db.insert_changesets_to_static_files(changesets.clone(), None)
716            .expect("insert changesets to static files");
717
718        let mut storage_indices: BTreeMap<(alloy_primitives::Address, B256), Vec<u64>> =
719            BTreeMap::new();
720        for (block, changeset) in changesets.iter().enumerate() {
721            for (address, _, storage_entries) in changeset {
722                for entry in storage_entries {
723                    storage_indices.entry((*address, entry.key)).or_default().push(block as u64);
724                }
725            }
726        }
727
728        {
729            let rocksdb = db.factory.rocksdb_provider();
730            let mut batch = rocksdb.batch();
731            for ((address, storage_key), block_numbers) in &storage_indices {
732                let shard = BlockNumberList::new_pre_sorted(block_numbers.clone());
733                batch
734                    .put::<tables::StoragesHistory>(
735                        StorageShardedKey::last(*address, *storage_key),
736                        &shard,
737                    )
738                    .expect("insert storage history shard");
739            }
740            batch.commit().expect("commit rocksdb batch");
741        }
742
743        {
744            let rocksdb = db.factory.rocksdb_provider();
745            for (address, storage_key) in storage_indices.keys() {
746                let shards = rocksdb.storage_history_shards(*address, *storage_key).unwrap();
747                assert!(!shards.is_empty(), "RocksDB should contain storage history before prune");
748            }
749        }
750
751        let to_block = 50u64;
752        let prune_mode = PruneMode::Before(to_block);
753        let input =
754            PruneInput { previous_checkpoint: None, to_block, limiter: PruneLimiter::default() };
755        let segment = StorageHistory::new(prune_mode);
756
757        let provider = db.factory.database_provider_rw().unwrap();
758        provider.set_storage_settings_cache(StorageSettings::v2());
759        let result = segment.prune(&provider, input).unwrap();
760        provider.commit().expect("commit");
761
762        assert_matches!(
763            result,
764            SegmentOutput { progress: PruneProgress::Finished, checkpoint: Some(_), .. }
765        );
766
767        {
768            let rocksdb = db.factory.rocksdb_provider();
769            for ((address, storage_key), block_numbers) in &storage_indices {
770                let shards = rocksdb.storage_history_shards(*address, *storage_key).unwrap();
771
772                let remaining_blocks: Vec<u64> =
773                    block_numbers.iter().copied().filter(|&b| b > to_block).collect();
774
775                if remaining_blocks.is_empty() {
776                    assert!(
777                        shards.is_empty(),
778                        "Shard for {:?}/{:?} should be deleted when all blocks pruned",
779                        address,
780                        storage_key
781                    );
782                } else {
783                    assert!(!shards.is_empty(), "Shard should exist with remaining blocks");
784                    let actual_blocks: Vec<u64> =
785                        shards.iter().flat_map(|(_, list)| list.iter()).collect();
786                    assert_eq!(
787                        actual_blocks, remaining_blocks,
788                        "RocksDB shard should only contain blocks > {}",
789                        to_block
790                    );
791                }
792            }
793        }
794    }
795
796    /// A block holding at least a whole run's budget of changesets must not stall pruning: the
797    /// walk deletes no changesets, so a checkpoint rewound below such a block would make every
798    /// later run reread it and never advance.
799    #[test]
800    fn dense_block_advances_rocksdb_checkpoint() {
801        use alloy_primitives::U256;
802        use reth_db_api::models::storage_sharded_key::StorageShardedKey;
803        use reth_primitives_traits::StorageEntry;
804        use reth_provider::RocksDBProviderFactory;
805        use reth_storage_api::StorageSettings;
806
807        let db = TestStageDB::default();
808        let mut rng = generators::rng();
809
810        let blocks = random_block_range(
811            &mut rng,
812            0..=20,
813            BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
814        );
815        db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
816
817        // Two storage changesets per block, so a budget of two makes every block "dense".
818        const ENTRIES_PER_BLOCK: usize = 2;
819        let (address, account) = random_eoa_accounts(&mut rng, 1).into_iter().next().unwrap();
820        let keys = [B256::with_last_byte(1), B256::with_last_byte(2)];
821        let changesets = (0..=20)
822            .map(|_| {
823                vec![(
824                    address,
825                    account,
826                    keys.iter()
827                        .map(|key| StorageEntry { key: *key, value: U256::from(1) })
828                        .collect(),
829                )]
830            })
831            .collect::<Vec<_>>();
832        db.insert_changesets_to_static_files(changesets, None)
833            .expect("insert changesets to static files");
834
835        {
836            let rocksdb = db.factory.rocksdb_provider();
837            let mut batch = rocksdb.batch();
838            for key in keys {
839                batch
840                    .put::<tables::StoragesHistory>(
841                        StorageShardedKey::last(address, key),
842                        &BlockNumberList::new_pre_sorted(0..=20),
843                    )
844                    .expect("insert storage history shard");
845            }
846            batch.commit().expect("commit rocksdb batch");
847        }
848
849        let to_block = 15u64;
850        let prune_mode = PruneMode::Before(to_block);
851        let segment = StorageHistory::new(prune_mode);
852
853        // Start from a checkpoint in the middle so a rewind can't be masked by block 0.
854        let mut checkpoint = PruneCheckpoint { block_number: Some(4), tx_number: None, prune_mode };
855
856        let run_prune = |checkpoint: PruneCheckpoint, limit: usize| {
857            let input = PruneInput {
858                previous_checkpoint: Some(checkpoint),
859                to_block,
860                limiter: PruneLimiter::default().set_deleted_entries_limit(limit),
861            };
862
863            let provider = db.factory.database_provider_rw().unwrap();
864            provider.set_storage_settings_cache(StorageSettings::v2());
865            let result = segment.prune(&provider, input).unwrap();
866            segment
867                .save_checkpoint(
868                    &provider,
869                    result.checkpoint.unwrap().as_prune_checkpoint(prune_mode),
870                )
871                .unwrap();
872            provider.commit().expect("commit");
873
874            let checkpoint = db
875                .factory
876                .provider()
877                .unwrap()
878                .get_prune_checkpoint(PruneSegment::StorageHistory)
879                .unwrap()
880                .unwrap();
881            (result, checkpoint)
882        };
883
884        // The RocksDB path does not halve the limit, so this budget is exactly one dense block.
885        for _ in 0..3 {
886            let previous = checkpoint.block_number;
887            let (result, next) = run_prune(checkpoint, ENTRIES_PER_BLOCK);
888            checkpoint = next;
889
890            assert!(
891                !result.progress.is_finished(),
892                "the range is longer than one run's budget allows"
893            );
894            assert!(
895                checkpoint.block_number > previous,
896                "checkpoint must advance past the dense block, got {:?} after {previous:?}",
897                checkpoint.block_number
898            );
899        }
900        assert_eq!(checkpoint.block_number, Some(7), "one dense block cleared per run");
901
902        // With enough budget the remainder of the range completes in one run.
903        let (result, checkpoint) = run_prune(checkpoint, 1000);
904        assert!(result.progress.is_finished());
905        assert_eq!(checkpoint.block_number, Some(to_block));
906    }
907}