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