Skip to main content

reth_provider/providers/database/
provider.rs

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