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