Skip to main content

reth_provider/
either_writer.rs

1//! Generic reader and writer abstractions for interacting with either database tables or static
2//! files.
3
4use std::{
5    collections::BTreeSet,
6    marker::PhantomData,
7    ops::{Range, RangeInclusive},
8};
9
10use crate::{
11    providers::{
12        history_info, rocksdb::RocksDBBatch, HistoryInfo, StaticFileProvider,
13        StaticFileProviderRWRefMut,
14    },
15    StaticFileProviderFactory,
16};
17use alloy_primitives::{map::HashMap, Address, BlockNumber, TxHash, TxNumber, B256};
18use rayon::slice::ParallelSliceMut;
19use reth_db::{
20    cursor::{DbCursorRO, DbDupCursorRW},
21    models::{AccountBeforeTx, StorageBeforeTx},
22    static_file::TransactionSenderMask,
23    table::Value,
24    transaction::{CursorMutTy, CursorTy, DbTx, DbTxMut, DupCursorMutTy, DupCursorTy},
25};
26use reth_db_api::{
27    cursor::DbCursorRW,
28    models::{storage_sharded_key::StorageShardedKey, BlockNumberAddress, ShardedKey},
29    tables,
30    tables::BlockNumberList,
31};
32use reth_errors::ProviderError;
33use reth_node_types::NodePrimitives;
34use reth_primitives_traits::{ReceiptTy, StorageEntry};
35use reth_static_file_types::StaticFileSegment;
36use reth_storage_api::{
37    ChangeSetReader, DBProvider, DbTxProvider, NodePrimitivesProvider, StorageSettingsCache,
38};
39use reth_storage_errors::provider::ProviderResult;
40use strum::{Display, EnumIs};
41
42/// Type alias for [`EitherReader`] constructors.
43type EitherReaderTy<'a, 'db, P, T> = EitherReader<
44    'a,
45    'db,
46    CursorTy<<P as DbTxProvider>::Tx, T>,
47    <P as NodePrimitivesProvider>::Primitives,
48>;
49
50/// Type alias for [`EitherReader`] constructors.
51type DupEitherReaderTy<'a, 'db, P, T> = EitherReader<
52    'a,
53    'db,
54    DupCursorTy<<P as DbTxProvider>::Tx, T>,
55    <P as NodePrimitivesProvider>::Primitives,
56>;
57
58/// Type alias for dup [`EitherWriter`] constructors.
59type DupEitherWriterTy<'a, P, T> = EitherWriter<
60    'a,
61    DupCursorMutTy<<P as DbTxProvider>::Tx, T>,
62    <P as NodePrimitivesProvider>::Primitives,
63>;
64
65/// Type alias for [`EitherWriter`] constructors.
66type EitherWriterTy<'a, P, T> = EitherWriter<
67    'a,
68    CursorMutTy<<P as DbTxProvider>::Tx, T>,
69    <P as NodePrimitivesProvider>::Primitives,
70>;
71
72/// Helper type for `RocksDB` batch argument in writer constructors.
73pub type RocksBatchArg<'a> = crate::providers::rocksdb::RocksDBBatch<'a>;
74
75/// The raw `RocksDB` batch type returned by [`EitherWriter::into_raw_rocksdb_batch`].
76pub type RawRocksDBBatch = rocksdb::WriteBatchWithTransaction<true>;
77
78/// Helper type for `RocksDB` snapshot argument in reader constructors.
79///
80/// The `Option` allows callers to skip `RocksDB` access when it isn't needed
81/// (e.g., on legacy MDBX-only nodes).
82pub type RocksDBRefArg<'a, 'db> = Option<&'a crate::providers::rocksdb::RocksReadSnapshot<'db>>;
83
84/// Represents a destination for writing data, either to database, static files, or `RocksDB`.
85#[derive(Debug, Display)]
86pub enum EitherWriter<'a, CURSOR, N> {
87    /// Write to database table via cursor
88    Database(CURSOR),
89    /// Write to static file
90    StaticFile(StaticFileProviderRWRefMut<'a, N>),
91    /// Write to `RocksDB` using a write-only batch (historical tables).
92    RocksDB(RocksDBBatch<'a>),
93}
94
95impl<'a> EitherWriter<'a, (), ()> {
96    /// Creates a new [`EitherWriter`] for receipts based on storage settings and prune modes.
97    pub fn new_receipts<P>(
98        provider: &'a P,
99        block_number: BlockNumber,
100    ) -> ProviderResult<EitherWriterTy<'a, P, tables::Receipts<ReceiptTy<P::Primitives>>>>
101    where
102        P: DBProvider + NodePrimitivesProvider + StorageSettingsCache + StaticFileProviderFactory,
103        P::Tx: DbTxMut,
104        ReceiptTy<P::Primitives>: Value,
105    {
106        if Self::receipts_destination(provider).is_static_file() {
107            Ok(EitherWriter::StaticFile(
108                provider.get_static_file_writer(block_number, StaticFileSegment::Receipts)?,
109            ))
110        } else {
111            Ok(EitherWriter::Database(
112                provider.tx_ref().cursor_write::<tables::Receipts<ReceiptTy<P::Primitives>>>()?,
113            ))
114        }
115    }
116
117    /// Creates a new [`EitherWriter`] for senders based on storage settings.
118    pub fn new_senders<P>(
119        provider: &'a P,
120        block_number: BlockNumber,
121    ) -> ProviderResult<EitherWriterTy<'a, P, tables::TransactionSenders>>
122    where
123        P: DBProvider + NodePrimitivesProvider + StorageSettingsCache + StaticFileProviderFactory,
124        P::Tx: DbTxMut,
125    {
126        if EitherWriterDestination::senders(provider).is_static_file() {
127            Ok(EitherWriter::StaticFile(
128                provider
129                    .get_static_file_writer(block_number, StaticFileSegment::TransactionSenders)?,
130            ))
131        } else {
132            Ok(EitherWriter::Database(
133                provider.tx_ref().cursor_write::<tables::TransactionSenders>()?,
134            ))
135        }
136    }
137
138    /// Creates a new [`EitherWriter`] for account changesets based on storage settings and prune
139    /// modes.
140    pub fn new_account_changesets<P>(
141        provider: &'a P,
142        block_number: BlockNumber,
143    ) -> ProviderResult<DupEitherWriterTy<'a, P, tables::AccountChangeSets>>
144    where
145        P: DBProvider + NodePrimitivesProvider + StorageSettingsCache + StaticFileProviderFactory,
146        P::Tx: DbTxMut,
147    {
148        if provider.cached_storage_settings().storage_v2 {
149            Ok(EitherWriter::StaticFile(
150                provider
151                    .get_static_file_writer(block_number, StaticFileSegment::AccountChangeSets)?,
152            ))
153        } else {
154            Ok(EitherWriter::Database(
155                provider.tx_ref().cursor_dup_write::<tables::AccountChangeSets>()?,
156            ))
157        }
158    }
159
160    /// Creates a new [`EitherWriter`] for storage changesets based on storage settings.
161    pub fn new_storage_changesets<P>(
162        provider: &'a P,
163        block_number: BlockNumber,
164    ) -> ProviderResult<DupEitherWriterTy<'a, P, tables::StorageChangeSets>>
165    where
166        P: DBProvider + NodePrimitivesProvider + StorageSettingsCache + StaticFileProviderFactory,
167        P::Tx: DbTxMut,
168    {
169        if provider.cached_storage_settings().storage_v2 {
170            Ok(EitherWriter::StaticFile(
171                provider
172                    .get_static_file_writer(block_number, StaticFileSegment::StorageChangeSets)?,
173            ))
174        } else {
175            Ok(EitherWriter::Database(
176                provider.tx_ref().cursor_dup_write::<tables::StorageChangeSets>()?,
177            ))
178        }
179    }
180
181    /// Returns the destination for writing receipts.
182    ///
183    /// The rules are as follows:
184    /// - If the node should not always write receipts to static files, and any receipt pruning is
185    ///   enabled, write to the database.
186    /// - If the node should always write receipts to static files, but receipt log filter pruning
187    ///   is enabled, write to the database.
188    /// - Otherwise, write to static files.
189    pub fn receipts_destination<P: DBProvider + StorageSettingsCache>(
190        provider: &P,
191    ) -> EitherWriterDestination {
192        let receipts_in_static_files = provider.cached_storage_settings().storage_v2;
193        let prune_modes = provider.prune_modes_ref();
194
195        if !receipts_in_static_files && prune_modes.has_receipts_pruning() ||
196            // TODO: support writing receipts to static files with log filter pruning enabled
197            receipts_in_static_files && !prune_modes.receipts_log_filter.is_empty()
198        {
199            EitherWriterDestination::Database
200        } else {
201            EitherWriterDestination::StaticFile
202        }
203    }
204
205    /// Returns the destination for writing account changesets.
206    ///
207    /// This determines the destination based solely on storage settings.
208    pub fn account_changesets_destination<P: DBProvider + StorageSettingsCache>(
209        provider: &P,
210    ) -> EitherWriterDestination {
211        if provider.cached_storage_settings().storage_v2 {
212            EitherWriterDestination::StaticFile
213        } else {
214            EitherWriterDestination::Database
215        }
216    }
217
218    /// Returns the destination for writing storage changesets.
219    ///
220    /// This determines the destination based solely on storage settings.
221    pub fn storage_changesets_destination<P: DBProvider + StorageSettingsCache>(
222        provider: &P,
223    ) -> EitherWriterDestination {
224        if provider.cached_storage_settings().storage_v2 {
225            EitherWriterDestination::StaticFile
226        } else {
227            EitherWriterDestination::Database
228        }
229    }
230
231    /// Creates a new [`EitherWriter`] for storages history based on storage settings.
232    pub fn new_storages_history<P>(
233        provider: &P,
234        _rocksdb_batch: RocksBatchArg<'a>,
235    ) -> ProviderResult<EitherWriterTy<'a, P, tables::StoragesHistory>>
236    where
237        P: DBProvider + NodePrimitivesProvider + StorageSettingsCache,
238        P::Tx: DbTxMut,
239    {
240        if provider.cached_storage_settings().storage_v2 {
241            return Ok(EitherWriter::RocksDB(_rocksdb_batch));
242        }
243
244        Ok(EitherWriter::Database(provider.tx_ref().cursor_write::<tables::StoragesHistory>()?))
245    }
246
247    /// Creates a new [`EitherWriter`] for transaction hash numbers based on storage settings.
248    pub fn new_transaction_hash_numbers<P>(
249        provider: &P,
250        _rocksdb_batch: RocksBatchArg<'a>,
251    ) -> ProviderResult<EitherWriterTy<'a, P, tables::TransactionHashNumbers>>
252    where
253        P: DBProvider + NodePrimitivesProvider + StorageSettingsCache,
254        P::Tx: DbTxMut,
255    {
256        if provider.cached_storage_settings().storage_v2 {
257            return Ok(EitherWriter::RocksDB(_rocksdb_batch));
258        }
259
260        Ok(EitherWriter::Database(
261            provider.tx_ref().cursor_write::<tables::TransactionHashNumbers>()?,
262        ))
263    }
264
265    /// Creates a new [`EitherWriter`] for account history based on storage settings.
266    pub fn new_accounts_history<P>(
267        provider: &P,
268        _rocksdb_batch: RocksBatchArg<'a>,
269    ) -> ProviderResult<EitherWriterTy<'a, P, tables::AccountsHistory>>
270    where
271        P: DBProvider + NodePrimitivesProvider + StorageSettingsCache,
272        P::Tx: DbTxMut,
273    {
274        if provider.cached_storage_settings().storage_v2 {
275            return Ok(EitherWriter::RocksDB(_rocksdb_batch));
276        }
277
278        Ok(EitherWriter::Database(provider.tx_ref().cursor_write::<tables::AccountsHistory>()?))
279    }
280}
281
282impl<'a, CURSOR, N: NodePrimitives> EitherWriter<'a, CURSOR, N> {
283    /// Extracts the raw `RocksDB` write batch from this writer, if it contains one.
284    ///
285    /// Returns `Some(WriteBatchWithTransaction)` for [`Self::RocksDB`] variant,
286    /// `None` for other variants.
287    ///
288    /// This is used to defer `RocksDB` commits to the provider level, ensuring all
289    /// storage commits (MDBX, static files, `RocksDB`) happen atomically in a single place.
290    pub fn into_raw_rocksdb_batch(self) -> Option<rocksdb::WriteBatchWithTransaction<true>> {
291        match self {
292            Self::Database(_) | Self::StaticFile(_) => None,
293            Self::RocksDB(batch) => Some(batch.into_inner()),
294        }
295    }
296
297    /// Increment the block number.
298    ///
299    /// Relevant only for [`Self::StaticFile`]. It is a no-op for [`Self::Database`].
300    pub fn increment_block(&mut self, expected_block_number: BlockNumber) -> ProviderResult<()> {
301        match self {
302            Self::Database(_) => Ok(()),
303            Self::StaticFile(writer) => writer.increment_block(expected_block_number),
304            Self::RocksDB(_) => Err(ProviderError::UnsupportedProvider),
305        }
306    }
307
308    /// Ensures that the writer is positioned at the specified block number.
309    ///
310    /// If the writer is positioned at a greater block number than the specified one, the writer
311    /// will NOT be unwound and the error will be returned.
312    ///
313    /// Relevant only for [`Self::StaticFile`]. It is a no-op for [`Self::Database`].
314    pub fn ensure_at_block(&mut self, block_number: BlockNumber) -> ProviderResult<()> {
315        match self {
316            Self::Database(_) => Ok(()),
317            Self::StaticFile(writer) => writer.ensure_at_block(block_number),
318            Self::RocksDB(_) => Err(ProviderError::UnsupportedProvider),
319        }
320    }
321}
322
323impl<'a, CURSOR, N: NodePrimitives> EitherWriter<'a, CURSOR, N>
324where
325    N::Receipt: Value,
326    CURSOR: DbCursorRW<tables::Receipts<N::Receipt>>,
327{
328    /// Append a transaction receipt.
329    pub fn append_receipt(&mut self, tx_num: TxNumber, receipt: &N::Receipt) -> ProviderResult<()> {
330        match self {
331            Self::Database(cursor) => Ok(cursor.append(tx_num, receipt)?),
332            Self::StaticFile(writer) => writer.append_receipt(tx_num, receipt),
333            Self::RocksDB(_) => Err(ProviderError::UnsupportedProvider),
334        }
335    }
336}
337
338impl<'a, CURSOR, N: NodePrimitives> EitherWriter<'a, CURSOR, N>
339where
340    CURSOR: DbCursorRW<tables::TransactionSenders>,
341{
342    /// Append a transaction sender to the destination
343    pub fn append_sender(&mut self, tx_num: TxNumber, sender: &Address) -> ProviderResult<()> {
344        match self {
345            Self::Database(cursor) => Ok(cursor.append(tx_num, sender)?),
346            Self::StaticFile(writer) => writer.append_transaction_sender(tx_num, sender),
347            Self::RocksDB(_) => Err(ProviderError::UnsupportedProvider),
348        }
349    }
350
351    /// Append transaction senders to the destination
352    pub fn append_senders<I>(&mut self, senders: I) -> ProviderResult<()>
353    where
354        I: Iterator<Item = (TxNumber, Address)>,
355    {
356        match self {
357            Self::Database(cursor) => {
358                for (tx_num, sender) in senders {
359                    cursor.append(tx_num, &sender)?;
360                }
361                Ok(())
362            }
363            Self::StaticFile(writer) => writer.append_transaction_senders(senders),
364            Self::RocksDB(_) => Err(ProviderError::UnsupportedProvider),
365        }
366    }
367
368    /// Removes all transaction senders above the given transaction number, and stops at the given
369    /// block number.
370    pub fn prune_senders(
371        &mut self,
372        unwind_tx_from: TxNumber,
373        block: BlockNumber,
374    ) -> ProviderResult<()>
375    where
376        CURSOR: DbCursorRO<tables::TransactionSenders>,
377    {
378        match self {
379            Self::Database(cursor) => {
380                let mut walker = cursor.walk_range(unwind_tx_from..)?;
381                while walker.next().transpose()?.is_some() {
382                    walker.delete_current()?;
383                }
384            }
385            Self::StaticFile(writer) => {
386                let static_file_transaction_sender_num = writer
387                    .reader()
388                    .get_highest_static_file_tx(StaticFileSegment::TransactionSenders);
389
390                let to_delete = static_file_transaction_sender_num
391                    .map(|static_num| (static_num + 1).saturating_sub(unwind_tx_from))
392                    .unwrap_or_default();
393
394                writer.prune_transaction_senders(to_delete, block)?;
395            }
396            Self::RocksDB(_) => return Err(ProviderError::UnsupportedProvider),
397        }
398
399        Ok(())
400    }
401}
402
403impl<'a, CURSOR, N: NodePrimitives> EitherWriter<'a, CURSOR, N>
404where
405    CURSOR: DbCursorRW<tables::TransactionHashNumbers> + DbCursorRO<tables::TransactionHashNumbers>,
406{
407    /// Puts a transaction hash number mapping.
408    ///
409    /// When `append_only` is true, uses `cursor.append()` which is significantly faster
410    /// but requires entries to be inserted in order and the table to be empty.
411    /// When false, uses `cursor.upsert()` which handles arbitrary insertion order and duplicates.
412    pub fn put_transaction_hash_number(
413        &mut self,
414        hash: TxHash,
415        tx_num: TxNumber,
416        append_only: bool,
417    ) -> ProviderResult<()> {
418        match self {
419            Self::Database(cursor) => {
420                if append_only {
421                    Ok(cursor.append(hash, &tx_num)?)
422                } else {
423                    Ok(cursor.upsert(hash, &tx_num)?)
424                }
425            }
426            Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
427            Self::RocksDB(batch) => batch.put::<tables::TransactionHashNumbers>(hash, &tx_num),
428        }
429    }
430
431    /// Puts multiple transaction hash number mappings in a batch.
432    ///
433    /// Accepts a vector of `(TxHash, TxNumber)` tuples and writes them all using the same cursor.
434    /// This is more efficient than calling `put_transaction_hash_number` repeatedly.
435    ///
436    /// When `append_only` is true, uses `cursor.append()` which requires entries to be
437    /// pre-sorted and the table to be empty or have only lower keys.
438    /// When false, uses `cursor.upsert()` which handles arbitrary insertion order.
439    pub fn put_transaction_hash_numbers_batch(
440        &mut self,
441        entries: Vec<(TxHash, TxNumber)>,
442        append_only: bool,
443    ) -> ProviderResult<()> {
444        match self {
445            Self::Database(cursor) => {
446                for (hash, tx_num) in entries {
447                    if append_only {
448                        cursor.append(hash, &tx_num)?;
449                    } else {
450                        cursor.upsert(hash, &tx_num)?;
451                    }
452                }
453                Ok(())
454            }
455            Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
456            Self::RocksDB(batch) => {
457                for (hash, tx_num) in entries {
458                    batch.put::<tables::TransactionHashNumbers>(hash, &tx_num)?;
459                }
460                Ok(())
461            }
462        }
463    }
464
465    /// Deletes a transaction hash number mapping.
466    pub fn delete_transaction_hash_number(&mut self, hash: TxHash) -> ProviderResult<()> {
467        match self {
468            Self::Database(cursor) => {
469                if cursor.seek_exact(hash)?.is_some() {
470                    cursor.delete_current()?;
471                }
472                Ok(())
473            }
474            Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
475            Self::RocksDB(batch) => batch.delete::<tables::TransactionHashNumbers>(hash),
476        }
477    }
478}
479
480impl<'a, CURSOR, N: NodePrimitives> EitherWriter<'a, CURSOR, N>
481where
482    CURSOR: DbCursorRW<tables::StoragesHistory> + DbCursorRO<tables::StoragesHistory>,
483{
484    /// Puts a storage history entry.
485    pub fn put_storage_history(
486        &mut self,
487        key: StorageShardedKey,
488        value: &BlockNumberList,
489    ) -> ProviderResult<()> {
490        match self {
491            Self::Database(cursor) => Ok(cursor.upsert(key, value)?),
492            Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
493            Self::RocksDB(batch) => batch.put::<tables::StoragesHistory>(key, value),
494        }
495    }
496
497    /// Deletes a storage history entry.
498    pub fn delete_storage_history(&mut self, key: StorageShardedKey) -> ProviderResult<()> {
499        match self {
500            Self::Database(cursor) => {
501                if cursor.seek_exact(key)?.is_some() {
502                    cursor.delete_current()?;
503                }
504                Ok(())
505            }
506            Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
507            Self::RocksDB(batch) => batch.delete::<tables::StoragesHistory>(key),
508        }
509    }
510
511    /// Appends a storage history entry (for first sync - more efficient).
512    pub fn append_storage_history(
513        &mut self,
514        key: StorageShardedKey,
515        value: &BlockNumberList,
516    ) -> ProviderResult<()> {
517        match self {
518            Self::Database(cursor) => Ok(cursor.append(key, value)?),
519            Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
520            Self::RocksDB(batch) => batch.put::<tables::StoragesHistory>(key, value),
521        }
522    }
523
524    /// Upserts a storage history entry (for incremental sync).
525    pub fn upsert_storage_history(
526        &mut self,
527        key: StorageShardedKey,
528        value: &BlockNumberList,
529    ) -> ProviderResult<()> {
530        match self {
531            Self::Database(cursor) => Ok(cursor.upsert(key, value)?),
532            Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
533            Self::RocksDB(batch) => batch.put::<tables::StoragesHistory>(key, value),
534        }
535    }
536
537    /// Gets the last shard for an address and storage key (keyed with `u64::MAX`).
538    pub fn get_last_storage_history_shard(
539        &mut self,
540        address: Address,
541        storage_key: B256,
542    ) -> ProviderResult<Option<BlockNumberList>> {
543        let key = StorageShardedKey::last(address, storage_key);
544        match self {
545            Self::Database(cursor) => Ok(cursor.seek_exact(key)?.map(|(_, v)| v)),
546            Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
547            Self::RocksDB(batch) => batch.get::<tables::StoragesHistory>(key),
548        }
549    }
550}
551
552impl<'a, CURSOR, N: NodePrimitives> EitherWriter<'a, CURSOR, N>
553where
554    CURSOR: DbCursorRW<tables::AccountsHistory> + DbCursorRO<tables::AccountsHistory>,
555{
556    /// Appends an account history entry (for first sync - more efficient).
557    pub fn append_account_history(
558        &mut self,
559        key: ShardedKey<Address>,
560        value: &BlockNumberList,
561    ) -> ProviderResult<()> {
562        match self {
563            Self::Database(cursor) => Ok(cursor.append(key, value)?),
564            Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
565            Self::RocksDB(batch) => batch.put::<tables::AccountsHistory>(key, value),
566        }
567    }
568
569    /// Upserts an account history entry (for incremental sync).
570    pub fn upsert_account_history(
571        &mut self,
572        key: ShardedKey<Address>,
573        value: &BlockNumberList,
574    ) -> ProviderResult<()> {
575        match self {
576            Self::Database(cursor) => Ok(cursor.upsert(key, value)?),
577            Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
578            Self::RocksDB(batch) => batch.put::<tables::AccountsHistory>(key, value),
579        }
580    }
581
582    /// Gets the last shard for an address (keyed with `u64::MAX`).
583    pub fn get_last_account_history_shard(
584        &mut self,
585        address: Address,
586    ) -> ProviderResult<Option<BlockNumberList>> {
587        match self {
588            Self::Database(cursor) => {
589                Ok(cursor.seek_exact(ShardedKey::last(address))?.map(|(_, v)| v))
590            }
591            Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
592            Self::RocksDB(batch) => batch.get::<tables::AccountsHistory>(ShardedKey::last(address)),
593        }
594    }
595
596    /// Deletes an account history entry.
597    pub fn delete_account_history(&mut self, key: ShardedKey<Address>) -> ProviderResult<()> {
598        match self {
599            Self::Database(cursor) => {
600                if cursor.seek_exact(key)?.is_some() {
601                    cursor.delete_current()?;
602                }
603                Ok(())
604            }
605            Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
606            Self::RocksDB(batch) => batch.delete::<tables::AccountsHistory>(key),
607        }
608    }
609}
610
611impl<'a, CURSOR, N: NodePrimitives> EitherWriter<'a, CURSOR, N>
612where
613    CURSOR: DbDupCursorRW<tables::AccountChangeSets>,
614{
615    /// Append account changeset for a block.
616    ///
617    /// NOTE: This _sorts_ the changesets by address before appending
618    pub fn append_account_changeset(
619        &mut self,
620        block_number: BlockNumber,
621        mut changeset: Vec<AccountBeforeTx>,
622    ) -> ProviderResult<()> {
623        // First sort the changesets
624        changeset.par_sort_by_key(|a| a.address);
625        match self {
626            Self::Database(cursor) => {
627                for change in changeset {
628                    cursor.append_dup(block_number, change)?;
629                }
630            }
631            Self::StaticFile(writer) => {
632                writer.append_account_changeset(changeset, block_number)?;
633            }
634            Self::RocksDB(_) => return Err(ProviderError::UnsupportedProvider),
635        }
636
637        Ok(())
638    }
639}
640
641impl<'a, CURSOR, N: NodePrimitives> EitherWriter<'a, CURSOR, N>
642where
643    CURSOR: DbDupCursorRW<tables::StorageChangeSets>,
644{
645    /// Append storage changeset for a block.
646    ///
647    /// NOTE: This _sorts_ the changesets by address and storage key before appending.
648    pub fn append_storage_changeset(
649        &mut self,
650        block_number: BlockNumber,
651        mut changeset: Vec<StorageBeforeTx>,
652    ) -> ProviderResult<()> {
653        changeset.par_sort_by_key(|change| (change.address, change.key));
654
655        match self {
656            Self::Database(cursor) => {
657                for change in changeset {
658                    let storage_id = BlockNumberAddress((block_number, change.address));
659                    cursor.append_dup(
660                        storage_id,
661                        StorageEntry { key: change.key, value: change.value },
662                    )?;
663                }
664            }
665            Self::StaticFile(writer) => {
666                writer.append_storage_changeset(changeset, block_number)?;
667            }
668            Self::RocksDB(_) => return Err(ProviderError::UnsupportedProvider),
669        }
670
671        Ok(())
672    }
673}
674
675/// Represents a source for reading data, either from database, static files, or `RocksDB`.
676#[derive(Debug, Display)]
677pub enum EitherReader<'a, 'db, CURSOR, N> {
678    /// Read from database table via cursor
679    Database(CURSOR, PhantomData<&'a ()>),
680    /// Read from static file
681    StaticFile(StaticFileProvider<N>, PhantomData<&'a ()>),
682    /// Read from `RocksDB` snapshot (works in both read-only and read-write modes)
683    RocksDB(&'a crate::providers::rocksdb::RocksReadSnapshot<'db>),
684}
685
686impl<'a, 'db> EitherReader<'a, 'db, (), ()> {
687    /// Creates a new [`EitherReader`] for senders based on storage settings.
688    pub fn new_senders<P>(
689        provider: &P,
690    ) -> ProviderResult<EitherReaderTy<'a, 'db, P, tables::TransactionSenders>>
691    where
692        P: DBProvider + NodePrimitivesProvider + StorageSettingsCache + StaticFileProviderFactory,
693        P::Tx: DbTx,
694    {
695        if EitherWriterDestination::senders(provider).is_static_file() {
696            Ok(EitherReader::StaticFile(provider.static_file_provider(), PhantomData))
697        } else {
698            Ok(EitherReader::Database(
699                provider.tx_ref().cursor_read::<tables::TransactionSenders>()?,
700                PhantomData,
701            ))
702        }
703    }
704
705    /// Creates a new [`EitherReader`] for storages history based on storage settings.
706    pub fn new_storages_history<P>(
707        provider: &P,
708        rocksdb: RocksDBRefArg<'a, 'db>,
709    ) -> ProviderResult<EitherReaderTy<'a, 'db, P, tables::StoragesHistory>>
710    where
711        P: DBProvider + NodePrimitivesProvider + StorageSettingsCache,
712        P::Tx: DbTx,
713    {
714        if provider.cached_storage_settings().storage_v2 {
715            return Ok(EitherReader::RocksDB(
716                rocksdb.expect("storages_history_in_rocksdb requires rocksdb snapshot"),
717            ));
718        }
719
720        Ok(EitherReader::Database(
721            provider.tx_ref().cursor_read::<tables::StoragesHistory>()?,
722            PhantomData,
723        ))
724    }
725
726    /// Creates a new [`EitherReader`] for transaction hash numbers based on storage settings.
727    pub fn new_transaction_hash_numbers<P>(
728        provider: &P,
729        rocksdb: RocksDBRefArg<'a, 'db>,
730    ) -> ProviderResult<EitherReaderTy<'a, 'db, P, tables::TransactionHashNumbers>>
731    where
732        P: DBProvider + NodePrimitivesProvider + StorageSettingsCache,
733        P::Tx: DbTx,
734    {
735        if provider.cached_storage_settings().storage_v2 {
736            return Ok(EitherReader::RocksDB(
737                rocksdb.expect("transaction_hash_numbers_in_rocksdb requires rocksdb snapshot"),
738            ));
739        }
740
741        Ok(EitherReader::Database(
742            provider.tx_ref().cursor_read::<tables::TransactionHashNumbers>()?,
743            PhantomData,
744        ))
745    }
746
747    /// Creates a new [`EitherReader`] for account history based on storage settings.
748    pub fn new_accounts_history<P>(
749        provider: &P,
750        rocksdb: RocksDBRefArg<'a, 'db>,
751    ) -> ProviderResult<EitherReaderTy<'a, 'db, P, tables::AccountsHistory>>
752    where
753        P: DBProvider + NodePrimitivesProvider + StorageSettingsCache,
754        P::Tx: DbTx,
755    {
756        if provider.cached_storage_settings().storage_v2 {
757            return Ok(EitherReader::RocksDB(
758                rocksdb.expect("account_history_in_rocksdb requires rocksdb snapshot"),
759            ));
760        }
761
762        Ok(EitherReader::Database(
763            provider.tx_ref().cursor_read::<tables::AccountsHistory>()?,
764            PhantomData,
765        ))
766    }
767
768    /// Creates a new [`EitherReader`] for account changesets based on storage settings.
769    pub fn new_account_changesets<P>(
770        provider: &P,
771    ) -> ProviderResult<DupEitherReaderTy<'a, 'db, P, tables::AccountChangeSets>>
772    where
773        P: DBProvider + NodePrimitivesProvider + StorageSettingsCache + StaticFileProviderFactory,
774        P::Tx: DbTx,
775    {
776        if EitherWriterDestination::account_changesets(provider).is_static_file() {
777            Ok(EitherReader::StaticFile(provider.static_file_provider(), PhantomData))
778        } else {
779            Ok(EitherReader::Database(
780                provider.tx_ref().cursor_dup_read::<tables::AccountChangeSets>()?,
781                PhantomData,
782            ))
783        }
784    }
785}
786
787impl<CURSOR, N: NodePrimitives> EitherReader<'_, '_, CURSOR, N>
788where
789    CURSOR: DbCursorRO<tables::TransactionSenders>,
790{
791    /// Fetches the senders for a range of transactions.
792    pub fn senders_by_tx_range(
793        &mut self,
794        range: Range<TxNumber>,
795    ) -> ProviderResult<HashMap<TxNumber, Address>> {
796        match self {
797            Self::Database(cursor, _) => cursor
798                .walk_range(range)?
799                .map(|result| result.map_err(ProviderError::from))
800                .collect::<ProviderResult<HashMap<_, _>>>(),
801            Self::StaticFile(provider, _) => range
802                .clone()
803                .zip(provider.fetch_range_iter(
804                    StaticFileSegment::TransactionSenders,
805                    range,
806                    |cursor, number| cursor.get_one::<TransactionSenderMask>(number.into()),
807                )?)
808                .filter_map(|(tx_num, sender)| {
809                    let result = sender.transpose()?;
810                    Some(result.map(|sender| (tx_num, sender)))
811                })
812                .collect::<ProviderResult<HashMap<_, _>>>(),
813            Self::RocksDB(_) => Err(ProviderError::UnsupportedProvider),
814        }
815    }
816}
817
818impl<CURSOR, N: NodePrimitives> EitherReader<'_, '_, CURSOR, N>
819where
820    CURSOR: DbCursorRO<tables::TransactionHashNumbers>,
821{
822    /// Gets a transaction number by its hash.
823    pub fn get_transaction_hash_number(
824        &mut self,
825        hash: TxHash,
826    ) -> ProviderResult<Option<TxNumber>> {
827        match self {
828            Self::Database(cursor, _) => Ok(cursor.seek_exact(hash)?.map(|(_, v)| v)),
829            Self::StaticFile(_, _) => Err(ProviderError::UnsupportedProvider),
830            Self::RocksDB(snapshot) => snapshot.get::<tables::TransactionHashNumbers>(hash),
831        }
832    }
833}
834
835impl<CURSOR, N: NodePrimitives> EitherReader<'_, '_, CURSOR, N>
836where
837    CURSOR: DbCursorRO<tables::StoragesHistory>,
838{
839    /// Gets a storage history shard entry for the given [`StorageShardedKey`], if present.
840    pub fn get_storage_history(
841        &mut self,
842        key: StorageShardedKey,
843    ) -> ProviderResult<Option<BlockNumberList>> {
844        match self {
845            Self::Database(cursor, _) => Ok(cursor.seek_exact(key)?.map(|(_, v)| v)),
846            Self::StaticFile(_, _) => Err(ProviderError::UnsupportedProvider),
847            Self::RocksDB(snapshot) => snapshot.get::<tables::StoragesHistory>(key),
848        }
849    }
850
851    /// Lookup storage history and return [`HistoryInfo`].
852    pub fn storage_history_info(
853        &mut self,
854        address: Address,
855        storage_key: alloy_primitives::B256,
856        block_number: BlockNumber,
857        lowest_available_block_number: Option<BlockNumber>,
858        visible_tip: BlockNumber,
859    ) -> ProviderResult<HistoryInfo> {
860        match self {
861            Self::Database(cursor, _) => {
862                let key = StorageShardedKey::new(address, storage_key, block_number);
863                history_info::<tables::StoragesHistory, _, _>(
864                    cursor,
865                    key,
866                    block_number,
867                    |k| k.address == address && k.sharded_key.key == storage_key,
868                    lowest_available_block_number,
869                )
870            }
871            Self::StaticFile(_, _) => Err(ProviderError::UnsupportedProvider),
872            Self::RocksDB(snapshot) => snapshot.storage_history_info(
873                address,
874                storage_key,
875                block_number,
876                lowest_available_block_number,
877                visible_tip,
878            ),
879        }
880    }
881}
882
883impl<CURSOR, N: NodePrimitives> EitherReader<'_, '_, CURSOR, N>
884where
885    CURSOR: DbCursorRO<tables::AccountsHistory>,
886{
887    /// Gets an account history shard entry for the given [`ShardedKey`], if present.
888    pub fn get_account_history(
889        &mut self,
890        key: ShardedKey<Address>,
891    ) -> ProviderResult<Option<BlockNumberList>> {
892        match self {
893            Self::Database(cursor, _) => Ok(cursor.seek_exact(key)?.map(|(_, v)| v)),
894            Self::StaticFile(_, _) => Err(ProviderError::UnsupportedProvider),
895            Self::RocksDB(snapshot) => snapshot.get::<tables::AccountsHistory>(key),
896        }
897    }
898
899    /// Lookup account history and return [`HistoryInfo`].
900    pub fn account_history_info(
901        &mut self,
902        address: Address,
903        block_number: BlockNumber,
904        lowest_available_block_number: Option<BlockNumber>,
905        visible_tip: BlockNumber,
906    ) -> ProviderResult<HistoryInfo> {
907        match self {
908            Self::Database(cursor, _) => {
909                let key = ShardedKey::new(address, block_number);
910                history_info::<tables::AccountsHistory, _, _>(
911                    cursor,
912                    key,
913                    block_number,
914                    |k| k.key == address,
915                    lowest_available_block_number,
916                )
917            }
918            Self::StaticFile(_, _) => Err(ProviderError::UnsupportedProvider),
919            Self::RocksDB(snapshot) => snapshot.account_history_info(
920                address,
921                block_number,
922                lowest_available_block_number,
923                visible_tip,
924            ),
925        }
926    }
927}
928
929impl<CURSOR, N: NodePrimitives> EitherReader<'_, '_, CURSOR, N>
930where
931    CURSOR: DbCursorRO<tables::AccountChangeSets>,
932{
933    /// Iterate over account changesets and return all account address that were changed.
934    pub fn changed_accounts_with_range(
935        &mut self,
936        range: RangeInclusive<BlockNumber>,
937    ) -> ProviderResult<BTreeSet<Address>> {
938        match self {
939            Self::StaticFile(provider, _) => {
940                let highest_static_block =
941                    provider.get_highest_static_file_block(StaticFileSegment::AccountChangeSets);
942
943                let Some(highest) = highest_static_block else {
944                    return Err(ProviderError::MissingHighestStaticFileBlock(
945                        StaticFileSegment::AccountChangeSets,
946                    ))
947                };
948
949                let start = *range.start();
950                let static_end = (*range.end()).min(highest);
951
952                let mut changed_accounts = BTreeSet::default();
953                if start <= static_end {
954                    for block in start..=static_end {
955                        let block_changesets = provider.account_block_changeset(block)?;
956                        for changeset in block_changesets {
957                            changed_accounts.insert(changeset.address);
958                        }
959                    }
960                }
961
962                Ok(changed_accounts)
963            }
964            Self::Database(provider, _) => provider
965                .walk_range(range)?
966                .map(|entry| {
967                    entry.map(|(_, account_before)| account_before.address).map_err(Into::into)
968                })
969                .collect(),
970            Self::RocksDB(_) => Err(ProviderError::UnsupportedProvider),
971        }
972    }
973}
974
975/// Destination for writing data.
976#[derive(Debug, EnumIs)]
977pub enum EitherWriterDestination {
978    /// Write to database table
979    Database,
980    /// Write to static file
981    StaticFile,
982    /// Write to `RocksDB`
983    RocksDB,
984}
985
986impl EitherWriterDestination {
987    /// Returns the destination for writing senders based on storage settings.
988    pub fn senders<P>(provider: &P) -> Self
989    where
990        P: StorageSettingsCache,
991    {
992        // Write senders to static files only if they're explicitly enabled
993        if provider.cached_storage_settings().storage_v2 {
994            Self::StaticFile
995        } else {
996            Self::Database
997        }
998    }
999
1000    /// Returns the destination for writing account changesets based on storage settings.
1001    pub fn account_changesets<P>(provider: &P) -> Self
1002    where
1003        P: StorageSettingsCache,
1004    {
1005        // Write account changesets to static files only if they're explicitly enabled
1006        if provider.cached_storage_settings().storage_v2 {
1007            Self::StaticFile
1008        } else {
1009            Self::Database
1010        }
1011    }
1012
1013    /// Returns the destination for writing storage changesets based on storage settings.
1014    pub fn storage_changesets<P>(provider: &P) -> Self
1015    where
1016        P: StorageSettingsCache,
1017    {
1018        // Write storage changesets to static files only if they're explicitly enabled
1019        if provider.cached_storage_settings().storage_v2 {
1020            Self::StaticFile
1021        } else {
1022            Self::Database
1023        }
1024    }
1025}
1026
1027#[cfg(test)]
1028mod tests {
1029    use crate::{test_utils::create_test_provider_factory, StaticFileWriter};
1030
1031    use super::*;
1032    use alloy_primitives::Address;
1033    use reth_db::models::AccountBeforeTx;
1034    use reth_static_file_types::StaticFileSegment;
1035    use reth_storage_api::{DatabaseProviderFactory, StorageSettings};
1036
1037    /// Verifies that `changed_accounts_with_range` correctly caps the query range to the
1038    /// static file tip when the requested range extends beyond it.
1039    ///
1040    /// This test documents the fix for an off-by-one bug where the code computed
1041    /// `static_end = range.end().min(highest + 1)` instead of `min(highest)`.
1042    /// The bug allowed iteration to attempt reading block `highest + 1` which doesn't
1043    /// exist (silently returning empty results due to `MissingStaticFileBlock` handling).
1044    ///
1045    /// While the bug was masked by error handling, it caused:
1046    /// 1. Unnecessary iteration/lookup for non-existent blocks
1047    /// 2. Potential overflow when `highest == u64::MAX`
1048    #[test]
1049    fn test_changed_accounts_with_range_caps_at_static_file_tip() {
1050        let factory = create_test_provider_factory();
1051        let highest_block = 5u64;
1052
1053        let addresses: Vec<Address> = (0..=highest_block)
1054            .map(|i| {
1055                let mut addr = Address::ZERO;
1056                addr.0[0] = i as u8;
1057                addr
1058            })
1059            .collect();
1060
1061        {
1062            let sf_provider = factory.static_file_provider();
1063            let mut writer =
1064                sf_provider.latest_writer(StaticFileSegment::AccountChangeSets).unwrap();
1065
1066            for block_num in 0..=highest_block {
1067                let changeset =
1068                    vec![AccountBeforeTx { address: addresses[block_num as usize], info: None }];
1069                writer.append_account_changeset(changeset, block_num).unwrap();
1070            }
1071            writer.commit().unwrap();
1072        }
1073
1074        factory.set_storage_settings_cache(StorageSettings::v2());
1075
1076        let provider = factory.database_provider_ro().unwrap();
1077
1078        let sf_tip = provider
1079            .static_file_provider()
1080            .get_highest_static_file_block(StaticFileSegment::AccountChangeSets);
1081        assert_eq!(sf_tip, Some(highest_block));
1082
1083        let mut reader = EitherReader::new_account_changesets(&provider).unwrap();
1084        assert!(matches!(reader, EitherReader::StaticFile(_, _)));
1085
1086        // Query range 0..=10 when tip is 5 - should return only accounts from blocks 0-5
1087        let result = reader.changed_accounts_with_range(0..=10).unwrap();
1088
1089        let expected: BTreeSet<Address> = addresses.into_iter().collect();
1090        assert_eq!(result, expected);
1091    }
1092
1093    #[test]
1094    fn test_reader_senders_by_tx_range() {
1095        let factory = create_test_provider_factory();
1096
1097        // Insert senders only from 1 to 4, but we will query from 0 to 5.
1098        let senders = [
1099            (1, Address::random()),
1100            (2, Address::random()),
1101            (3, Address::random()),
1102            (4, Address::random()),
1103        ];
1104
1105        for transaction_senders_in_static_files in [false, true] {
1106            factory.set_storage_settings_cache(if transaction_senders_in_static_files {
1107                StorageSettings::v2()
1108            } else {
1109                StorageSettings::v1()
1110            });
1111
1112            let provider = factory.database_provider_rw().unwrap();
1113            let mut writer = EitherWriter::new_senders(&provider, 0).unwrap();
1114            if transaction_senders_in_static_files {
1115                assert!(matches!(writer, EitherWriter::StaticFile(_)));
1116            } else {
1117                assert!(matches!(writer, EitherWriter::Database(_)));
1118            }
1119
1120            writer.increment_block(0).unwrap();
1121            writer.append_senders(senders.iter().copied()).unwrap();
1122            drop(writer);
1123            provider.commit().unwrap();
1124
1125            let provider = factory.database_provider_ro().unwrap();
1126            let mut reader = EitherReader::new_senders(&provider).unwrap();
1127            if transaction_senders_in_static_files {
1128                assert!(matches!(reader, EitherReader::StaticFile(_, _)));
1129            } else {
1130                assert!(matches!(reader, EitherReader::Database(_, _)));
1131            }
1132
1133            assert_eq!(
1134                reader.senders_by_tx_range(0..6).unwrap(),
1135                senders.iter().copied().collect::<HashMap<_, _>>(),
1136                "{reader}"
1137            );
1138        }
1139    }
1140}
1141
1142#[cfg(test)]
1143mod rocksdb_tests {
1144    use super::*;
1145    use crate::{
1146        providers::rocksdb::{RocksDBBuilder, RocksDBProvider},
1147        test_utils::create_test_provider_factory,
1148        RocksDBProviderFactory,
1149    };
1150    use alloy_primitives::{Address, B256};
1151    use reth_db_api::{
1152        models::{storage_sharded_key::StorageShardedKey, IntegerList, ShardedKey},
1153        tables,
1154        transaction::DbTxMut,
1155    };
1156    use reth_ethereum_primitives::EthPrimitives;
1157    use reth_storage_api::{DatabaseProviderFactory, StorageSettings};
1158    use std::marker::PhantomData;
1159    use tempfile::TempDir;
1160
1161    fn create_rocksdb_provider() -> (TempDir, RocksDBProvider) {
1162        let temp_dir = TempDir::new().unwrap();
1163        let provider = RocksDBBuilder::new(temp_dir.path())
1164            .with_table::<tables::TransactionHashNumbers>()
1165            .with_table::<tables::StoragesHistory>()
1166            .with_table::<tables::AccountsHistory>()
1167            .build()
1168            .unwrap();
1169        (temp_dir, provider)
1170    }
1171
1172    /// Test that `EitherWriter::new_transaction_hash_numbers` creates a `RocksDB` writer
1173    /// when the storage setting is enabled, and that put operations followed by commit
1174    /// persist the data to `RocksDB`.
1175    #[test]
1176    fn test_either_writer_transaction_hash_numbers_with_rocksdb() {
1177        let factory = create_test_provider_factory();
1178
1179        // Enable RocksDB for transaction hash numbers
1180        factory.set_storage_settings_cache(StorageSettings::v2());
1181
1182        let hash1 = B256::from([1u8; 32]);
1183        let hash2 = B256::from([2u8; 32]);
1184        let tx_num1 = 100u64;
1185        let tx_num2 = 200u64;
1186
1187        // Get the RocksDB batch from the provider
1188        let rocksdb = factory.rocksdb_provider();
1189        let batch = rocksdb.batch();
1190
1191        // Create EitherWriter with RocksDB
1192        let provider = factory.database_provider_rw().unwrap();
1193        let mut writer = EitherWriter::new_transaction_hash_numbers(&provider, batch).unwrap();
1194
1195        // Verify we got a RocksDB writer
1196        assert!(matches!(writer, EitherWriter::RocksDB(_)));
1197
1198        // Write transaction hash numbers (append_only=false since we're using RocksDB)
1199        writer.put_transaction_hash_number(hash1, tx_num1, false).unwrap();
1200        writer.put_transaction_hash_number(hash2, tx_num2, false).unwrap();
1201
1202        // Extract the batch and register with provider for commit
1203        if let Some(batch) = writer.into_raw_rocksdb_batch() {
1204            provider.set_pending_rocksdb_batch(batch);
1205        }
1206
1207        // Commit via provider - this commits RocksDB batch too
1208        provider.commit().unwrap();
1209
1210        // Verify data was written to RocksDB
1211        let rocksdb = factory.rocksdb_provider();
1212        assert_eq!(rocksdb.get::<tables::TransactionHashNumbers>(hash1).unwrap(), Some(tx_num1));
1213        assert_eq!(rocksdb.get::<tables::TransactionHashNumbers>(hash2).unwrap(), Some(tx_num2));
1214    }
1215
1216    /// Test that `EitherWriter::delete_transaction_hash_number` works with `RocksDB`.
1217    #[test]
1218    fn test_either_writer_delete_transaction_hash_number_with_rocksdb() {
1219        let factory = create_test_provider_factory();
1220
1221        // Enable RocksDB for transaction hash numbers
1222        factory.set_storage_settings_cache(StorageSettings::v2());
1223
1224        let hash = B256::from([1u8; 32]);
1225        let tx_num = 100u64;
1226
1227        // First, write a value directly to RocksDB
1228        let rocksdb = factory.rocksdb_provider();
1229        rocksdb.put::<tables::TransactionHashNumbers>(hash, &tx_num).unwrap();
1230        assert_eq!(rocksdb.get::<tables::TransactionHashNumbers>(hash).unwrap(), Some(tx_num));
1231
1232        // Now delete using EitherWriter
1233        let batch = rocksdb.batch();
1234        let provider = factory.database_provider_rw().unwrap();
1235        let mut writer = EitherWriter::new_transaction_hash_numbers(&provider, batch).unwrap();
1236        writer.delete_transaction_hash_number(hash).unwrap();
1237
1238        // Extract the batch and commit via provider
1239        if let Some(batch) = writer.into_raw_rocksdb_batch() {
1240            provider.set_pending_rocksdb_batch(batch);
1241        }
1242        provider.commit().unwrap();
1243
1244        // Verify deletion
1245        let rocksdb = factory.rocksdb_provider();
1246        assert_eq!(rocksdb.get::<tables::TransactionHashNumbers>(hash).unwrap(), None);
1247    }
1248
1249    #[test]
1250    fn test_rocksdb_batch_transaction_hash_numbers() {
1251        let (_temp_dir, provider) = create_rocksdb_provider();
1252
1253        let hash1 = B256::from([1u8; 32]);
1254        let hash2 = B256::from([2u8; 32]);
1255        let tx_num1 = 100u64;
1256        let tx_num2 = 200u64;
1257
1258        // Write via RocksDBBatch (same as EitherWriter::RocksDB would use internally)
1259        let mut batch = provider.batch();
1260        batch.put::<tables::TransactionHashNumbers>(hash1, &tx_num1).unwrap();
1261        batch.put::<tables::TransactionHashNumbers>(hash2, &tx_num2).unwrap();
1262        batch.commit().unwrap();
1263
1264        // Read via RocksTx (same as EitherReader::RocksDB would use internally)
1265        let tx = provider.tx();
1266        assert_eq!(tx.get::<tables::TransactionHashNumbers>(hash1).unwrap(), Some(tx_num1));
1267        assert_eq!(tx.get::<tables::TransactionHashNumbers>(hash2).unwrap(), Some(tx_num2));
1268
1269        // Test missing key
1270        let missing_hash = B256::from([99u8; 32]);
1271        assert_eq!(tx.get::<tables::TransactionHashNumbers>(missing_hash).unwrap(), None);
1272    }
1273
1274    #[test]
1275    fn test_rocksdb_batch_storage_history() {
1276        let (_temp_dir, provider) = create_rocksdb_provider();
1277
1278        let address = Address::random();
1279        let storage_key = B256::from([1u8; 32]);
1280        let key = StorageShardedKey::new(address, storage_key, 1000);
1281        let value = IntegerList::new([1, 5, 10, 50]).unwrap();
1282
1283        // Write via RocksDBBatch
1284        let mut batch = provider.batch();
1285        batch.put::<tables::StoragesHistory>(key.clone(), &value).unwrap();
1286        batch.commit().unwrap();
1287
1288        // Read via RocksTx
1289        let tx = provider.tx();
1290        let result = tx.get::<tables::StoragesHistory>(key).unwrap();
1291        assert_eq!(result, Some(value));
1292
1293        // Test missing key
1294        let missing_key = StorageShardedKey::new(Address::random(), B256::ZERO, 0);
1295        assert_eq!(tx.get::<tables::StoragesHistory>(missing_key).unwrap(), None);
1296    }
1297
1298    #[test]
1299    fn test_rocksdb_batch_account_history() {
1300        let (_temp_dir, provider) = create_rocksdb_provider();
1301
1302        let address = Address::random();
1303        let key = ShardedKey::new(address, 1000);
1304        let value = IntegerList::new([1, 10, 100, 500]).unwrap();
1305
1306        // Write via RocksDBBatch
1307        let mut batch = provider.batch();
1308        batch.put::<tables::AccountsHistory>(key.clone(), &value).unwrap();
1309        batch.commit().unwrap();
1310
1311        // Read via RocksTx
1312        let tx = provider.tx();
1313        let result = tx.get::<tables::AccountsHistory>(key).unwrap();
1314        assert_eq!(result, Some(value));
1315
1316        // Test missing key
1317        let missing_key = ShardedKey::new(Address::random(), 0);
1318        assert_eq!(tx.get::<tables::AccountsHistory>(missing_key).unwrap(), None);
1319    }
1320
1321    #[test]
1322    fn test_rocksdb_batch_delete_transaction_hash_number() {
1323        let (_temp_dir, provider) = create_rocksdb_provider();
1324
1325        let hash = B256::from([1u8; 32]);
1326        let tx_num = 100u64;
1327
1328        // First write
1329        provider.put::<tables::TransactionHashNumbers>(hash, &tx_num).unwrap();
1330        assert_eq!(provider.get::<tables::TransactionHashNumbers>(hash).unwrap(), Some(tx_num));
1331
1332        // Delete via RocksDBBatch
1333        let mut batch = provider.batch();
1334        batch.delete::<tables::TransactionHashNumbers>(hash).unwrap();
1335        batch.commit().unwrap();
1336
1337        // Verify deletion
1338        assert_eq!(provider.get::<tables::TransactionHashNumbers>(hash).unwrap(), None);
1339    }
1340
1341    #[test]
1342    fn test_rocksdb_batch_delete_storage_history() {
1343        let (_temp_dir, provider) = create_rocksdb_provider();
1344
1345        let address = Address::random();
1346        let storage_key = B256::from([1u8; 32]);
1347        let key = StorageShardedKey::new(address, storage_key, 1000);
1348        let value = IntegerList::new([1, 5, 10]).unwrap();
1349
1350        // First write
1351        provider.put::<tables::StoragesHistory>(key.clone(), &value).unwrap();
1352        assert!(provider.get::<tables::StoragesHistory>(key.clone()).unwrap().is_some());
1353
1354        // Delete via RocksDBBatch
1355        let mut batch = provider.batch();
1356        batch.delete::<tables::StoragesHistory>(key.clone()).unwrap();
1357        batch.commit().unwrap();
1358
1359        // Verify deletion
1360        assert_eq!(provider.get::<tables::StoragesHistory>(key).unwrap(), None);
1361    }
1362
1363    #[test]
1364    fn test_rocksdb_batch_delete_account_history() {
1365        let (_temp_dir, provider) = create_rocksdb_provider();
1366
1367        let address = Address::random();
1368        let key = ShardedKey::new(address, 1000);
1369        let value = IntegerList::new([1, 10, 100]).unwrap();
1370
1371        // First write
1372        provider.put::<tables::AccountsHistory>(key.clone(), &value).unwrap();
1373        assert!(provider.get::<tables::AccountsHistory>(key.clone()).unwrap().is_some());
1374
1375        // Delete via RocksDBBatch
1376        let mut batch = provider.batch();
1377        batch.delete::<tables::AccountsHistory>(key.clone()).unwrap();
1378        batch.commit().unwrap();
1379
1380        // Verify deletion
1381        assert_eq!(provider.get::<tables::AccountsHistory>(key).unwrap(), None);
1382    }
1383
1384    // ==================== Parametrized Backend Equivalence Tests ====================
1385    //
1386    // These tests verify that MDBX and RocksDB produce identical results for history lookups.
1387    // Each scenario sets up the same data in both backends and asserts identical HistoryInfo.
1388
1389    /// Query parameters for a history lookup test case.
1390    struct HistoryQuery {
1391        block_number: BlockNumber,
1392        lowest_available: Option<BlockNumber>,
1393        expected: HistoryInfo,
1394    }
1395
1396    // Type aliases for cursor types (needed for EitherWriter/EitherReader type inference)
1397    type AccountsHistoryWriteCursor =
1398        reth_db::mdbx::cursor::Cursor<reth_db::mdbx::RW, tables::AccountsHistory>;
1399    type StoragesHistoryWriteCursor =
1400        reth_db::mdbx::cursor::Cursor<reth_db::mdbx::RW, tables::StoragesHistory>;
1401    type AccountsHistoryReadCursor =
1402        reth_db::mdbx::cursor::Cursor<reth_db::mdbx::RO, tables::AccountsHistory>;
1403    type StoragesHistoryReadCursor =
1404        reth_db::mdbx::cursor::Cursor<reth_db::mdbx::RO, tables::StoragesHistory>;
1405
1406    /// Runs the same account history queries against both MDBX and `RocksDB` backends,
1407    /// asserting they produce identical results.
1408    fn run_account_history_scenario(
1409        scenario_name: &str,
1410        address: Address,
1411        shards: &[(BlockNumber, Vec<BlockNumber>)], // (shard_highest_block, blocks_in_shard)
1412        queries: &[HistoryQuery],
1413    ) {
1414        // Setup MDBX and RocksDB with identical data using EitherWriter
1415        let factory = create_test_provider_factory();
1416        let mdbx_provider = factory.database_provider_rw().unwrap();
1417        let (temp_dir, rocks_provider) = create_rocksdb_provider();
1418
1419        // Create writers for both backends
1420        let mut mdbx_writer: EitherWriter<'_, AccountsHistoryWriteCursor, EthPrimitives> =
1421            EitherWriter::Database(
1422                mdbx_provider.tx_ref().cursor_write::<tables::AccountsHistory>().unwrap(),
1423            );
1424        let mut rocks_writer: EitherWriter<'_, AccountsHistoryWriteCursor, EthPrimitives> =
1425            EitherWriter::RocksDB(rocks_provider.batch());
1426
1427        // Write identical data to both backends in a single loop
1428        for (highest_block, blocks) in shards {
1429            let key = ShardedKey::new(address, *highest_block);
1430            let value = IntegerList::new(blocks.clone()).unwrap();
1431            mdbx_writer.upsert_account_history(key.clone(), &value).unwrap();
1432            rocks_writer.upsert_account_history(key, &value).unwrap();
1433        }
1434
1435        // Commit both backends
1436        drop(mdbx_writer);
1437        mdbx_provider.commit().unwrap();
1438        if let EitherWriter::RocksDB(batch) = rocks_writer {
1439            batch.commit().unwrap();
1440        }
1441
1442        // Run queries against both backends using EitherReader
1443        let mdbx_ro = factory.database_provider_ro().unwrap();
1444        let rocks_snapshot = rocks_provider.snapshot();
1445
1446        for (i, query) in queries.iter().enumerate() {
1447            // MDBX query via EitherReader
1448            let mut mdbx_reader: EitherReader<'_, '_, AccountsHistoryReadCursor, EthPrimitives> =
1449                EitherReader::Database(
1450                    mdbx_ro.tx_ref().cursor_read::<tables::AccountsHistory>().unwrap(),
1451                    PhantomData,
1452                );
1453            let mdbx_result = mdbx_reader
1454                .account_history_info(address, query.block_number, query.lowest_available, u64::MAX)
1455                .unwrap();
1456
1457            // RocksDB query via EitherReader — reuse snapshot for consistent view
1458            let rocks_result = rocks_snapshot
1459                .account_history_info(address, query.block_number, query.lowest_available, u64::MAX)
1460                .unwrap();
1461
1462            // Assert both backends produce identical results
1463            assert_eq!(
1464                mdbx_result,
1465                rocks_result,
1466                "Backend mismatch in scenario '{}' query {}: block={}, lowest={:?}\n\
1467                 MDBX: {:?}, RocksDB: {:?}",
1468                scenario_name,
1469                i,
1470                query.block_number,
1471                query.lowest_available,
1472                mdbx_result,
1473                rocks_result
1474            );
1475
1476            // Also verify against expected result
1477            assert_eq!(
1478                mdbx_result,
1479                query.expected,
1480                "Unexpected result in scenario '{}' query {}: block={}, lowest={:?}\n\
1481                 Got: {:?}, Expected: {:?}",
1482                scenario_name,
1483                i,
1484                query.block_number,
1485                query.lowest_available,
1486                mdbx_result,
1487                query.expected
1488            );
1489        }
1490
1491        drop(temp_dir);
1492    }
1493
1494    /// Runs the same storage history queries against both MDBX and `RocksDB` backends,
1495    /// asserting they produce identical results.
1496    fn run_storage_history_scenario(
1497        scenario_name: &str,
1498        address: Address,
1499        storage_key: B256,
1500        shards: &[(BlockNumber, Vec<BlockNumber>)], // (shard_highest_block, blocks_in_shard)
1501        queries: &[HistoryQuery],
1502    ) {
1503        // Setup MDBX and RocksDB with identical data using EitherWriter
1504        let factory = create_test_provider_factory();
1505        let mdbx_provider = factory.database_provider_rw().unwrap();
1506        let (temp_dir, rocks_provider) = create_rocksdb_provider();
1507
1508        // Create writers for both backends
1509        let mut mdbx_writer: EitherWriter<'_, StoragesHistoryWriteCursor, EthPrimitives> =
1510            EitherWriter::Database(
1511                mdbx_provider.tx_ref().cursor_write::<tables::StoragesHistory>().unwrap(),
1512            );
1513        let mut rocks_writer: EitherWriter<'_, StoragesHistoryWriteCursor, EthPrimitives> =
1514            EitherWriter::RocksDB(rocks_provider.batch());
1515
1516        // Write identical data to both backends in a single loop
1517        for (highest_block, blocks) in shards {
1518            let key = StorageShardedKey::new(address, storage_key, *highest_block);
1519            let value = IntegerList::new(blocks.clone()).unwrap();
1520            mdbx_writer.put_storage_history(key.clone(), &value).unwrap();
1521            rocks_writer.put_storage_history(key, &value).unwrap();
1522        }
1523
1524        // Commit both backends
1525        drop(mdbx_writer);
1526        mdbx_provider.commit().unwrap();
1527        if let EitherWriter::RocksDB(batch) = rocks_writer {
1528            batch.commit().unwrap();
1529        }
1530
1531        // Run queries against both backends using EitherReader
1532        let mdbx_ro = factory.database_provider_ro().unwrap();
1533        let rocks_snapshot = rocks_provider.snapshot();
1534
1535        for (i, query) in queries.iter().enumerate() {
1536            // MDBX query via EitherReader
1537            let mut mdbx_reader: EitherReader<'_, '_, StoragesHistoryReadCursor, EthPrimitives> =
1538                EitherReader::Database(
1539                    mdbx_ro.tx_ref().cursor_read::<tables::StoragesHistory>().unwrap(),
1540                    PhantomData,
1541                );
1542            let mdbx_result = mdbx_reader
1543                .storage_history_info(
1544                    address,
1545                    storage_key,
1546                    query.block_number,
1547                    query.lowest_available,
1548                    u64::MAX,
1549                )
1550                .unwrap();
1551
1552            // RocksDB query via snapshot — reuse for consistent view
1553            let rocks_result = rocks_snapshot
1554                .storage_history_info(
1555                    address,
1556                    storage_key,
1557                    query.block_number,
1558                    query.lowest_available,
1559                    u64::MAX,
1560                )
1561                .unwrap();
1562
1563            // Assert both backends produce identical results
1564            assert_eq!(
1565                mdbx_result,
1566                rocks_result,
1567                "Backend mismatch in scenario '{}' query {}: block={}, lowest={:?}\n\
1568                 MDBX: {:?}, RocksDB: {:?}",
1569                scenario_name,
1570                i,
1571                query.block_number,
1572                query.lowest_available,
1573                mdbx_result,
1574                rocks_result
1575            );
1576
1577            // Also verify against expected result
1578            assert_eq!(
1579                mdbx_result,
1580                query.expected,
1581                "Unexpected result in scenario '{}' query {}: block={}, lowest={:?}\n\
1582                 Got: {:?}, Expected: {:?}",
1583                scenario_name,
1584                i,
1585                query.block_number,
1586                query.lowest_available,
1587                mdbx_result,
1588                query.expected
1589            );
1590        }
1591
1592        drop(temp_dir);
1593    }
1594
1595    /// Tests account history lookups across both MDBX and `RocksDB` backends.
1596    ///
1597    /// Covers the following scenarios from PR2's `RocksDB`-only tests:
1598    /// 1. Single shard - basic lookups within one shard
1599    /// 2. Multiple shards - `prev()` shard detection and transitions
1600    /// 3. No history - query address with no entries
1601    /// 4. Pruning boundary - `lowest_available` boundary behavior (block at/after boundary)
1602    #[test]
1603    fn test_account_history_info_both_backends() {
1604        let address = Address::from([0x42; 20]);
1605
1606        // Scenario 1: Single shard with blocks [100, 200, 300]
1607        run_account_history_scenario(
1608            "single_shard",
1609            address,
1610            &[(u64::MAX, vec![100, 200, 300])],
1611            &[
1612                // Before first entry -> NotYetWritten
1613                HistoryQuery {
1614                    block_number: 50,
1615                    lowest_available: None,
1616                    expected: HistoryInfo::NotYetWritten,
1617                },
1618                // Between entries -> InChangeset(next_write)
1619                HistoryQuery {
1620                    block_number: 150,
1621                    lowest_available: None,
1622                    expected: HistoryInfo::InChangeset(200),
1623                },
1624                // Exact match on entry -> InChangeset(same_block)
1625                HistoryQuery {
1626                    block_number: 300,
1627                    lowest_available: None,
1628                    expected: HistoryInfo::InChangeset(300),
1629                },
1630                // After last entry in last shard -> InPlainState
1631                HistoryQuery {
1632                    block_number: 500,
1633                    lowest_available: None,
1634                    expected: HistoryInfo::InPlainState,
1635                },
1636            ],
1637        );
1638
1639        // Scenario 2: Multiple shards - tests prev() shard detection
1640        run_account_history_scenario(
1641            "multiple_shards",
1642            address,
1643            &[
1644                (500, vec![100, 200, 300, 400, 500]), // First shard ends at 500
1645                (u64::MAX, vec![600, 700, 800]),      // Last shard
1646            ],
1647            &[
1648                // Before first shard, no prev -> NotYetWritten
1649                HistoryQuery {
1650                    block_number: 50,
1651                    lowest_available: None,
1652                    expected: HistoryInfo::NotYetWritten,
1653                },
1654                // Within first shard
1655                HistoryQuery {
1656                    block_number: 150,
1657                    lowest_available: None,
1658                    expected: HistoryInfo::InChangeset(200),
1659                },
1660                // Between shards - prev() should find first shard
1661                HistoryQuery {
1662                    block_number: 550,
1663                    lowest_available: None,
1664                    expected: HistoryInfo::InChangeset(600),
1665                },
1666                // After all entries
1667                HistoryQuery {
1668                    block_number: 900,
1669                    lowest_available: None,
1670                    expected: HistoryInfo::InPlainState,
1671                },
1672            ],
1673        );
1674
1675        // Scenario 3: No history for address
1676        let address_without_history = Address::from([0x43; 20]);
1677        run_account_history_scenario(
1678            "no_history",
1679            address_without_history,
1680            &[], // No shards for this address
1681            &[HistoryQuery {
1682                block_number: 150,
1683                lowest_available: None,
1684                expected: HistoryInfo::NotYetWritten,
1685            }],
1686        );
1687
1688        // Scenario 4: Query at pruning boundary
1689        // We test block >= lowest_available because state queries reject blocks below the
1690        // pruning boundary before doing the lookup. The RocksDB implementation doesn't have
1691        // this check at the same level.
1692        // This tests that when pruning IS available, both backends agree.
1693        run_account_history_scenario(
1694            "with_pruning_boundary",
1695            address,
1696            &[(u64::MAX, vec![100, 200, 300])],
1697            &[
1698                // At pruning boundary -> InChangeset(first entry after block)
1699                HistoryQuery {
1700                    block_number: 100,
1701                    lowest_available: Some(100),
1702                    expected: HistoryInfo::InChangeset(100),
1703                },
1704                // After pruning boundary, between entries
1705                HistoryQuery {
1706                    block_number: 150,
1707                    lowest_available: Some(100),
1708                    expected: HistoryInfo::InChangeset(200),
1709                },
1710            ],
1711        );
1712    }
1713
1714    /// Tests storage history lookups across both MDBX and `RocksDB` backends.
1715    #[test]
1716    fn test_storage_history_info_both_backends() {
1717        let address = Address::from([0x42; 20]);
1718        let storage_key = B256::from([0x01; 32]);
1719        let other_storage_key = B256::from([0x02; 32]);
1720
1721        // Single shard with blocks [100, 200, 300]
1722        run_storage_history_scenario(
1723            "storage_single_shard",
1724            address,
1725            storage_key,
1726            &[(u64::MAX, vec![100, 200, 300])],
1727            &[
1728                // Before first entry -> NotYetWritten
1729                HistoryQuery {
1730                    block_number: 50,
1731                    lowest_available: None,
1732                    expected: HistoryInfo::NotYetWritten,
1733                },
1734                // Between entries -> InChangeset(next_write)
1735                HistoryQuery {
1736                    block_number: 150,
1737                    lowest_available: None,
1738                    expected: HistoryInfo::InChangeset(200),
1739                },
1740                // After last entry -> InPlainState
1741                HistoryQuery {
1742                    block_number: 500,
1743                    lowest_available: None,
1744                    expected: HistoryInfo::InPlainState,
1745                },
1746            ],
1747        );
1748
1749        // No history for different storage key
1750        run_storage_history_scenario(
1751            "storage_no_history",
1752            address,
1753            other_storage_key,
1754            &[], // No shards for this storage key
1755            &[HistoryQuery {
1756                block_number: 150,
1757                lowest_available: None,
1758                expected: HistoryInfo::NotYetWritten,
1759            }],
1760        );
1761    }
1762
1763    /// Test that `RocksDB` batches created via `EitherWriter` are only made visible when
1764    /// `provider.commit()` is called, not when the writer is dropped.
1765    #[test]
1766    fn test_rocksdb_commits_at_provider_level() {
1767        let factory = create_test_provider_factory();
1768
1769        // Enable RocksDB for transaction hash numbers
1770        factory.set_storage_settings_cache(StorageSettings::v2());
1771
1772        let hash1 = B256::from([1u8; 32]);
1773        let hash2 = B256::from([2u8; 32]);
1774        let tx_num1 = 100u64;
1775        let tx_num2 = 200u64;
1776
1777        // Get the RocksDB batch from the provider
1778        let rocksdb = factory.rocksdb_provider();
1779        let batch = rocksdb.batch();
1780
1781        // Create provider and EitherWriter
1782        let provider = factory.database_provider_rw().unwrap();
1783        let mut writer = EitherWriter::new_transaction_hash_numbers(&provider, batch).unwrap();
1784
1785        // Write transaction hash numbers (append_only=false since we're using RocksDB)
1786        writer.put_transaction_hash_number(hash1, tx_num1, false).unwrap();
1787        writer.put_transaction_hash_number(hash2, tx_num2, false).unwrap();
1788
1789        // Extract the raw batch from the writer and register it with the provider
1790        let raw_batch = writer.into_raw_rocksdb_batch();
1791        if let Some(batch) = raw_batch {
1792            provider.set_pending_rocksdb_batch(batch);
1793        }
1794
1795        // Data should NOT be visible yet (batch not committed)
1796        let rocksdb = factory.rocksdb_provider();
1797        assert_eq!(
1798            rocksdb.get::<tables::TransactionHashNumbers>(hash1).unwrap(),
1799            None,
1800            "Data should not be visible before provider.commit()"
1801        );
1802
1803        // Commit the provider - this should commit both MDBX and RocksDB
1804        provider.commit().unwrap();
1805
1806        // Now data should be visible in RocksDB
1807        let rocksdb = factory.rocksdb_provider();
1808        assert_eq!(
1809            rocksdb.get::<tables::TransactionHashNumbers>(hash1).unwrap(),
1810            Some(tx_num1),
1811            "Data should be visible after provider.commit()"
1812        );
1813        assert_eq!(
1814            rocksdb.get::<tables::TransactionHashNumbers>(hash2).unwrap(),
1815            Some(tx_num2),
1816            "Data should be visible after provider.commit()"
1817        );
1818    }
1819
1820    /// Test that `EitherReader::new_accounts_history` panics when settings require
1821    /// `RocksDB` but no snapshot is given (`None`). This is an invariant violation that
1822    /// indicates a bug - `with_rocksdb_snapshot` should always provide a snapshot when needed.
1823    #[test]
1824    #[should_panic(expected = "account_history_in_rocksdb requires rocksdb snapshot")]
1825    fn test_settings_mismatch_panics() {
1826        let factory = create_test_provider_factory();
1827
1828        factory.set_storage_settings_cache(StorageSettings::v2());
1829
1830        let provider = factory.database_provider_ro().unwrap();
1831        let _ = EitherReader::<(), ()>::new_accounts_history(&provider, None);
1832    }
1833}