Skip to main content

reth_provider/providers/database/
provider.rs

1use super::SaveBlocksInput;
2use crate::{
3    changesets_utils::StorageRevertsIter,
4    providers::{
5        database::{chain::ChainStorage, metrics, DatabaseProviderMetrics},
6        rocksdb::{PendingRocksDBBatches, RocksDBProvider, RocksDBWriteCtx},
7        static_file::{StaticFileWriteCtx, StaticFileWriter},
8        NodeTypesForProvider, StaticFileProvider,
9    },
10    to_range,
11    traits::{
12        AccountExtReader, BlockSource, ChangeSetReader, ReceiptProvider, StageCheckpointWriter,
13    },
14    AccountReader, BlockBodyWriter, BlockExecutionWriter, BlockHashReader, BlockNumReader,
15    BlockReader, BlockWriter, BundleStateInit, ChainStateBlockReader, ChainStateBlockWriter,
16    DBProvider, DbTxProvider, EitherReader, EitherWriter, EitherWriterDestination, HashingWriter,
17    HeaderProvider, HeaderSyncGapProvider, HistoricalStateProvider, HistoricalStateProviderRef,
18    HistoryWriter, LatestStateProvider, LatestStateProviderRef, OriginalValuesKnown,
19    PersistenceFrontiers, ProviderError, PruneCheckpointReader, PruneCheckpointWriter,
20    RawRocksDBBatch, RevertsInit, RocksBatchArg, RocksDBProviderFactory, StageCheckpointReader,
21    StateProviderBox, StateWriter, StaticFileProviderFactory, StatsReader, StorageReader,
22    StorageTrieWriter, TransactionVariant, TransactionsProvider, TransactionsProviderExt,
23    TrieWriter,
24};
25use alloy_consensus::{
26    transaction::{SignerRecoverable, TransactionMeta, TxHashRef},
27    BlockHeader, TxReceipt,
28};
29use alloy_eips::BlockHashOrNumber;
30use alloy_primitives::{
31    keccak256,
32    map::{hash_map, AddressSet, B256Map, HashMap},
33    Address, BlockHash, BlockNumber, StorageKey, StorageValue, TxHash, TxNumber, B256,
34};
35use itertools::Itertools;
36use parking_lot::RwLock;
37use rayon::slice::ParallelSliceMut;
38use reth_chain_state::ExecutedBlock;
39use reth_chainspec::{ChainInfo, ChainSpecProvider, EthChainSpec};
40use reth_db_api::{
41    cursor::{DbCursorRO, DbCursorRW, DbDupCursorRO, DbDupCursorRW},
42    database::{Database, ReaderTxnTracker},
43    models::{
44        sharded_key, storage_sharded_key::StorageShardedKey, AccountBeforeTx, BlockNumberAddress,
45        BlockNumberAddressRange, ShardedKey, StorageBeforeTx, StorageSettings,
46        StoredBlockBodyIndices,
47    },
48    table::Table,
49    tables,
50    transaction::{DbTx, DbTxMut},
51    BlockNumberList,
52};
53use reth_execution_types::{BlockExecutionOutput, BlockExecutionResult, Chain, ExecutionOutcome};
54use reth_node_types::{BlockTy, BodyTy, HeaderTy, NodeTypes, ReceiptTy, TxTy};
55use reth_primitives_traits::{
56    Account, Block as _, BlockBody as _, Bytecode, FastInstant as Instant, RecoveredBlock,
57    SealedHeader, StorageEntry,
58};
59use reth_prune_types::{
60    PruneCheckpoint, PruneMode, PruneModes, PruneSegment, MINIMUM_UNWIND_SAFE_DISTANCE,
61};
62use reth_stages_types::{FinishCheckpoint, StageCheckpoint, StageId};
63use reth_static_file_types::StaticFileSegment;
64use reth_storage_api::{
65    BlockBodyIndicesProvider, BlockBodyReader, MetadataProvider, MetadataWriter,
66    NodePrimitivesProvider, StateProvider, StateReader, StateWriteConfig, StorageChangeSetReader,
67    StoragePath, StorageSettingsCache, TryIntoHistoricalStateProvider, WriteStateInput,
68};
69use reth_storage_errors::provider::{ProviderResult, StaticFileWriterError};
70use reth_storage_overlay::OverlayManager;
71use reth_trie::{
72    updates::{StorageTrieUpdatesSorted, TrieUpdatesSorted},
73    ComputedTrieData, HashedPostStateSorted,
74};
75use reth_trie_db::{DatabaseStorageTrieCursor, TrieTableAdapter};
76use revm::database::states::{
77    PlainStateReverts, PlainStorageChangeset, PlainStorageRevert, StateChangeset,
78};
79use smallvec::SmallVec;
80use std::{
81    cmp::Ordering,
82    collections::{BTreeMap, BTreeSet},
83    fmt::Debug,
84    ops::{Deref, DerefMut, Range, RangeBounds, RangeInclusive},
85    path::PathBuf,
86    sync::Arc,
87};
88use tracing::{debug, instrument, trace};
89
90/// Determines the commit order for database operations.
91#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
92pub enum CommitOrder {
93    /// Normal commit order: static files first, then `RocksDB`, then MDBX.
94    #[default]
95    Normal,
96    /// Unwind commit order: MDBX first, then `RocksDB`, then static files.
97    /// Used for unwind operations to allow recovery by truncating static files on restart.
98    Unwind,
99}
100
101impl CommitOrder {
102    /// Returns true if this is unwind commit order.
103    pub const fn is_unwind(&self) -> bool {
104        matches!(self, Self::Unwind)
105    }
106}
107
108/// A [`DatabaseProvider`] that holds a read-only database transaction.
109pub type DatabaseProviderRO<DB, N> = DatabaseProvider<<DB as Database>::TX, N>;
110
111/// A [`DatabaseProvider`] that holds a read-write database transaction.
112///
113/// Ideally this would be an alias type. However, there's some weird compiler error (<https://github.com/rust-lang/rust/issues/102211>), that forces us to wrap this in a struct instead.
114/// Once that issue is solved, we can probably revert back to being an alias type.
115#[derive(Debug)]
116pub struct DatabaseProviderRW<DB: Database, N: NodeTypes>(
117    pub DatabaseProvider<<DB as Database>::TXMut, N>,
118);
119
120impl<DB: Database, N: NodeTypes> Deref for DatabaseProviderRW<DB, N> {
121    type Target = DatabaseProvider<<DB as Database>::TXMut, N>;
122
123    fn deref(&self) -> &Self::Target {
124        &self.0
125    }
126}
127
128impl<DB: Database, N: NodeTypes> DerefMut for DatabaseProviderRW<DB, N> {
129    fn deref_mut(&mut self) -> &mut Self::Target {
130        &mut self.0
131    }
132}
133
134impl<DB: Database, N: NodeTypes> AsRef<DatabaseProvider<<DB as Database>::TXMut, N>>
135    for DatabaseProviderRW<DB, N>
136{
137    fn as_ref(&self) -> &DatabaseProvider<<DB as Database>::TXMut, N> {
138        &self.0
139    }
140}
141
142impl<DB: Database, N: NodeTypes + 'static> DatabaseProviderRW<DB, N> {
143    /// Commit database transaction and static file if it exists.
144    pub fn commit(self) -> ProviderResult<()> {
145        self.0.commit()
146    }
147
148    /// Consume `DbTx` or `DbTxMut`.
149    pub fn into_tx(self) -> <DB as Database>::TXMut {
150        self.0.into_tx()
151    }
152
153    /// Override the minimum pruning distance for testing purposes.
154    #[cfg(any(test, feature = "test-utils"))]
155    pub const fn with_minimum_pruning_distance(mut self, distance: u64) -> Self {
156        self.0.minimum_pruning_distance = distance;
157        self
158    }
159}
160
161impl<DB: Database, N: NodeTypes> From<DatabaseProviderRW<DB, N>>
162    for DatabaseProvider<<DB as Database>::TXMut, N>
163{
164    fn from(provider: DatabaseProviderRW<DB, N>) -> Self {
165        provider.0
166    }
167}
168
169/// Mode for [`DatabaseProvider::save_blocks_inner`].
170#[derive(Debug, Clone, Copy, PartialEq, Eq)]
171enum SaveBlocksMode {
172    /// Full mode: write block structure + receipts + state + trie.
173    /// Used by engine/production code.
174    Full,
175    /// Blocks only: write block structure (headers, txs, senders, indices).
176    /// Receipts/state/trie are skipped - they may come later via separate calls.
177    /// Used by `insert_block`.
178    BlocksOnly,
179}
180
181impl SaveBlocksMode {
182    /// Returns `true` if this is [`SaveBlocksMode::Full`].
183    const fn with_state(self) -> bool {
184        matches!(self, Self::Full)
185    }
186}
187
188/// A provider struct that fetches data from the database.
189/// Wrapper around [`DbTx`] and [`DbTxMut`]. Example: [`HeaderProvider`] [`BlockHashReader`]
190pub struct DatabaseProvider<TX, N: NodeTypes> {
191    /// Database transaction.
192    tx: TX,
193    /// Chain spec
194    chain_spec: Arc<N::ChainSpec>,
195    /// Static File provider
196    static_file_provider: StaticFileProvider<N::Primitives>,
197    /// Pruning configuration
198    prune_modes: PruneModes,
199    /// Node storage handler.
200    storage: Arc<N::Storage>,
201    /// Storage configuration settings for this node
202    storage_settings: Arc<RwLock<StorageSettings>>,
203    /// `RocksDB` provider
204    rocksdb_provider: RocksDBProvider,
205    /// Manager for state trie overlays and cached changesets.
206    overlay_manager: OverlayManager<N::Primitives>,
207    /// Task runtime for spawning parallel I/O work.
208    runtime: reth_tasks::Runtime,
209    /// Path to the database directory.
210    db_path: PathBuf,
211    /// Pending `RocksDB` batches to be committed at provider commit time.
212    pending_rocksdb_batches: PendingRocksDBBatches,
213    /// Commit order for database operations.
214    commit_order: CommitOrder,
215    /// Minimum distance from tip required for pruning
216    minimum_pruning_distance: u64,
217    /// Database provider metrics
218    metrics: Arc<DatabaseProviderMetrics>,
219    /// Database handle used to inspect active MDBX readers during unwind commits.
220    reader_txn_tracker: Option<Arc<dyn ReaderTxnTracker>>,
221}
222
223impl<TX: Debug, N: NodeTypes> Debug for DatabaseProvider<TX, N> {
224    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
225        let mut s = f.debug_struct("DatabaseProvider");
226        s.field("tx", &self.tx)
227            .field("chain_spec", &self.chain_spec)
228            .field("static_file_provider", &self.static_file_provider)
229            .field("prune_modes", &self.prune_modes)
230            .field("storage", &self.storage)
231            .field("storage_settings", &self.storage_settings)
232            .field("rocksdb_provider", &self.rocksdb_provider)
233            .field("overlay_manager", &self.overlay_manager)
234            .field("runtime", &self.runtime)
235            .field("pending_rocksdb_batches", &"<pending batches>")
236            .field("commit_order", &self.commit_order)
237            .field("minimum_pruning_distance", &self.minimum_pruning_distance)
238            .field("reader_txn_tracker", &"<reader txn tracker>")
239            .finish()
240    }
241}
242
243impl<TX, N: NodeTypes> DatabaseProvider<TX, N> {
244    /// Returns reference to prune modes.
245    pub const fn prune_modes_ref(&self) -> &PruneModes {
246        &self.prune_modes
247    }
248
249    /// Sets the minimum pruning distance.
250    pub const fn with_minimum_pruning_distance(mut self, distance: u64) -> Self {
251        self.minimum_pruning_distance = distance;
252        self
253    }
254
255    /// Attaches reader tracking so unwind commits can wait on active readers.
256    pub(crate) fn with_reader_txn_tracker<T>(mut self, reader_txn_tracker: T) -> Self
257    where
258        T: ReaderTxnTracker + 'static,
259    {
260        self.reader_txn_tracker = Some(Arc::new(reader_txn_tracker));
261        self
262    }
263}
264
265impl<TX: DbTx + 'static, N: NodeTypes> DatabaseProvider<TX, N> {
266    /// Commits unwind writes in MDBX -> `RocksDB` -> static-file order.
267    ///
268    /// This keeps MDBX as the first durable step so an interrupted unwind can be recovered by
269    /// truncating static files from checkpoints on the next startup.
270    ///
271    /// This waits after the MDBX commit so readers holding older MDBX-visible views cannot overlap
272    /// later cross-store unwind steps.
273    ///
274    /// Historical `storage_v2` reads ignore `RocksDB` history entries above their MDBX-visible tip,
275    /// so no additional post-`RocksDB` wait is needed before static-file commit.
276    fn commit_unwind(self) -> ProviderResult<()> {
277        let storage_v2 = self.cached_storage_settings().storage_v2;
278        let reader_txn_tracker = self.reader_txn_tracker.clone();
279        self.tx.commit()?;
280
281        if let Some(reader_txn_tracker) = reader_txn_tracker.as_ref() {
282            reader_txn_tracker.wait_for_pre_commit_readers();
283        }
284
285        if storage_v2 {
286            let batches = std::mem::take(&mut *self.pending_rocksdb_batches.lock());
287            for batch in batches {
288                self.rocksdb_provider.commit_batch(batch)?;
289            }
290        }
291
292        self.static_file_provider.commit()?;
293        Ok(())
294    }
295
296    /// State provider for latest state
297    pub fn latest<'a>(&'a self) -> Box<dyn StateProvider + 'a> {
298        trace!(target: "providers::db", "Returning latest state provider");
299        Box::new(LatestStateProviderRef::new(self))
300    }
301
302    /// Storage provider for state at that given block hash
303    pub fn history_by_block_hash<'a>(
304        &'a self,
305        block_hash: BlockHash,
306    ) -> ProviderResult<Box<dyn StateProvider + 'a>> {
307        let block_number =
308            self.block_number(block_hash)?.ok_or(ProviderError::BlockHashNotFound(block_hash))?;
309        self.history_by_block_number(block_number)
310    }
311
312    /// Storage provider for state at that given block number
313    pub fn history_by_block_number<'a>(
314        &'a self,
315        mut block_number: BlockNumber,
316    ) -> ProviderResult<Box<dyn StateProvider + 'a>> {
317        if block_number == self.best_block_number().unwrap_or_default() &&
318            block_number == self.last_block_number().unwrap_or_default()
319        {
320            return Ok(Box::new(LatestStateProviderRef::new(self)))
321        }
322
323        // +1 as the changeset that we want is the one that was applied after this block.
324        block_number += 1;
325
326        let account_history_prune_checkpoint =
327            self.get_prune_checkpoint(PruneSegment::AccountHistory)?;
328        let storage_history_prune_checkpoint =
329            self.get_prune_checkpoint(PruneSegment::StorageHistory)?;
330
331        let mut state_provider =
332            HistoricalStateProviderRef::new(self, block_number, self.overlay_manager.clone());
333        // If we pruned account or storage history, we can't return state on every historical block.
334        // Instead, we should cap it at the latest prune checkpoint for corresponding prune segment.
335        if let Some(prune_checkpoint_block_number) =
336            account_history_prune_checkpoint.and_then(|checkpoint| checkpoint.block_number)
337        {
338            state_provider = state_provider.with_lowest_available_account_history_block_number(
339                prune_checkpoint_block_number + 1,
340            );
341        }
342        if let Some(prune_checkpoint_block_number) =
343            storage_history_prune_checkpoint.and_then(|checkpoint| checkpoint.block_number)
344        {
345            state_provider = state_provider.with_lowest_available_storage_history_block_number(
346                prune_checkpoint_block_number + 1,
347            );
348        }
349
350        Ok(Box::new(state_provider))
351    }
352
353    #[cfg(feature = "test-utils")]
354    /// Sets the prune modes for provider.
355    pub fn set_prune_modes(&mut self, prune_modes: PruneModes) {
356        self.prune_modes = prune_modes;
357    }
358}
359
360impl<TX, N: NodeTypes> NodePrimitivesProvider for DatabaseProvider<TX, N> {
361    type Primitives = N::Primitives;
362}
363
364impl<TX, N: NodeTypes> StaticFileProviderFactory for DatabaseProvider<TX, N> {
365    /// Returns a static file provider
366    fn static_file_provider(&self) -> StaticFileProvider<Self::Primitives> {
367        self.static_file_provider.clone()
368    }
369
370    fn get_static_file_writer(
371        &self,
372        block: BlockNumber,
373        segment: StaticFileSegment,
374    ) -> ProviderResult<crate::providers::StaticFileProviderRWRefMut<'_, Self::Primitives>> {
375        self.static_file_provider.get_writer(block, segment)
376    }
377}
378
379impl<TX, N: NodeTypes> RocksDBProviderFactory for DatabaseProvider<TX, N> {
380    /// Returns the `RocksDB` provider.
381    fn rocksdb_provider(&self) -> RocksDBProvider {
382        self.rocksdb_provider.clone()
383    }
384
385    fn set_pending_rocksdb_batch(&self, batch: rocksdb::WriteBatchWithTransaction<true>) {
386        self.pending_rocksdb_batches.lock().push(batch);
387    }
388
389    fn commit_pending_rocksdb_batches(&self) -> ProviderResult<()> {
390        let batches = std::mem::take(&mut *self.pending_rocksdb_batches.lock());
391        for batch in batches {
392            self.rocksdb_provider.commit_batch(batch)?;
393        }
394        Ok(())
395    }
396}
397
398impl<TX: Debug + Send, N: NodeTypes<ChainSpec: EthChainSpec + 'static>> ChainSpecProvider
399    for DatabaseProvider<TX, N>
400{
401    type ChainSpec = N::ChainSpec;
402
403    fn chain_spec(&self) -> Arc<Self::ChainSpec> {
404        self.chain_spec.clone()
405    }
406}
407
408impl<TX: DbTxMut, N: NodeTypes> DatabaseProvider<TX, N> {
409    /// Creates a provider with an inner read-write transaction.
410    #[expect(clippy::too_many_arguments)]
411    fn new_rw_inner(
412        tx: TX,
413        chain_spec: Arc<N::ChainSpec>,
414        static_file_provider: StaticFileProvider<N::Primitives>,
415        prune_modes: PruneModes,
416        storage: Arc<N::Storage>,
417        storage_settings: Arc<RwLock<StorageSettings>>,
418        rocksdb_provider: RocksDBProvider,
419        overlay_manager: OverlayManager<N::Primitives>,
420        runtime: reth_tasks::Runtime,
421        db_path: PathBuf,
422        commit_order: CommitOrder,
423        metrics: Arc<DatabaseProviderMetrics>,
424    ) -> Self {
425        Self {
426            tx,
427            chain_spec,
428            static_file_provider,
429            prune_modes,
430            storage,
431            storage_settings,
432            rocksdb_provider,
433            overlay_manager,
434            runtime,
435            db_path,
436            pending_rocksdb_batches: Default::default(),
437            commit_order,
438            minimum_pruning_distance: MINIMUM_UNWIND_SAFE_DISTANCE,
439            metrics,
440            reader_txn_tracker: None,
441        }
442    }
443
444    /// Creates a provider with an inner read-write transaction using normal commit order.
445    #[expect(clippy::too_many_arguments)]
446    pub fn new_rw(
447        tx: TX,
448        chain_spec: Arc<N::ChainSpec>,
449        static_file_provider: StaticFileProvider<N::Primitives>,
450        prune_modes: PruneModes,
451        storage: Arc<N::Storage>,
452        storage_settings: Arc<RwLock<StorageSettings>>,
453        rocksdb_provider: RocksDBProvider,
454        overlay_manager: OverlayManager<N::Primitives>,
455        runtime: reth_tasks::Runtime,
456        db_path: PathBuf,
457        metrics: Arc<DatabaseProviderMetrics>,
458    ) -> Self {
459        Self::new_rw_inner(
460            tx,
461            chain_spec,
462            static_file_provider,
463            prune_modes,
464            storage,
465            storage_settings,
466            rocksdb_provider,
467            overlay_manager,
468            runtime,
469            db_path,
470            CommitOrder::Normal,
471            metrics,
472        )
473    }
474
475    /// Creates a provider with an inner read-write transaction using unwind commit order.
476    #[expect(clippy::too_many_arguments)]
477    pub fn new_unwind_rw(
478        tx: TX,
479        chain_spec: Arc<N::ChainSpec>,
480        static_file_provider: StaticFileProvider<N::Primitives>,
481        prune_modes: PruneModes,
482        storage: Arc<N::Storage>,
483        storage_settings: Arc<RwLock<StorageSettings>>,
484        rocksdb_provider: RocksDBProvider,
485        overlay_manager: OverlayManager<N::Primitives>,
486        runtime: reth_tasks::Runtime,
487        db_path: PathBuf,
488        metrics: Arc<DatabaseProviderMetrics>,
489    ) -> Self {
490        Self::new_rw_inner(
491            tx,
492            chain_spec,
493            static_file_provider,
494            prune_modes,
495            storage,
496            storage_settings,
497            rocksdb_provider,
498            overlay_manager,
499            runtime,
500            db_path,
501            CommitOrder::Unwind,
502            metrics,
503        )
504    }
505}
506
507impl<TX, N: NodeTypes> AsRef<Self> for DatabaseProvider<TX, N> {
508    fn as_ref(&self) -> &Self {
509        self
510    }
511}
512
513impl<TX: DbTx + DbTxMut + 'static, N: NodeTypesForProvider> DatabaseProvider<TX, N> {
514    /// Executes a closure with a `RocksDB` batch, automatically registering it for commit.
515    ///
516    /// This helper encapsulates all the cfg-gated `RocksDB` batch handling.
517    pub fn with_rocksdb_batch<F, R>(&self, f: F) -> ProviderResult<R>
518    where
519        F: FnOnce(RocksBatchArg<'_>) -> ProviderResult<(R, Option<RawRocksDBBatch>)>,
520    {
521        let rocksdb = self.rocksdb_provider();
522        let rocksdb_batch = rocksdb.batch();
523
524        let (result, raw_batch) = f(rocksdb_batch)?;
525
526        if let Some(batch) = raw_batch {
527            self.set_pending_rocksdb_batch(batch);
528        }
529        let _ = raw_batch; // silence unused warning when rocksdb feature is disabled
530
531        Ok(result)
532    }
533
534    /// Creates the context for static file writes.
535    fn static_file_write_ctx(
536        &self,
537        save_mode: SaveBlocksMode,
538        first_block: BlockNumber,
539        last_block: BlockNumber,
540    ) -> ProviderResult<StaticFileWriteCtx> {
541        let tip = self.last_block_number()?.max(last_block);
542        Ok(StaticFileWriteCtx {
543            write_senders: EitherWriterDestination::senders(self).is_static_file() &&
544                self.prune_modes.sender_recovery.is_none_or(|m| !m.is_full()),
545            write_receipts: save_mode.with_state() &&
546                EitherWriter::receipts_destination(self).is_static_file(),
547            write_account_changesets: save_mode.with_state() &&
548                EitherWriterDestination::account_changesets(self).is_static_file(),
549            write_storage_changesets: save_mode.with_state() &&
550                EitherWriterDestination::storage_changesets(self).is_static_file(),
551            tip,
552            receipts_prune_mode: self.prune_modes.receipts,
553            // Receipts are prunable if no receipts exist in SF yet and within pruning distance
554            receipts_prunable: self
555                .static_file_provider
556                .get_highest_static_file_tx(StaticFileSegment::Receipts)
557                .is_none() &&
558                PruneMode::Distance(self.minimum_pruning_distance)
559                    .should_prune(first_block, tip),
560        })
561    }
562
563    /// Creates the context for `RocksDB` writes.
564    fn rocksdb_write_ctx(&self, first_block: BlockNumber) -> RocksDBWriteCtx {
565        RocksDBWriteCtx {
566            first_block_number: first_block,
567            prune_tx_lookup: self.prune_modes.transaction_lookup,
568            storage_settings: self.cached_storage_settings(),
569            pending_batches: self.pending_rocksdb_batches.clone(),
570        }
571    }
572
573    /// Advances the independent persistence frontiers described by [`SaveBlocksInput`].
574    ///
575    /// Ordinary block data and hashed-state/trie updates advance independently according to the
576    /// ranges derived by the input.
577    #[instrument(level = "debug", target = "providers::db", skip_all, fields(block_count = input.persist_rest_blocks().len()))]
578    pub fn save_blocks(&self, input: &SaveBlocksInput<N::Primitives>) -> ProviderResult<()> {
579        let (db_tip, partial_state_trie) = self
580            .get_stage_checkpoint(StageId::Finish)?
581            .map(|checkpoint| {
582                let partial_state_trie = checkpoint
583                    .finish_stage_checkpoint()
584                    .and_then(|finish| finish.partial_state_trie())
585                    .unwrap_or(checkpoint.block_number);
586                (checkpoint.block_number, partial_state_trie)
587            })
588            .unwrap_or_default();
589
590        if db_tip != input.prev_db_tip() || partial_state_trie != input.prev_partial_state_trie() {
591            return Err(ProviderError::other(std::io::Error::other(format!(
592                "persistence frontiers do not match Finish checkpoint: expected database/state-trie tips #{}/{}, got #{}/{}",
593                input.prev_db_tip(),
594                input.prev_partial_state_trie(),
595                db_tip,
596                partial_state_trie,
597            ))))
598        }
599
600        self.save_blocks_inner(
601            input.persist_rest_blocks(),
602            input.state_trie_blocks(),
603            input.state_trie_masking_blocks(),
604            (input.new_partial_state_trie() < input.new_db_tip())
605                .then_some(input.new_partial_state_trie()),
606            SaveBlocksMode::Full,
607        )
608    }
609
610    fn save_blocks_inner(
611        &self,
612        blocks: &[ExecutedBlock<N::Primitives>],
613        state_trie_blocks: &[ExecutedBlock<N::Primitives>],
614        state_trie_masking_blocks: &[ExecutedBlock<N::Primitives>],
615        partial_state_trie: Option<BlockNumber>,
616        save_mode: SaveBlocksMode,
617    ) -> ProviderResult<()> {
618        let total_start = Instant::now();
619        let block_count = blocks.len() as u64;
620        // With no new block data, the masking suffix still ends at the database tip. If the
621        // masking suffix is empty, the state/trie range ends there instead.
622        let last_block_number = blocks
623            .last()
624            .or_else(|| state_trie_masking_blocks.last())
625            .or_else(|| state_trie_blocks.last())
626            .expect("at least one persistence range must be non-empty")
627            .recovered_block()
628            .number();
629        let first_number = blocks.first().map(|block| block.recovered_block().number());
630
631        debug!(target: "providers::db", block_count, "Writing blocks and execution data to storage");
632
633        // Compute tx_nums upfront (both threads need these)
634        let mut tx_nums: SmallVec<[TxNumber; 4]> = SmallVec::with_capacity(blocks.len());
635        if !blocks.is_empty() {
636            let first_tx_num = self
637                .tx
638                .cursor_read::<tables::TransactionBlocks>()?
639                .last()?
640                .map(|(n, _)| n + 1)
641                .unwrap_or_default();
642            let mut current = first_tx_num;
643            for block in blocks {
644                tx_nums.push(current);
645                current += block.recovered_block().body().transaction_count() as u64;
646            }
647        }
648
649        let mut timings =
650            metrics::SaveBlocksTimings { batch_size: block_count, ..Default::default() };
651
652        // avoid capturing &self.tx in scope below.
653        let sf_provider = &self.static_file_provider;
654        let rocksdb_provider = &self.rocksdb_provider;
655        let sf_ctx = first_number
656            .map(|first_number| {
657                self.static_file_write_ctx(save_mode, first_number, last_block_number)
658            })
659            .transpose()?;
660        let rocksdb_ctx = first_number.map(|first_number| self.rocksdb_write_ctx(first_number));
661        let rocksdb_enabled =
662            rocksdb_ctx.as_ref().is_some_and(|ctx| ctx.storage_settings.storage_v2);
663
664        let mut sf_result = None;
665        let mut rocksdb_result = None;
666
667        // Write to all backends in parallel.
668        let runtime = &self.runtime;
669        // Propagate tracing context into rayon-spawned threads so that static file
670        // and RocksDB write spans appear as children of save_blocks in traces.
671        let span = tracing::Span::current();
672        runtime.storage_pool().in_place_scope(|s| {
673            // SF writes
674            if sf_ctx.is_some() {
675                s.spawn(|_| {
676                    let _guard = span.enter();
677                    let start = Instant::now();
678                    let sf_ctx =
679                        sf_ctx.expect("static file context exists when blocks are persisted");
680                    sf_result = Some(
681                        sf_provider
682                            .write_blocks_data(blocks, &tx_nums, sf_ctx, runtime)
683                            .map(|()| start.elapsed()),
684                    );
685                });
686            }
687
688            // RocksDB writes
689            if rocksdb_enabled {
690                s.spawn(|_| {
691                    let _guard = span.enter();
692                    let start = Instant::now();
693                    let rocksdb_ctx =
694                        rocksdb_ctx.clone().expect("RocksDB context exists when enabled");
695                    rocksdb_result = Some(
696                        rocksdb_provider
697                            .write_blocks_data(blocks, &tx_nums, rocksdb_ctx, runtime)
698                            .map(|()| start.elapsed()),
699                    );
700                });
701            }
702
703            // MDBX writes
704            let mdbx_start = Instant::now();
705
706            // Collect all transaction hashes across all blocks, sort them, and write in batch
707            if !blocks.is_empty() &&
708                !self.cached_storage_settings().storage_v2 &&
709                self.prune_modes.transaction_lookup.is_none_or(|m| !m.is_full())
710            {
711                let start = Instant::now();
712                let total_tx_count: usize =
713                    blocks.iter().map(|b| b.recovered_block().body().transaction_count()).sum();
714                let mut all_tx_hashes = Vec::with_capacity(total_tx_count);
715                for (i, block) in blocks.iter().enumerate() {
716                    let recovered_block = block.recovered_block();
717                    for (tx_num, transaction) in
718                        (tx_nums[i]..).zip(recovered_block.body().transactions_iter())
719                    {
720                        all_tx_hashes.push((*transaction.tx_hash(), tx_num));
721                    }
722                }
723
724                // Sort by hash for optimal MDBX insertion performance
725                all_tx_hashes.sort_unstable_by_key(|(hash, _)| *hash);
726
727                // Write all transaction hash numbers in a single batch
728                self.with_rocksdb_batch(|batch| {
729                    let mut tx_hash_writer =
730                        EitherWriter::new_transaction_hash_numbers(self, batch)?;
731                    tx_hash_writer.put_transaction_hash_numbers_batch(all_tx_hashes, false)?;
732                    let raw_batch = tx_hash_writer.into_raw_rocksdb_batch();
733                    Ok(((), raw_batch))
734                })?;
735                self.metrics.record_duration(
736                    metrics::Action::InsertTransactionHashNumbers,
737                    start.elapsed(),
738                );
739            }
740
741            for (i, block) in blocks.iter().enumerate() {
742                let recovered_block = block.recovered_block();
743
744                let start = Instant::now();
745                self.insert_block_mdbx_only(recovered_block, tx_nums[i])?;
746                timings.insert_block += start.elapsed();
747
748                if save_mode.with_state() {
749                    let execution_output = block.execution_outcome();
750                    let sf_ctx =
751                        sf_ctx.expect("static file context exists when blocks are persisted");
752
753                    // Write state and changesets to the database.
754                    // Must be written after blocks because of the receipt lookup.
755                    // Skip receipts/account changesets if they're being written to static files.
756                    let start = Instant::now();
757                    self.write_state(
758                        WriteStateInput::Single {
759                            outcome: execution_output,
760                            block: recovered_block.number(),
761                        },
762                        OriginalValuesKnown::No,
763                        StateWriteConfig {
764                            write_receipts: !sf_ctx.write_receipts,
765                            write_account_changesets: !sf_ctx.write_account_changesets,
766                            write_storage_changesets: !sf_ctx.write_storage_changesets,
767                        },
768                    )?;
769                    timings.write_state += start.elapsed();
770                }
771            }
772
773            // Write all hashed state and trie updates in single batches.
774            // This reduces cursor open/close overhead from N calls to 1.
775            if save_mode.with_state() && !state_trie_blocks.is_empty() {
776                let start = Instant::now();
777                let batch = ExecutedBlock::hashed_state_refs(state_trie_blocks);
778                let mask = ExecutedBlock::hashed_state_refs(state_trie_masking_blocks);
779                let merged_hashed_state =
780                    HashedPostStateSorted::disjointed_merge_batch(&batch, &mask);
781                if !merged_hashed_state.is_empty() {
782                    self.write_hashed_state(&merged_hashed_state)?;
783                }
784                timings.write_hashed_state += start.elapsed();
785
786                let start = Instant::now();
787                let batch = ExecutedBlock::trie_updates_refs(state_trie_blocks);
788                let mask = ExecutedBlock::trie_updates_refs(state_trie_masking_blocks);
789                // Sparse-trie streaming emits node updates, while the serial state-root path,
790                // including fallback, can emit a whole storage-trie wipe. `disjointed_merge_batch`
791                // rejects wipes, so use the ordered merge. Persisting updates shadowed by the
792                // masking suffix is redundant but safe: the in-memory overlay still takes
793                // precedence.
794                //
795                // With a revm release containing bluealloy/revm#3863, post-Cancun selfdestructs
796                // will no longer result in `storage.is_deleted` in serial trie updates. The flag
797                // remains valid `save_blocks` input when processing pre-Cancun historical data.
798                let contains_storage_wipe = batch.iter().chain(&mask).any(|updates| {
799                    updates.storage_tries_ref().values().any(|storage| storage.is_deleted)
800                });
801                let merged_trie = if contains_storage_wipe {
802                    TrieUpdatesSorted::merge_batch(
803                        state_trie_blocks.iter().rev().map(|block| block.trie_updates()),
804                    )
805                } else {
806                    Arc::new(TrieUpdatesSorted::disjointed_merge_batch(&batch, &mask))
807                };
808                if !merged_trie.is_empty() {
809                    self.write_trie_updates_sorted(&merged_trie)?;
810                }
811                timings.write_trie_updates += start.elapsed();
812            }
813
814            // Full mode: update history indices
815            if save_mode.with_state() &&
816                let Some(first_number) = first_number
817            {
818                let start = Instant::now();
819                self.update_history_indices(first_number..=last_block_number)?;
820                timings.update_history_indices = start.elapsed();
821            }
822
823            // Update pipeline progress
824            let start = Instant::now();
825            if !blocks.is_empty() {
826                self.update_pipeline_stages(last_block_number, false)?;
827            }
828            if save_mode.with_state() {
829                let checkpoint = match partial_state_trie {
830                    Some(partial_state_trie) => StageCheckpoint::new(last_block_number)
831                        .with_finish_stage_checkpoint(FinishCheckpoint {
832                            partial_state_trie: Some(partial_state_trie),
833                        }),
834                    None => StageCheckpoint::new(last_block_number),
835                };
836                self.save_stage_checkpoint(StageId::Finish, checkpoint)?;
837            }
838            timings.update_pipeline_stages = start.elapsed();
839
840            timings.mdbx = mdbx_start.elapsed();
841
842            Ok::<_, ProviderError>(())
843        })?;
844
845        // Collect results from spawned tasks
846        if !blocks.is_empty() {
847            timings.sf = sf_result.ok_or(StaticFileWriterError::ThreadPanic("static file"))??;
848        }
849
850        if rocksdb_enabled {
851            timings.rocksdb = rocksdb_result.ok_or_else(|| {
852                ProviderError::Database(reth_db_api::DatabaseError::Other(
853                    "RocksDB thread panicked".into(),
854                ))
855            })??;
856        }
857
858        timings.total = total_start.elapsed();
859
860        self.metrics.record_save_blocks(&timings);
861        if let Some(first_number) = first_number {
862            debug!(target: "providers::db", range = ?first_number..=last_block_number, "Appended block data");
863        }
864
865        Ok(())
866    }
867
868    /// Writes MDBX-only data for a block (indices, lookups, and senders if configured for MDBX).
869    ///
870    /// SF data (headers, transactions, senders if SF, receipts if SF) must be written separately.
871    #[instrument(level = "debug", target = "providers::db", skip_all)]
872    fn insert_block_mdbx_only(
873        &self,
874        block: &RecoveredBlock<BlockTy<N>>,
875        first_tx_num: TxNumber,
876    ) -> ProviderResult<StoredBlockBodyIndices> {
877        if self.prune_modes.sender_recovery.is_none_or(|m| !m.is_full()) &&
878            EitherWriterDestination::senders(self).is_database()
879        {
880            let start = Instant::now();
881            let tx_nums_iter = std::iter::successors(Some(first_tx_num), |n| Some(n + 1));
882            let mut cursor = self.tx.cursor_write::<tables::TransactionSenders>()?;
883            for (tx_num, sender) in tx_nums_iter.zip(block.senders_iter().copied()) {
884                cursor.append(tx_num, &sender)?;
885            }
886            self.metrics
887                .record_duration(metrics::Action::InsertTransactionSenders, start.elapsed());
888        }
889
890        let block_number = block.number();
891        let tx_count = block.body().transaction_count() as u64;
892
893        let start = Instant::now();
894        self.tx.put::<tables::HeaderNumbers>(block.hash(), block_number)?;
895        self.metrics.record_duration(metrics::Action::InsertHeaderNumbers, start.elapsed());
896
897        self.write_block_body_indices(block_number, block.body(), first_tx_num, tx_count)?;
898
899        Ok(StoredBlockBodyIndices { first_tx_num, tx_count })
900    }
901
902    /// Writes MDBX block body indices (`BlockBodyIndices`, `TransactionBlocks`,
903    /// `Ommers`/`Withdrawals`).
904    fn write_block_body_indices(
905        &self,
906        block_number: BlockNumber,
907        body: &BodyTy<N>,
908        first_tx_num: TxNumber,
909        tx_count: u64,
910    ) -> ProviderResult<()> {
911        // MDBX: BlockBodyIndices
912        let start = Instant::now();
913        self.tx
914            .cursor_write::<tables::BlockBodyIndices>()?
915            .append(block_number, &StoredBlockBodyIndices { first_tx_num, tx_count })?;
916        self.metrics.record_duration(metrics::Action::InsertBlockBodyIndices, start.elapsed());
917
918        // MDBX: TransactionBlocks (last tx -> block mapping)
919        if tx_count > 0 {
920            let start = Instant::now();
921            self.tx
922                .cursor_write::<tables::TransactionBlocks>()?
923                .append(first_tx_num + tx_count - 1, &block_number)?;
924            self.metrics.record_duration(metrics::Action::InsertTransactionBlocks, start.elapsed());
925        }
926
927        // MDBX: Ommers/Withdrawals
928        self.storage.writer().write_block_bodies(self, vec![(block_number, Some(body))])?;
929
930        Ok(())
931    }
932
933    /// Unwinds trie state starting at and including the given block.
934    ///
935    /// This includes calculating the resulted state root and comparing it with the parent block
936    /// state root.
937    pub fn unwind_trie_state_from(&self, from: BlockNumber) -> ProviderResult<()> {
938        let changed_accounts = self.account_changesets_range(from..)?;
939
940        // Unwind account hashes.
941        self.unwind_account_hashing(changed_accounts.iter())?;
942
943        // Unwind account history indices.
944        self.unwind_account_history_indices(changed_accounts.iter())?;
945
946        let changed_storages = self.storage_changesets_range(from..)?;
947
948        // Unwind storage hashes.
949        self.unwind_storage_hashing(changed_storages.iter().copied())?;
950
951        // Unwind storage history indices.
952        self.unwind_storage_history_indices(changed_storages.iter().copied())?;
953
954        // Unwind accounts/storages trie tables using the revert.
955        // Get the database tip block number
956        let db_tip_block = self
957            .get_stage_checkpoint(reth_stages_types::StageId::Finish)?
958            .as_ref()
959            .map(|chk| chk.block_number)
960            .ok_or_else(|| ProviderError::InsufficientChangesets {
961                requested: from,
962                available: 0..=0,
963            })?;
964
965        let trie_revert = self
966            .overlay_manager
967            .get_or_compute_cached_changesets_range(self, from..=db_tip_block)?;
968        self.write_trie_updates_sorted(&trie_revert)?;
969
970        Ok(())
971    }
972
973    /// Removes receipts from all transactions starting with provided number (inclusive).
974    fn remove_receipts_from(
975        &self,
976        from_tx: TxNumber,
977        last_block: BlockNumber,
978    ) -> ProviderResult<()> {
979        // iterate over block body and remove receipts
980        self.remove::<tables::Receipts<ReceiptTy<N>>>(from_tx..)?;
981
982        if EitherWriter::receipts_destination(self).is_static_file() {
983            let static_file_receipt_num =
984                self.static_file_provider.get_highest_static_file_tx(StaticFileSegment::Receipts);
985
986            let to_delete = static_file_receipt_num
987                .map(|static_num| (static_num + 1).saturating_sub(from_tx))
988                .unwrap_or_default();
989
990            self.static_file_provider
991                .latest_writer(StaticFileSegment::Receipts)?
992                .prune_receipts(to_delete, last_block)?;
993        }
994
995        Ok(())
996    }
997
998    /// Writes bytecodes to MDBX.
999    fn write_bytecodes(
1000        &self,
1001        bytecodes: impl IntoIterator<Item = (B256, Bytecode)>,
1002    ) -> ProviderResult<()> {
1003        let mut bytecodes_cursor = self.tx_ref().cursor_write::<tables::Bytecodes>()?;
1004        for (hash, bytecode) in bytecodes {
1005            bytecodes_cursor.upsert(hash, &bytecode)?;
1006        }
1007        Ok(())
1008    }
1009}
1010
1011impl<TX: DbTx + 'static, N: NodeTypes> TryIntoHistoricalStateProvider for DatabaseProvider<TX, N> {
1012    fn try_into_history_at_block(
1013        self,
1014        mut block_number: BlockNumber,
1015    ) -> ProviderResult<StateProviderBox> {
1016        let best_block = self.best_block_number().unwrap_or_default();
1017
1018        // Reject requests for blocks beyond the best block
1019        if block_number > best_block {
1020            return Err(ProviderError::BlockNotExecuted {
1021                requested: block_number,
1022                executed: best_block,
1023            });
1024        }
1025
1026        // If requesting state at the best block, use the latest state provider
1027        if block_number == best_block {
1028            return Ok(Box::new(LatestStateProvider::new(self)));
1029        }
1030
1031        // +1 as the changeset that we want is the one that was applied after this block.
1032        block_number += 1;
1033
1034        let account_history_prune_checkpoint =
1035            self.get_prune_checkpoint(PruneSegment::AccountHistory)?;
1036        let storage_history_prune_checkpoint =
1037            self.get_prune_checkpoint(PruneSegment::StorageHistory)?;
1038        let overlay_manager = self.overlay_manager.clone();
1039
1040        let mut state_provider = HistoricalStateProvider::new(self, block_number, overlay_manager);
1041
1042        // If we pruned account or storage history, we can't return state on every historical block.
1043        // Instead, we should cap it at the latest prune checkpoint for corresponding prune segment.
1044        if let Some(prune_checkpoint_block_number) =
1045            account_history_prune_checkpoint.and_then(|checkpoint| checkpoint.block_number)
1046        {
1047            state_provider = state_provider.with_lowest_available_account_history_block_number(
1048                prune_checkpoint_block_number + 1,
1049            );
1050        }
1051        if let Some(prune_checkpoint_block_number) =
1052            storage_history_prune_checkpoint.and_then(|checkpoint| checkpoint.block_number)
1053        {
1054            state_provider = state_provider.with_lowest_available_storage_history_block_number(
1055                prune_checkpoint_block_number + 1,
1056            );
1057        }
1058
1059        Ok(Box::new(state_provider))
1060    }
1061}
1062
1063/// For a given key, unwind all history shards that contain block numbers at or above the given
1064/// block number.
1065///
1066/// S - Sharded key subtype.
1067/// T - Table to walk over.
1068/// C - Cursor implementation.
1069///
1070/// This function walks the entries from the given start key and deletes all shards that belong to
1071/// the key and contain block numbers at or above the given block number. Shards entirely below
1072/// the block number are preserved.
1073///
1074/// The boundary shard (the shard that spans across the block number) is removed from the database.
1075/// Any indices that are below the block number are filtered out and returned for reinsertion.
1076/// The boundary shard is returned for reinsertion (if it's not empty).
1077fn unwind_history_shards<S, T, C>(
1078    cursor: &mut C,
1079    start_key: T::Key,
1080    block_number: BlockNumber,
1081    mut shard_belongs_to_key: impl FnMut(&T::Key) -> bool,
1082) -> ProviderResult<Vec<u64>>
1083where
1084    T: Table<Value = BlockNumberList>,
1085    T::Key: AsRef<ShardedKey<S>>,
1086    C: DbCursorRO<T> + DbCursorRW<T>,
1087{
1088    // Start from the given key and iterate through shards
1089    let mut item = cursor.seek_exact(start_key)?;
1090    while let Some((sharded_key, list)) = item {
1091        // If the shard does not belong to the key, break.
1092        if !shard_belongs_to_key(&sharded_key) {
1093            break
1094        }
1095
1096        // Always delete the current shard from the database first
1097        // We'll decide later what (if anything) to reinsert
1098        cursor.delete_current()?;
1099
1100        // Get the first (lowest) block number in this shard
1101        // All block numbers in a shard are sorted in ascending order
1102        let first = list.iter().next().expect("List can't be empty");
1103
1104        // Case 1: Entire shard is at or above the unwinding point
1105        // Keep it deleted (don't return anything for reinsertion)
1106        if first >= block_number {
1107            item = cursor.prev()?;
1108            continue
1109        }
1110        // Case 2: This is a boundary shard (spans across the unwinding point)
1111        // The shard contains some blocks below and some at/above the unwinding point
1112        else if block_number <= sharded_key.as_ref().highest_block_number {
1113            // Return only the block numbers that are below the unwinding point
1114            // These will be reinserted to preserve the historical data
1115            return Ok(list.iter().take_while(|i| *i < block_number).collect::<Vec<_>>())
1116        }
1117        // Case 3: Entire shard is below the unwinding point
1118        // Return all block numbers for reinsertion (preserve entire shard)
1119        return Ok(list.iter().collect::<Vec<_>>())
1120    }
1121
1122    // No shards found or all processed
1123    Ok(Vec::new())
1124}
1125
1126impl<TX: DbTx + 'static, N: NodeTypesForProvider> DatabaseProvider<TX, N> {
1127    /// Creates a provider with an inner read-only transaction.
1128    #[expect(clippy::too_many_arguments)]
1129    pub fn new(
1130        tx: TX,
1131        chain_spec: Arc<N::ChainSpec>,
1132        static_file_provider: StaticFileProvider<N::Primitives>,
1133        prune_modes: PruneModes,
1134        storage: Arc<N::Storage>,
1135        storage_settings: Arc<RwLock<StorageSettings>>,
1136        rocksdb_provider: RocksDBProvider,
1137        overlay_manager: OverlayManager<N::Primitives>,
1138        runtime: reth_tasks::Runtime,
1139        db_path: PathBuf,
1140        metrics: Arc<DatabaseProviderMetrics>,
1141    ) -> Self {
1142        Self {
1143            tx,
1144            chain_spec,
1145            static_file_provider,
1146            prune_modes,
1147            storage,
1148            storage_settings,
1149            rocksdb_provider,
1150            overlay_manager,
1151            runtime,
1152            db_path,
1153            pending_rocksdb_batches: Default::default(),
1154            commit_order: CommitOrder::Normal,
1155            minimum_pruning_distance: MINIMUM_UNWIND_SAFE_DISTANCE,
1156            metrics,
1157            reader_txn_tracker: None,
1158        }
1159    }
1160
1161    /// Consume `DbTx` or `DbTxMut`.
1162    pub fn into_tx(self) -> TX {
1163        self.tx
1164    }
1165
1166    /// Pass `DbTx` or `DbTxMut` mutable reference.
1167    pub const fn tx_mut(&mut self) -> &mut TX {
1168        &mut self.tx
1169    }
1170
1171    /// Pass `DbTx` or `DbTxMut` immutable reference.
1172    pub const fn tx_ref(&self) -> &TX {
1173        &self.tx
1174    }
1175
1176    /// Returns a reference to the chain specification.
1177    pub fn chain_spec(&self) -> &N::ChainSpec {
1178        &self.chain_spec
1179    }
1180}
1181
1182impl<TX: DbTx + 'static, N: NodeTypesForProvider> DatabaseProvider<TX, N> {
1183    fn recovered_block<H, HF, B, BF>(
1184        &self,
1185        id: BlockHashOrNumber,
1186        _transaction_kind: TransactionVariant,
1187        header_by_number: HF,
1188        construct_block: BF,
1189    ) -> ProviderResult<Option<B>>
1190    where
1191        H: AsRef<HeaderTy<N>>,
1192        HF: FnOnce(BlockNumber) -> ProviderResult<Option<H>>,
1193        BF: FnOnce(H, BodyTy<N>, Vec<Address>) -> ProviderResult<Option<B>>,
1194    {
1195        let Some(block_number) = self.convert_hash_or_number(id)? else { return Ok(None) };
1196        let earliest_available = self.static_file_provider.earliest_history_height();
1197        if block_number < earliest_available {
1198            return Err(ProviderError::BlockExpired { requested: block_number, earliest_available })
1199        }
1200        let Some(header) = header_by_number(block_number)? else { return Ok(None) };
1201
1202        // Get the block body
1203        //
1204        // If the body indices are not found, this means that the transactions either do not exist
1205        // in the database yet, or they do exit but are not indexed. If they exist but are not
1206        // indexed, we don't have enough information to return the block anyways, so we return
1207        // `None`.
1208        let Some(body) = self.block_body_indices(block_number)? else { return Ok(None) };
1209
1210        let tx_range = body.tx_num_range();
1211
1212        let transactions = if tx_range.is_empty() {
1213            vec![]
1214        } else {
1215            self.transactions_by_tx_range(tx_range.clone())?
1216        };
1217
1218        let body = self
1219            .storage
1220            .reader()
1221            .read_block_bodies(self, vec![(header.as_ref(), transactions)])?
1222            .pop()
1223            .ok_or(ProviderError::InvalidStorageOutput)?;
1224
1225        let senders = if tx_range.is_empty() {
1226            vec![]
1227        } else {
1228            let known_senders: HashMap<TxNumber, Address> =
1229                EitherReader::new_senders(self)?.senders_by_tx_range(tx_range.clone())?;
1230
1231            let mut senders = Vec::with_capacity(body.transactions().len());
1232            for (tx_num, tx) in tx_range.zip(body.transactions()) {
1233                match known_senders.get(&tx_num) {
1234                    None => {
1235                        let sender = tx.recover_signer_unchecked()?;
1236                        senders.push(sender);
1237                    }
1238                    Some(sender) => senders.push(*sender),
1239                }
1240            }
1241            senders
1242        };
1243
1244        construct_block(header, body, senders)
1245    }
1246
1247    /// Returns a range of blocks from the database.
1248    ///
1249    /// Uses the provided `headers_range` to get the headers for the range, and `assemble_block` to
1250    /// construct blocks from the following inputs:
1251    ///     – Header
1252    ///     - Range of transaction numbers
1253    ///     – Ommers
1254    ///     – Withdrawals
1255    ///     – Senders
1256    fn block_range<F, H, HF, R>(
1257        &self,
1258        range: RangeInclusive<BlockNumber>,
1259        headers_range: HF,
1260        mut assemble_block: F,
1261    ) -> ProviderResult<Vec<R>>
1262    where
1263        H: AsRef<HeaderTy<N>>,
1264        HF: FnOnce(RangeInclusive<BlockNumber>) -> ProviderResult<Vec<H>>,
1265        F: FnMut(H, BodyTy<N>, Range<TxNumber>) -> ProviderResult<R>,
1266    {
1267        if range.is_empty() {
1268            return Ok(Vec::new())
1269        }
1270
1271        let len = range.end().saturating_sub(*range.start()) as usize + 1;
1272        let mut blocks = Vec::with_capacity(len);
1273
1274        let headers = headers_range(range.clone())?;
1275
1276        // If the body indices are not found, this means that the transactions either do
1277        // not exist in the database yet, or they do exit but are
1278        // not indexed. If they exist but are not indexed, we don't
1279        // have enough information to return the block anyways, so
1280        // we skip the block.
1281        let present_headers = self
1282            .block_body_indices_range(range)?
1283            .into_iter()
1284            .map(|b| b.tx_num_range())
1285            .zip(headers)
1286            .collect::<Vec<_>>();
1287
1288        let mut inputs = Vec::with_capacity(present_headers.len());
1289        for (tx_range, header) in &present_headers {
1290            let transactions = if tx_range.is_empty() {
1291                Vec::new()
1292            } else {
1293                self.transactions_by_tx_range(tx_range.clone())?
1294            };
1295
1296            inputs.push((header.as_ref(), transactions));
1297        }
1298
1299        let bodies = self.storage.reader().read_block_bodies(self, inputs)?;
1300
1301        for ((tx_range, header), body) in present_headers.into_iter().zip(bodies) {
1302            blocks.push(assemble_block(header, body, tx_range)?);
1303        }
1304
1305        Ok(blocks)
1306    }
1307
1308    /// Returns a range of blocks from the database, along with the senders of each
1309    /// transaction in the blocks.
1310    ///
1311    /// Uses the provided `headers_range` to get the headers for the range, and `assemble_block` to
1312    /// construct blocks from the following inputs:
1313    ///     – Header
1314    ///     - Transactions
1315    ///     – Ommers
1316    ///     – Withdrawals
1317    ///     – Senders
1318    fn block_with_senders_range<H, HF, B, BF>(
1319        &self,
1320        range: RangeInclusive<BlockNumber>,
1321        headers_range: HF,
1322        assemble_block: BF,
1323    ) -> ProviderResult<Vec<B>>
1324    where
1325        H: AsRef<HeaderTy<N>>,
1326        HF: Fn(RangeInclusive<BlockNumber>) -> ProviderResult<Vec<H>>,
1327        BF: Fn(H, BodyTy<N>, Vec<Address>) -> ProviderResult<B>,
1328    {
1329        self.block_range(range, headers_range, |header, body, tx_range| {
1330            let senders = if tx_range.is_empty() {
1331                Vec::new()
1332            } else {
1333                let known_senders: HashMap<TxNumber, Address> =
1334                    EitherReader::new_senders(self)?.senders_by_tx_range(tx_range.clone())?;
1335
1336                let mut senders = Vec::with_capacity(body.transactions().len());
1337                for (tx_num, tx) in tx_range.zip(body.transactions()) {
1338                    match known_senders.get(&tx_num) {
1339                        None => {
1340                            // recover the sender from the transaction if not found
1341                            let sender = tx.recover_signer_unchecked()?;
1342                            senders.push(sender);
1343                        }
1344                        Some(sender) => senders.push(*sender),
1345                    }
1346                }
1347
1348                senders
1349            };
1350
1351            assemble_block(header, body, senders)
1352        })
1353    }
1354
1355    /// Populate a [`BundleStateInit`] and [`RevertsInit`] using cursors over the
1356    /// [`tables::PlainAccountState`] and [`tables::PlainStorageState`] tables, based on the given
1357    /// storage and account changesets.
1358    fn populate_bundle_state(
1359        &self,
1360        account_changeset: Vec<(u64, AccountBeforeTx)>,
1361        storage_changeset: Vec<(BlockNumberAddress, StorageEntry)>,
1362        mut get_account: impl FnMut(Address) -> ProviderResult<Option<Account>>,
1363        mut get_storage: impl FnMut(Address, StorageKey) -> ProviderResult<Option<StorageValue>>,
1364    ) -> ProviderResult<(BundleStateInit, RevertsInit)> {
1365        // iterate previous value and get plain state value to create changeset
1366        // Double option around Account represent if Account state is know (first option) and
1367        // account is removed (Second Option)
1368        let mut state: BundleStateInit = HashMap::default();
1369
1370        // This is not working for blocks that are not at tip. as plain state is not the last
1371        // state of end range. We should rename the functions or add support to access
1372        // History state. Accessing history state can be tricky but we are not gaining
1373        // anything.
1374
1375        let mut reverts: RevertsInit = HashMap::default();
1376
1377        // add account changeset changes
1378        for (block_number, account_before) in account_changeset.into_iter().rev() {
1379            let AccountBeforeTx { info: old_info, address } = account_before;
1380            match state.entry(address) {
1381                hash_map::Entry::Vacant(entry) => {
1382                    let new_info = get_account(address)?;
1383                    entry.insert((old_info, new_info, HashMap::default()));
1384                }
1385                hash_map::Entry::Occupied(mut entry) => {
1386                    // overwrite old account state.
1387                    entry.get_mut().0 = old_info;
1388                }
1389            }
1390            // insert old info into reverts.
1391            reverts.entry(block_number).or_default().entry(address).or_default().0 = Some(old_info);
1392        }
1393
1394        // add storage changeset changes
1395        for (block_and_address, old_storage) in storage_changeset.into_iter().rev() {
1396            let BlockNumberAddress((block_number, address)) = block_and_address;
1397            // get account state or insert from plain state.
1398            let account_state = match state.entry(address) {
1399                hash_map::Entry::Vacant(entry) => {
1400                    let present_info = get_account(address)?;
1401                    entry.insert((present_info, present_info, HashMap::default()))
1402                }
1403                hash_map::Entry::Occupied(entry) => entry.into_mut(),
1404            };
1405
1406            // match storage.
1407            match account_state.2.entry(old_storage.key) {
1408                hash_map::Entry::Vacant(entry) => {
1409                    let new_storage = get_storage(address, old_storage.key)?.unwrap_or_default();
1410                    entry.insert((old_storage.value, new_storage));
1411                }
1412                hash_map::Entry::Occupied(mut entry) => {
1413                    entry.get_mut().0 = old_storage.value;
1414                }
1415            };
1416
1417            reverts
1418                .entry(block_number)
1419                .or_default()
1420                .entry(address)
1421                .or_default()
1422                .1
1423                .push(old_storage);
1424        }
1425
1426        Ok((state, reverts))
1427    }
1428
1429    /// Invokes [`populate_bundle_state`](Self::populate_bundle_state) with the given plain state
1430    /// cursors.
1431    fn populate_bundle_state_plain(
1432        &self,
1433        account_changeset: Vec<(u64, AccountBeforeTx)>,
1434        storage_changeset: Vec<(BlockNumberAddress, StorageEntry)>,
1435        plain_accounts_cursor: &mut impl DbCursorRO<tables::PlainAccountState>,
1436        plain_storage_cursor: &mut impl DbDupCursorRO<tables::PlainStorageState>,
1437    ) -> ProviderResult<(BundleStateInit, RevertsInit)> {
1438        self.populate_bundle_state(
1439            account_changeset,
1440            storage_changeset,
1441            |address| Ok(plain_accounts_cursor.seek_exact(address)?.map(|kv| kv.1)),
1442            |address, storage_key| {
1443                Ok(plain_storage_cursor
1444                    .seek_by_key_subkey(address, storage_key)?
1445                    .filter(|s| s.key == storage_key)
1446                    .map(|s| s.value))
1447            },
1448        )
1449    }
1450
1451    /// Like [`populate_bundle_state`](Self::populate_bundle_state), but reads current values from
1452    /// `HashedAccounts`/`HashedStorages`. Addresses and storage keys are hashed via `keccak256`
1453    /// for DB lookups. The output `BundleStateInit`/`RevertsInit` structures remain keyed by
1454    /// plain address and plain storage key.
1455    fn populate_bundle_state_hashed(
1456        &self,
1457        account_changeset: Vec<(u64, AccountBeforeTx)>,
1458        storage_changeset: Vec<(BlockNumberAddress, StorageEntry)>,
1459        hashed_accounts_cursor: &mut impl DbCursorRO<tables::HashedAccounts>,
1460        hashed_storage_cursor: &mut impl DbDupCursorRO<tables::HashedStorages>,
1461    ) -> ProviderResult<(BundleStateInit, RevertsInit)> {
1462        self.populate_bundle_state(
1463            account_changeset,
1464            storage_changeset,
1465            |address| Ok(hashed_accounts_cursor.seek_exact(keccak256(address))?.map(|kv| kv.1)),
1466            |address, storage_key| {
1467                let hashed_storage_key = keccak256(storage_key);
1468                Ok(hashed_storage_cursor
1469                    .seek_by_key_subkey(keccak256(address), hashed_storage_key)?
1470                    .filter(|s| s.key == hashed_storage_key)
1471                    .map(|s| s.value))
1472            },
1473        )
1474    }
1475
1476    fn populate_bundle_state_with_provider(
1477        &self,
1478        account_changeset: Vec<(u64, AccountBeforeTx)>,
1479        storage_changeset: Vec<(BlockNumberAddress, StorageEntry)>,
1480        state_provider: impl StateProvider,
1481    ) -> ProviderResult<(BundleStateInit, RevertsInit)> {
1482        self.populate_bundle_state(
1483            account_changeset,
1484            storage_changeset,
1485            |address| state_provider.basic_account(&address),
1486            |address, storage_key| state_provider.storage(address, storage_key),
1487        )
1488    }
1489}
1490
1491impl<TX: DbTxMut + DbTx + 'static, N: NodeTypes> DatabaseProvider<TX, N> {
1492    /// Insert history index to the database.
1493    ///
1494    /// For each updated partial key, this function retrieves the last shard from the database
1495    /// (if any), appends the new indices to it, chunks the resulting list if needed, and upserts
1496    /// the shards back into the database.
1497    ///
1498    /// This function is used by history indexing stages.
1499    fn append_history_index<P, T>(
1500        &self,
1501        index_updates: impl IntoIterator<Item = (P, impl IntoIterator<Item = u64>)>,
1502        mut sharded_key_factory: impl FnMut(P, BlockNumber) -> T::Key,
1503    ) -> ProviderResult<()>
1504    where
1505        P: Copy,
1506        T: Table<Value = BlockNumberList>,
1507    {
1508        // This function cannot be used with DUPSORT tables because `upsert` on DUPSORT tables
1509        // will append duplicate entries instead of updating existing ones, causing data corruption.
1510        assert!(!T::DUPSORT, "append_history_index cannot be used with DUPSORT tables");
1511
1512        let mut cursor = self.tx.cursor_write::<T>()?;
1513
1514        for (partial_key, indices) in index_updates {
1515            let last_key = sharded_key_factory(partial_key, u64::MAX);
1516            let mut last_shard = cursor
1517                .seek_exact(last_key.clone())?
1518                .map(|(_, list)| list)
1519                .unwrap_or_else(BlockNumberList::empty);
1520
1521            last_shard.append(indices).map_err(ProviderError::other)?;
1522
1523            // fast path: all indices fit in one shard
1524            if last_shard.len() <= sharded_key::NUM_OF_INDICES_IN_SHARD as u64 {
1525                cursor.upsert(last_key, &last_shard)?;
1526                continue;
1527            }
1528
1529            // slow path: rechunk into multiple shards
1530            let chunks = last_shard.iter().chunks(sharded_key::NUM_OF_INDICES_IN_SHARD);
1531            let mut chunks_peekable = chunks.into_iter().peekable();
1532
1533            while let Some(chunk) = chunks_peekable.next() {
1534                let shard = BlockNumberList::new_pre_sorted(chunk);
1535                let highest_block_number = if chunks_peekable.peek().is_some() {
1536                    shard.iter().next_back().expect("`chunks` does not return empty list")
1537                } else {
1538                    // Insert last list with `u64::MAX`.
1539                    u64::MAX
1540                };
1541
1542                cursor.upsert(sharded_key_factory(partial_key, highest_block_number), &shard)?;
1543            }
1544        }
1545
1546        Ok(())
1547    }
1548}
1549
1550impl<TX: DbTx, N: NodeTypes> AccountReader for DatabaseProvider<TX, N> {
1551    fn basic_account(&self, address: &Address) -> ProviderResult<Option<Account>> {
1552        if self.cached_storage_settings().use_hashed_state() {
1553            let hashed_address = keccak256(address);
1554            Ok(self.tx.get_by_encoded_key::<tables::HashedAccounts>(&hashed_address)?)
1555        } else {
1556            Ok(self.tx.get_by_encoded_key::<tables::PlainAccountState>(address)?)
1557        }
1558    }
1559}
1560
1561impl<TX: DbTx + 'static, N: NodeTypes> AccountExtReader for DatabaseProvider<TX, N> {
1562    fn changed_accounts_with_range(
1563        &self,
1564        range: RangeInclusive<BlockNumber>,
1565    ) -> ProviderResult<BTreeSet<Address>> {
1566        let mut reader = EitherReader::new_account_changesets(self)?;
1567
1568        reader.changed_accounts_with_range(range)
1569    }
1570
1571    fn basic_accounts(
1572        &self,
1573        iter: impl IntoIterator<Item = Address>,
1574    ) -> ProviderResult<Vec<(Address, Option<Account>)>> {
1575        if self.cached_storage_settings().use_hashed_state() {
1576            let mut hashed_accounts = self.tx.cursor_read::<tables::HashedAccounts>()?;
1577            Ok(iter
1578                .into_iter()
1579                .map(|address| {
1580                    let hashed_address = keccak256(address);
1581                    hashed_accounts.seek_exact(hashed_address).map(|a| (address, a.map(|(_, v)| v)))
1582                })
1583                .collect::<Result<Vec<_>, _>>()?)
1584        } else {
1585            let mut plain_accounts = self.tx.cursor_read::<tables::PlainAccountState>()?;
1586            Ok(iter
1587                .into_iter()
1588                .map(|address| {
1589                    plain_accounts.seek_exact(address).map(|a| (address, a.map(|(_, v)| v)))
1590                })
1591                .collect::<Result<Vec<_>, _>>()?)
1592        }
1593    }
1594
1595    fn changed_accounts_and_blocks_with_range(
1596        &self,
1597        range: RangeInclusive<BlockNumber>,
1598    ) -> ProviderResult<BTreeMap<Address, Vec<u64>>> {
1599        let highest_static_block = self
1600            .static_file_provider
1601            .get_highest_static_file_block(StaticFileSegment::AccountChangeSets);
1602
1603        if let Some(highest) = highest_static_block &&
1604            self.cached_storage_settings().storage_v2
1605        {
1606            let start = *range.start();
1607            let static_end = (*range.end()).min(highest);
1608
1609            let mut changed_accounts_and_blocks: BTreeMap<_, Vec<u64>> = BTreeMap::default();
1610            if start <= static_end {
1611                for block in start..=static_end {
1612                    let block_changesets = self.account_block_changeset(block)?;
1613                    for changeset in block_changesets {
1614                        changed_accounts_and_blocks
1615                            .entry(changeset.address)
1616                            .or_default()
1617                            .push(block);
1618                    }
1619                }
1620            }
1621
1622            Ok(changed_accounts_and_blocks)
1623        } else {
1624            let mut changeset_cursor = self.tx.cursor_read::<tables::AccountChangeSets>()?;
1625
1626            let account_transitions = changeset_cursor.walk_range(range)?.try_fold(
1627                BTreeMap::new(),
1628                |mut accounts: BTreeMap<Address, Vec<u64>>, entry| -> ProviderResult<_> {
1629                    let (index, account) = entry?;
1630                    accounts.entry(account.address).or_default().push(index);
1631                    Ok(accounts)
1632                },
1633            )?;
1634
1635            Ok(account_transitions)
1636        }
1637    }
1638}
1639
1640impl<TX: DbTx, N: NodeTypes> StorageChangeSetReader for DatabaseProvider<TX, N> {
1641    fn storage_changeset(
1642        &self,
1643        block_number: BlockNumber,
1644    ) -> ProviderResult<Vec<(BlockNumberAddress, StorageEntry)>> {
1645        if self.cached_storage_settings().storage_v2 {
1646            self.static_file_provider.storage_changeset(block_number)
1647        } else {
1648            let range = block_number..=block_number;
1649            let storage_range = BlockNumberAddress::range(range);
1650            self.tx
1651                .cursor_dup_read::<tables::StorageChangeSets>()?
1652                .walk_range(storage_range)?
1653                .map(|r| {
1654                    let (bna, entry) = r?;
1655                    Ok((bna, entry))
1656                })
1657                .collect()
1658        }
1659    }
1660
1661    fn get_storage_before_block(
1662        &self,
1663        block_number: BlockNumber,
1664        address: Address,
1665        storage_key: B256,
1666    ) -> ProviderResult<Option<StorageEntry>> {
1667        if self.cached_storage_settings().storage_v2 {
1668            self.static_file_provider.get_storage_before_block(block_number, address, storage_key)
1669        } else {
1670            Ok(self
1671                .tx
1672                .cursor_dup_read::<tables::StorageChangeSets>()?
1673                .seek_by_key_subkey(BlockNumberAddress((block_number, address)), storage_key)?
1674                .filter(|entry| entry.key == storage_key))
1675        }
1676    }
1677
1678    fn storage_changesets_range(
1679        &self,
1680        range: impl RangeBounds<BlockNumber>,
1681    ) -> ProviderResult<Vec<(BlockNumberAddress, StorageEntry)>> {
1682        if self.cached_storage_settings().storage_v2 {
1683            self.static_file_provider.storage_changesets_range(range)
1684        } else {
1685            self.tx
1686                .cursor_dup_read::<tables::StorageChangeSets>()?
1687                .walk_range(BlockNumberAddressRange::from(range))?
1688                .map(|r| {
1689                    let (bna, entry) = r?;
1690                    Ok((bna, entry))
1691                })
1692                .collect()
1693        }
1694    }
1695}
1696
1697impl<TX: DbTx, N: NodeTypes> ChangeSetReader for DatabaseProvider<TX, N> {
1698    fn account_block_changeset(
1699        &self,
1700        block_number: BlockNumber,
1701    ) -> ProviderResult<Vec<AccountBeforeTx>> {
1702        if self.cached_storage_settings().storage_v2 {
1703            let static_changesets =
1704                self.static_file_provider.account_block_changeset(block_number)?;
1705            Ok(static_changesets)
1706        } else {
1707            let range = block_number..=block_number;
1708            self.tx
1709                .cursor_read::<tables::AccountChangeSets>()?
1710                .walk_range(range)?
1711                .map(|result| -> ProviderResult<_> {
1712                    let (_, account_before) = result?;
1713                    Ok(account_before)
1714                })
1715                .collect()
1716        }
1717    }
1718
1719    fn get_account_before_block(
1720        &self,
1721        block_number: BlockNumber,
1722        address: Address,
1723    ) -> ProviderResult<Option<AccountBeforeTx>> {
1724        if self.cached_storage_settings().storage_v2 {
1725            Ok(self.static_file_provider.get_account_before_block(block_number, address)?)
1726        } else {
1727            self.tx
1728                .cursor_dup_read::<tables::AccountChangeSets>()?
1729                .seek_by_key_subkey(block_number, address)?
1730                .filter(|acc| acc.address == address)
1731                .map(Ok)
1732                .transpose()
1733        }
1734    }
1735
1736    fn account_changesets_range(
1737        &self,
1738        range: impl core::ops::RangeBounds<BlockNumber>,
1739    ) -> ProviderResult<Vec<(BlockNumber, AccountBeforeTx)>> {
1740        if self.cached_storage_settings().storage_v2 {
1741            self.static_file_provider.account_changesets_range(range)
1742        } else {
1743            self.tx
1744                .cursor_read::<tables::AccountChangeSets>()?
1745                .walk_range(to_range(range))?
1746                .map(|r| r.map_err(Into::into))
1747                .collect()
1748        }
1749    }
1750}
1751
1752impl<Tx: DbTx + 'static, N: NodeTypesForProvider> StateReader for DatabaseProvider<Tx, N> {
1753    type Receipt = ReceiptTy<N>;
1754
1755    fn get_state(
1756        &self,
1757        block: BlockNumber,
1758    ) -> ProviderResult<Option<ExecutionOutcome<Self::Receipt>>> {
1759        let Some(block_body) = self.block_body_indices(block)? else { return Ok(None) };
1760
1761        let from_transaction_num = block_body.first_tx_num();
1762        let to_transaction_num = block_body.last_tx_num();
1763
1764        let account_changeset = self.account_changesets_range(block..=block)?;
1765        let storage_changeset = self.storage_changeset(block)?;
1766
1767        let Some(block_hash) = self.block_hash(block)? else { return Ok(None) };
1768        let state_provider = self.history_by_block_hash(block_hash)?;
1769        let (state, reverts) = self.populate_bundle_state_with_provider(
1770            account_changeset,
1771            storage_changeset,
1772            state_provider,
1773        )?;
1774
1775        let receipts = self.receipts_by_tx_range(from_transaction_num..=to_transaction_num)?;
1776
1777        Ok(Some(ExecutionOutcome::new_init(
1778            state,
1779            reverts,
1780            // We skip new contracts since we never delete them from the database
1781            Vec::new(),
1782            vec![receipts],
1783            block,
1784            Vec::new(),
1785        )))
1786    }
1787}
1788
1789impl<TX: DbTx + 'static, N: NodeTypesForProvider> HeaderSyncGapProvider
1790    for DatabaseProvider<TX, N>
1791{
1792    type Header = HeaderTy<N>;
1793
1794    fn local_tip_header(
1795        &self,
1796        highest_uninterrupted_block: BlockNumber,
1797    ) -> ProviderResult<SealedHeader<Self::Header>> {
1798        let static_file_provider = self.static_file_provider();
1799
1800        // Make sure Headers static file is at the same height. If it's further, this
1801        // input execution was interrupted previously and we need to unwind the static file.
1802        let next_static_file_block_num = static_file_provider
1803            .get_highest_static_file_block(StaticFileSegment::Headers)
1804            .map(|id| id + 1)
1805            .unwrap_or_default();
1806        let next_block = highest_uninterrupted_block + 1;
1807
1808        match next_static_file_block_num.cmp(&next_block) {
1809            // The node shutdown between an executed static file commit and before the database
1810            // commit, so we need to unwind the static files.
1811            Ordering::Greater => {
1812                let mut static_file_producer =
1813                    static_file_provider.latest_writer(StaticFileSegment::Headers)?;
1814                static_file_producer.prune_headers(next_static_file_block_num - next_block)?;
1815                // Since this is a database <-> static file inconsistency, we commit the change
1816                // straight away.
1817                static_file_producer.commit()?
1818            }
1819            Ordering::Less => {
1820                // There's either missing or corrupted files.
1821                return Err(ProviderError::HeaderNotFound(next_static_file_block_num.into()))
1822            }
1823            Ordering::Equal => {}
1824        }
1825
1826        let local_head = static_file_provider
1827            .sealed_header(highest_uninterrupted_block)?
1828            .ok_or_else(|| ProviderError::HeaderNotFound(highest_uninterrupted_block.into()))?;
1829
1830        Ok(local_head)
1831    }
1832}
1833
1834impl<TX: DbTx + 'static, N: NodeTypesForProvider> HeaderProvider for DatabaseProvider<TX, N> {
1835    type Header = HeaderTy<N>;
1836
1837    fn header(&self, block_hash: BlockHash) -> ProviderResult<Option<Self::Header>> {
1838        if let Some(num) = self.block_number(block_hash)? {
1839            Ok(self.header_by_number(num)?)
1840        } else {
1841            Ok(None)
1842        }
1843    }
1844
1845    fn header_by_number(&self, num: BlockNumber) -> ProviderResult<Option<Self::Header>> {
1846        self.static_file_provider.header_by_number(num)
1847    }
1848
1849    fn headers_range(
1850        &self,
1851        range: impl RangeBounds<BlockNumber>,
1852    ) -> ProviderResult<Vec<Self::Header>> {
1853        self.static_file_provider.headers_range(range)
1854    }
1855
1856    fn sealed_header(
1857        &self,
1858        number: BlockNumber,
1859    ) -> ProviderResult<Option<SealedHeader<Self::Header>>> {
1860        self.static_file_provider.sealed_header(number)
1861    }
1862
1863    fn sealed_headers_while(
1864        &self,
1865        range: impl RangeBounds<BlockNumber>,
1866        predicate: impl FnMut(&SealedHeader<Self::Header>) -> bool,
1867    ) -> ProviderResult<Vec<SealedHeader<Self::Header>>> {
1868        self.static_file_provider.sealed_headers_while(range, predicate)
1869    }
1870}
1871
1872impl<TX: DbTx + 'static, N: NodeTypes> BlockHashReader for DatabaseProvider<TX, N> {
1873    fn block_hash(&self, number: u64) -> ProviderResult<Option<B256>> {
1874        self.static_file_provider.block_hash(number)
1875    }
1876
1877    fn canonical_hashes_range(
1878        &self,
1879        start: BlockNumber,
1880        end: BlockNumber,
1881    ) -> ProviderResult<Vec<B256>> {
1882        self.static_file_provider.canonical_hashes_range(start, end)
1883    }
1884}
1885
1886impl<TX: DbTx + 'static, N: NodeTypes> BlockNumReader for DatabaseProvider<TX, N> {
1887    fn chain_info(&self) -> ProviderResult<ChainInfo> {
1888        let best_number = self.best_block_number()?;
1889        let best_hash = self.block_hash(best_number)?.unwrap_or_default();
1890        Ok(ChainInfo { best_hash, best_number })
1891    }
1892
1893    fn best_block_number(&self) -> ProviderResult<BlockNumber> {
1894        // The best block number is tracked via the finished stage which gets updated in the same tx
1895        // when new blocks committed
1896        Ok(self
1897            .get_stage_checkpoint(StageId::Finish)?
1898            .map(|checkpoint| checkpoint.block_number)
1899            .unwrap_or_default())
1900    }
1901
1902    fn last_block_number(&self) -> ProviderResult<BlockNumber> {
1903        self.static_file_provider.last_block_number()
1904    }
1905
1906    fn block_number(&self, hash: B256) -> ProviderResult<Option<BlockNumber>> {
1907        Ok(self.tx.get::<tables::HeaderNumbers>(hash)?)
1908    }
1909}
1910
1911impl<TX: DbTx + 'static, N: NodeTypesForProvider> BlockReader for DatabaseProvider<TX, N> {
1912    type Block = BlockTy<N>;
1913
1914    fn find_block_by_hash(
1915        &self,
1916        hash: B256,
1917        source: BlockSource,
1918    ) -> ProviderResult<Option<Self::Block>> {
1919        if source.is_canonical() {
1920            self.block(hash.into())
1921        } else {
1922            Ok(None)
1923        }
1924    }
1925
1926    /// Returns the block with matching number from database.
1927    ///
1928    /// If the header for this block is not found, this returns `None`.
1929    /// If the header is found, but the transactions either do not exist, or are not indexed, this
1930    /// will return None.
1931    ///
1932    /// Returns an error if the requested block is below the earliest available history.
1933    fn block(&self, id: BlockHashOrNumber) -> ProviderResult<Option<Self::Block>> {
1934        if let Some(number) = self.convert_hash_or_number(id)? {
1935            let earliest_available = self.static_file_provider.earliest_history_height();
1936            if number < earliest_available {
1937                return Err(ProviderError::BlockExpired { requested: number, earliest_available })
1938            }
1939
1940            let Some(header) = self.header_by_number(number)? else { return Ok(None) };
1941
1942            // If the body indices are not found, this means that the transactions either do not
1943            // exist in the database yet, or they do exit but are not indexed.
1944            // If they exist but are not indexed, we don't have enough
1945            // information to return the block anyways, so we return `None`.
1946            let Some(transactions) = self.transactions_by_block(number.into())? else {
1947                return Ok(None)
1948            };
1949
1950            let body = self
1951                .storage
1952                .reader()
1953                .read_block_bodies(self, vec![(&header, transactions)])?
1954                .pop()
1955                .ok_or(ProviderError::InvalidStorageOutput)?;
1956
1957            return Ok(Some(Self::Block::new(header, body)))
1958        }
1959
1960        Ok(None)
1961    }
1962
1963    fn pending_block(&self) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
1964        Ok(None)
1965    }
1966
1967    fn pending_block_and_receipts(
1968        &self,
1969    ) -> ProviderResult<Option<(RecoveredBlock<Self::Block>, Vec<Self::Receipt>)>> {
1970        Ok(None)
1971    }
1972
1973    /// Returns the block with senders with matching number or hash from database.
1974    ///
1975    /// **NOTE: The transactions have invalid hashes, since they would need to be calculated on the
1976    /// spot, and we want fast querying.**
1977    ///
1978    /// If the header for this block is not found, this returns `None`.
1979    /// If the header is found, but the transactions either do not exist, or are not indexed, this
1980    /// will return None.
1981    fn recovered_block(
1982        &self,
1983        id: BlockHashOrNumber,
1984        transaction_kind: TransactionVariant,
1985    ) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
1986        self.recovered_block(
1987            id,
1988            transaction_kind,
1989            |block_number| self.header_by_number(block_number),
1990            |header, body, senders| {
1991                Self::Block::new(header, body)
1992                    // Note: we're using unchecked here because we know the block contains valid txs
1993                    // wrt to its height and can ignore the s value check so pre
1994                    // EIP-2 txs are allowed
1995                    .try_into_recovered_unchecked(senders)
1996                    .map(Some)
1997                    .map_err(|_| ProviderError::SenderRecoveryError)
1998            },
1999        )
2000    }
2001
2002    fn sealed_block_with_senders(
2003        &self,
2004        id: BlockHashOrNumber,
2005        transaction_kind: TransactionVariant,
2006    ) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
2007        self.recovered_block(
2008            id,
2009            transaction_kind,
2010            |block_number| self.sealed_header(block_number),
2011            |header, body, senders| {
2012                Self::Block::new_sealed(header, body)
2013                    // Note: we're using unchecked here because we know the block contains valid txs
2014                    // wrt to its height and can ignore the s value check so pre
2015                    // EIP-2 txs are allowed
2016                    .try_with_senders_unchecked(senders)
2017                    .map(Some)
2018                    .map_err(|_| ProviderError::SenderRecoveryError)
2019            },
2020        )
2021    }
2022
2023    fn block_range(&self, range: RangeInclusive<BlockNumber>) -> ProviderResult<Vec<Self::Block>> {
2024        self.block_range(
2025            range,
2026            |range| self.headers_range(range),
2027            |header, body, _| Ok(Self::Block::new(header, body)),
2028        )
2029    }
2030
2031    fn block_with_senders_range(
2032        &self,
2033        range: RangeInclusive<BlockNumber>,
2034    ) -> ProviderResult<Vec<RecoveredBlock<Self::Block>>> {
2035        self.block_with_senders_range(
2036            range,
2037            |range| self.headers_range(range),
2038            |header, body, senders| {
2039                Self::Block::new(header, body)
2040                    .try_into_recovered_unchecked(senders)
2041                    .map_err(|_| ProviderError::SenderRecoveryError)
2042            },
2043        )
2044    }
2045
2046    fn recovered_block_range(
2047        &self,
2048        range: RangeInclusive<BlockNumber>,
2049    ) -> ProviderResult<Vec<RecoveredBlock<Self::Block>>> {
2050        self.block_with_senders_range(
2051            range,
2052            |range| self.sealed_headers_range(range),
2053            |header, body, senders| {
2054                Self::Block::new_sealed(header, body)
2055                    .try_with_senders(senders)
2056                    .map_err(|_| ProviderError::SenderRecoveryError)
2057            },
2058        )
2059    }
2060
2061    fn block_by_transaction_id(&self, id: TxNumber) -> ProviderResult<Option<BlockNumber>> {
2062        Ok(self
2063            .tx
2064            .cursor_read::<tables::TransactionBlocks>()?
2065            .seek(id)
2066            .map(|b| b.map(|(_, bn)| bn))?)
2067    }
2068}
2069
2070impl<TX: DbTx + 'static, N: NodeTypesForProvider> TransactionsProviderExt
2071    for DatabaseProvider<TX, N>
2072{
2073    /// Recovers transaction hashes by walking through `Transactions` table and
2074    /// calculating them in a parallel manner. Returned unsorted.
2075    fn transaction_hashes_by_range(
2076        &self,
2077        tx_range: Range<TxNumber>,
2078    ) -> ProviderResult<Vec<(TxHash, TxNumber)>> {
2079        self.static_file_provider.transaction_hashes_by_range(tx_range)
2080    }
2081}
2082
2083// Calculates the hash of the given transaction
2084impl<TX: DbTx + 'static, N: NodeTypesForProvider> TransactionsProvider for DatabaseProvider<TX, N> {
2085    type Transaction = TxTy<N>;
2086
2087    fn transaction_id(&self, tx_hash: TxHash) -> ProviderResult<Option<TxNumber>> {
2088        self.with_rocksdb_snapshot(|rocksdb_ref| {
2089            let mut reader = EitherReader::new_transaction_hash_numbers(self, rocksdb_ref)?;
2090            reader.get_transaction_hash_number(tx_hash)
2091        })
2092    }
2093
2094    fn transaction_by_id(&self, id: TxNumber) -> ProviderResult<Option<Self::Transaction>> {
2095        self.static_file_provider.transaction_by_id(id)
2096    }
2097
2098    fn transaction_by_id_unhashed(
2099        &self,
2100        id: TxNumber,
2101    ) -> ProviderResult<Option<Self::Transaction>> {
2102        self.static_file_provider.transaction_by_id_unhashed(id)
2103    }
2104
2105    fn transaction_by_hash(&self, hash: TxHash) -> ProviderResult<Option<Self::Transaction>> {
2106        if let Some(id) = self.transaction_id(hash)? {
2107            Ok(self.transaction_by_id_unhashed(id)?)
2108        } else {
2109            Ok(None)
2110        }
2111    }
2112
2113    fn transaction_by_hash_with_meta(
2114        &self,
2115        tx_hash: TxHash,
2116    ) -> ProviderResult<Option<(Self::Transaction, TransactionMeta)>> {
2117        if let Some(transaction_id) = self.transaction_id(tx_hash)? &&
2118            let Some(transaction) = self.transaction_by_id_unhashed(transaction_id)? &&
2119            let Some(block_number) = self.block_by_transaction_id(transaction_id)? &&
2120            let Some(sealed_header) = self.sealed_header(block_number)?
2121        {
2122            let (header, block_hash) = sealed_header.split();
2123            if let Some(block_body) = self.block_body_indices(block_number)? {
2124                // the index of the tx in the block is the offset:
2125                // len([start..tx_id])
2126                // NOTE: `transaction_id` is always `>=` the block's first
2127                // index
2128                let index = transaction_id - block_body.first_tx_num();
2129
2130                let meta = TransactionMeta {
2131                    tx_hash,
2132                    index,
2133                    block_hash,
2134                    block_number,
2135                    base_fee: header.base_fee_per_gas(),
2136                    excess_blob_gas: header.excess_blob_gas(),
2137                    timestamp: header.timestamp(),
2138                };
2139
2140                return Ok(Some((transaction, meta)))
2141            }
2142        }
2143
2144        Ok(None)
2145    }
2146
2147    fn transactions_by_block(
2148        &self,
2149        id: BlockHashOrNumber,
2150    ) -> ProviderResult<Option<Vec<Self::Transaction>>> {
2151        if let Some(block_number) = self.convert_hash_or_number(id)? &&
2152            let Some(body) = self.block_body_indices(block_number)?
2153        {
2154            let tx_range = body.tx_num_range();
2155            return if tx_range.is_empty() {
2156                Ok(Some(Vec::new()))
2157            } else {
2158                self.transactions_by_tx_range(tx_range).map(Some)
2159            }
2160        }
2161        Ok(None)
2162    }
2163
2164    fn transactions_by_block_range(
2165        &self,
2166        range: impl RangeBounds<BlockNumber>,
2167    ) -> ProviderResult<Vec<Vec<Self::Transaction>>> {
2168        let range = to_range(range);
2169
2170        self.block_body_indices_range(range.start..=range.end.saturating_sub(1))?
2171            .into_iter()
2172            .map(|body| {
2173                let tx_num_range = body.tx_num_range();
2174                if tx_num_range.is_empty() {
2175                    Ok(Vec::new())
2176                } else {
2177                    self.transactions_by_tx_range(tx_num_range)
2178                }
2179            })
2180            .collect()
2181    }
2182
2183    fn transactions_by_tx_range(
2184        &self,
2185        range: impl RangeBounds<TxNumber>,
2186    ) -> ProviderResult<Vec<Self::Transaction>> {
2187        self.static_file_provider.transactions_by_tx_range(range)
2188    }
2189
2190    fn senders_by_tx_range(
2191        &self,
2192        range: impl RangeBounds<TxNumber>,
2193    ) -> ProviderResult<Vec<Address>> {
2194        if EitherWriterDestination::senders(self).is_static_file() {
2195            self.static_file_provider.senders_by_tx_range(range)
2196        } else {
2197            self.cursor_read_collect::<tables::TransactionSenders>(range)
2198        }
2199    }
2200
2201    fn transaction_sender(&self, id: TxNumber) -> ProviderResult<Option<Address>> {
2202        if EitherWriterDestination::senders(self).is_static_file() {
2203            self.static_file_provider.transaction_sender(id)
2204        } else {
2205            Ok(self.tx.get::<tables::TransactionSenders>(id)?)
2206        }
2207    }
2208}
2209
2210impl<TX: DbTx + 'static, N: NodeTypesForProvider> ReceiptProvider for DatabaseProvider<TX, N> {
2211    type Receipt = ReceiptTy<N>;
2212
2213    fn receipt(&self, id: TxNumber) -> ProviderResult<Option<Self::Receipt>> {
2214        self.static_file_provider.get_with_static_file_or_database(
2215            StaticFileSegment::Receipts,
2216            id,
2217            |static_file| static_file.receipt(id),
2218            || Ok(self.tx.get::<tables::Receipts<Self::Receipt>>(id)?),
2219        )
2220    }
2221
2222    fn receipt_by_hash(&self, hash: TxHash) -> ProviderResult<Option<Self::Receipt>> {
2223        if let Some(id) = self.transaction_id(hash)? {
2224            self.receipt(id)
2225        } else {
2226            Ok(None)
2227        }
2228    }
2229
2230    fn receipts_by_block(
2231        &self,
2232        block: BlockHashOrNumber,
2233    ) -> ProviderResult<Option<Vec<Self::Receipt>>> {
2234        if let Some(number) = self.convert_hash_or_number(block)? &&
2235            let Some(body) = self.block_body_indices(number)?
2236        {
2237            let tx_range = body.tx_num_range();
2238            return if tx_range.is_empty() {
2239                Ok(Some(Vec::new()))
2240            } else {
2241                let receipts = self.receipts_by_tx_range(tx_range)?;
2242
2243                if receipts.len() != body.tx_count as usize {
2244                    return Ok(None)
2245                }
2246
2247                Ok(Some(receipts))
2248            }
2249        }
2250        Ok(None)
2251    }
2252
2253    fn receipts_by_tx_range(
2254        &self,
2255        range: impl RangeBounds<TxNumber>,
2256    ) -> ProviderResult<Vec<Self::Receipt>> {
2257        self.static_file_provider.get_range_with_static_file_or_database(
2258            StaticFileSegment::Receipts,
2259            to_range(range),
2260            |static_file, range, _| static_file.receipts_by_tx_range(range),
2261            |range, _| self.cursor_read_collect::<tables::Receipts<Self::Receipt>>(range),
2262            |_| true,
2263        )
2264    }
2265
2266    fn receipts_by_block_range(
2267        &self,
2268        block_range: RangeInclusive<BlockNumber>,
2269    ) -> ProviderResult<Vec<Vec<Self::Receipt>>> {
2270        if block_range.is_empty() {
2271            return Ok(Vec::new());
2272        }
2273
2274        // collect block body indices for each block in the range
2275        let range_len = block_range.end().saturating_sub(*block_range.start()) as usize + 1;
2276        let mut block_body_indices = Vec::with_capacity(range_len);
2277        for block_num in block_range {
2278            if let Some(indices) = self.block_body_indices(block_num)? {
2279                block_body_indices.push(indices);
2280            } else {
2281                // use default indices for missing blocks (empty block)
2282                block_body_indices.push(StoredBlockBodyIndices::default());
2283            }
2284        }
2285
2286        if block_body_indices.is_empty() {
2287            return Ok(Vec::new());
2288        }
2289
2290        // find blocks with transactions to determine transaction range
2291        let non_empty_blocks: Vec<_> =
2292            block_body_indices.iter().filter(|indices| indices.tx_count > 0).collect();
2293
2294        if non_empty_blocks.is_empty() {
2295            // all blocks are empty
2296            return Ok(vec![Vec::new(); block_body_indices.len()]);
2297        }
2298
2299        // calculate the overall transaction range
2300        let first_tx = non_empty_blocks[0].first_tx_num();
2301        let last_tx = non_empty_blocks[non_empty_blocks.len() - 1].last_tx_num();
2302
2303        // fetch all receipts in the transaction range
2304        let all_receipts = self.receipts_by_tx_range(first_tx..=last_tx)?;
2305        let mut receipts_iter = all_receipts.into_iter();
2306
2307        // distribute receipts to their respective blocks
2308        let mut result = Vec::with_capacity(block_body_indices.len());
2309        for indices in &block_body_indices {
2310            if indices.tx_count == 0 {
2311                result.push(Vec::new());
2312            } else {
2313                let block_receipts =
2314                    receipts_iter.by_ref().take(indices.tx_count as usize).collect();
2315                result.push(block_receipts);
2316            }
2317        }
2318
2319        Ok(result)
2320    }
2321}
2322
2323impl<TX: DbTx + 'static, N: NodeTypesForProvider> BlockBodyIndicesProvider
2324    for DatabaseProvider<TX, N>
2325{
2326    fn block_body_indices(&self, num: u64) -> ProviderResult<Option<StoredBlockBodyIndices>> {
2327        Ok(self.tx.get::<tables::BlockBodyIndices>(num)?)
2328    }
2329
2330    fn block_body_indices_range(
2331        &self,
2332        range: RangeInclusive<BlockNumber>,
2333    ) -> ProviderResult<Vec<StoredBlockBodyIndices>> {
2334        self.cursor_read_collect::<tables::BlockBodyIndices>(range)
2335    }
2336}
2337
2338impl<TX: DbTx, N: NodeTypes> StageCheckpointReader for DatabaseProvider<TX, N> {
2339    fn get_stage_checkpoint(&self, id: StageId) -> ProviderResult<Option<StageCheckpoint>> {
2340        Ok(if let Some(encoded) = id.get_pre_encoded() {
2341            self.tx.get_by_encoded_key::<tables::StageCheckpoints>(encoded)?
2342        } else {
2343            self.tx.get::<tables::StageCheckpoints>(id.to_string())?
2344        })
2345    }
2346
2347    /// Get stage checkpoint progress.
2348    fn get_stage_checkpoint_progress(&self, id: StageId) -> ProviderResult<Option<Vec<u8>>> {
2349        Ok(self.tx.get::<tables::StageCheckpointProgresses>(id.to_string())?)
2350    }
2351
2352    fn get_all_checkpoints(&self) -> ProviderResult<Vec<(String, StageCheckpoint)>> {
2353        self.tx
2354            .cursor_read::<tables::StageCheckpoints>()?
2355            .walk(None)?
2356            .collect::<Result<Vec<(String, StageCheckpoint)>, _>>()
2357            .map_err(ProviderError::Database)
2358    }
2359}
2360
2361impl<TX: DbTxMut, N: NodeTypes> StageCheckpointWriter for DatabaseProvider<TX, N> {
2362    /// Save stage checkpoint.
2363    fn save_stage_checkpoint(
2364        &self,
2365        id: StageId,
2366        checkpoint: StageCheckpoint,
2367    ) -> ProviderResult<()> {
2368        Ok(self.tx.put::<tables::StageCheckpoints>(id.to_string(), checkpoint)?)
2369    }
2370
2371    /// Save stage checkpoint progress.
2372    fn save_stage_checkpoint_progress(
2373        &self,
2374        id: StageId,
2375        checkpoint: Vec<u8>,
2376    ) -> ProviderResult<()> {
2377        Ok(self.tx.put::<tables::StageCheckpointProgresses>(id.to_string(), checkpoint)?)
2378    }
2379
2380    #[instrument(level = "debug", target = "providers::db", skip_all)]
2381    fn update_pipeline_stages(
2382        &self,
2383        block_number: BlockNumber,
2384        drop_stage_checkpoint: bool,
2385    ) -> ProviderResult<()> {
2386        // iterate over all existing stages in the table and update its progress.
2387        let mut cursor = self.tx.cursor_write::<tables::StageCheckpoints>()?;
2388        for stage_id in StageId::ALL {
2389            let (_, checkpoint) = cursor.seek_exact(stage_id.to_string())?.unwrap_or_default();
2390            cursor.upsert(
2391                stage_id.to_string(),
2392                &StageCheckpoint {
2393                    block_number,
2394                    ..if drop_stage_checkpoint { Default::default() } else { checkpoint }
2395                },
2396            )?;
2397        }
2398
2399        Ok(())
2400    }
2401}
2402
2403impl<TX: DbTxMut + DbTx, N: NodeTypes> DatabaseProvider<TX, N> {
2404    /// Updates pipeline checkpoints after an unwind while preserving an explicitly lagging
2405    /// state/trie frontier.
2406    fn update_pipeline_stages_after_unwind(
2407        &self,
2408        block_number: BlockNumber,
2409    ) -> ProviderResult<PersistenceFrontiers> {
2410        let partial_state_trie = self
2411            .get_stage_checkpoint(StageId::Finish)?
2412            .map(|checkpoint| {
2413                checkpoint
2414                    .finish_stage_checkpoint()
2415                    .and_then(|finish| finish.partial_state_trie())
2416                    .unwrap_or(checkpoint.block_number)
2417            })
2418            .unwrap_or(block_number)
2419            .min(block_number);
2420
2421        self.update_pipeline_stages(block_number, true)?;
2422        if partial_state_trie < block_number {
2423            self.save_stage_checkpoint(
2424                StageId::Finish,
2425                StageCheckpoint::new(block_number).with_finish_stage_checkpoint(FinishCheckpoint {
2426                    partial_state_trie: Some(partial_state_trie),
2427                }),
2428            )?;
2429        }
2430
2431        Ok(PersistenceFrontiers { db_tip: block_number, partial_state_trie })
2432    }
2433}
2434
2435impl<TX: DbTx + 'static, N: NodeTypes> StorageReader for DatabaseProvider<TX, N> {
2436    fn plain_state_storages(
2437        &self,
2438        addresses_with_keys: impl IntoIterator<Item = (Address, impl IntoIterator<Item = B256>)>,
2439    ) -> ProviderResult<Vec<(Address, Vec<StorageEntry>)>> {
2440        if self.cached_storage_settings().use_hashed_state() {
2441            let mut hashed_storage = self.tx.cursor_dup_read::<tables::HashedStorages>()?;
2442
2443            addresses_with_keys
2444                .into_iter()
2445                .map(|(address, storage)| {
2446                    let hashed_address = keccak256(address);
2447                    storage
2448                        .into_iter()
2449                        .map(|key| -> ProviderResult<_> {
2450                            let hashed_key = keccak256(key);
2451                            let value = hashed_storage
2452                                .seek_by_key_subkey(hashed_address, hashed_key)?
2453                                .filter(|v| v.key == hashed_key)
2454                                .map(|v| v.value)
2455                                .unwrap_or_default();
2456                            Ok(StorageEntry { key, value })
2457                        })
2458                        .collect::<ProviderResult<Vec<_>>>()
2459                        .map(|storage| (address, storage))
2460                })
2461                .collect::<ProviderResult<Vec<(_, _)>>>()
2462        } else {
2463            let mut plain_storage = self.tx.cursor_dup_read::<tables::PlainStorageState>()?;
2464
2465            addresses_with_keys
2466                .into_iter()
2467                .map(|(address, storage)| {
2468                    storage
2469                        .into_iter()
2470                        .map(|key| -> ProviderResult<_> {
2471                            Ok(plain_storage
2472                                .seek_by_key_subkey(address, key)?
2473                                .filter(|v| v.key == key)
2474                                .unwrap_or_else(|| StorageEntry { key, value: Default::default() }))
2475                        })
2476                        .collect::<ProviderResult<Vec<_>>>()
2477                        .map(|storage| (address, storage))
2478                })
2479                .collect::<ProviderResult<Vec<(_, _)>>>()
2480        }
2481    }
2482
2483    fn changed_storages_with_range(
2484        &self,
2485        range: RangeInclusive<BlockNumber>,
2486    ) -> ProviderResult<BTreeMap<Address, BTreeSet<B256>>> {
2487        if self.cached_storage_settings().storage_v2 {
2488            self.storage_changesets_range(range)?.into_iter().try_fold(
2489                BTreeMap::new(),
2490                |mut accounts: BTreeMap<Address, BTreeSet<B256>>, entry| {
2491                    let (BlockNumberAddress((_, address)), storage_entry) = entry;
2492                    accounts.entry(address).or_default().insert(storage_entry.key);
2493                    Ok(accounts)
2494                },
2495            )
2496        } else {
2497            self.tx
2498                .cursor_read::<tables::StorageChangeSets>()?
2499                .walk_range(BlockNumberAddress::range(range))?
2500                // fold all storages and save its old state so we can remove it from HashedStorage
2501                // it is needed as it is dup table.
2502                .try_fold(
2503                    BTreeMap::new(),
2504                    |mut accounts: BTreeMap<Address, BTreeSet<B256>>, entry| {
2505                        let (BlockNumberAddress((_, address)), storage_entry) = entry?;
2506                        accounts.entry(address).or_default().insert(storage_entry.key);
2507                        Ok(accounts)
2508                    },
2509                )
2510        }
2511    }
2512
2513    fn changed_storages_and_blocks_with_range(
2514        &self,
2515        range: RangeInclusive<BlockNumber>,
2516    ) -> ProviderResult<BTreeMap<(Address, B256), Vec<u64>>> {
2517        if self.cached_storage_settings().storage_v2 {
2518            self.storage_changesets_range(range)?.into_iter().try_fold(
2519                BTreeMap::new(),
2520                |mut storages: BTreeMap<(Address, B256), Vec<u64>>, (index, storage)| {
2521                    storages
2522                        .entry((index.address(), storage.key))
2523                        .or_default()
2524                        .push(index.block_number());
2525                    Ok(storages)
2526                },
2527            )
2528        } else {
2529            let mut changeset_cursor = self.tx.cursor_read::<tables::StorageChangeSets>()?;
2530
2531            let storage_changeset_lists =
2532                changeset_cursor.walk_range(BlockNumberAddress::range(range))?.try_fold(
2533                    BTreeMap::new(),
2534                    |mut storages: BTreeMap<(Address, B256), Vec<u64>>,
2535                     entry|
2536                     -> ProviderResult<_> {
2537                        let (index, storage) = entry?;
2538                        storages
2539                            .entry((index.address(), storage.key))
2540                            .or_default()
2541                            .push(index.block_number());
2542                        Ok(storages)
2543                    },
2544                )?;
2545
2546            Ok(storage_changeset_lists)
2547        }
2548    }
2549}
2550
2551impl<TX: DbTxMut + DbTx + 'static, N: NodeTypesForProvider> StateWriter
2552    for DatabaseProvider<TX, N>
2553{
2554    type Receipt = ReceiptTy<N>;
2555
2556    #[instrument(level = "debug", target = "providers::db", skip_all)]
2557    fn write_state<'a>(
2558        &self,
2559        execution_outcome: impl Into<WriteStateInput<'a, Self::Receipt>>,
2560        is_value_known: OriginalValuesKnown,
2561        config: StateWriteConfig,
2562    ) -> ProviderResult<()> {
2563        let execution_outcome = execution_outcome.into();
2564
2565        if self.cached_storage_settings().use_hashed_state() &&
2566            !config.write_receipts &&
2567            !config.write_account_changesets &&
2568            !config.write_storage_changesets
2569        {
2570            // In storage v2 with all outputs directed to static files, plain state and changesets
2571            // are written elsewhere. Only bytecodes need MDBX writes, so skip the expensive
2572            // to_plain_state_and_reverts conversion that iterates all accounts and storage.
2573            self.write_bytecodes(
2574                execution_outcome.state().contracts.iter().map(|(h, b)| (*h, Bytecode(b.clone()))),
2575            )?;
2576            return Ok(());
2577        }
2578
2579        let first_block = execution_outcome.first_block();
2580        let (plain_state, reverts) =
2581            execution_outcome.state().to_plain_state_and_reverts(is_value_known);
2582
2583        self.write_state_reverts(reverts, first_block, config)?;
2584        self.write_state_changes(plain_state)?;
2585
2586        if !config.write_receipts {
2587            return Ok(());
2588        }
2589
2590        let block_count = execution_outcome.len() as u64;
2591        let last_block = execution_outcome.last_block();
2592        let block_range = first_block..=last_block;
2593
2594        let tip = self.last_block_number()?.max(last_block);
2595
2596        // Fetch the first transaction number for each block in the range
2597        let block_indices: Vec<_> = self
2598            .block_body_indices_range(block_range)?
2599            .into_iter()
2600            .map(|b| b.first_tx_num)
2601            .collect();
2602
2603        // Ensure all expected blocks are present.
2604        if block_indices.len() < block_count as usize {
2605            let missing_blocks = block_count - block_indices.len() as u64;
2606            return Err(ProviderError::BlockBodyIndicesNotFound(
2607                last_block.saturating_sub(missing_blocks - 1),
2608            ));
2609        }
2610
2611        let mut receipts_writer = EitherWriter::new_receipts(self, first_block)?;
2612
2613        let has_contract_log_filter = !self.prune_modes.receipts_log_filter.is_empty();
2614        let contract_log_pruner = self.prune_modes.receipts_log_filter.group_by_block(tip, None)?;
2615
2616        // All receipts from the last 128 blocks are required for blockchain tree, even with
2617        // [`PruneSegment::ContractLogs`].
2618        //
2619        // Receipts can only be skipped if we're dealing with legacy nodes that write them to
2620        // Database, OR if receipts_in_static_files is enabled but no receipts exist in static
2621        // files yet. Once receipts exist in static files, we must continue writing to maintain
2622        // continuity and have no gaps.
2623        let prunable_receipts = (EitherWriter::receipts_destination(self).is_database() ||
2624            self.static_file_provider()
2625                .get_highest_static_file_tx(StaticFileSegment::Receipts)
2626                .is_none()) &&
2627            PruneMode::Distance(self.minimum_pruning_distance).should_prune(first_block, tip);
2628
2629        // Prepare set of addresses which logs should not be pruned.
2630        let mut allowed_addresses: AddressSet = AddressSet::default();
2631        for (_, addresses) in contract_log_pruner.range(..first_block) {
2632            allowed_addresses.extend(addresses.iter().copied());
2633        }
2634
2635        for (idx, (receipts, first_tx_index)) in
2636            execution_outcome.receipts().zip(block_indices).enumerate()
2637        {
2638            let block_number = first_block + idx as u64;
2639
2640            // Increment block number for receipts static file writer
2641            receipts_writer.increment_block(block_number)?;
2642
2643            // Skip writing receipts if pruning configuration requires us to.
2644            if prunable_receipts &&
2645                self.prune_modes
2646                    .receipts
2647                    .is_some_and(|mode| mode.should_prune(block_number, tip))
2648            {
2649                continue
2650            }
2651
2652            // If there are new addresses to retain after this block number, track them
2653            if let Some(new_addresses) = contract_log_pruner.get(&block_number) {
2654                allowed_addresses.extend(new_addresses.iter().copied());
2655            }
2656
2657            for (idx, receipt) in receipts.iter().enumerate() {
2658                let receipt_idx = first_tx_index + idx as u64;
2659                // Skip writing receipt if log filter is active and it does not have any logs to
2660                // retain
2661                if prunable_receipts &&
2662                    has_contract_log_filter &&
2663                    !receipt.logs().iter().any(|log| allowed_addresses.contains(&log.address))
2664                {
2665                    continue
2666                }
2667
2668                receipts_writer.append_receipt(receipt_idx, receipt)?;
2669            }
2670        }
2671
2672        Ok(())
2673    }
2674
2675    fn write_state_reverts(
2676        &self,
2677        reverts: PlainStateReverts,
2678        first_block: BlockNumber,
2679        config: StateWriteConfig,
2680    ) -> ProviderResult<()> {
2681        // Write storage changes
2682        if config.write_storage_changesets {
2683            tracing::trace!("Writing storage changes");
2684            let mut storages_cursor =
2685                self.tx_ref().cursor_dup_write::<tables::PlainStorageState>()?;
2686            for (block_index, mut storage_changes) in reverts.storage.into_iter().enumerate() {
2687                let block_number = first_block + block_index as BlockNumber;
2688
2689                tracing::trace!(block_number, "Writing block change");
2690                // sort changes by address.
2691                storage_changes.par_sort_unstable_by_key(|a| a.address);
2692                let total_changes =
2693                    storage_changes.iter().map(|change| change.storage_revert.len()).sum();
2694                let mut changeset = Vec::with_capacity(total_changes);
2695                for PlainStorageRevert { address, wiped, storage_revert } in storage_changes {
2696                    let mut storage = storage_revert
2697                        .into_iter()
2698                        .map(|(k, v)| (B256::from(k.to_be_bytes()), v))
2699                        .collect::<Vec<_>>();
2700                    // sort storage slots by key.
2701                    storage.par_sort_unstable_by_key(|a| a.0);
2702
2703                    // If we are writing the primary storage wipe transition, the pre-existing
2704                    // storage state has to be taken from the database and written to storage
2705                    // history. See [StorageWipe::Primary] for more details.
2706                    //
2707                    // TODO(mediocregopher): This could be rewritten in a way which doesn't
2708                    // require collecting wiped entries into a Vec like this, see
2709                    // `write_storage_trie_changesets`.
2710                    let mut wiped_storage = Vec::new();
2711                    if wiped {
2712                        tracing::trace!(?address, "Wiping storage");
2713                        if let Some((_, entry)) = storages_cursor.seek_exact(address)? {
2714                            wiped_storage.push((entry.key, entry.value));
2715                            while let Some(entry) = storages_cursor.next_dup_val()? {
2716                                wiped_storage.push((entry.key, entry.value))
2717                            }
2718                        }
2719                    }
2720
2721                    tracing::trace!(?address, ?storage, "Writing storage reverts");
2722                    for (key, value) in StorageRevertsIter::new(storage, wiped_storage) {
2723                        changeset.push(StorageBeforeTx { address, key, value });
2724                    }
2725                }
2726
2727                let mut storage_changesets_writer =
2728                    EitherWriter::new_storage_changesets(self, block_number)?;
2729                storage_changesets_writer.append_storage_changeset(block_number, changeset)?;
2730            }
2731        }
2732
2733        if !config.write_account_changesets {
2734            return Ok(());
2735        }
2736
2737        // Write account changes
2738        tracing::trace!(?first_block, "Writing account changes");
2739        for (block_index, account_block_reverts) in reverts.accounts.into_iter().enumerate() {
2740            let block_number = first_block + block_index as BlockNumber;
2741            let changeset = account_block_reverts
2742                .into_iter()
2743                .map(|(address, info)| AccountBeforeTx { address, info: info.map(Into::into) })
2744                .collect::<Vec<_>>();
2745            let mut account_changesets_writer =
2746                EitherWriter::new_account_changesets(self, block_number)?;
2747
2748            account_changesets_writer.append_account_changeset(block_number, changeset)?;
2749        }
2750
2751        Ok(())
2752    }
2753
2754    fn write_state_changes(&self, mut changes: StateChangeset) -> ProviderResult<()> {
2755        // sort all entries so they can be written to database in more performant way.
2756        // and take smaller memory footprint.
2757        changes.accounts.par_sort_by_key(|a| a.0);
2758        changes.storage.par_sort_by_key(|a| a.address);
2759        changes.contracts.par_sort_by_key(|a| a.0);
2760
2761        if !self.cached_storage_settings().use_hashed_state() {
2762            // Write new account state
2763            tracing::trace!(len = changes.accounts.len(), "Writing new account state");
2764            let mut accounts_cursor = self.tx_ref().cursor_write::<tables::PlainAccountState>()?;
2765            // write account to database.
2766            for (address, account) in changes.accounts {
2767                if let Some(account) = account {
2768                    tracing::trace!(?address, "Updating plain state account");
2769                    accounts_cursor.upsert(address, &account.into())?;
2770                } else if accounts_cursor.seek_exact(address)?.is_some() {
2771                    tracing::trace!(?address, "Deleting plain state account");
2772                    accounts_cursor.delete_current()?;
2773                }
2774            }
2775
2776            // Write new storage state and wipe storage if needed.
2777            tracing::trace!(len = changes.storage.len(), "Writing new storage state");
2778            let mut storages_cursor =
2779                self.tx_ref().cursor_dup_write::<tables::PlainStorageState>()?;
2780            for PlainStorageChangeset { address, wipe_storage, storage } in changes.storage {
2781                // Wiping of storage.
2782                if wipe_storage && storages_cursor.seek_exact(address)?.is_some() {
2783                    storages_cursor.delete_current_duplicates()?;
2784                }
2785                // cast storages to B256.
2786                let mut storage = storage
2787                    .into_iter()
2788                    .map(|(k, value)| StorageEntry { key: k.into(), value })
2789                    .collect::<Vec<_>>();
2790                // sort storage slots by key.
2791                storage.par_sort_unstable_by_key(|a| a.key);
2792
2793                for entry in storage {
2794                    tracing::trace!(?address, ?entry.key, "Updating plain state storage");
2795                    if let Some(db_entry) =
2796                        storages_cursor.seek_by_key_subkey(address, entry.key)? &&
2797                        db_entry.key == entry.key
2798                    {
2799                        storages_cursor.delete_current()?;
2800                    }
2801
2802                    if !entry.value.is_zero() {
2803                        storages_cursor.upsert(address, &entry)?;
2804                    }
2805                }
2806            }
2807        }
2808
2809        // Write bytecode
2810        tracing::trace!(len = changes.contracts.len(), "Writing bytecodes");
2811        self.write_bytecodes(
2812            changes.contracts.into_iter().map(|(hash, bytecode)| (hash, Bytecode(bytecode))),
2813        )?;
2814
2815        Ok(())
2816    }
2817
2818    #[instrument(level = "debug", target = "providers::db", skip_all)]
2819    fn write_hashed_state(&self, hashed_state: &HashedPostStateSorted) -> ProviderResult<()> {
2820        // Write hashed account updates.
2821        let mut hashed_accounts_cursor = self.tx_ref().cursor_write::<tables::HashedAccounts>()?;
2822        for (hashed_address, account) in hashed_state.accounts() {
2823            if let Some(account) = account {
2824                hashed_accounts_cursor.upsert(*hashed_address, account)?;
2825            } else if hashed_accounts_cursor.seek_exact(*hashed_address)?.is_some() {
2826                hashed_accounts_cursor.delete_current()?;
2827            }
2828        }
2829
2830        // Write hashed storage changes.
2831        let sorted_storages = hashed_state.account_storages().iter().sorted_by_key(|(key, _)| *key);
2832        let mut hashed_storage_cursor =
2833            self.tx_ref().cursor_dup_write::<tables::HashedStorages>()?;
2834        for (hashed_address, storage) in sorted_storages {
2835            if storage.is_wiped() && hashed_storage_cursor.seek_exact(*hashed_address)?.is_some() {
2836                hashed_storage_cursor.delete_current_duplicates()?;
2837            }
2838
2839            for (hashed_slot, value) in storage.storage_slots_ref() {
2840                let entry = StorageEntry { key: *hashed_slot, value: *value };
2841
2842                if let Some(db_entry) =
2843                    hashed_storage_cursor.seek_by_key_subkey(*hashed_address, entry.key)? &&
2844                    db_entry.key == entry.key
2845                {
2846                    hashed_storage_cursor.delete_current()?;
2847                }
2848
2849                if !entry.value.is_zero() {
2850                    hashed_storage_cursor.upsert(*hashed_address, &entry)?;
2851                }
2852            }
2853        }
2854
2855        Ok(())
2856    }
2857
2858    /// Remove the last N blocks of state.
2859    ///
2860    /// The latest state will be unwound
2861    ///
2862    /// 1. Read the retained block's [`BlockBodyIndices`][tables::BlockBodyIndices] entry to get the
2863    ///    first transaction id to remove.
2864    /// 2. Iterate over the [`StorageChangeSets`][tables::StorageChangeSets] table and the
2865    ///    [`AccountChangeSets`][tables::AccountChangeSets] tables in reverse order to reconstruct
2866    ///    the changesets.
2867    ///    - In order to have both the old and new values in the changesets, we also access the
2868    ///      plain state tables.
2869    /// 3. While iterating over the changeset tables, if we encounter a new account or storage slot,
2870    ///    we:
2871    ///     1. Take the old value from the changeset
2872    ///     2. Take the new value from the plain state
2873    ///     3. Save the old value to the local state
2874    /// 4. While iterating over the changeset tables, if we encounter an account/storage slot we
2875    ///    have seen before we:
2876    ///     1. Take the old value from the changeset
2877    ///     2. Take the new value from the local state
2878    ///     3. Set the local state to the value in the changeset
2879    fn remove_state_above(&self, block: BlockNumber) -> ProviderResult<()> {
2880        let range = block + 1..=self.last_block_number()?;
2881
2882        if range.is_empty() {
2883            return Ok(());
2884        }
2885
2886        // get transaction receipts
2887        let from_transaction_num = self
2888            .block_body_indices(block)?
2889            .map(|b| b.next_tx_num())
2890            .ok_or(ProviderError::BlockBodyIndicesNotFound(block))?;
2891
2892        let storage_range = BlockNumberAddress::range(range.clone());
2893        let storage_changeset = if self.cached_storage_settings().storage_v2 {
2894            let changesets = self.storage_changesets_range(range.clone())?;
2895            let mut changeset_writer =
2896                self.static_file_provider.latest_writer(StaticFileSegment::StorageChangeSets)?;
2897            changeset_writer.prune_storage_changesets(block)?;
2898            changesets
2899        } else {
2900            self.take::<tables::StorageChangeSets>(storage_range)?.into_iter().collect()
2901        };
2902        let account_changeset = if self.cached_storage_settings().storage_v2 {
2903            let changesets = self.account_changesets_range(range)?;
2904            let mut changeset_writer =
2905                self.static_file_provider.latest_writer(StaticFileSegment::AccountChangeSets)?;
2906            changeset_writer.prune_account_changesets(block)?;
2907            changesets
2908        } else {
2909            self.take::<tables::AccountChangeSets>(range)?
2910        };
2911
2912        if self.cached_storage_settings().use_hashed_state() {
2913            let mut hashed_accounts_cursor = self.tx.cursor_write::<tables::HashedAccounts>()?;
2914            let mut hashed_storage_cursor = self.tx.cursor_dup_write::<tables::HashedStorages>()?;
2915
2916            let (state, _) = self.populate_bundle_state_hashed(
2917                account_changeset,
2918                storage_changeset,
2919                &mut hashed_accounts_cursor,
2920                &mut hashed_storage_cursor,
2921            )?;
2922
2923            for (address, (old_account, new_account, storage)) in &state {
2924                if old_account != new_account {
2925                    let hashed_address = keccak256(address);
2926                    let existing_entry = hashed_accounts_cursor.seek_exact(hashed_address)?;
2927                    if let Some(account) = old_account {
2928                        hashed_accounts_cursor.upsert(hashed_address, account)?;
2929                    } else if existing_entry.is_some() {
2930                        hashed_accounts_cursor.delete_current()?;
2931                    }
2932                }
2933
2934                for (storage_key, (old_storage_value, _new_storage_value)) in storage {
2935                    let hashed_address = keccak256(address);
2936                    let hashed_storage_key = keccak256(storage_key);
2937                    let storage_entry =
2938                        StorageEntry { key: hashed_storage_key, value: *old_storage_value };
2939                    if hashed_storage_cursor
2940                        .seek_by_key_subkey(hashed_address, hashed_storage_key)?
2941                        .is_some_and(|s| s.key == hashed_storage_key)
2942                    {
2943                        hashed_storage_cursor.delete_current()?
2944                    }
2945
2946                    if !old_storage_value.is_zero() {
2947                        hashed_storage_cursor.upsert(hashed_address, &storage_entry)?;
2948                    }
2949                }
2950            }
2951        } else {
2952            // This is not working for blocks that are not at tip. as plain state is not the last
2953            // state of end range. We should rename the functions or add support to access
2954            // History state. Accessing history state can be tricky but we are not gaining
2955            // anything.
2956            let mut plain_accounts_cursor = self.tx.cursor_write::<tables::PlainAccountState>()?;
2957            let mut plain_storage_cursor =
2958                self.tx.cursor_dup_write::<tables::PlainStorageState>()?;
2959
2960            let (state, _) = self.populate_bundle_state_plain(
2961                account_changeset,
2962                storage_changeset,
2963                &mut plain_accounts_cursor,
2964                &mut plain_storage_cursor,
2965            )?;
2966
2967            for (address, (old_account, new_account, storage)) in &state {
2968                if old_account != new_account {
2969                    let existing_entry = plain_accounts_cursor.seek_exact(*address)?;
2970                    if let Some(account) = old_account {
2971                        plain_accounts_cursor.upsert(*address, account)?;
2972                    } else if existing_entry.is_some() {
2973                        plain_accounts_cursor.delete_current()?;
2974                    }
2975                }
2976
2977                for (storage_key, (old_storage_value, _new_storage_value)) in storage {
2978                    let storage_entry =
2979                        StorageEntry { key: *storage_key, value: *old_storage_value };
2980                    if plain_storage_cursor
2981                        .seek_by_key_subkey(*address, *storage_key)?
2982                        .is_some_and(|s| s.key == *storage_key)
2983                    {
2984                        plain_storage_cursor.delete_current()?
2985                    }
2986
2987                    if !old_storage_value.is_zero() {
2988                        plain_storage_cursor.upsert(*address, &storage_entry)?;
2989                    }
2990                }
2991            }
2992        }
2993
2994        self.remove_receipts_from(from_transaction_num, block)?;
2995
2996        Ok(())
2997    }
2998
2999    /// Take the last N blocks of state, recreating the [`ExecutionOutcome`].
3000    ///
3001    /// The latest state will be unwound and returned back with all the blocks
3002    ///
3003    /// 1. Iterate over the [`BlockBodyIndices`][tables::BlockBodyIndices] table to get all the
3004    ///    transaction ids.
3005    /// 2. Iterate over the [`StorageChangeSets`][tables::StorageChangeSets] table and the
3006    ///    [`AccountChangeSets`][tables::AccountChangeSets] tables in reverse order to reconstruct
3007    ///    the changesets.
3008    ///    - In order to have both the old and new values in the changesets, we also access the
3009    ///      plain state tables.
3010    /// 3. While iterating over the changeset tables, if we encounter a new account or storage slot,
3011    ///    we:
3012    ///     1. Take the old value from the changeset
3013    ///     2. Take the new value from the plain state
3014    ///     3. Save the old value to the local state
3015    /// 4. While iterating over the changeset tables, if we encounter an account/storage slot we
3016    ///    have seen before we:
3017    ///     1. Take the old value from the changeset
3018    ///     2. Take the new value from the local state
3019    ///     3. Set the local state to the value in the changeset
3020    fn take_state_above(
3021        &self,
3022        block: BlockNumber,
3023    ) -> ProviderResult<ExecutionOutcome<Self::Receipt>> {
3024        let range = block + 1..=self.last_block_number()?;
3025
3026        if range.is_empty() {
3027            return Ok(ExecutionOutcome::default())
3028        }
3029        let start_block_number = *range.start();
3030
3031        // We are not removing block meta as it is used to get block changesets.
3032        let block_bodies = self.block_body_indices_range(range.clone())?;
3033
3034        // get transaction receipts
3035        let from_transaction_num =
3036            block_bodies.first().expect("already checked if there are blocks").first_tx_num();
3037        let to_transaction_num =
3038            block_bodies.last().expect("already checked if there are blocks").last_tx_num();
3039
3040        let storage_range = BlockNumberAddress::range(range.clone());
3041        let storage_changeset = if let Some(highest_block) = self
3042            .static_file_provider
3043            .get_highest_static_file_block(StaticFileSegment::StorageChangeSets) &&
3044            self.cached_storage_settings().storage_v2
3045        {
3046            let changesets = self.storage_changesets_range(block + 1..=highest_block)?;
3047            let mut changeset_writer =
3048                self.static_file_provider.latest_writer(StaticFileSegment::StorageChangeSets)?;
3049            changeset_writer.prune_storage_changesets(block)?;
3050            changesets
3051        } else {
3052            self.take::<tables::StorageChangeSets>(storage_range)?.into_iter().collect()
3053        };
3054
3055        // if there are static files for this segment, prune them.
3056        let highest_changeset_block = self
3057            .static_file_provider
3058            .get_highest_static_file_block(StaticFileSegment::AccountChangeSets);
3059        let account_changeset = if let Some(highest_block) = highest_changeset_block &&
3060            self.cached_storage_settings().storage_v2
3061        {
3062            // TODO: add a `take` method that removes and returns the items instead of doing this
3063            let changesets = self.account_changesets_range(block + 1..highest_block + 1)?;
3064            let mut changeset_writer =
3065                self.static_file_provider.latest_writer(StaticFileSegment::AccountChangeSets)?;
3066            changeset_writer.prune_account_changesets(block)?;
3067
3068            changesets
3069        } else {
3070            // Have to remove from static files if they exist, otherwise remove using `take` for the
3071            // changeset tables
3072            self.take::<tables::AccountChangeSets>(range)?
3073        };
3074
3075        let (state, reverts) = if self.cached_storage_settings().use_hashed_state() {
3076            let mut hashed_accounts_cursor = self.tx.cursor_write::<tables::HashedAccounts>()?;
3077            let mut hashed_storage_cursor = self.tx.cursor_dup_write::<tables::HashedStorages>()?;
3078
3079            let (state, reverts) = self.populate_bundle_state_hashed(
3080                account_changeset,
3081                storage_changeset,
3082                &mut hashed_accounts_cursor,
3083                &mut hashed_storage_cursor,
3084            )?;
3085
3086            for (address, (old_account, new_account, storage)) in &state {
3087                if old_account != new_account {
3088                    let hashed_address = keccak256(address);
3089                    let existing_entry = hashed_accounts_cursor.seek_exact(hashed_address)?;
3090                    if let Some(account) = old_account {
3091                        hashed_accounts_cursor.upsert(hashed_address, account)?;
3092                    } else if existing_entry.is_some() {
3093                        hashed_accounts_cursor.delete_current()?;
3094                    }
3095                }
3096
3097                for (storage_key, (old_storage_value, _new_storage_value)) in storage {
3098                    let hashed_address = keccak256(address);
3099                    let hashed_storage_key = keccak256(storage_key);
3100                    let storage_entry =
3101                        StorageEntry { key: hashed_storage_key, value: *old_storage_value };
3102                    if hashed_storage_cursor
3103                        .seek_by_key_subkey(hashed_address, hashed_storage_key)?
3104                        .is_some_and(|s| s.key == hashed_storage_key)
3105                    {
3106                        hashed_storage_cursor.delete_current()?
3107                    }
3108
3109                    if !old_storage_value.is_zero() {
3110                        hashed_storage_cursor.upsert(hashed_address, &storage_entry)?;
3111                    }
3112                }
3113            }
3114
3115            (state, reverts)
3116        } else {
3117            // This is not working for blocks that are not at tip. as plain state is not the last
3118            // state of end range. We should rename the functions or add support to access
3119            // History state. Accessing history state can be tricky but we are not gaining
3120            // anything.
3121            let mut plain_accounts_cursor = self.tx.cursor_write::<tables::PlainAccountState>()?;
3122            let mut plain_storage_cursor =
3123                self.tx.cursor_dup_write::<tables::PlainStorageState>()?;
3124
3125            let (state, reverts) = self.populate_bundle_state_plain(
3126                account_changeset,
3127                storage_changeset,
3128                &mut plain_accounts_cursor,
3129                &mut plain_storage_cursor,
3130            )?;
3131
3132            for (address, (old_account, new_account, storage)) in &state {
3133                if old_account != new_account {
3134                    let existing_entry = plain_accounts_cursor.seek_exact(*address)?;
3135                    if let Some(account) = old_account {
3136                        plain_accounts_cursor.upsert(*address, account)?;
3137                    } else if existing_entry.is_some() {
3138                        plain_accounts_cursor.delete_current()?;
3139                    }
3140                }
3141
3142                for (storage_key, (old_storage_value, _new_storage_value)) in storage {
3143                    let storage_entry =
3144                        StorageEntry { key: *storage_key, value: *old_storage_value };
3145                    if plain_storage_cursor
3146                        .seek_by_key_subkey(*address, *storage_key)?
3147                        .is_some_and(|s| s.key == *storage_key)
3148                    {
3149                        plain_storage_cursor.delete_current()?
3150                    }
3151
3152                    if !old_storage_value.is_zero() {
3153                        plain_storage_cursor.upsert(*address, &storage_entry)?;
3154                    }
3155                }
3156            }
3157
3158            (state, reverts)
3159        };
3160
3161        // Collect receipts into tuples (tx_num, receipt) to correctly handle pruned receipts
3162        let mut receipts_iter = self
3163            .static_file_provider
3164            .get_range_with_static_file_or_database(
3165                StaticFileSegment::Receipts,
3166                from_transaction_num..to_transaction_num + 1,
3167                |static_file, range, _| {
3168                    static_file
3169                        .receipts_by_tx_range(range.clone())
3170                        .map(|r| range.into_iter().zip(r).collect())
3171                },
3172                |range, _| {
3173                    self.tx
3174                        .cursor_read::<tables::Receipts<Self::Receipt>>()?
3175                        .walk_range(range)?
3176                        .map(|r| r.map_err(Into::into))
3177                        .collect()
3178                },
3179                |_| true,
3180            )?
3181            .into_iter()
3182            .peekable();
3183
3184        let mut receipts = Vec::with_capacity(block_bodies.len());
3185        // loop break if we are at the end of the blocks.
3186        for block_body in block_bodies {
3187            let mut block_receipts = Vec::with_capacity(block_body.tx_count as usize);
3188            for num in block_body.tx_num_range() {
3189                if receipts_iter.peek().is_some_and(|(n, _)| *n == num) {
3190                    block_receipts.push(receipts_iter.next().unwrap().1);
3191                }
3192            }
3193            receipts.push(block_receipts);
3194        }
3195
3196        self.remove_receipts_from(from_transaction_num, block)?;
3197
3198        Ok(ExecutionOutcome::new_init(
3199            state,
3200            reverts,
3201            Vec::new(),
3202            receipts,
3203            start_block_number,
3204            Vec::new(),
3205        ))
3206    }
3207}
3208
3209impl<TX: DbTxMut + DbTx + 'static, N: NodeTypes> DatabaseProvider<TX, N> {
3210    fn write_account_trie_updates<A: TrieTableAdapter>(
3211        tx: &TX,
3212        trie_updates: &TrieUpdatesSorted,
3213        num_entries: &mut usize,
3214    ) -> ProviderResult<()>
3215    where
3216        TX: DbTxMut,
3217    {
3218        let mut account_trie_cursor = tx.cursor_write::<A::AccountTrieTable>()?;
3219        // Process sorted account nodes
3220        for (key, updated_node) in trie_updates.account_nodes_ref() {
3221            let nibbles = A::AccountKey::from(*key);
3222            match updated_node {
3223                Some(node) => {
3224                    if !key.is_empty() {
3225                        *num_entries += 1;
3226                        account_trie_cursor.upsert(nibbles, node)?;
3227                    }
3228                }
3229                None => {
3230                    *num_entries += 1;
3231                    if account_trie_cursor.seek_exact(nibbles)?.is_some() {
3232                        account_trie_cursor.delete_current()?;
3233                    }
3234                }
3235            }
3236        }
3237        Ok(())
3238    }
3239
3240    fn write_storage_tries<A: TrieTableAdapter>(
3241        tx: &TX,
3242        storage_tries: Vec<(&B256, &StorageTrieUpdatesSorted)>,
3243        num_entries: &mut usize,
3244    ) -> ProviderResult<()>
3245    where
3246        TX: DbTxMut,
3247    {
3248        let mut cursor = tx.cursor_dup_write::<A::StorageTrieTable>()?;
3249        for (hashed_address, storage_trie_updates) in storage_tries {
3250            let mut db_storage_trie_cursor: DatabaseStorageTrieCursor<_, A> =
3251                DatabaseStorageTrieCursor::new(cursor, *hashed_address);
3252            *num_entries +=
3253                db_storage_trie_cursor.write_storage_trie_updates_sorted(storage_trie_updates)?;
3254            cursor = db_storage_trie_cursor.cursor;
3255        }
3256        Ok(())
3257    }
3258}
3259
3260impl<TX: DbTxMut + DbTx + 'static, N: NodeTypes> TrieWriter for DatabaseProvider<TX, N> {
3261    /// Writes trie updates to the database with already sorted updates.
3262    ///
3263    /// Returns the number of entries modified.
3264    #[instrument(level = "debug", target = "providers::db", skip_all)]
3265    fn write_trie_updates_sorted(&self, trie_updates: &TrieUpdatesSorted) -> ProviderResult<usize> {
3266        if trie_updates.is_empty() {
3267            return Ok(0)
3268        }
3269
3270        // Track the number of inserted entries.
3271        let mut num_entries = 0;
3272
3273        reth_trie_db::with_adapter!(self, |A| {
3274            Self::write_account_trie_updates::<A>(self.tx_ref(), trie_updates, &mut num_entries)?;
3275        });
3276
3277        num_entries +=
3278            self.write_storage_trie_updates_sorted(trie_updates.storage_tries_ref().iter())?;
3279
3280        Ok(num_entries)
3281    }
3282}
3283
3284impl<TX: DbTxMut + DbTx + 'static, N: NodeTypes> StorageTrieWriter for DatabaseProvider<TX, N> {
3285    /// Writes storage trie updates from the given storage trie map with already sorted updates.
3286    ///
3287    /// Expects the storage trie updates to already be sorted by the hashed address key.
3288    ///
3289    /// Returns the number of entries modified.
3290    fn write_storage_trie_updates_sorted<'a>(
3291        &self,
3292        storage_tries: impl Iterator<Item = (&'a B256, &'a StorageTrieUpdatesSorted)>,
3293    ) -> ProviderResult<usize> {
3294        let mut num_entries = 0;
3295        let mut storage_tries = storage_tries.collect::<Vec<_>>();
3296        storage_tries.sort_unstable_by(|a, b| a.0.cmp(b.0));
3297        reth_trie_db::with_adapter!(self, |A| {
3298            Self::write_storage_tries::<A>(self.tx_ref(), storage_tries, &mut num_entries)?;
3299        });
3300        Ok(num_entries)
3301    }
3302}
3303
3304impl<TX: DbTxMut + DbTx + 'static, N: NodeTypes> HashingWriter for DatabaseProvider<TX, N> {
3305    fn unwind_account_hashing<'a>(
3306        &self,
3307        changesets: impl Iterator<Item = &'a (BlockNumber, AccountBeforeTx)>,
3308    ) -> ProviderResult<BTreeMap<B256, Option<Account>>> {
3309        // Aggregate all block changesets and make a list of accounts that have been changed.
3310        // Note that collecting and then reversing the order is necessary to ensure that the
3311        // changes are applied in the correct order.
3312        let hashed_accounts = changesets
3313            .into_iter()
3314            .map(|(_, e)| (keccak256(e.address), e.info))
3315            .collect::<Vec<_>>()
3316            .into_iter()
3317            .rev()
3318            .collect::<BTreeMap<_, _>>();
3319
3320        // Apply values to HashedState, and remove the account if it's None.
3321        let mut hashed_accounts_cursor = self.tx.cursor_write::<tables::HashedAccounts>()?;
3322        for (hashed_address, account) in &hashed_accounts {
3323            if let Some(account) = account {
3324                hashed_accounts_cursor.upsert(*hashed_address, account)?;
3325            } else if hashed_accounts_cursor.seek_exact(*hashed_address)?.is_some() {
3326                hashed_accounts_cursor.delete_current()?;
3327            }
3328        }
3329
3330        Ok(hashed_accounts)
3331    }
3332
3333    fn unwind_account_hashing_range(
3334        &self,
3335        range: impl RangeBounds<BlockNumber>,
3336    ) -> ProviderResult<BTreeMap<B256, Option<Account>>> {
3337        let changesets = self.account_changesets_range(range)?;
3338        self.unwind_account_hashing(changesets.iter())
3339    }
3340
3341    fn insert_account_for_hashing(
3342        &self,
3343        changesets: impl IntoIterator<Item = (Address, Option<Account>)>,
3344    ) -> ProviderResult<BTreeMap<B256, Option<Account>>> {
3345        let mut hashed_accounts_cursor = self.tx.cursor_write::<tables::HashedAccounts>()?;
3346        let hashed_accounts =
3347            changesets.into_iter().map(|(ad, ac)| (keccak256(ad), ac)).collect::<BTreeMap<_, _>>();
3348        for (hashed_address, account) in &hashed_accounts {
3349            if let Some(account) = account {
3350                hashed_accounts_cursor.upsert(*hashed_address, account)?;
3351            } else if hashed_accounts_cursor.seek_exact(*hashed_address)?.is_some() {
3352                hashed_accounts_cursor.delete_current()?;
3353            }
3354        }
3355        Ok(hashed_accounts)
3356    }
3357
3358    fn unwind_storage_hashing(
3359        &self,
3360        changesets: impl Iterator<Item = (BlockNumberAddress, StorageEntry)>,
3361    ) -> ProviderResult<B256Map<BTreeSet<B256>>> {
3362        // Aggregate all block changesets and make list of accounts that have been changed.
3363        let mut hashed_storages = changesets
3364            .into_iter()
3365            .map(|(BlockNumberAddress((_, address)), storage_entry)| {
3366                let hashed_key = keccak256(storage_entry.key);
3367                (keccak256(address), hashed_key, storage_entry.value)
3368            })
3369            .collect::<Vec<_>>();
3370        hashed_storages.sort_by_key(|(ha, hk, _)| (*ha, *hk));
3371
3372        // Apply values to HashedState, and remove the account if it's None.
3373        let mut hashed_storage_keys: B256Map<BTreeSet<B256>> =
3374            B256Map::with_capacity_and_hasher(hashed_storages.len(), Default::default());
3375        let mut hashed_storage = self.tx.cursor_dup_write::<tables::HashedStorages>()?;
3376        for (hashed_address, key, value) in hashed_storages.into_iter().rev() {
3377            hashed_storage_keys.entry(hashed_address).or_default().insert(key);
3378
3379            if hashed_storage
3380                .seek_by_key_subkey(hashed_address, key)?
3381                .is_some_and(|entry| entry.key == key)
3382            {
3383                hashed_storage.delete_current()?;
3384            }
3385
3386            if !value.is_zero() {
3387                hashed_storage.upsert(hashed_address, &StorageEntry { key, value })?;
3388            }
3389        }
3390        Ok(hashed_storage_keys)
3391    }
3392
3393    fn unwind_storage_hashing_range(
3394        &self,
3395        range: impl RangeBounds<BlockNumber>,
3396    ) -> ProviderResult<B256Map<BTreeSet<B256>>> {
3397        let changesets = self.storage_changesets_range(range)?;
3398        self.unwind_storage_hashing(changesets.into_iter())
3399    }
3400
3401    fn insert_storage_for_hashing(
3402        &self,
3403        storages: impl IntoIterator<Item = (Address, impl IntoIterator<Item = StorageEntry>)>,
3404    ) -> ProviderResult<B256Map<BTreeSet<B256>>> {
3405        // hash values
3406        let hashed_storages =
3407            storages.into_iter().fold(BTreeMap::new(), |mut map, (address, storage)| {
3408                let storage = storage.into_iter().fold(BTreeMap::new(), |mut map, entry| {
3409                    map.insert(keccak256(entry.key), entry.value);
3410                    map
3411                });
3412                map.insert(keccak256(address), storage);
3413                map
3414            });
3415
3416        let hashed_storage_keys = hashed_storages
3417            .iter()
3418            .map(|(hashed_address, entries)| (*hashed_address, entries.keys().copied().collect()))
3419            .collect();
3420
3421        let mut hashed_storage_cursor = self.tx.cursor_dup_write::<tables::HashedStorages>()?;
3422        // Hash the address and key and apply them to HashedStorage (if Storage is None
3423        // just remove it);
3424        hashed_storages.into_iter().try_for_each(|(hashed_address, storage)| {
3425            storage.into_iter().try_for_each(|(key, value)| -> ProviderResult<()> {
3426                if hashed_storage_cursor
3427                    .seek_by_key_subkey(hashed_address, key)?
3428                    .is_some_and(|entry| entry.key == key)
3429                {
3430                    hashed_storage_cursor.delete_current()?;
3431                }
3432
3433                if !value.is_zero() {
3434                    hashed_storage_cursor.upsert(hashed_address, &StorageEntry { key, value })?;
3435                }
3436                Ok(())
3437            })
3438        })?;
3439
3440        Ok(hashed_storage_keys)
3441    }
3442}
3443
3444impl<TX: DbTxMut + DbTx + 'static, N: NodeTypes> HistoryWriter for DatabaseProvider<TX, N> {
3445    fn unwind_account_history_indices<'a>(
3446        &self,
3447        changesets: impl Iterator<Item = &'a (BlockNumber, AccountBeforeTx)>,
3448    ) -> ProviderResult<usize> {
3449        let mut last_indices = changesets
3450            .into_iter()
3451            .map(|(index, account)| (account.address, *index))
3452            .collect::<Vec<_>>();
3453        last_indices.sort_unstable_by_key(|(a, _)| *a);
3454
3455        if self.cached_storage_settings().storage_v2 {
3456            let batch = self.rocksdb_provider.unwind_account_history_indices(&last_indices)?;
3457            self.pending_rocksdb_batches.lock().push(batch);
3458        } else {
3459            // Unwind the account history index in MDBX.
3460            let mut cursor = self.tx.cursor_write::<tables::AccountsHistory>()?;
3461            for &(address, rem_index) in &last_indices {
3462                let partial_shard = unwind_history_shards::<_, tables::AccountsHistory, _>(
3463                    &mut cursor,
3464                    ShardedKey::last(address),
3465                    rem_index,
3466                    |sharded_key| sharded_key.key == address,
3467                )?;
3468
3469                // Check the last returned partial shard.
3470                // If it's not empty, the shard needs to be reinserted.
3471                if !partial_shard.is_empty() {
3472                    cursor.insert(
3473                        ShardedKey::last(address),
3474                        &BlockNumberList::new_pre_sorted(partial_shard),
3475                    )?;
3476                }
3477            }
3478        }
3479
3480        let changesets = last_indices.len();
3481        Ok(changesets)
3482    }
3483
3484    fn unwind_account_history_indices_range(
3485        &self,
3486        range: impl RangeBounds<BlockNumber>,
3487    ) -> ProviderResult<usize> {
3488        let changesets = self.account_changesets_range(range)?;
3489        self.unwind_account_history_indices(changesets.iter())
3490    }
3491
3492    fn insert_account_history_index(
3493        &self,
3494        account_transitions: impl IntoIterator<Item = (Address, impl IntoIterator<Item = u64>)>,
3495    ) -> ProviderResult<()> {
3496        self.append_history_index::<_, tables::AccountsHistory>(
3497            account_transitions,
3498            ShardedKey::new,
3499        )
3500    }
3501
3502    fn unwind_storage_history_indices(
3503        &self,
3504        changesets: impl Iterator<Item = (BlockNumberAddress, StorageEntry)>,
3505    ) -> ProviderResult<usize> {
3506        let mut storage_changesets = changesets
3507            .into_iter()
3508            .map(|(BlockNumberAddress((bn, address)), storage)| (address, storage.key, bn))
3509            .collect::<Vec<_>>();
3510        storage_changesets.sort_unstable_by_key(|(address, key, _)| (*address, *key));
3511
3512        if self.cached_storage_settings().storage_v2 {
3513            let batch =
3514                self.rocksdb_provider.unwind_storage_history_indices(&storage_changesets)?;
3515            self.pending_rocksdb_batches.lock().push(batch);
3516        } else {
3517            // Unwind the storage history index in MDBX.
3518            let mut cursor = self.tx.cursor_write::<tables::StoragesHistory>()?;
3519            for &(address, storage_key, rem_index) in &storage_changesets {
3520                let partial_shard = unwind_history_shards::<_, tables::StoragesHistory, _>(
3521                    &mut cursor,
3522                    StorageShardedKey::last(address, storage_key),
3523                    rem_index,
3524                    |storage_sharded_key| {
3525                        storage_sharded_key.address == address &&
3526                            storage_sharded_key.sharded_key.key == storage_key
3527                    },
3528                )?;
3529
3530                // Check the last returned partial shard.
3531                // If it's not empty, the shard needs to be reinserted.
3532                if !partial_shard.is_empty() {
3533                    cursor.insert(
3534                        StorageShardedKey::last(address, storage_key),
3535                        &BlockNumberList::new_pre_sorted(partial_shard),
3536                    )?;
3537                }
3538            }
3539        }
3540
3541        let changesets = storage_changesets.len();
3542        Ok(changesets)
3543    }
3544
3545    fn unwind_storage_history_indices_range(
3546        &self,
3547        range: impl RangeBounds<BlockNumber>,
3548    ) -> ProviderResult<usize> {
3549        let changesets = self.storage_changesets_range(range)?;
3550        self.unwind_storage_history_indices(changesets.into_iter())
3551    }
3552
3553    fn insert_storage_history_index(
3554        &self,
3555        storage_transitions: impl IntoIterator<Item = ((Address, B256), impl IntoIterator<Item = u64>)>,
3556    ) -> ProviderResult<()> {
3557        self.append_history_index::<_, tables::StoragesHistory>(
3558            storage_transitions,
3559            |(address, storage_key), highest_block_number| {
3560                StorageShardedKey::new(address, storage_key, highest_block_number)
3561            },
3562        )
3563    }
3564
3565    #[instrument(level = "debug", target = "providers::db", skip_all)]
3566    fn update_history_indices(&self, range: RangeInclusive<BlockNumber>) -> ProviderResult<()> {
3567        let storage_settings = self.cached_storage_settings();
3568        if !storage_settings.storage_v2 {
3569            let indices = self.changed_accounts_and_blocks_with_range(range.clone())?;
3570            self.insert_account_history_index(indices)?;
3571        }
3572
3573        if !storage_settings.storage_v2 {
3574            let indices = self.changed_storages_and_blocks_with_range(range)?;
3575            self.insert_storage_history_index(indices)?;
3576        }
3577
3578        Ok(())
3579    }
3580}
3581
3582impl<TX: DbTxMut + DbTx + 'static, N: NodeTypesForProvider> BlockExecutionWriter
3583    for DatabaseProvider<TX, N>
3584{
3585    fn take_block_and_execution_above(
3586        &self,
3587        block: BlockNumber,
3588    ) -> ProviderResult<Chain<Self::Primitives>> {
3589        let range = block + 1..=self.last_block_number()?;
3590
3591        self.unwind_trie_state_from(block + 1)?;
3592
3593        // get execution res
3594        let execution_state = self.take_state_above(block)?;
3595
3596        let blocks = self.recovered_block_range(range)?;
3597
3598        // remove block bodies it is needed for both get block range and get block execution results
3599        // that is why it is deleted afterwards.
3600        self.remove_blocks_above(block)?;
3601
3602        // Update pipeline progress
3603        self.update_pipeline_stages_after_unwind(block)?;
3604
3605        Ok(Chain::new(blocks, execution_state, BTreeMap::new()))
3606    }
3607
3608    fn remove_block_and_execution_above(
3609        &self,
3610        block: BlockNumber,
3611    ) -> ProviderResult<PersistenceFrontiers> {
3612        self.unwind_trie_state_from(block + 1)?;
3613
3614        // remove execution res
3615        self.remove_state_above(block)?;
3616
3617        // remove block bodies it is needed for both get block range and get block execution results
3618        // that is why it is deleted afterwards.
3619        self.remove_blocks_above(block)?;
3620
3621        // Update pipeline progress
3622        self.update_pipeline_stages_after_unwind(block)
3623    }
3624}
3625
3626impl<TX: DbTxMut + DbTx + 'static, N: NodeTypesForProvider> BlockWriter
3627    for DatabaseProvider<TX, N>
3628{
3629    type Block = BlockTy<N>;
3630    type Receipt = ReceiptTy<N>;
3631
3632    /// Inserts the block into the database, writing to both static files and MDBX.
3633    ///
3634    /// This is a convenience method primarily used in tests. For production use,
3635    /// prefer [`Self::save_blocks`] which handles execution output and trie data.
3636    fn insert_block(
3637        &self,
3638        block: &RecoveredBlock<Self::Block>,
3639    ) -> ProviderResult<StoredBlockBodyIndices> {
3640        let block_number = block.number();
3641
3642        // Wrap block in ExecutedBlock with empty execution output (no receipts/state/trie)
3643        let executed_block = ExecutedBlock::new(
3644            Arc::new(block.clone()),
3645            Arc::new(BlockExecutionOutput {
3646                result: BlockExecutionResult {
3647                    receipts: Default::default(),
3648                    requests: Default::default(),
3649                    gas_used: 0,
3650                    blob_gas_used: 0,
3651                },
3652                state: Default::default(),
3653            }),
3654            ComputedTrieData::default(),
3655        );
3656
3657        self.save_blocks_inner(
3658            std::slice::from_ref(&executed_block),
3659            &[],
3660            &[],
3661            None,
3662            SaveBlocksMode::BlocksOnly,
3663        )?;
3664
3665        // Return the body indices
3666        self.block_body_indices(block_number)?
3667            .ok_or(ProviderError::BlockBodyIndicesNotFound(block_number))
3668    }
3669
3670    fn append_block_bodies(
3671        &self,
3672        bodies: Vec<(BlockNumber, Option<&BodyTy<N>>)>,
3673    ) -> ProviderResult<()> {
3674        let Some(from_block) = bodies.first().map(|(block, _)| *block) else { return Ok(()) };
3675
3676        // Initialize writer if we will be writing transactions to staticfiles
3677        let mut tx_writer =
3678            self.static_file_provider.get_writer(from_block, StaticFileSegment::Transactions)?;
3679
3680        let mut block_indices_cursor = self.tx.cursor_write::<tables::BlockBodyIndices>()?;
3681        let mut tx_block_cursor = self.tx.cursor_write::<tables::TransactionBlocks>()?;
3682
3683        // Get id for the next tx_num or zero if there are no transactions.
3684        let mut next_tx_num = tx_block_cursor.last()?.map(|(id, _)| id + 1).unwrap_or_default();
3685
3686        for (block_number, body) in &bodies {
3687            // Increment block on static file header.
3688            tx_writer.increment_block(*block_number)?;
3689
3690            let tx_count = body.as_ref().map(|b| b.transactions().len() as u64).unwrap_or_default();
3691            let block_indices = StoredBlockBodyIndices { first_tx_num: next_tx_num, tx_count };
3692
3693            let mut durations_recorder = metrics::DurationsRecorder::new(&self.metrics);
3694
3695            // insert block meta
3696            block_indices_cursor.append(*block_number, &block_indices)?;
3697
3698            durations_recorder.record_relative(metrics::Action::InsertBlockBodyIndices);
3699
3700            let Some(body) = body else { continue };
3701
3702            // write transaction block index
3703            if !body.transactions().is_empty() {
3704                tx_block_cursor.append(block_indices.last_tx_num(), block_number)?;
3705                durations_recorder.record_relative(metrics::Action::InsertTransactionBlocks);
3706            }
3707
3708            // write transactions
3709            for transaction in body.transactions() {
3710                tx_writer.append_transaction(next_tx_num, transaction)?;
3711
3712                // Increment transaction id for each transaction.
3713                next_tx_num += 1;
3714            }
3715        }
3716
3717        self.storage.writer().write_block_bodies(self, bodies)?;
3718
3719        Ok(())
3720    }
3721
3722    fn remove_blocks_above(&self, block: BlockNumber) -> ProviderResult<()> {
3723        let last_block_number = self.last_block_number()?;
3724        // Clean up HeaderNumbers for blocks being removed, we must clear all indexes from MDBX.
3725        for hash in self.canonical_hashes_range(block + 1, last_block_number + 1)? {
3726            self.tx.delete::<tables::HeaderNumbers>(hash, None)?;
3727        }
3728
3729        // Get highest static file block for the total block range
3730        let highest_static_file_block = self
3731            .static_file_provider()
3732            .get_highest_static_file_block(StaticFileSegment::Headers)
3733            .expect("todo: error handling, headers should exist");
3734
3735        // IMPORTANT: we use `highest_static_file_block.saturating_sub(block_number)` to make sure
3736        // we remove only what is ABOVE the block.
3737        //
3738        // i.e., if the highest static file block is 8, we want to remove above block 5 only, we
3739        // will have three blocks to remove, which will be block 8, 7, and 6.
3740        debug!(target: "providers::db", ?block, "Removing static file blocks above block_number");
3741        self.static_file_provider()
3742            .get_writer(block, StaticFileSegment::Headers)?
3743            .prune_headers(highest_static_file_block.saturating_sub(block))?;
3744
3745        // First transaction to be removed
3746        let unwind_tx_from = self
3747            .block_body_indices(block)?
3748            .map(|b| b.next_tx_num())
3749            .ok_or(ProviderError::BlockBodyIndicesNotFound(block))?;
3750
3751        // Last transaction to be removed
3752        let unwind_tx_to = self
3753            .tx
3754            .cursor_read::<tables::BlockBodyIndices>()?
3755            .last()?
3756            // shouldn't happen because this was OK above
3757            .ok_or(ProviderError::BlockBodyIndicesNotFound(block))?
3758            .1
3759            .last_tx_num();
3760
3761        if unwind_tx_from <= unwind_tx_to {
3762            let hashes = self.transaction_hashes_by_range(unwind_tx_from..(unwind_tx_to + 1))?;
3763            self.with_rocksdb_batch(|batch| {
3764                let mut writer = EitherWriter::new_transaction_hash_numbers(self, batch)?;
3765                for (hash, _) in hashes {
3766                    writer.delete_transaction_hash_number(hash)?;
3767                }
3768                Ok(((), writer.into_raw_rocksdb_batch()))
3769            })?;
3770        }
3771
3772        // Skip sender pruning when sender_recovery is fully pruned, since no sender data
3773        // exists in static files or the database.
3774        if self.prune_modes.sender_recovery.is_none_or(|m| !m.is_full()) {
3775            EitherWriter::new_senders(self, last_block_number)?
3776                .prune_senders(unwind_tx_from, block)?;
3777        }
3778
3779        self.remove_bodies_above(block)?;
3780
3781        Ok(())
3782    }
3783
3784    fn remove_bodies_above(&self, block: BlockNumber) -> ProviderResult<()> {
3785        self.storage.writer().remove_block_bodies_above(self, block)?;
3786
3787        // First transaction to be removed
3788        let unwind_tx_from = self
3789            .block_body_indices(block)?
3790            .map(|b| b.next_tx_num())
3791            .ok_or(ProviderError::BlockBodyIndicesNotFound(block))?;
3792
3793        self.remove::<tables::BlockBodyIndices>(block + 1..)?;
3794        self.remove::<tables::TransactionBlocks>(unwind_tx_from..)?;
3795
3796        let static_file_tx_num =
3797            self.static_file_provider.get_highest_static_file_tx(StaticFileSegment::Transactions);
3798
3799        let to_delete = static_file_tx_num
3800            .map(|static_tx| (static_tx + 1).saturating_sub(unwind_tx_from))
3801            .unwrap_or_default();
3802
3803        self.static_file_provider
3804            .latest_writer(StaticFileSegment::Transactions)?
3805            .prune_transactions(to_delete, block)?;
3806
3807        Ok(())
3808    }
3809
3810    /// Appends blocks with their execution state to the database.
3811    ///
3812    /// **Note:** This function is only used in tests.
3813    ///
3814    /// History indices are written to the appropriate backend based on storage settings:
3815    /// MDBX when `*_history_in_rocksdb` is false, `RocksDB` when true.
3816    ///
3817    /// TODO(joshie): this fn should be moved to `UnifiedStorageWriter` eventually
3818    fn append_blocks_with_state(
3819        &self,
3820        blocks: Vec<RecoveredBlock<Self::Block>>,
3821        execution_outcome: &ExecutionOutcome<Self::Receipt>,
3822        hashed_state: HashedPostStateSorted,
3823    ) -> ProviderResult<()> {
3824        if blocks.is_empty() {
3825            debug!(target: "providers::db", "Attempted to append empty block range");
3826            return Ok(())
3827        }
3828
3829        // Blocks are not empty, so no need to handle the case of `blocks.first()` being
3830        // `None`.
3831        let first_number = blocks[0].number();
3832
3833        // Blocks are not empty, so no need to handle the case of `blocks.last()` being
3834        // `None`.
3835        let last_block_number = blocks[blocks.len() - 1].number();
3836
3837        let mut durations_recorder = metrics::DurationsRecorder::new(&self.metrics);
3838
3839        // Extract account and storage transitions from the bundle reverts BEFORE writing state.
3840        // This is necessary because with edge storage, changesets are written to static files
3841        // whose index isn't updated until commit, making them invisible to subsequent reads
3842        // within the same transaction.
3843        let (account_transitions, storage_transitions) = {
3844            let mut account_transitions: BTreeMap<Address, Vec<u64>> = BTreeMap::new();
3845            let mut storage_transitions: BTreeMap<(Address, B256), Vec<u64>> = BTreeMap::new();
3846            for (block_idx, block_reverts) in execution_outcome.bundle.reverts.iter().enumerate() {
3847                let block_number = first_number + block_idx as u64;
3848                for (address, account_revert) in block_reverts {
3849                    account_transitions.entry(*address).or_default().push(block_number);
3850                    for storage_key in account_revert.storage.keys() {
3851                        let key = B256::from(storage_key.to_be_bytes());
3852                        storage_transitions.entry((*address, key)).or_default().push(block_number);
3853                    }
3854                }
3855            }
3856            (account_transitions, storage_transitions)
3857        };
3858
3859        // Insert the blocks
3860        for block in blocks {
3861            self.insert_block(&block)?;
3862            durations_recorder.record_relative(metrics::Action::InsertBlock);
3863        }
3864
3865        self.write_state(execution_outcome, OriginalValuesKnown::No, StateWriteConfig::default())?;
3866        durations_recorder.record_relative(metrics::Action::InsertState);
3867
3868        // insert hashes and intermediate merkle nodes
3869        self.write_hashed_state(&hashed_state)?;
3870        durations_recorder.record_relative(metrics::Action::InsertHashes);
3871
3872        // Use pre-computed transitions for history indices since static file
3873        // writes aren't visible until commit.
3874        // Note: For MDBX we use insert_*_history_index. For RocksDB we use
3875        // append_*_history_shard which handles read-merge-write internally.
3876        let storage_settings = self.cached_storage_settings();
3877        if storage_settings.storage_v2 {
3878            self.with_rocksdb_batch(|mut batch| {
3879                for (address, blocks) in account_transitions {
3880                    batch.append_account_history_shard(address, blocks)?;
3881                }
3882                Ok(((), Some(batch.into_inner())))
3883            })?;
3884        } else {
3885            self.insert_account_history_index(account_transitions)?;
3886        }
3887        if storage_settings.storage_v2 {
3888            self.with_rocksdb_batch(|mut batch| {
3889                for ((address, key), blocks) in storage_transitions {
3890                    batch.append_storage_history_shard(address, key, blocks)?;
3891                }
3892                Ok(((), Some(batch.into_inner())))
3893            })?;
3894        } else {
3895            self.insert_storage_history_index(storage_transitions)?;
3896        }
3897        durations_recorder.record_relative(metrics::Action::InsertHistoryIndices);
3898
3899        // Update pipeline progress
3900        self.update_pipeline_stages(last_block_number, false)?;
3901        durations_recorder.record_relative(metrics::Action::UpdatePipelineStages);
3902
3903        debug!(target: "providers::db", range = ?first_number..=last_block_number, actions = ?durations_recorder.actions, "Appended blocks");
3904
3905        Ok(())
3906    }
3907}
3908
3909impl<TX: DbTx + 'static, N: NodeTypes> PruneCheckpointReader for DatabaseProvider<TX, N> {
3910    fn get_prune_checkpoint(
3911        &self,
3912        segment: PruneSegment,
3913    ) -> ProviderResult<Option<PruneCheckpoint>> {
3914        Ok(self.tx.get::<tables::PruneCheckpoints>(segment)?)
3915    }
3916
3917    fn get_prune_checkpoints(&self) -> ProviderResult<Vec<(PruneSegment, PruneCheckpoint)>> {
3918        Ok(PruneSegment::variants()
3919            .filter_map(|segment| {
3920                self.tx
3921                    .get::<tables::PruneCheckpoints>(segment)
3922                    .transpose()
3923                    .map(|chk| chk.map(|chk| (segment, chk)))
3924            })
3925            .collect::<Result<_, _>>()?)
3926    }
3927}
3928
3929impl<TX: DbTxMut, N: NodeTypes> PruneCheckpointWriter for DatabaseProvider<TX, N> {
3930    fn save_prune_checkpoint(
3931        &self,
3932        segment: PruneSegment,
3933        checkpoint: PruneCheckpoint,
3934    ) -> ProviderResult<()> {
3935        Ok(self.tx.put::<tables::PruneCheckpoints>(segment, checkpoint)?)
3936    }
3937}
3938
3939impl<TX: DbTx + 'static, N: NodeTypesForProvider> StatsReader for DatabaseProvider<TX, N> {
3940    fn count_entries<T: Table>(&self) -> ProviderResult<usize> {
3941        let db_entries = self.tx.entries::<T>()?;
3942        let static_file_entries = match self.static_file_provider.count_entries::<T>() {
3943            Ok(entries) => entries,
3944            Err(ProviderError::UnsupportedProvider) => 0,
3945            Err(err) => return Err(err),
3946        };
3947
3948        Ok(db_entries + static_file_entries)
3949    }
3950}
3951
3952impl<TX: DbTx + 'static, N: NodeTypes> ChainStateBlockReader for DatabaseProvider<TX, N> {
3953    fn last_finalized_block_number(&self) -> ProviderResult<Option<BlockNumber>> {
3954        let mut finalized_blocks = self
3955            .tx
3956            .cursor_read::<tables::ChainState>()?
3957            .walk(Some(tables::ChainStateKey::LastFinalizedBlock))?
3958            .take(1)
3959            .collect::<Result<BTreeMap<tables::ChainStateKey, BlockNumber>, _>>()?;
3960
3961        let last_finalized_block_number = finalized_blocks.pop_first().map(|pair| pair.1);
3962        Ok(last_finalized_block_number)
3963    }
3964
3965    fn last_safe_block_number(&self) -> ProviderResult<Option<BlockNumber>> {
3966        let mut finalized_blocks = self
3967            .tx
3968            .cursor_read::<tables::ChainState>()?
3969            .walk(Some(tables::ChainStateKey::LastSafeBlock))?
3970            .take(1)
3971            .collect::<Result<BTreeMap<tables::ChainStateKey, BlockNumber>, _>>()?;
3972
3973        let last_finalized_block_number = finalized_blocks.pop_first().map(|pair| pair.1);
3974        Ok(last_finalized_block_number)
3975    }
3976}
3977
3978impl<TX: DbTxMut, N: NodeTypes> ChainStateBlockWriter for DatabaseProvider<TX, N> {
3979    fn save_finalized_block_number(&self, block_number: BlockNumber) -> ProviderResult<()> {
3980        Ok(self
3981            .tx
3982            .put::<tables::ChainState>(tables::ChainStateKey::LastFinalizedBlock, block_number)?)
3983    }
3984
3985    fn save_safe_block_number(&self, block_number: BlockNumber) -> ProviderResult<()> {
3986        Ok(self.tx.put::<tables::ChainState>(tables::ChainStateKey::LastSafeBlock, block_number)?)
3987    }
3988}
3989
3990impl<TX: DbTx + 'static, N: NodeTypes + 'static> DbTxProvider for DatabaseProvider<TX, N> {
3991    type Tx = TX;
3992
3993    fn tx(&self) -> &Self::Tx {
3994        &self.tx
3995    }
3996}
3997
3998impl<TX: DbTx + 'static, N: NodeTypes + 'static> DBProvider for DatabaseProvider<TX, N> {
3999    fn tx_mut(&mut self) -> &mut Self::Tx {
4000        &mut self.tx
4001    }
4002
4003    fn into_tx(self) -> Self::Tx {
4004        self.tx
4005    }
4006
4007    fn prune_modes_ref(&self) -> &PruneModes {
4008        self.prune_modes_ref()
4009    }
4010
4011    /// Commit database transaction, static files, and pending `RocksDB` batches.
4012    #[instrument(
4013        name = "DatabaseProvider::commit",
4014        level = "debug",
4015        target = "providers::db",
4016        skip_all
4017    )]
4018    fn commit(self) -> ProviderResult<()> {
4019        if self.static_file_provider.has_unwind_queued() || self.commit_order.is_unwind() {
4020            self.commit_unwind()?;
4021        } else {
4022            // Normal path: finalize() will call sync_all() if not already synced
4023            let mut timings = metrics::CommitTimings::default();
4024
4025            let start = Instant::now();
4026            self.static_file_provider.finalize()?;
4027            timings.sf = start.elapsed();
4028
4029            let start = Instant::now();
4030            let batches = std::mem::take(&mut *self.pending_rocksdb_batches.lock());
4031            for batch in batches {
4032                self.rocksdb_provider.commit_batch(batch)?;
4033            }
4034            timings.rocksdb = start.elapsed();
4035
4036            let start = Instant::now();
4037            self.tx.commit()?;
4038            timings.mdbx = start.elapsed();
4039
4040            self.metrics.record_commit(&timings);
4041        }
4042
4043        Ok(())
4044    }
4045}
4046
4047impl<TX: DbTx, N: NodeTypes> MetadataProvider for DatabaseProvider<TX, N> {
4048    fn get_metadata(&self, key: &str) -> ProviderResult<Option<Vec<u8>>> {
4049        self.tx.get::<tables::Metadata>(key.to_string()).map_err(Into::into)
4050    }
4051}
4052
4053impl<TX: DbTxMut, N: NodeTypes> MetadataWriter for DatabaseProvider<TX, N> {
4054    fn write_metadata(&self, key: &str, value: Vec<u8>) -> ProviderResult<()> {
4055        self.tx.put::<tables::Metadata>(key.to_string(), value).map_err(Into::into)
4056    }
4057
4058    fn delete_metadata(&self, key: &str) -> ProviderResult<()> {
4059        self.tx.delete::<tables::Metadata>(key.to_string(), None)?;
4060        Ok(())
4061    }
4062}
4063
4064impl<TX: Send, N: NodeTypes> StorageSettingsCache for DatabaseProvider<TX, N> {
4065    fn cached_storage_settings(&self) -> StorageSettings {
4066        *self.storage_settings.read()
4067    }
4068
4069    fn set_storage_settings_cache(&self, settings: StorageSettings) {
4070        *self.storage_settings.write() = settings;
4071    }
4072}
4073
4074impl<TX: Send, N: NodeTypes> StoragePath for DatabaseProvider<TX, N> {
4075    fn storage_path(&self) -> PathBuf {
4076        self.db_path.clone()
4077    }
4078}
4079
4080#[cfg(test)]
4081mod tests {
4082    use super::*;
4083    use crate::{
4084        test_utils::{blocks::BlockchainTestData, create_test_provider_factory},
4085        BlockWriter,
4086    };
4087    use alloy_consensus::Header;
4088    use alloy_primitives::{
4089        map::{AddressMap, B256Map},
4090        U256,
4091    };
4092    #[cfg(feature = "partial-persistence")]
4093    use reth_chain_state::test_utils::TestBlockBuilder;
4094    use reth_chain_state::ExecutedBlock;
4095    use reth_db_api::models::StorageSettings;
4096    use reth_ethereum_primitives::Receipt;
4097    use reth_execution_types::{AccountRevertInit, BlockExecutionOutput, BlockExecutionResult};
4098    use reth_primitives_traits::SealedBlock;
4099    use reth_storage_api::{MetadataProvider, MetadataWriter};
4100    use reth_testing_utils::generators::{self, random_block, BlockParams};
4101    use reth_trie::{
4102        HashedPostState, KeccakKeyHasher, Nibbles, SortedTrieData, StoredNibbles,
4103        StoredNibblesSubKey,
4104    };
4105    use revm::{database::BundleState, state::AccountInfo};
4106    use std::{sync::mpsc, time::Duration};
4107
4108    /// Seeds block zero through the writer core because [`SaveBlocksInput`] only describes
4109    /// advancing an existing persistence frontier.
4110    fn save_genesis<TX, N>(
4111        provider: &DatabaseProvider<TX, N>,
4112        genesis: &ExecutedBlock<N::Primitives>,
4113    ) -> ProviderResult<()>
4114    where
4115        TX: DbTx + DbTxMut + 'static,
4116        N: NodeTypesForProvider,
4117    {
4118        assert_eq!(genesis.recovered_block().number(), 0);
4119        provider.save_blocks_inner(
4120            std::slice::from_ref(genesis),
4121            std::slice::from_ref(genesis),
4122            &[],
4123            None,
4124            SaveBlocksMode::Full,
4125        )
4126    }
4127
4128    #[test]
4129    fn test_receipts_by_block_range_empty_range() {
4130        let factory = create_test_provider_factory();
4131        let provider = factory.provider().unwrap();
4132
4133        // empty range should return empty vec
4134        let start = 10u64;
4135        let end = 9u64;
4136        let result = provider.receipts_by_block_range(start..=end).unwrap();
4137        assert_eq!(result, Vec::<Vec<reth_ethereum_primitives::Receipt>>::new());
4138    }
4139
4140    #[test]
4141    fn metadata_can_be_deleted() {
4142        let factory = create_test_provider_factory();
4143        let key = "metadata-delete-test";
4144
4145        let provider_rw = factory.provider_rw().unwrap();
4146        provider_rw.write_metadata(key, vec![1]).unwrap();
4147        provider_rw.commit().unwrap();
4148        assert_eq!(factory.provider().unwrap().get_metadata(key).unwrap(), Some(vec![1]));
4149
4150        let provider_rw = factory.provider_rw().unwrap();
4151        provider_rw.delete_metadata(key).unwrap();
4152        provider_rw.commit().unwrap();
4153        assert_eq!(factory.provider().unwrap().get_metadata(key).unwrap(), None);
4154    }
4155
4156    #[test]
4157    fn unwind_commit_waits_for_pre_commit_readers() {
4158        let factory = create_test_provider_factory();
4159
4160        let reader = factory.provider().unwrap();
4161        let provider_rw = factory.unwind_provider_rw().unwrap();
4162        provider_rw.write_metadata("unwind-wait-test", vec![1]).unwrap();
4163        let (done_tx, done_rx) = mpsc::channel();
4164
4165        let handle = std::thread::spawn(move || {
4166            let result = provider_rw.commit();
4167            done_tx.send(result).unwrap();
4168        });
4169
4170        assert!(
4171            done_rx.recv_timeout(Duration::from_millis(50)).is_err(),
4172            "unwind commit should wait while an older read transaction is still open"
4173        );
4174
4175        drop(reader);
4176
4177        done_rx.recv_timeout(Duration::from_secs(1)).unwrap().unwrap();
4178        handle.join().unwrap();
4179    }
4180
4181    #[test]
4182    fn test_receipts_by_block_range_nonexistent_blocks() {
4183        let factory = create_test_provider_factory();
4184        let provider = factory.provider().unwrap();
4185
4186        // non-existent blocks should return empty vecs for each block
4187        let result = provider.receipts_by_block_range(10..=12).unwrap();
4188        assert_eq!(result, vec![vec![], vec![], vec![]]);
4189    }
4190
4191    #[test]
4192    fn test_receipts_by_block_range_single_block() {
4193        let factory = create_test_provider_factory();
4194        let data = BlockchainTestData::default();
4195
4196        let provider_rw = factory.provider_rw().unwrap();
4197        provider_rw.insert_block(&data.genesis.try_recover().unwrap()).unwrap();
4198        provider_rw
4199            .write_state(
4200                &ExecutionOutcome { first_block: 0, receipts: vec![vec![]], ..Default::default() },
4201                crate::OriginalValuesKnown::No,
4202                StateWriteConfig::default(),
4203            )
4204            .unwrap();
4205        provider_rw.insert_block(&data.blocks[0].0).unwrap();
4206        provider_rw
4207            .write_state(
4208                &data.blocks[0].1,
4209                crate::OriginalValuesKnown::No,
4210                StateWriteConfig::default(),
4211            )
4212            .unwrap();
4213        provider_rw.commit().unwrap();
4214
4215        let provider = factory.provider().unwrap();
4216        let result = provider.receipts_by_block_range(1..=1).unwrap();
4217
4218        // should have one vec with one receipt
4219        assert_eq!(result.len(), 1);
4220        assert_eq!(result[0].len(), 1);
4221        assert_eq!(result[0][0], data.blocks[0].1.receipts()[0][0]);
4222    }
4223
4224    #[test]
4225    fn test_receipts_by_block_range_multiple_blocks() {
4226        let factory = create_test_provider_factory();
4227        let data = BlockchainTestData::default();
4228
4229        let provider_rw = factory.provider_rw().unwrap();
4230        provider_rw.insert_block(&data.genesis.try_recover().unwrap()).unwrap();
4231        provider_rw
4232            .write_state(
4233                &ExecutionOutcome { first_block: 0, receipts: vec![vec![]], ..Default::default() },
4234                crate::OriginalValuesKnown::No,
4235                StateWriteConfig::default(),
4236            )
4237            .unwrap();
4238        for i in 0..3 {
4239            provider_rw.insert_block(&data.blocks[i].0).unwrap();
4240            provider_rw
4241                .write_state(
4242                    &data.blocks[i].1,
4243                    crate::OriginalValuesKnown::No,
4244                    StateWriteConfig::default(),
4245                )
4246                .unwrap();
4247        }
4248        provider_rw.commit().unwrap();
4249
4250        let provider = factory.provider().unwrap();
4251        let result = provider.receipts_by_block_range(1..=3).unwrap();
4252
4253        // should have 3 vecs, each with one receipt
4254        assert_eq!(result.len(), 3);
4255        for (i, block_receipts) in result.iter().enumerate() {
4256            assert_eq!(block_receipts.len(), 1);
4257            assert_eq!(block_receipts[0], data.blocks[i].1.receipts()[0][0]);
4258        }
4259    }
4260
4261    #[test]
4262    fn test_receipts_by_block_range_blocks_with_varying_tx_counts() {
4263        let factory = create_test_provider_factory();
4264        let data = BlockchainTestData::default();
4265
4266        let provider_rw = factory.provider_rw().unwrap();
4267        provider_rw.insert_block(&data.genesis.try_recover().unwrap()).unwrap();
4268        provider_rw
4269            .write_state(
4270                &ExecutionOutcome { first_block: 0, receipts: vec![vec![]], ..Default::default() },
4271                crate::OriginalValuesKnown::No,
4272                StateWriteConfig::default(),
4273            )
4274            .unwrap();
4275
4276        // insert blocks 1-3 with receipts
4277        for i in 0..3 {
4278            provider_rw.insert_block(&data.blocks[i].0).unwrap();
4279            provider_rw
4280                .write_state(
4281                    &data.blocks[i].1,
4282                    crate::OriginalValuesKnown::No,
4283                    StateWriteConfig::default(),
4284                )
4285                .unwrap();
4286        }
4287        provider_rw.commit().unwrap();
4288
4289        let provider = factory.provider().unwrap();
4290        let result = provider.receipts_by_block_range(1..=3).unwrap();
4291
4292        // verify each block has one receipt
4293        assert_eq!(result.len(), 3);
4294        for block_receipts in &result {
4295            assert_eq!(block_receipts.len(), 1);
4296        }
4297    }
4298
4299    #[test]
4300    fn test_receipts_by_block_range_partial_range() {
4301        let factory = create_test_provider_factory();
4302        let data = BlockchainTestData::default();
4303
4304        let provider_rw = factory.provider_rw().unwrap();
4305        provider_rw.insert_block(&data.genesis.try_recover().unwrap()).unwrap();
4306        provider_rw
4307            .write_state(
4308                &ExecutionOutcome { first_block: 0, receipts: vec![vec![]], ..Default::default() },
4309                crate::OriginalValuesKnown::No,
4310                StateWriteConfig::default(),
4311            )
4312            .unwrap();
4313        for i in 0..3 {
4314            provider_rw.insert_block(&data.blocks[i].0).unwrap();
4315            provider_rw
4316                .write_state(
4317                    &data.blocks[i].1,
4318                    crate::OriginalValuesKnown::No,
4319                    StateWriteConfig::default(),
4320                )
4321                .unwrap();
4322        }
4323        provider_rw.commit().unwrap();
4324
4325        let provider = factory.provider().unwrap();
4326
4327        // request range that includes both existing and non-existing blocks
4328        let result = provider.receipts_by_block_range(2..=5).unwrap();
4329        assert_eq!(result.len(), 4);
4330
4331        // blocks 2-3 should have receipts, blocks 4-5 should be empty
4332        assert_eq!(result[0].len(), 1); // block 2
4333        assert_eq!(result[1].len(), 1); // block 3
4334        assert_eq!(result[2].len(), 0); // block 4 (doesn't exist)
4335        assert_eq!(result[3].len(), 0); // block 5 (doesn't exist)
4336
4337        assert_eq!(result[0][0], data.blocks[1].1.receipts()[0][0]);
4338        assert_eq!(result[1][0], data.blocks[2].1.receipts()[0][0]);
4339    }
4340
4341    #[test]
4342    fn test_receipts_by_block_range_all_empty_blocks() {
4343        let factory = create_test_provider_factory();
4344        let mut rng = generators::rng();
4345
4346        // create blocks with no transactions
4347        let mut blocks = Vec::new();
4348        for i in 0..3 {
4349            let block =
4350                random_block(&mut rng, i, BlockParams { tx_count: Some(0), ..Default::default() });
4351            blocks.push(block);
4352        }
4353
4354        let provider_rw = factory.provider_rw().unwrap();
4355        for block in blocks {
4356            provider_rw.insert_block(&block.try_recover().unwrap()).unwrap();
4357        }
4358        provider_rw.commit().unwrap();
4359
4360        let provider = factory.provider().unwrap();
4361        let result = provider.receipts_by_block_range(1..=3).unwrap();
4362
4363        assert_eq!(result.len(), 3);
4364        for block_receipts in result {
4365            assert_eq!(block_receipts.len(), 0);
4366        }
4367    }
4368
4369    #[test]
4370    fn test_receipts_by_block_range_consistency_with_individual_calls() {
4371        let factory = create_test_provider_factory();
4372        let data = BlockchainTestData::default();
4373
4374        let provider_rw = factory.provider_rw().unwrap();
4375        provider_rw.insert_block(&data.genesis.try_recover().unwrap()).unwrap();
4376        provider_rw
4377            .write_state(
4378                &ExecutionOutcome { first_block: 0, receipts: vec![vec![]], ..Default::default() },
4379                crate::OriginalValuesKnown::No,
4380                StateWriteConfig::default(),
4381            )
4382            .unwrap();
4383        for i in 0..3 {
4384            provider_rw.insert_block(&data.blocks[i].0).unwrap();
4385            provider_rw
4386                .write_state(
4387                    &data.blocks[i].1,
4388                    crate::OriginalValuesKnown::No,
4389                    StateWriteConfig::default(),
4390                )
4391                .unwrap();
4392        }
4393        provider_rw.commit().unwrap();
4394
4395        let provider = factory.provider().unwrap();
4396
4397        // get receipts using block range method
4398        let range_result = provider.receipts_by_block_range(1..=3).unwrap();
4399
4400        // get receipts using individual block calls
4401        let mut individual_results = Vec::new();
4402        for block_num in 1..=3 {
4403            let receipts =
4404                provider.receipts_by_block(block_num.into()).unwrap().unwrap_or_default();
4405            individual_results.push(receipts);
4406        }
4407
4408        assert_eq!(range_result, individual_results);
4409    }
4410
4411    #[test]
4412    fn test_receipts_by_block_returns_none_for_missing_unpruned_receipts() {
4413        let factory = create_test_provider_factory();
4414        let data = BlockchainTestData::default();
4415
4416        let provider_rw = factory.provider_rw().unwrap();
4417        provider_rw.insert_block(&data.genesis.try_recover().unwrap()).unwrap();
4418        provider_rw.insert_block(&data.blocks[0].0).unwrap();
4419        provider_rw.commit().unwrap();
4420
4421        let provider = factory.provider().unwrap();
4422        assert!(provider.receipts_by_block(1.into()).unwrap().is_none());
4423    }
4424
4425    #[test]
4426    fn test_write_trie_updates_sorted() {
4427        use reth_trie::{
4428            updates::{StorageTrieUpdatesSorted, TrieUpdatesSorted},
4429            BranchNodeCompact, StorageTrieEntry,
4430        };
4431
4432        let factory = create_test_provider_factory();
4433        let provider_rw = factory.provider_rw().unwrap();
4434
4435        // Pre-populate account trie with data that will be deleted
4436        {
4437            let tx = provider_rw.tx_ref();
4438            let mut cursor = tx.cursor_write::<tables::AccountsTrie>().unwrap();
4439
4440            // Add account node that will be deleted
4441            let to_delete = StoredNibbles(Nibbles::from_nibbles([0x3, 0x4]));
4442            cursor
4443                .upsert(
4444                    to_delete,
4445                    &BranchNodeCompact::new(
4446                        0b1010_1010_1010_1010, // state_mask
4447                        0b0000_0000_0000_0000, // tree_mask
4448                        0b0000_0000_0000_0000, // hash_mask
4449                        vec![],
4450                        None,
4451                    ),
4452                )
4453                .unwrap();
4454
4455            // Add account node that will be updated
4456            let to_update = StoredNibbles(Nibbles::from_nibbles([0x1, 0x2]));
4457            cursor
4458                .upsert(
4459                    to_update,
4460                    &BranchNodeCompact::new(
4461                        0b0101_0101_0101_0101, // old state_mask (will be updated)
4462                        0b0000_0000_0000_0000, // tree_mask
4463                        0b0000_0000_0000_0000, // hash_mask
4464                        vec![],
4465                        None,
4466                    ),
4467                )
4468                .unwrap();
4469        }
4470
4471        // Pre-populate storage tries with data
4472        let storage_address1 = B256::from([1u8; 32]);
4473        let storage_address2 = B256::from([2u8; 32]);
4474        {
4475            let tx = provider_rw.tx_ref();
4476            let mut storage_cursor = tx.cursor_dup_write::<tables::StoragesTrie>().unwrap();
4477
4478            // Add storage nodes for address1 (one will be deleted)
4479            storage_cursor
4480                .upsert(
4481                    storage_address1,
4482                    &StorageTrieEntry {
4483                        nibbles: StoredNibblesSubKey(Nibbles::from_nibbles([0x2, 0x0])),
4484                        node: BranchNodeCompact::new(
4485                            0b0011_0011_0011_0011, // will be deleted
4486                            0b0000_0000_0000_0000,
4487                            0b0000_0000_0000_0000,
4488                            vec![],
4489                            None,
4490                        ),
4491                    },
4492                )
4493                .unwrap();
4494
4495            // Add storage nodes for address2 (will be wiped)
4496            storage_cursor
4497                .upsert(
4498                    storage_address2,
4499                    &StorageTrieEntry {
4500                        nibbles: StoredNibblesSubKey(Nibbles::from_nibbles([0xa, 0xb])),
4501                        node: BranchNodeCompact::new(
4502                            0b1100_1100_1100_1100, // will be wiped
4503                            0b0000_0000_0000_0000,
4504                            0b0000_0000_0000_0000,
4505                            vec![],
4506                            None,
4507                        ),
4508                    },
4509                )
4510                .unwrap();
4511            storage_cursor
4512                .upsert(
4513                    storage_address2,
4514                    &StorageTrieEntry {
4515                        nibbles: StoredNibblesSubKey(Nibbles::from_nibbles([0xc, 0xd])),
4516                        node: BranchNodeCompact::new(
4517                            0b0011_1100_0011_1100, // will be wiped
4518                            0b0000_0000_0000_0000,
4519                            0b0000_0000_0000_0000,
4520                            vec![],
4521                            None,
4522                        ),
4523                    },
4524                )
4525                .unwrap();
4526        }
4527
4528        // Create sorted account trie updates
4529        let account_nodes = vec![
4530            (
4531                Nibbles::from_nibbles([0x1, 0x2]),
4532                Some(BranchNodeCompact::new(
4533                    0b1111_1111_1111_1111, // state_mask (updated)
4534                    0b0000_0000_0000_0000, // tree_mask
4535                    0b0000_0000_0000_0000, // hash_mask (no hashes)
4536                    vec![],
4537                    None,
4538                )),
4539            ),
4540            (Nibbles::from_nibbles([0x3, 0x4]), None), // Deletion
4541            (
4542                Nibbles::from_nibbles([0x5, 0x6]),
4543                Some(BranchNodeCompact::new(
4544                    0b1111_1111_1111_1111, // state_mask
4545                    0b0000_0000_0000_0000, // tree_mask
4546                    0b0000_0000_0000_0000, // hash_mask (no hashes)
4547                    vec![],
4548                    None,
4549                )),
4550            ),
4551        ];
4552
4553        // Create sorted storage trie updates
4554        let storage_trie1 = StorageTrieUpdatesSorted {
4555            is_deleted: false,
4556            storage_nodes: vec![
4557                (
4558                    Nibbles::from_nibbles([0x1, 0x0]),
4559                    Some(BranchNodeCompact::new(
4560                        0b1111_0000_0000_0000, // state_mask
4561                        0b0000_0000_0000_0000, // tree_mask
4562                        0b0000_0000_0000_0000, // hash_mask (no hashes)
4563                        vec![],
4564                        None,
4565                    )),
4566                ),
4567                (Nibbles::from_nibbles([0x2, 0x0]), None), // Deletion of existing node
4568            ],
4569        };
4570
4571        let storage_trie2 = StorageTrieUpdatesSorted {
4572            is_deleted: true, // Wipe all storage for this address
4573            storage_nodes: vec![],
4574        };
4575
4576        let mut storage_tries = B256Map::default();
4577        storage_tries.insert(storage_address1, storage_trie1);
4578        storage_tries.insert(storage_address2, storage_trie2);
4579
4580        let trie_updates = TrieUpdatesSorted::new(account_nodes, storage_tries);
4581
4582        // Write the sorted trie updates
4583        let num_entries = provider_rw.write_trie_updates_sorted(&trie_updates).unwrap();
4584
4585        // We should have 2 account insertions + 1 account deletion + 1 storage insertion + 1
4586        // storage deletion = 5
4587        assert_eq!(num_entries, 5);
4588
4589        // Verify account trie updates were written correctly
4590        let tx = provider_rw.tx_ref();
4591        let mut cursor = tx.cursor_read::<tables::AccountsTrie>().unwrap();
4592
4593        // Check first account node was updated
4594        let nibbles1 = StoredNibbles(Nibbles::from_nibbles([0x1, 0x2]));
4595        let entry1 = cursor.seek_exact(nibbles1).unwrap();
4596        assert!(entry1.is_some(), "Updated account node should exist");
4597        let expected_mask = reth_trie::TrieMask::new(0b1111_1111_1111_1111);
4598        assert_eq!(
4599            entry1.unwrap().1.state_mask,
4600            expected_mask,
4601            "Account node should have updated state_mask"
4602        );
4603
4604        // Check deleted account node no longer exists
4605        let nibbles2 = StoredNibbles(Nibbles::from_nibbles([0x3, 0x4]));
4606        let entry2 = cursor.seek_exact(nibbles2).unwrap();
4607        assert!(entry2.is_none(), "Deleted account node should not exist");
4608
4609        // Check new account node exists
4610        let nibbles3 = StoredNibbles(Nibbles::from_nibbles([0x5, 0x6]));
4611        let entry3 = cursor.seek_exact(nibbles3).unwrap();
4612        assert!(entry3.is_some(), "New account node should exist");
4613
4614        // Verify storage trie updates were written correctly
4615        let mut storage_cursor = tx.cursor_dup_read::<tables::StoragesTrie>().unwrap();
4616
4617        // Check storage for address1
4618        let storage_entries1: Vec<_> = storage_cursor
4619            .walk_dup(Some(storage_address1), None)
4620            .unwrap()
4621            .collect::<Result<Vec<_>, _>>()
4622            .unwrap();
4623        assert_eq!(
4624            storage_entries1.len(),
4625            1,
4626            "Storage address1 should have 1 entry after deletion"
4627        );
4628        assert_eq!(
4629            storage_entries1[0].1.nibbles.0,
4630            Nibbles::from_nibbles([0x1, 0x0]),
4631            "Remaining entry should be [0x1, 0x0]"
4632        );
4633
4634        // Check storage for address2 was wiped
4635        let storage_entries2: Vec<_> = storage_cursor
4636            .walk_dup(Some(storage_address2), None)
4637            .unwrap()
4638            .collect::<Result<Vec<_>, _>>()
4639            .unwrap();
4640        assert_eq!(storage_entries2.len(), 0, "Storage address2 should be empty after wipe");
4641
4642        provider_rw.commit().unwrap();
4643    }
4644
4645    #[cfg(feature = "partial-persistence")]
4646    #[test]
4647    fn test_save_blocks_only_masks_trie_with_deferred_blocks() {
4648        use reth_trie::{
4649            updates::{StorageTrieUpdatesSorted, TrieUpdatesSorted},
4650            BranchNodeCompact, HashedPostStateSorted, HashedStorageSorted,
4651        };
4652
4653        fn branch(mask: u16) -> BranchNodeCompact {
4654            BranchNodeCompact::new(mask, 0, 0, vec![], None)
4655        }
4656
4657        let factory = create_test_provider_factory();
4658        factory.set_storage_settings_cache(StorageSettings::v1());
4659
4660        let mut test_block_builder = TestBlockBuilder::eth().with_state();
4661        let genesis = test_block_builder.get_executed_blocks(0..1).next().unwrap();
4662        let blocks: Vec<_> = test_block_builder.get_executed_blocks(1..3).collect();
4663
4664        let provider_rw = factory.provider_rw().unwrap();
4665        save_genesis(&provider_rw, &genesis).unwrap();
4666        provider_rw.commit().unwrap();
4667
4668        let kept_account = B256::with_last_byte(0x11);
4669        let masked_account = B256::with_last_byte(0x12);
4670        let kept_storage = B256::with_last_byte(0x21);
4671        let masked_storage = B256::with_last_byte(0x22);
4672        let kept_slot = B256::with_last_byte(0x31);
4673        let masked_slot = B256::with_last_byte(0x32);
4674        let kept_account_node = Nibbles::from_nibbles([0x1, 0x2]);
4675        let masked_account_node = Nibbles::from_nibbles([0x1, 0x3]);
4676        let kept_storage_node = Nibbles::from_nibbles([0x2, 0x1]);
4677        let masked_storage_node = Nibbles::from_nibbles([0x2, 0x2]);
4678        let full_persist_base = &blocks[0];
4679        let deferred_trie_base = &blocks[1];
4680
4681        let full_persist_hashed_state = HashedPostStateSorted::new(
4682            vec![
4683                (kept_account, Some(Account::default())),
4684                (masked_account, Some(Account { nonce: 1, ..Default::default() })),
4685            ],
4686            B256Map::from_iter([
4687                (
4688                    kept_storage,
4689                    HashedStorageSorted {
4690                        wiped: false,
4691                        storage_slots: vec![(kept_slot, U256::from(1))],
4692                    },
4693                ),
4694                (
4695                    masked_storage,
4696                    HashedStorageSorted {
4697                        wiped: false,
4698                        storage_slots: vec![(masked_slot, U256::from(2))],
4699                    },
4700                ),
4701            ]),
4702        );
4703        let full_persist_trie_updates = TrieUpdatesSorted::new(
4704            vec![
4705                (kept_account_node, Some(branch(0b0000_1111_0000_1111))),
4706                (masked_account_node, Some(branch(0b1111_0000_1111_0000))),
4707            ],
4708            B256Map::from_iter([
4709                (
4710                    kept_storage,
4711                    StorageTrieUpdatesSorted {
4712                        is_deleted: false,
4713                        storage_nodes: vec![(kept_storage_node, Some(branch(0b1010)))],
4714                    },
4715                ),
4716                (
4717                    masked_storage,
4718                    StorageTrieUpdatesSorted {
4719                        is_deleted: false,
4720                        storage_nodes: vec![(masked_storage_node, Some(branch(0b0101)))],
4721                    },
4722                ),
4723            ]),
4724        );
4725
4726        let full_persist_block = ExecutedBlock::new(
4727            Arc::clone(&full_persist_base.recovered_block),
4728            Arc::clone(&full_persist_base.execution_output),
4729            ComputedTrieData::new(
4730                Arc::new(full_persist_hashed_state),
4731                Arc::new(full_persist_trie_updates),
4732            ),
4733        );
4734
4735        let deferred_trie_hashed_state = HashedPostStateSorted::new(
4736            vec![(masked_account, Some(Account { nonce: 3, ..Default::default() }))],
4737            B256Map::from_iter([(
4738                masked_storage,
4739                HashedStorageSorted {
4740                    wiped: false,
4741                    storage_slots: vec![(masked_slot, U256::from(4))],
4742                },
4743            )]),
4744        );
4745        let deferred_trie_updates = TrieUpdatesSorted::new(
4746            vec![(masked_account_node, Some(branch(0b0011_0011)))],
4747            B256Map::from_iter([(
4748                masked_storage,
4749                StorageTrieUpdatesSorted {
4750                    is_deleted: false,
4751                    storage_nodes: vec![(masked_storage_node, Some(branch(0b1100)))],
4752                },
4753            )]),
4754        );
4755        let deferred_trie_block = ExecutedBlock::new(
4756            Arc::clone(&deferred_trie_base.recovered_block),
4757            Arc::clone(&deferred_trie_base.execution_output),
4758            ComputedTrieData::new(
4759                Arc::new(deferred_trie_hashed_state),
4760                Arc::new(deferred_trie_updates),
4761            ),
4762        );
4763
4764        let provider_rw = factory.provider_rw().unwrap();
4765        let input = SaveBlocksInput::new(
4766            vec![full_persist_block.clone(), deferred_trie_block.clone()],
4767            0,
4768            0,
4769            2,
4770            0,
4771        );
4772        provider_rw.save_blocks(&input).unwrap();
4773        provider_rw.commit().unwrap();
4774
4775        let provider_rw = factory.provider_rw().unwrap();
4776        let input =
4777            SaveBlocksInput::new(vec![full_persist_block, deferred_trie_block.clone()], 2, 0, 2, 1);
4778        assert!(input.persist_rest_blocks().is_empty());
4779        provider_rw.save_blocks(&input).unwrap();
4780        provider_rw.commit().unwrap();
4781
4782        let provider = factory.provider().unwrap();
4783        let tx = provider.tx_ref();
4784        let finish_checkpoint = provider.get_stage_checkpoint(StageId::Finish).unwrap().unwrap();
4785        assert_eq!(finish_checkpoint.block_number, 2);
4786        assert_eq!(
4787            finish_checkpoint.finish_stage_checkpoint().unwrap().partial_state_trie,
4788            Some(1)
4789        );
4790        assert!(provider.block_hash(2).unwrap().is_some());
4791
4792        let mut hashed_accounts = tx.cursor_read::<tables::HashedAccounts>().unwrap();
4793        assert!(hashed_accounts.seek_exact(kept_account).unwrap().is_some());
4794        assert!(hashed_accounts.seek_exact(masked_account).unwrap().is_none());
4795
4796        let mut hashed_storages = tx.cursor_dup_read::<tables::HashedStorages>().unwrap();
4797        assert!(hashed_storages.seek_by_key_subkey(kept_storage, kept_slot).unwrap().is_some());
4798        assert!(hashed_storages
4799            .walk_dup(Some(masked_storage), None)
4800            .unwrap()
4801            .next()
4802            .transpose()
4803            .unwrap()
4804            .is_none());
4805
4806        let mut account_trie = tx.cursor_read::<tables::AccountsTrie>().unwrap();
4807        assert!(account_trie.seek_exact(StoredNibbles(kept_account_node)).unwrap().is_some());
4808        assert!(account_trie.seek_exact(StoredNibbles(masked_account_node)).unwrap().is_none());
4809
4810        let mut storage_trie = tx.cursor_dup_read::<tables::StoragesTrie>().unwrap();
4811        let kept_entries: Vec<_> = storage_trie
4812            .walk_dup(Some(kept_storage), None)
4813            .unwrap()
4814            .collect::<Result<Vec<_>, _>>()
4815            .unwrap();
4816        assert_eq!(kept_entries.len(), 1);
4817        assert_eq!(kept_entries[0].1.nibbles.0, kept_storage_node);
4818
4819        let masked_entries: Vec<_> = storage_trie
4820            .walk_dup(Some(masked_storage), None)
4821            .unwrap()
4822            .collect::<Result<Vec<_>, _>>()
4823            .unwrap();
4824        assert!(masked_entries.is_empty());
4825
4826        drop(storage_trie);
4827        drop(account_trie);
4828        drop(hashed_storages);
4829        drop(hashed_accounts);
4830        drop(provider);
4831
4832        let provider_rw = factory.provider_rw().unwrap();
4833        let input = SaveBlocksInput::new(vec![deferred_trie_block], 2, 1, 2, 2);
4834        assert!(input.persist_rest_blocks().is_empty());
4835        provider_rw.save_blocks(&input).unwrap();
4836        provider_rw.commit().unwrap();
4837
4838        let provider = factory.provider().unwrap();
4839        let finish_checkpoint = provider.get_stage_checkpoint(StageId::Finish).unwrap().unwrap();
4840        assert_eq!(finish_checkpoint.block_number, 2);
4841        assert!(finish_checkpoint.finish_stage_checkpoint().is_none());
4842
4843        let mut hashed_accounts =
4844            provider.tx_ref().cursor_read::<tables::HashedAccounts>().unwrap();
4845        let (_, account) = hashed_accounts.seek_exact(masked_account).unwrap().unwrap();
4846        assert_eq!(account.nonce, 3);
4847
4848        let mut hashed_storages =
4849            provider.tx_ref().cursor_dup_read::<tables::HashedStorages>().unwrap();
4850        let storage =
4851            hashed_storages.seek_by_key_subkey(masked_storage, masked_slot).unwrap().unwrap();
4852        assert_eq!(storage.value, U256::from(4));
4853
4854        let mut account_trie = provider.tx_ref().cursor_read::<tables::AccountsTrie>().unwrap();
4855        assert!(account_trie.seek_exact(StoredNibbles(masked_account_node)).unwrap().is_some());
4856
4857        let mut storage_trie = provider.tx_ref().cursor_dup_read::<tables::StoragesTrie>().unwrap();
4858        let masked_entries: Vec<_> = storage_trie
4859            .walk_dup(Some(masked_storage), None)
4860            .unwrap()
4861            .collect::<Result<Vec<_>, _>>()
4862            .unwrap();
4863        assert_eq!(masked_entries.len(), 1);
4864        assert_eq!(masked_entries[0].1.nibbles.0, masked_storage_node);
4865    }
4866
4867    #[cfg(feature = "partial-persistence")]
4868    #[test]
4869    fn test_save_blocks_merges_storage_wipe_in_multi_block_batch() {
4870        use reth_trie::{
4871            updates::{StorageTrieUpdatesSorted, TrieUpdatesSorted},
4872            BranchNodeCompact, HashedPostStateSorted,
4873        };
4874
4875        let factory = create_test_provider_factory();
4876        factory.set_storage_settings_cache(StorageSettings::v2());
4877
4878        let mut test_block_builder = TestBlockBuilder::eth().with_state();
4879        let genesis = test_block_builder.get_executed_blocks(0..1).next().unwrap();
4880        let mut blocks: Vec<_> = test_block_builder.get_executed_blocks(1..4).collect();
4881        let address = B256::with_last_byte(1);
4882        let storage_path = Nibbles::from_nibbles([0x1, 0x2]);
4883
4884        let provider_rw = factory.provider_rw().unwrap();
4885        save_genesis(&provider_rw, &genesis).unwrap();
4886        provider_rw
4887            .write_trie_updates_sorted(&TrieUpdatesSorted::new(
4888                vec![],
4889                B256Map::from_iter([(
4890                    address,
4891                    StorageTrieUpdatesSorted {
4892                        is_deleted: false,
4893                        storage_nodes: vec![(storage_path, Some(BranchNodeCompact::default()))],
4894                    },
4895                )]),
4896            ))
4897            .unwrap();
4898        provider_rw.commit().unwrap();
4899
4900        let wipe = TrieUpdatesSorted::new(
4901            vec![],
4902            B256Map::from_iter([(
4903                address,
4904                StorageTrieUpdatesSorted { is_deleted: true, storage_nodes: vec![] },
4905            )]),
4906        );
4907        let wipe_block = &blocks[2];
4908        blocks[2] = ExecutedBlock::new(
4909            Arc::clone(&wipe_block.recovered_block),
4910            Arc::clone(&wipe_block.execution_output),
4911            ComputedTrieData::new(Arc::new(HashedPostStateSorted::default()), Arc::new(wipe)),
4912        );
4913
4914        let provider_rw = factory.provider_rw().unwrap();
4915        let input = SaveBlocksInput::new(blocks, 0, 0, 3, 3);
4916        provider_rw.save_blocks(&input).unwrap();
4917        provider_rw.commit().unwrap();
4918
4919        let provider = factory.provider().unwrap();
4920        let checkpoint = provider.get_stage_checkpoint(StageId::Finish).unwrap().unwrap();
4921        assert_eq!(checkpoint.block_number, 3);
4922        let storage_entries = provider
4923            .tx_ref()
4924            .cursor_dup_read::<tables::StoragesTrie>()
4925            .unwrap()
4926            .walk_dup(Some(address), None)
4927            .unwrap()
4928            .collect::<Result<Vec<_>, _>>()
4929            .unwrap();
4930        assert!(storage_entries.is_empty());
4931    }
4932
4933    #[cfg(feature = "partial-persistence")]
4934    #[test]
4935    fn test_save_blocks_batches_transient_storage_wipe() {
4936        use alloy_primitives::map::B256Set;
4937        use reth_trie::{updates::TrieUpdates, HashBuilder, HashedPostStateSorted};
4938
4939        let factory = create_test_provider_factory();
4940        factory.set_storage_settings_cache(StorageSettings::v2());
4941
4942        let mut test_block_builder = TestBlockBuilder::eth().with_state();
4943        let genesis = test_block_builder.get_executed_blocks(0..1).next().unwrap();
4944        let blocks: Vec<_> = test_block_builder.get_executed_blocks(1..4).collect();
4945
4946        let provider_rw = factory.provider_rw().unwrap();
4947        save_genesis(&provider_rw, &genesis).unwrap();
4948        provider_rw.commit().unwrap();
4949
4950        // A contract created and self-destructed in the same transaction leaves an empty whole-
4951        // storage wipe in the trie updates, even though the account never reaches persisted state.
4952        let ephemeral_account = B256::with_last_byte(0x57);
4953        let hashed_state = HashedPostStateSorted::default();
4954        let mut trie_updates = TrieUpdates::default();
4955        trie_updates.finalize(
4956            HashBuilder::default(),
4957            Default::default(),
4958            B256Set::from_iter([ephemeral_account]),
4959        );
4960        let trie_updates = trie_updates.into_sorted();
4961        assert!(trie_updates.storage_tries_ref()[&ephemeral_account].is_deleted);
4962        let selfdestruct_block = ExecutedBlock::new(
4963            Arc::clone(&blocks[2].recovered_block),
4964            Arc::clone(&blocks[2].execution_output),
4965            ComputedTrieData::new(Arc::new(hashed_state), Arc::new(trie_updates)),
4966        );
4967
4968        let provider_rw = factory.provider_rw().unwrap();
4969        let input = SaveBlocksInput::new(
4970            vec![blocks[0].clone(), blocks[1].clone(), selfdestruct_block],
4971            0,
4972            0,
4973            3,
4974            3,
4975        );
4976        provider_rw.save_blocks(&input).unwrap();
4977        provider_rw.commit().unwrap();
4978
4979        let checkpoint =
4980            factory.provider().unwrap().get_stage_checkpoint(StageId::Finish).unwrap().unwrap();
4981        assert_eq!(checkpoint.block_number, 3);
4982    }
4983
4984    #[cfg(feature = "partial-persistence")]
4985    #[test]
4986    fn test_save_blocks_partial_cycles_do_not_duplicate_static_file_writes() {
4987        let factory = create_test_provider_factory();
4988        let mut test_block_builder = TestBlockBuilder::eth().with_state();
4989
4990        let genesis = test_block_builder.get_executed_blocks(0..1).next().unwrap();
4991        let blocks: Vec<_> = test_block_builder.get_executed_blocks(1..5).collect();
4992
4993        let provider_rw = factory.provider_rw().unwrap();
4994        save_genesis(&provider_rw, &genesis).unwrap();
4995        provider_rw.commit().unwrap();
4996
4997        let provider_rw = factory.provider_rw().unwrap();
4998        let input = SaveBlocksInput::new(blocks[..2].to_vec(), 0, 0, 2, 2);
4999        provider_rw.save_blocks(&input).unwrap();
5000        provider_rw.commit().unwrap();
5001
5002        let provider_rw = factory.provider_rw().unwrap();
5003        let input = SaveBlocksInput::new(blocks[2..].to_vec(), 2, 2, 4, 2);
5004        provider_rw.save_blocks(&input).unwrap();
5005        provider_rw.commit().unwrap();
5006
5007        let provider_rw = factory.provider_rw().unwrap();
5008        let stale_input = SaveBlocksInput::new(vec![blocks[0].clone()], 0, 0, 1, 1);
5009        let err = provider_rw.save_blocks(&stale_input).unwrap_err();
5010        assert!(err.to_string().contains("persistence frontiers do not match Finish checkpoint"));
5011        drop(provider_rw);
5012
5013        let provider = factory.provider().unwrap();
5014        let finish_checkpoint = provider.get_stage_checkpoint(StageId::Finish).unwrap().unwrap();
5015        assert_eq!(finish_checkpoint.block_number, 4);
5016        assert_eq!(
5017            finish_checkpoint.finish_stage_checkpoint().unwrap().partial_state_trie,
5018            Some(2)
5019        );
5020
5021        let static_files = factory.static_file_provider();
5022        assert_eq!(static_files.get_highest_static_file_block(StaticFileSegment::Headers), Some(4));
5023        assert_eq!(
5024            static_files.get_highest_static_file_block(StaticFileSegment::Transactions),
5025            Some(4)
5026        );
5027        assert_eq!(
5028            static_files.get_highest_static_file_block(StaticFileSegment::Receipts),
5029            Some(4)
5030        );
5031    }
5032
5033    #[cfg(feature = "partial-persistence")]
5034    #[test]
5035    fn remove_block_and_execution_above_returns_persistence_frontiers() {
5036        let factory = create_test_provider_factory();
5037        let mut test_block_builder = TestBlockBuilder::eth().with_state();
5038
5039        let genesis = test_block_builder.get_executed_blocks(0..1).next().unwrap();
5040        let blocks: Vec<_> = test_block_builder.get_executed_blocks(1..5).collect();
5041
5042        let provider_rw = factory.provider_rw().unwrap();
5043        save_genesis(&provider_rw, &genesis).unwrap();
5044        provider_rw.commit().unwrap();
5045
5046        let provider_rw = factory.provider_rw().unwrap();
5047        let input = SaveBlocksInput::new(blocks, 0, 0, 4, 2);
5048        provider_rw.save_blocks(&input).unwrap();
5049        provider_rw.commit().unwrap();
5050
5051        let provider_rw = factory.provider_rw().unwrap();
5052        let frontiers = provider_rw.remove_block_and_execution_above(3).unwrap();
5053        assert_eq!(frontiers, PersistenceFrontiers { db_tip: 3, partial_state_trie: 2 });
5054        provider_rw.commit().unwrap();
5055
5056        let provider = factory.provider().unwrap();
5057        let checkpoint = provider.get_stage_checkpoint(StageId::Finish).unwrap().unwrap();
5058        assert_eq!(
5059            checkpoint.finish_stage_checkpoint().and_then(|finish| finish.partial_state_trie()),
5060            Some(2)
5061        );
5062    }
5063
5064    #[test]
5065    fn test_prunable_receipts_logic() {
5066        let insert_blocks =
5067            |provider_rw: &DatabaseProviderRW<_, _>, tip_block: u64, tx_count: u8| {
5068                let mut rng = generators::rng();
5069                for block_num in 0..=tip_block {
5070                    let block = random_block(
5071                        &mut rng,
5072                        block_num,
5073                        BlockParams { tx_count: Some(tx_count), ..Default::default() },
5074                    );
5075                    provider_rw.insert_block(&block.try_recover().unwrap()).unwrap();
5076                }
5077            };
5078
5079        let write_receipts = |provider_rw: DatabaseProviderRW<_, _>, block: u64| {
5080            let outcome = ExecutionOutcome {
5081                first_block: block,
5082                receipts: vec![vec![Receipt {
5083                    tx_type: Default::default(),
5084                    success: true,
5085                    cumulative_gas_used: block, // identifier to assert against
5086                    logs: vec![],
5087                }]],
5088                ..Default::default()
5089            };
5090            provider_rw
5091                .write_state(&outcome, crate::OriginalValuesKnown::No, StateWriteConfig::default())
5092                .unwrap();
5093            provider_rw.commit().unwrap();
5094        };
5095
5096        // Legacy mode (receipts in DB) - should be prunable
5097        {
5098            let factory = create_test_provider_factory();
5099            let storage_settings = StorageSettings::v1();
5100            factory.set_storage_settings_cache(storage_settings);
5101            let factory = factory.with_prune_modes(PruneModes {
5102                receipts: Some(PruneMode::Before(100)),
5103                ..Default::default()
5104            });
5105
5106            let tip_block = 200u64;
5107            let first_block = 1u64;
5108
5109            // create chain
5110            let provider_rw = factory.provider_rw().unwrap();
5111            insert_blocks(&provider_rw, tip_block, 1);
5112            provider_rw.commit().unwrap();
5113
5114            write_receipts(
5115                factory.provider_rw().unwrap().with_minimum_pruning_distance(100),
5116                first_block,
5117            );
5118            write_receipts(
5119                factory.provider_rw().unwrap().with_minimum_pruning_distance(100),
5120                tip_block - 1,
5121            );
5122
5123            let provider = factory.provider().unwrap();
5124
5125            assert!(provider.receipts_by_block(0.into()).unwrap().is_none());
5126            assert!(provider
5127                .receipts_by_block((tip_block - 1).into())
5128                .unwrap()
5129                .is_some_and(|r| r.len() == 1));
5130        }
5131
5132        // Static files mode
5133        {
5134            let factory = create_test_provider_factory();
5135            let storage_settings = StorageSettings::v2();
5136            factory.set_storage_settings_cache(storage_settings);
5137            let factory = factory.with_prune_modes(PruneModes {
5138                receipts: Some(PruneMode::Before(2)),
5139                ..Default::default()
5140            });
5141
5142            let tip_block = 200u64;
5143
5144            // create chain
5145            let provider_rw = factory.provider_rw().unwrap();
5146            insert_blocks(&provider_rw, tip_block, 1);
5147            provider_rw.commit().unwrap();
5148
5149            // Attempt to write receipts for block 0 and 1 (should be skipped)
5150            write_receipts(factory.provider_rw().unwrap().with_minimum_pruning_distance(100), 0);
5151            write_receipts(factory.provider_rw().unwrap().with_minimum_pruning_distance(100), 1);
5152
5153            assert!(factory
5154                .static_file_provider()
5155                .get_highest_static_file_tx(StaticFileSegment::Receipts)
5156                .is_none(),);
5157            assert!(factory
5158                .static_file_provider()
5159                .get_highest_static_file_block(StaticFileSegment::Receipts)
5160                .is_some_and(|b| b == 1),);
5161
5162            // Since we have prune mode Before(2), the next receipt (block 2) should be written to
5163            // static files.
5164            write_receipts(factory.provider_rw().unwrap().with_minimum_pruning_distance(100), 2);
5165            assert!(factory
5166                .static_file_provider()
5167                .get_highest_static_file_tx(StaticFileSegment::Receipts)
5168                .is_some_and(|num| num == 2),);
5169
5170            // After having a receipt already in static files, attempt to skip the next receipt by
5171            // changing the prune mode. It should NOT skip it and should still write the receipt,
5172            // since static files do not support gaps.
5173            let factory = factory.with_prune_modes(PruneModes {
5174                receipts: Some(PruneMode::Before(100)),
5175                ..Default::default()
5176            });
5177            let provider_rw = factory.provider_rw().unwrap().with_minimum_pruning_distance(1);
5178            assert!(PruneMode::Distance(1).should_prune(3, tip_block));
5179            write_receipts(provider_rw, 3);
5180
5181            // Ensure we can only fetch the 2 last receipts.
5182            //
5183            // Test setup only has 1 tx per block and each receipt has its cumulative_gas_used set
5184            // to the block number it belongs to easily identify and assert.
5185            let provider = factory.provider().unwrap();
5186            assert!(EitherWriter::receipts_destination(&provider).is_static_file());
5187            for (num, has_receipt) in [(0, false), (1, false), (2, true), (3, true)] {
5188                let receipts = provider.receipts_by_block(num.into()).unwrap();
5189                if has_receipt {
5190                    assert!(receipts.is_some_and(|r| r.len() == 1));
5191                } else {
5192                    assert!(receipts.is_none());
5193                }
5194
5195                let receipt = provider.receipt(num).unwrap();
5196                if has_receipt {
5197                    assert!(receipt.is_some_and(|r| r.cumulative_gas_used == num));
5198                } else {
5199                    assert!(receipt.is_none());
5200                }
5201            }
5202        }
5203    }
5204
5205    #[test]
5206    fn test_try_into_history_rejects_unexecuted_blocks() {
5207        use reth_storage_api::TryIntoHistoricalStateProvider;
5208
5209        let factory = create_test_provider_factory();
5210
5211        // Insert genesis block to have some data
5212        let data = BlockchainTestData::default();
5213        let provider_rw = factory.provider_rw().unwrap();
5214        provider_rw.insert_block(&data.genesis.try_recover().unwrap()).unwrap();
5215        provider_rw
5216            .write_state(
5217                &ExecutionOutcome { first_block: 0, receipts: vec![vec![]], ..Default::default() },
5218                crate::OriginalValuesKnown::No,
5219                StateWriteConfig::default(),
5220            )
5221            .unwrap();
5222        provider_rw.commit().unwrap();
5223
5224        // Get a fresh provider - Execution checkpoint is 0, no receipts written beyond genesis
5225        let provider = factory.provider().unwrap();
5226
5227        // Requesting historical state for block 0 (executed) should succeed
5228        let result = provider.try_into_history_at_block(0);
5229        assert!(result.is_ok(), "Block 0 should be available");
5230
5231        // Get another provider and request state for block 100 (not executed)
5232        let provider = factory.provider().unwrap();
5233        let result = provider.try_into_history_at_block(100);
5234
5235        // Should fail with BlockNotExecuted error
5236        match result {
5237            Err(ProviderError::BlockNotExecuted { requested: 100, .. }) => {}
5238            Err(e) => panic!("Expected BlockNotExecuted error, got: {e:?}"),
5239            Ok(_) => panic!("Expected error, got Ok"),
5240        }
5241    }
5242
5243    #[test]
5244    fn test_unwind_storage_hashing_with_hashed_state() {
5245        let factory = create_test_provider_factory();
5246        let storage_settings = StorageSettings::v2();
5247        factory.set_storage_settings_cache(storage_settings);
5248
5249        let address = Address::random();
5250        let hashed_address = keccak256(address);
5251
5252        let plain_slot = B256::random();
5253        let hashed_slot = keccak256(plain_slot);
5254
5255        let current_value = U256::from(100);
5256        let old_value = U256::from(42);
5257
5258        let provider_rw = factory.provider_rw().unwrap();
5259        provider_rw
5260            .tx
5261            .cursor_dup_write::<tables::HashedStorages>()
5262            .unwrap()
5263            .upsert(hashed_address, &StorageEntry { key: hashed_slot, value: current_value })
5264            .unwrap();
5265
5266        let changesets = vec![(
5267            BlockNumberAddress((1, address)),
5268            StorageEntry { key: plain_slot, value: old_value },
5269        )];
5270
5271        let result = provider_rw.unwind_storage_hashing(changesets.into_iter()).unwrap();
5272
5273        assert_eq!(result.len(), 1);
5274        assert!(result.contains_key(&hashed_address));
5275        assert!(result[&hashed_address].contains(&hashed_slot));
5276
5277        let mut cursor = provider_rw.tx.cursor_dup_read::<tables::HashedStorages>().unwrap();
5278        let entry = cursor
5279            .seek_by_key_subkey(hashed_address, hashed_slot)
5280            .unwrap()
5281            .expect("entry should exist");
5282        assert_eq!(entry.key, hashed_slot);
5283        assert_eq!(entry.value, old_value);
5284    }
5285
5286    #[test]
5287    fn test_write_and_remove_state_roundtrip_legacy() {
5288        let factory = create_test_provider_factory();
5289        let storage_settings = StorageSettings::v1();
5290        assert!(!storage_settings.use_hashed_state());
5291        factory.set_storage_settings_cache(storage_settings);
5292
5293        let address = Address::with_last_byte(1);
5294        let hashed_address = keccak256(address);
5295        let slot = U256::from(5);
5296        let slot_key = B256::from(slot);
5297        let hashed_slot = keccak256(slot_key);
5298
5299        let mut rng = generators::rng();
5300        let block0 =
5301            random_block(&mut rng, 0, BlockParams { tx_count: Some(0), ..Default::default() });
5302        let block1 =
5303            random_block(&mut rng, 1, BlockParams { tx_count: Some(0), ..Default::default() });
5304
5305        {
5306            let provider_rw = factory.provider_rw().unwrap();
5307            provider_rw.insert_block(&block0.try_recover().unwrap()).unwrap();
5308            provider_rw.insert_block(&block1.try_recover().unwrap()).unwrap();
5309            provider_rw
5310                .tx
5311                .cursor_write::<tables::PlainAccountState>()
5312                .unwrap()
5313                .upsert(address, &Account { nonce: 0, balance: U256::ZERO, bytecode_hash: None })
5314                .unwrap();
5315            provider_rw.commit().unwrap();
5316        }
5317
5318        let provider_rw = factory.provider_rw().unwrap();
5319
5320        let mut state_init: BundleStateInit = AddressMap::default();
5321        let mut storage_map: B256Map<(U256, U256)> = B256Map::default();
5322        storage_map.insert(slot_key, (U256::ZERO, U256::from(10)));
5323        state_init.insert(
5324            address,
5325            (
5326                Some(Account { nonce: 0, balance: U256::ZERO, bytecode_hash: None }),
5327                Some(Account { nonce: 1, balance: U256::ZERO, bytecode_hash: None }),
5328                storage_map,
5329            ),
5330        );
5331
5332        let mut reverts_init: RevertsInit = HashMap::default();
5333        let mut block_reverts: AddressMap<AccountRevertInit> = AddressMap::default();
5334        block_reverts.insert(
5335            address,
5336            (
5337                Some(Some(Account { nonce: 0, balance: U256::ZERO, bytecode_hash: None })),
5338                vec![StorageEntry { key: slot_key, value: U256::ZERO }],
5339            ),
5340        );
5341        reverts_init.insert(1, block_reverts);
5342
5343        let execution_outcome =
5344            ExecutionOutcome::new_init(state_init, reverts_init, [], vec![vec![]], 1, vec![]);
5345
5346        provider_rw
5347            .write_state(
5348                &execution_outcome,
5349                OriginalValuesKnown::Yes,
5350                StateWriteConfig {
5351                    write_receipts: false,
5352                    write_account_changesets: true,
5353                    write_storage_changesets: true,
5354                },
5355            )
5356            .unwrap();
5357
5358        let hashed_state =
5359            execution_outcome.hash_state_slow::<reth_trie::KeccakKeyHasher>().into_sorted();
5360        provider_rw.write_hashed_state(&hashed_state).unwrap();
5361
5362        let account = provider_rw
5363            .tx
5364            .cursor_read::<tables::PlainAccountState>()
5365            .unwrap()
5366            .seek_exact(address)
5367            .unwrap()
5368            .unwrap()
5369            .1;
5370        assert_eq!(account.nonce, 1);
5371
5372        let storage_entry = provider_rw
5373            .tx
5374            .cursor_dup_read::<tables::PlainStorageState>()
5375            .unwrap()
5376            .seek_by_key_subkey(address, slot_key)
5377            .unwrap()
5378            .unwrap();
5379        assert_eq!(storage_entry.key, slot_key);
5380        assert_eq!(storage_entry.value, U256::from(10));
5381
5382        let hashed_entry = provider_rw
5383            .tx
5384            .cursor_dup_read::<tables::HashedStorages>()
5385            .unwrap()
5386            .seek_by_key_subkey(hashed_address, hashed_slot)
5387            .unwrap()
5388            .unwrap();
5389        assert_eq!(hashed_entry.key, hashed_slot);
5390        assert_eq!(hashed_entry.value, U256::from(10));
5391
5392        let account_cs_entries = provider_rw
5393            .tx
5394            .cursor_dup_read::<tables::AccountChangeSets>()
5395            .unwrap()
5396            .walk(Some(1))
5397            .unwrap()
5398            .collect::<Result<Vec<_>, _>>()
5399            .unwrap();
5400        assert!(!account_cs_entries.is_empty());
5401
5402        let storage_cs_entries = provider_rw
5403            .tx
5404            .cursor_read::<tables::StorageChangeSets>()
5405            .unwrap()
5406            .walk(Some(BlockNumberAddress((1, address))))
5407            .unwrap()
5408            .collect::<Result<Vec<_>, _>>()
5409            .unwrap();
5410        assert!(!storage_cs_entries.is_empty());
5411        assert_eq!(storage_cs_entries[0].1.key, slot_key);
5412
5413        provider_rw.remove_state_above(0).unwrap();
5414
5415        let restored_account = provider_rw
5416            .tx
5417            .cursor_read::<tables::PlainAccountState>()
5418            .unwrap()
5419            .seek_exact(address)
5420            .unwrap()
5421            .unwrap()
5422            .1;
5423        assert_eq!(restored_account.nonce, 0);
5424
5425        let storage_gone = provider_rw
5426            .tx
5427            .cursor_dup_read::<tables::PlainStorageState>()
5428            .unwrap()
5429            .seek_by_key_subkey(address, slot_key)
5430            .unwrap();
5431        assert!(storage_gone.is_none() || storage_gone.unwrap().key != slot_key);
5432
5433        let account_cs_after = provider_rw
5434            .tx
5435            .cursor_dup_read::<tables::AccountChangeSets>()
5436            .unwrap()
5437            .walk(Some(1))
5438            .unwrap()
5439            .collect::<Result<Vec<_>, _>>()
5440            .unwrap();
5441        assert!(account_cs_after.is_empty());
5442
5443        let storage_cs_after = provider_rw
5444            .tx
5445            .cursor_read::<tables::StorageChangeSets>()
5446            .unwrap()
5447            .walk(Some(BlockNumberAddress((1, address))))
5448            .unwrap()
5449            .collect::<Result<Vec<_>, _>>()
5450            .unwrap();
5451        assert!(storage_cs_after.is_empty());
5452    }
5453
5454    #[test]
5455    fn test_unwind_storage_hashing_legacy() {
5456        let factory = create_test_provider_factory();
5457        let storage_settings = StorageSettings::v1();
5458        assert!(!storage_settings.use_hashed_state());
5459        factory.set_storage_settings_cache(storage_settings);
5460
5461        let address = Address::random();
5462        let hashed_address = keccak256(address);
5463
5464        let plain_slot = B256::random();
5465        let hashed_slot = keccak256(plain_slot);
5466
5467        let current_value = U256::from(100);
5468        let old_value = U256::from(42);
5469
5470        let provider_rw = factory.provider_rw().unwrap();
5471        provider_rw
5472            .tx
5473            .cursor_dup_write::<tables::HashedStorages>()
5474            .unwrap()
5475            .upsert(hashed_address, &StorageEntry { key: hashed_slot, value: current_value })
5476            .unwrap();
5477
5478        let changesets = vec![(
5479            BlockNumberAddress((1, address)),
5480            StorageEntry { key: plain_slot, value: old_value },
5481        )];
5482
5483        let result = provider_rw.unwind_storage_hashing(changesets.into_iter()).unwrap();
5484
5485        assert_eq!(result.len(), 1);
5486        assert!(result.contains_key(&hashed_address));
5487        assert!(result[&hashed_address].contains(&hashed_slot));
5488
5489        let mut cursor = provider_rw.tx.cursor_dup_read::<tables::HashedStorages>().unwrap();
5490        let entry = cursor
5491            .seek_by_key_subkey(hashed_address, hashed_slot)
5492            .unwrap()
5493            .expect("entry should exist");
5494        assert_eq!(entry.key, hashed_slot);
5495        assert_eq!(entry.value, old_value);
5496    }
5497
5498    #[test]
5499    fn test_write_state_and_historical_read_hashed() {
5500        use reth_storage_api::StateProvider;
5501        use reth_trie::{HashedPostState, KeccakKeyHasher};
5502        use revm::{database::BundleState, state::AccountInfo};
5503
5504        let factory = create_test_provider_factory();
5505        factory.set_storage_settings_cache(StorageSettings::v2());
5506
5507        let address = Address::with_last_byte(1);
5508        let slot = U256::from(5);
5509        let slot_key = B256::from(slot);
5510        let hashed_address = keccak256(address);
5511        let hashed_slot = keccak256(slot_key);
5512
5513        {
5514            let sf = factory.static_file_provider();
5515            let mut hw = sf.latest_writer(StaticFileSegment::Headers).unwrap();
5516            let h0 = alloy_consensus::Header { number: 0, ..Default::default() };
5517            hw.append_header(&h0, &B256::ZERO).unwrap();
5518            let h1 = alloy_consensus::Header { number: 1, ..Default::default() };
5519            hw.append_header(&h1, &B256::ZERO).unwrap();
5520            hw.commit().unwrap();
5521
5522            let mut aw = sf.latest_writer(StaticFileSegment::AccountChangeSets).unwrap();
5523            aw.append_account_changeset(vec![], 0).unwrap();
5524            aw.commit().unwrap();
5525
5526            let mut sw = sf.latest_writer(StaticFileSegment::StorageChangeSets).unwrap();
5527            sw.append_storage_changeset(vec![], 0).unwrap();
5528            sw.commit().unwrap();
5529        }
5530
5531        let provider_rw = factory.provider_rw().unwrap();
5532
5533        let bundle = BundleState::builder(1..=1)
5534            .state_present_account_info(
5535                address,
5536                AccountInfo { nonce: 1, balance: U256::from(10), ..Default::default() },
5537            )
5538            .state_storage(address, HashMap::from_iter([(slot, (U256::ZERO, U256::from(10)))]))
5539            .revert_account_info(1, address, Some(None))
5540            .revert_storage(1, address, vec![(slot, U256::ZERO)])
5541            .build();
5542
5543        let execution_outcome = ExecutionOutcome::new(bundle.clone(), vec![vec![]], 1, Vec::new());
5544
5545        provider_rw
5546            .tx
5547            .put::<tables::BlockBodyIndices>(
5548                1,
5549                StoredBlockBodyIndices { first_tx_num: 0, tx_count: 0 },
5550            )
5551            .unwrap();
5552
5553        provider_rw
5554            .write_state(
5555                &execution_outcome,
5556                OriginalValuesKnown::Yes,
5557                StateWriteConfig {
5558                    write_receipts: false,
5559                    write_account_changesets: true,
5560                    write_storage_changesets: true,
5561                },
5562            )
5563            .unwrap();
5564
5565        let hashed_state =
5566            HashedPostState::from_bundle_state::<KeccakKeyHasher>(bundle.state()).into_sorted();
5567        provider_rw.write_hashed_state(&hashed_state).unwrap();
5568
5569        let plain_storage_entries = provider_rw
5570            .tx
5571            .cursor_dup_read::<tables::PlainStorageState>()
5572            .unwrap()
5573            .walk(None)
5574            .unwrap()
5575            .collect::<Result<Vec<_>, _>>()
5576            .unwrap();
5577        assert!(plain_storage_entries.is_empty());
5578
5579        let hashed_entry = provider_rw
5580            .tx
5581            .cursor_dup_read::<tables::HashedStorages>()
5582            .unwrap()
5583            .seek_by_key_subkey(hashed_address, hashed_slot)
5584            .unwrap()
5585            .unwrap();
5586        assert_eq!(hashed_entry.key, hashed_slot);
5587        assert_eq!(hashed_entry.value, U256::from(10));
5588
5589        provider_rw.static_file_provider().commit().unwrap();
5590
5591        let sf = factory.static_file_provider();
5592        let storage_cs = sf.storage_changeset(1).unwrap();
5593        assert!(!storage_cs.is_empty());
5594        assert_eq!(storage_cs[0].1.key, slot_key);
5595
5596        let account_cs = sf.account_block_changeset(1).unwrap();
5597        assert!(!account_cs.is_empty());
5598        assert_eq!(account_cs[0].address, address);
5599
5600        let historical_value =
5601            HistoricalStateProviderRef::new(&*provider_rw, 0, OverlayManager::default())
5602                .storage(address, slot_key)
5603                .unwrap();
5604        assert_eq!(historical_value, None);
5605    }
5606
5607    #[derive(Debug, Clone, Copy, PartialEq, Eq)]
5608    enum StorageMode {
5609        V1,
5610        V2,
5611    }
5612
5613    fn run_save_blocks_and_verify(mode: StorageMode) {
5614        use alloy_primitives::map::{FbBuildHasher, HashMap};
5615
5616        let factory = create_test_provider_factory();
5617
5618        match mode {
5619            StorageMode::V1 => factory.set_storage_settings_cache(StorageSettings::v1()),
5620            StorageMode::V2 => factory.set_storage_settings_cache(StorageSettings::v2()),
5621        }
5622
5623        let num_blocks = 3u64;
5624        let accounts_per_block = 5usize;
5625        let slots_per_account = 3usize;
5626
5627        let genesis = SealedBlock::<reth_ethereum_primitives::Block>::from_sealed_parts(
5628            SealedHeader::new(
5629                Header { number: 0, difficulty: U256::from(1), ..Default::default() },
5630                B256::ZERO,
5631            ),
5632            Default::default(),
5633        );
5634
5635        let genesis_executed: ExecutedBlock = ExecutedBlock::new(
5636            Arc::new(genesis.try_recover().unwrap()),
5637            Arc::new(BlockExecutionOutput {
5638                result: BlockExecutionResult {
5639                    receipts: vec![],
5640                    requests: Default::default(),
5641                    gas_used: 0,
5642                    blob_gas_used: 0,
5643                },
5644                state: Default::default(),
5645            }),
5646            ComputedTrieData::default(),
5647        );
5648        let provider_rw = factory.provider_rw().unwrap();
5649        save_genesis(&provider_rw, &genesis_executed).unwrap();
5650        provider_rw.commit().unwrap();
5651
5652        let mut blocks: Vec<ExecutedBlock> = Vec::new();
5653        let mut parent_hash = B256::ZERO;
5654
5655        for block_num in 1..=num_blocks {
5656            let mut builder = BundleState::builder(block_num..=block_num);
5657
5658            for acct_idx in 0..accounts_per_block {
5659                let address = Address::with_last_byte((block_num * 10 + acct_idx as u64) as u8);
5660                let info = AccountInfo {
5661                    nonce: block_num,
5662                    balance: U256::from(block_num * 100 + acct_idx as u64),
5663                    ..Default::default()
5664                };
5665
5666                let storage: HashMap<U256, (U256, U256), FbBuildHasher<32>> = (1..=
5667                    slots_per_account as u64)
5668                    .map(|s| {
5669                        (
5670                            U256::from(s + acct_idx as u64 * 100),
5671                            (U256::ZERO, U256::from(block_num * 1000 + s)),
5672                        )
5673                    })
5674                    .collect();
5675
5676                let revert_storage: Vec<(U256, U256)> = (1..=slots_per_account as u64)
5677                    .map(|s| (U256::from(s + acct_idx as u64 * 100), U256::ZERO))
5678                    .collect();
5679
5680                builder = builder
5681                    .state_present_account_info(address, info)
5682                    .revert_account_info(block_num, address, Some(None))
5683                    .state_storage(address, storage)
5684                    .revert_storage(block_num, address, revert_storage);
5685            }
5686
5687            let bundle = builder.build();
5688
5689            let hashed_state =
5690                HashedPostState::from_bundle_state::<KeccakKeyHasher>(bundle.state()).into_sorted();
5691
5692            let header = Header {
5693                number: block_num,
5694                parent_hash,
5695                difficulty: U256::from(1),
5696                ..Default::default()
5697            };
5698            let block = SealedBlock::<reth_ethereum_primitives::Block>::seal_parts(
5699                header,
5700                Default::default(),
5701            );
5702            parent_hash = block.hash();
5703
5704            let executed = ExecutedBlock::new(
5705                Arc::new(block.try_recover().unwrap()),
5706                Arc::new(BlockExecutionOutput {
5707                    result: BlockExecutionResult {
5708                        receipts: vec![],
5709                        requests: Default::default(),
5710                        gas_used: 0,
5711                        blob_gas_used: 0,
5712                    },
5713                    state: bundle,
5714                }),
5715                ComputedTrieData {
5716                    sorted: SortedTrieData::new(Arc::new(hashed_state), Default::default()),
5717                },
5718            );
5719            blocks.push(executed);
5720        }
5721
5722        let provider_rw = factory.provider_rw().unwrap();
5723        let input = SaveBlocksInput::new(blocks, 0, 0, num_blocks, num_blocks);
5724        provider_rw.save_blocks(&input).unwrap();
5725        provider_rw.commit().unwrap();
5726
5727        let provider = factory.provider().unwrap();
5728
5729        for block_num in 1..=num_blocks {
5730            for acct_idx in 0..accounts_per_block {
5731                let address = Address::with_last_byte((block_num * 10 + acct_idx as u64) as u8);
5732                let hashed_address = keccak256(address);
5733
5734                let ha_entry = provider
5735                    .tx_ref()
5736                    .cursor_read::<tables::HashedAccounts>()
5737                    .unwrap()
5738                    .seek_exact(hashed_address)
5739                    .unwrap();
5740                assert!(
5741                    ha_entry.is_some(),
5742                    "HashedAccounts missing for block {block_num} acct {acct_idx}"
5743                );
5744
5745                for s in 1..=slots_per_account as u64 {
5746                    let slot = U256::from(s + acct_idx as u64 * 100);
5747                    let slot_key = B256::from(slot);
5748                    let hashed_slot = keccak256(slot_key);
5749
5750                    let hs_entry = provider
5751                        .tx_ref()
5752                        .cursor_dup_read::<tables::HashedStorages>()
5753                        .unwrap()
5754                        .seek_by_key_subkey(hashed_address, hashed_slot)
5755                        .unwrap();
5756                    assert!(
5757                        hs_entry.is_some(),
5758                        "HashedStorages missing for block {block_num} acct {acct_idx} slot {s}"
5759                    );
5760                    let entry = hs_entry.unwrap();
5761                    assert_eq!(entry.key, hashed_slot);
5762                    assert_eq!(entry.value, U256::from(block_num * 1000 + s));
5763                }
5764            }
5765        }
5766
5767        for block_num in 1..=num_blocks {
5768            let header = provider.header_by_number(block_num).unwrap();
5769            assert!(header.is_some(), "Header missing for block {block_num}");
5770
5771            let indices = provider.block_body_indices(block_num).unwrap();
5772            assert!(indices.is_some(), "BlockBodyIndices missing for block {block_num}");
5773        }
5774
5775        let plain_accounts = provider.tx_ref().entries::<tables::PlainAccountState>().unwrap();
5776        let plain_storage = provider.tx_ref().entries::<tables::PlainStorageState>().unwrap();
5777
5778        if mode == StorageMode::V2 {
5779            assert_eq!(plain_accounts, 0, "v2: PlainAccountState should be empty");
5780            assert_eq!(plain_storage, 0, "v2: PlainStorageState should be empty");
5781
5782            let mdbx_account_cs = provider.tx_ref().entries::<tables::AccountChangeSets>().unwrap();
5783            assert_eq!(mdbx_account_cs, 0, "v2: AccountChangeSets in MDBX should be empty");
5784
5785            let mdbx_storage_cs = provider.tx_ref().entries::<tables::StorageChangeSets>().unwrap();
5786            assert_eq!(mdbx_storage_cs, 0, "v2: StorageChangeSets in MDBX should be empty");
5787
5788            provider.static_file_provider().commit().unwrap();
5789            let sf = factory.static_file_provider();
5790
5791            for block_num in 1..=num_blocks {
5792                let account_cs = sf.account_block_changeset(block_num).unwrap();
5793                assert!(
5794                    !account_cs.is_empty(),
5795                    "v2: static file AccountChangeSets should exist for block {block_num}"
5796                );
5797
5798                let storage_cs = sf.storage_changeset(block_num).unwrap();
5799                assert!(
5800                    !storage_cs.is_empty(),
5801                    "v2: static file StorageChangeSets should exist for block {block_num}"
5802                );
5803
5804                for (_, entry) in &storage_cs {
5805                    assert!(
5806                        entry.key != keccak256(entry.key),
5807                        "v2: static file storage changeset should have plain slot keys"
5808                    );
5809                }
5810            }
5811
5812            let rocksdb = factory.rocksdb_provider();
5813            for block_num in 1..=num_blocks {
5814                for acct_idx in 0..accounts_per_block {
5815                    let address = Address::with_last_byte((block_num * 10 + acct_idx as u64) as u8);
5816                    let shards = rocksdb.account_history_shards(address).unwrap();
5817                    assert!(
5818                        !shards.is_empty(),
5819                        "v2: RocksDB AccountsHistory missing for block {block_num} acct {acct_idx}"
5820                    );
5821
5822                    for s in 1..=slots_per_account as u64 {
5823                        let slot = U256::from(s + acct_idx as u64 * 100);
5824                        let slot_key = B256::from(slot);
5825                        let shards = rocksdb.storage_history_shards(address, slot_key).unwrap();
5826                        assert!(
5827                            !shards.is_empty(),
5828                            "v2: RocksDB StoragesHistory missing for block {block_num} acct {acct_idx} slot {s}"
5829                        );
5830                    }
5831                }
5832            }
5833        } else {
5834            assert!(plain_accounts > 0, "v1: PlainAccountState should not be empty");
5835            assert!(plain_storage > 0, "v1: PlainStorageState should not be empty");
5836
5837            let mdbx_account_cs = provider.tx_ref().entries::<tables::AccountChangeSets>().unwrap();
5838            assert!(mdbx_account_cs > 0, "v1: AccountChangeSets in MDBX should not be empty");
5839
5840            let mdbx_storage_cs = provider.tx_ref().entries::<tables::StorageChangeSets>().unwrap();
5841            assert!(mdbx_storage_cs > 0, "v1: StorageChangeSets in MDBX should not be empty");
5842
5843            for block_num in 1..=num_blocks {
5844                let storage_entries: Vec<_> = provider
5845                    .tx_ref()
5846                    .cursor_dup_read::<tables::StorageChangeSets>()
5847                    .unwrap()
5848                    .walk_range(BlockNumberAddress::range(block_num..=block_num))
5849                    .unwrap()
5850                    .collect::<Result<Vec<_>, _>>()
5851                    .unwrap();
5852                assert!(
5853                    !storage_entries.is_empty(),
5854                    "v1: MDBX StorageChangeSets should have entries for block {block_num}"
5855                );
5856
5857                for (_, entry) in &storage_entries {
5858                    let slot_key = B256::from(entry.key);
5859                    assert!(
5860                        slot_key != keccak256(slot_key),
5861                        "v1: storage changeset keys should be plain (not hashed)"
5862                    );
5863                }
5864            }
5865
5866            let mdbx_account_history =
5867                provider.tx_ref().entries::<tables::AccountsHistory>().unwrap();
5868            assert!(mdbx_account_history > 0, "v1: AccountsHistory in MDBX should not be empty");
5869
5870            let mdbx_storage_history =
5871                provider.tx_ref().entries::<tables::StoragesHistory>().unwrap();
5872            assert!(mdbx_storage_history > 0, "v1: StoragesHistory in MDBX should not be empty");
5873        }
5874    }
5875
5876    #[test]
5877    fn test_save_blocks_v1_table_assertions() {
5878        run_save_blocks_and_verify(StorageMode::V1);
5879    }
5880
5881    #[test]
5882    fn test_save_blocks_v2_table_assertions() {
5883        run_save_blocks_and_verify(StorageMode::V2);
5884    }
5885
5886    #[test]
5887    fn test_write_and_remove_state_roundtrip_v2() {
5888        let factory = create_test_provider_factory();
5889        let storage_settings = StorageSettings::v2();
5890        assert!(storage_settings.use_hashed_state());
5891        factory.set_storage_settings_cache(storage_settings);
5892
5893        let address = Address::with_last_byte(1);
5894        let hashed_address = keccak256(address);
5895        let slot = U256::from(5);
5896        let slot_key = B256::from(slot);
5897        let hashed_slot = keccak256(slot_key);
5898
5899        {
5900            let sf = factory.static_file_provider();
5901            let mut hw = sf.latest_writer(StaticFileSegment::Headers).unwrap();
5902            let h0 = alloy_consensus::Header { number: 0, ..Default::default() };
5903            hw.append_header(&h0, &B256::ZERO).unwrap();
5904            let h1 = alloy_consensus::Header { number: 1, ..Default::default() };
5905            hw.append_header(&h1, &B256::ZERO).unwrap();
5906            hw.commit().unwrap();
5907
5908            let mut aw = sf.latest_writer(StaticFileSegment::AccountChangeSets).unwrap();
5909            aw.append_account_changeset(vec![], 0).unwrap();
5910            aw.commit().unwrap();
5911
5912            let mut sw = sf.latest_writer(StaticFileSegment::StorageChangeSets).unwrap();
5913            sw.append_storage_changeset(vec![], 0).unwrap();
5914            sw.commit().unwrap();
5915        }
5916
5917        {
5918            let provider_rw = factory.provider_rw().unwrap();
5919            provider_rw
5920                .tx
5921                .put::<tables::BlockBodyIndices>(
5922                    0,
5923                    StoredBlockBodyIndices { first_tx_num: 0, tx_count: 0 },
5924                )
5925                .unwrap();
5926            provider_rw
5927                .tx
5928                .put::<tables::BlockBodyIndices>(
5929                    1,
5930                    StoredBlockBodyIndices { first_tx_num: 0, tx_count: 0 },
5931                )
5932                .unwrap();
5933            provider_rw
5934                .tx
5935                .cursor_write::<tables::HashedAccounts>()
5936                .unwrap()
5937                .upsert(
5938                    hashed_address,
5939                    &Account { nonce: 0, balance: U256::ZERO, bytecode_hash: None },
5940                )
5941                .unwrap();
5942            provider_rw.commit().unwrap();
5943        }
5944
5945        let provider_rw = factory.provider_rw().unwrap();
5946
5947        let bundle = BundleState::builder(1..=1)
5948            .state_present_account_info(
5949                address,
5950                AccountInfo { nonce: 1, balance: U256::from(10), ..Default::default() },
5951            )
5952            .state_storage(address, HashMap::from_iter([(slot, (U256::ZERO, U256::from(10)))]))
5953            .revert_account_info(1, address, Some(None))
5954            .revert_storage(1, address, vec![(slot, U256::ZERO)])
5955            .build();
5956
5957        let execution_outcome = ExecutionOutcome::new(bundle.clone(), vec![vec![]], 1, Vec::new());
5958
5959        provider_rw
5960            .write_state(
5961                &execution_outcome,
5962                OriginalValuesKnown::Yes,
5963                StateWriteConfig {
5964                    write_receipts: false,
5965                    write_account_changesets: true,
5966                    write_storage_changesets: true,
5967                },
5968            )
5969            .unwrap();
5970
5971        let hashed_state =
5972            HashedPostState::from_bundle_state::<KeccakKeyHasher>(bundle.state()).into_sorted();
5973        provider_rw.write_hashed_state(&hashed_state).unwrap();
5974
5975        let hashed_account = provider_rw
5976            .tx
5977            .cursor_read::<tables::HashedAccounts>()
5978            .unwrap()
5979            .seek_exact(hashed_address)
5980            .unwrap()
5981            .unwrap()
5982            .1;
5983        assert_eq!(hashed_account.nonce, 1);
5984
5985        let hashed_entry = provider_rw
5986            .tx
5987            .cursor_dup_read::<tables::HashedStorages>()
5988            .unwrap()
5989            .seek_by_key_subkey(hashed_address, hashed_slot)
5990            .unwrap()
5991            .unwrap();
5992        assert_eq!(hashed_entry.key, hashed_slot);
5993        assert_eq!(hashed_entry.value, U256::from(10));
5994
5995        let plain_accounts = provider_rw.tx.entries::<tables::PlainAccountState>().unwrap();
5996        assert_eq!(plain_accounts, 0, "v2: PlainAccountState should be empty");
5997
5998        let plain_storage = provider_rw.tx.entries::<tables::PlainStorageState>().unwrap();
5999        assert_eq!(plain_storage, 0, "v2: PlainStorageState should be empty");
6000
6001        provider_rw.static_file_provider().commit().unwrap();
6002
6003        let sf = factory.static_file_provider();
6004        let storage_cs = sf.storage_changeset(1).unwrap();
6005        assert!(!storage_cs.is_empty(), "v2: storage changesets should be in static files");
6006        assert_eq!(storage_cs[0].1.key, slot_key, "v2: changeset key should be plain");
6007
6008        provider_rw.remove_state_above(0).unwrap();
6009
6010        let restored_account = provider_rw
6011            .tx
6012            .cursor_read::<tables::HashedAccounts>()
6013            .unwrap()
6014            .seek_exact(hashed_address)
6015            .unwrap();
6016        assert!(
6017            restored_account.is_none(),
6018            "v2: account should be removed (didn't exist before block 1)"
6019        );
6020
6021        let storage_gone = provider_rw
6022            .tx
6023            .cursor_dup_read::<tables::HashedStorages>()
6024            .unwrap()
6025            .seek_by_key_subkey(hashed_address, hashed_slot)
6026            .unwrap();
6027        assert!(
6028            storage_gone.is_none() || storage_gone.unwrap().key != hashed_slot,
6029            "v2: storage should be reverted (removed or different key)"
6030        );
6031
6032        let mdbx_storage_cs = provider_rw.tx.entries::<tables::StorageChangeSets>().unwrap();
6033        assert_eq!(mdbx_storage_cs, 0, "v2: MDBX StorageChangeSets should remain empty");
6034
6035        let mdbx_account_cs = provider_rw.tx.entries::<tables::AccountChangeSets>().unwrap();
6036        assert_eq!(mdbx_account_cs, 0, "v2: MDBX AccountChangeSets should remain empty");
6037    }
6038
6039    #[test]
6040    fn test_unwind_storage_history_indices_v2() {
6041        let factory = create_test_provider_factory();
6042        factory.set_storage_settings_cache(StorageSettings::v2());
6043
6044        let address = Address::with_last_byte(1);
6045        let slot_key = B256::from(U256::from(42));
6046
6047        {
6048            let rocksdb = factory.rocksdb_provider();
6049            let mut batch = rocksdb.batch();
6050            batch.append_storage_history_shard(address, slot_key, vec![3u64, 7, 10]).unwrap();
6051            batch.commit().unwrap();
6052
6053            let shards = rocksdb.storage_history_shards(address, slot_key).unwrap();
6054            assert!(!shards.is_empty(), "history should be written to rocksdb");
6055        }
6056
6057        let provider_rw = factory.provider_rw().unwrap();
6058
6059        let changesets = vec![
6060            (
6061                BlockNumberAddress((7, address)),
6062                StorageEntry { key: slot_key, value: U256::from(5) },
6063            ),
6064            (
6065                BlockNumberAddress((10, address)),
6066                StorageEntry { key: slot_key, value: U256::from(8) },
6067            ),
6068        ];
6069
6070        let count = provider_rw.unwind_storage_history_indices(changesets.into_iter()).unwrap();
6071        assert_eq!(count, 2);
6072
6073        provider_rw.commit().unwrap();
6074
6075        let rocksdb = factory.rocksdb_provider();
6076        let shards = rocksdb.storage_history_shards(address, slot_key).unwrap();
6077
6078        assert!(
6079            !shards.is_empty(),
6080            "history shards should still exist with block 3 after partial unwind"
6081        );
6082
6083        let all_blocks: Vec<u64> = shards.iter().flat_map(|(_, list)| list.iter()).collect();
6084        assert!(all_blocks.contains(&3), "block 3 should remain");
6085        assert!(!all_blocks.contains(&7), "block 7 should be unwound");
6086        assert!(!all_blocks.contains(&10), "block 10 should be unwound");
6087    }
6088}