Skip to main content

reth_provider/providers/rocksdb/
provider.rs

1use super::metrics::{RocksDBMetrics, RocksDBOperation, ROCKSDB_TABLES};
2use crate::providers::{compute_history_rank, needs_prev_shard_check, HistoryInfo};
3use alloy_consensus::transaction::TxHashRef;
4use alloy_primitives::{
5    map::{AddressMap, HashMap},
6    Address, BlockNumber, TxNumber, B256,
7};
8use itertools::Itertools;
9use metrics::Label;
10use parking_lot::Mutex;
11use rayon::prelude::*;
12use reth_chain_state::ExecutedBlock;
13use reth_db_api::{
14    database_metrics::DatabaseMetrics,
15    models::{
16        sharded_key::NUM_OF_INDICES_IN_SHARD, storage_sharded_key::StorageShardedKey, ShardedKey,
17        StorageSettings,
18    },
19    table::{Compress, Decode, Decompress, Encode, Table},
20    tables, BlockNumberList, DatabaseError,
21};
22use reth_primitives_traits::{BlockBody as _, FastInstant as Instant};
23use reth_prune_types::PruneMode;
24use reth_storage_errors::{
25    db::{DatabaseErrorInfo, DatabaseWriteError, DatabaseWriteOperation, LogLevel},
26    provider::{ProviderError, ProviderResult},
27};
28use rocksdb::{
29    statistics::StatsLevel, BlockBasedOptions, Cache, ColumnFamilyDescriptor, CompactionPri,
30    DBCompressionType, DBRawIteratorWithThreadMode, IteratorMode, OptimisticTransactionDB,
31    OptimisticTransactionOptions, Options, ReadOptions, SnapshotWithThreadMode, Transaction,
32    WriteBatchWithTransaction, WriteBufferManager, WriteOptions, DB, DEFAULT_COLUMN_FAMILY_NAME,
33};
34use std::{
35    collections::BTreeMap,
36    fmt,
37    path::{Path, PathBuf},
38    sync::Arc,
39};
40use tracing::instrument;
41
42/// Returns [`WriteOptions`] with WAL sync enabled for crash durability.
43fn synced_write_options() -> WriteOptions {
44    let mut opts = WriteOptions::default();
45    opts.set_sync(true);
46    opts
47}
48
49/// Pending `RocksDB` batches type alias.
50pub(crate) type PendingRocksDBBatches = Arc<Mutex<Vec<WriteBatchWithTransaction<true>>>>;
51
52/// Raw key-value result from a `RocksDB` iterator.
53type RawKVResult = Result<(Box<[u8]>, Box<[u8]>), rocksdb::Error>;
54
55/// Statistics for a single `RocksDB` table (column family).
56#[derive(Debug, Clone)]
57pub struct RocksDBTableStats {
58    /// Size of SST files on disk in bytes.
59    pub sst_size_bytes: u64,
60    /// Size of memtables in memory in bytes.
61    pub memtable_size_bytes: u64,
62    /// Name of the table/column family.
63    pub name: String,
64    /// Estimated number of keys in the table.
65    pub estimated_num_keys: u64,
66    /// Estimated size of live data in bytes (SST files + memtables).
67    pub estimated_size_bytes: u64,
68    /// Estimated bytes pending compaction (reclaimable space).
69    pub pending_compaction_bytes: u64,
70}
71
72/// Database-level statistics for `RocksDB`.
73///
74/// Contains both per-table statistics and DB-level metrics like WAL size.
75#[derive(Debug, Clone)]
76pub struct RocksDBStats {
77    /// Statistics for each table (column family).
78    pub tables: Vec<RocksDBTableStats>,
79    /// Total size of WAL (Write-Ahead Log) files in bytes.
80    ///
81    /// WAL is shared across all tables and not included in per-table metrics.
82    pub wal_size_bytes: u64,
83}
84
85/// Context for `RocksDB` block writes.
86#[derive(Clone)]
87pub(crate) struct RocksDBWriteCtx {
88    /// The first block number being written.
89    pub first_block_number: BlockNumber,
90    /// The prune mode for transaction lookup, if any.
91    pub prune_tx_lookup: Option<PruneMode>,
92    /// Storage settings determining what goes to `RocksDB`.
93    pub storage_settings: StorageSettings,
94    /// Pending batches to push to after writing.
95    pub pending_batches: PendingRocksDBBatches,
96}
97
98impl fmt::Debug for RocksDBWriteCtx {
99    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
100        f.debug_struct("RocksDBWriteCtx")
101            .field("first_block_number", &self.first_block_number)
102            .field("prune_tx_lookup", &self.prune_tx_lookup)
103            .field("storage_settings", &self.storage_settings)
104            .field("pending_batches", &"<pending batches>")
105            .finish()
106    }
107}
108
109/// Default cache size for `RocksDB` block cache (128 MB).
110const DEFAULT_CACHE_SIZE: usize = 128 << 20;
111
112/// Default block size for `RocksDB` tables (16 KB).
113const DEFAULT_BLOCK_SIZE: usize = 16 * 1024;
114
115/// Default max background jobs for `RocksDB` compaction and flushing.
116const DEFAULT_MAX_BACKGROUND_JOBS: i32 = 6;
117
118/// Max open file descriptors for `RocksDB` when the process limit is low.
119///
120/// Caps the number of SST file handles `RocksDB` keeps open simultaneously.
121/// Set to 512 to stay within the common default OS `ulimit -n` of 1024,
122/// leaving headroom for MDBX, static files, and other I/O.
123const LIMITED_MAX_OPEN_FILES: i32 = 512;
124
125/// Keeps all `RocksDB` files open and avoids table cache lookups.
126const KEEP_ALL_FILES_OPEN: i32 = -1;
127
128/// Minimum process file descriptor limit for keeping all `RocksDB` files open.
129///
130/// Mature archive databases can use tens of thousands of file descriptors. This threshold keeps
131/// unlimited mode for hosts configured for that workload while retaining a bounded table cache on
132/// systems with low limits.
133const HIGH_FILE_DESCRIPTOR_LIMIT: u64 = 128 * 1024;
134
135/// Default bytes per sync for `RocksDB` WAL writes (1 MB).
136const DEFAULT_BYTES_PER_SYNC: u64 = 1_048_576;
137
138/// Default write buffer size for `RocksDB` memtables (128 MB).
139///
140/// Larger memtables reduce flush frequency during burst writes, providing more consistent
141/// tail latency. Benchmarks showed 128 MB reduces p99 latency variance by ~80% compared
142/// to 64 MB default, with negligible impact on mean throughput.
143const DEFAULT_WRITE_BUFFER_SIZE: usize = 128 << 20;
144
145/// Default total `RocksDB` memtable memory budget across column families (4 GiB).
146///
147/// This is a soft limit; with write stalls enabled, `RocksDB` waits for flushes once
148/// memtable arena usage exceeds the budget.
149const DEFAULT_WRITE_BUFFER_MANAGER_SIZE: usize = 4 * 1024 * 1024 * 1024;
150
151/// Default buffer capacity for compression in batches.
152/// 4 KiB matches common block/page sizes and comfortably holds typical history values,
153/// reducing the first few reallocations without over-allocating.
154const DEFAULT_COMPRESS_BUF_CAPACITY: usize = 4096;
155
156/// Default auto-commit threshold for batch writes (512 MiB).
157///
158/// When a batch exceeds this size, it is automatically committed to prevent OOM
159/// during large bulk writes. Keep this below the `RocksDB` write buffer manager
160/// budget so stalls can recover without waiting on a single large flush.
161/// The consistency check on startup heals any crash that occurs between auto-commits.
162const DEFAULT_AUTO_COMMIT_THRESHOLD: usize = 512 * 1024 * 1024;
163
164/// Minimum BAL value size stored in `BlobDB` files.
165///
166/// Smaller BALs stay inline. Larger payloads avoid regular LSM value compaction.
167const DEFAULT_BAL_MIN_BLOB_SIZE: u64 = 4 * 1024;
168
169/// Target BAL blob file size.
170const DEFAULT_BAL_BLOB_FILE_SIZE: u64 = 256 * 1024 * 1024;
171
172/// Builder for [`RocksDBProvider`].
173pub struct RocksDBBuilder {
174    path: PathBuf,
175    column_families: Vec<String>,
176    enable_metrics: bool,
177    enable_statistics: bool,
178    log_level: rocksdb::LogLevel,
179    block_cache: Cache,
180    read_only: bool,
181}
182
183impl fmt::Debug for RocksDBBuilder {
184    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
185        f.debug_struct("RocksDBBuilder")
186            .field("path", &self.path)
187            .field("column_families", &self.column_families)
188            .field("enable_metrics", &self.enable_metrics)
189            .finish()
190    }
191}
192
193impl RocksDBBuilder {
194    /// Creates a new builder with optimized default options.
195    pub fn new(path: impl AsRef<Path>) -> Self {
196        let cache = Cache::new_lru_cache(DEFAULT_CACHE_SIZE);
197        Self {
198            path: path.as_ref().to_path_buf(),
199            column_families: Vec::new(),
200            enable_metrics: false,
201            enable_statistics: false,
202            log_level: rocksdb::LogLevel::Info,
203            block_cache: cache,
204            read_only: false,
205        }
206    }
207
208    /// Creates default table options with shared block cache.
209    fn default_table_options(cache: &Cache) -> BlockBasedOptions {
210        let mut table_options = BlockBasedOptions::default();
211        table_options.set_block_size(DEFAULT_BLOCK_SIZE);
212        table_options.set_cache_index_and_filter_blocks(true);
213        table_options.set_pin_l0_filter_and_index_blocks_in_cache(true);
214        // Shared block cache for all column families.
215        table_options.set_block_cache(cache);
216        table_options
217    }
218
219    /// Creates optimized `RocksDB` options per `RocksDB` wiki recommendations.
220    fn default_options(
221        log_level: rocksdb::LogLevel,
222        cache: &Cache,
223        enable_statistics: bool,
224    ) -> Options {
225        // Follow recommend tuning guide from RocksDB wiki, see https://github.com/facebook/rocksdb/wiki/Setup-Options-and-Basic-Tuning
226        let table_options = Self::default_table_options(cache);
227
228        let mut options = Options::default();
229        options.set_block_based_table_factory(&table_options);
230        options.create_if_missing(true);
231        options.create_missing_column_families(true);
232        options.set_max_background_jobs(DEFAULT_MAX_BACKGROUND_JOBS);
233        options.set_bytes_per_sync(DEFAULT_BYTES_PER_SYNC);
234        let write_buffer_manager =
235            WriteBufferManager::new_write_buffer_manager(DEFAULT_WRITE_BUFFER_MANAGER_SIZE, true);
236        options.set_write_buffer_manager(&write_buffer_manager);
237
238        options.set_bottommost_compression_type(DBCompressionType::Zstd);
239        options.set_bottommost_zstd_max_train_bytes(0, true);
240        options.set_compression_type(DBCompressionType::Lz4);
241        options.set_compaction_pri(CompactionPri::MinOverlappingRatio);
242
243        options.set_log_level(log_level);
244
245        options.set_max_open_files(select_max_open_files());
246
247        // Delete obsolete WAL files immediately after all column families have flushed.
248        // Both set to 0 means "delete ASAP, no archival".
249        options.set_wal_ttl_seconds(0);
250        options.set_wal_size_limit_mb(0);
251
252        // Statistics can view from RocksDB log file
253        if enable_statistics {
254            options.enable_statistics();
255            // Nothing reads the statistics programmatically, they only end up in the periodic
256            // LOG dump, so collect the counters but not the timer histograms: at the default
257            // level every `Get`, `Seek` and write wraps itself in a `StopWatch` that reads the
258            // clock twice, while the tickers are plain atomic increments.
259            //
260            // Note: the discriminants of `rocksdb::StatsLevel` are one higher than the levels the
261            // C API defines, so this variant selects "except timers". Should that ever be
262            // corrected upstream it selects "except histogram or timers", which also keeps the
263            // tickers and drops the timers.
264            options.set_statistics_level(StatsLevel::ExceptHistogramOrTimers);
265        }
266
267        options
268    }
269
270    /// Creates optimized column family options.
271    fn default_column_family_options(cache: &Cache) -> Options {
272        // Follow recommend tuning guide from RocksDB wiki, see https://github.com/facebook/rocksdb/wiki/Setup-Options-and-Basic-Tuning
273        let table_options = Self::default_table_options(cache);
274
275        let mut cf_options = Options::default();
276        cf_options.set_block_based_table_factory(&table_options);
277        cf_options.set_level_compaction_dynamic_level_bytes(true);
278        // Recommend to use Zstd for bottommost compression and Lz4 for other levels, see https://github.com/facebook/rocksdb/wiki/Compression#configuration
279        cf_options.set_compression_type(DBCompressionType::Lz4);
280        cf_options.set_bottommost_compression_type(DBCompressionType::Zstd);
281        // Only use Zstd compression, disable dictionary training
282        cf_options.set_bottommost_zstd_max_train_bytes(0, true);
283        cf_options.set_write_buffer_size(DEFAULT_WRITE_BUFFER_SIZE);
284
285        cf_options
286    }
287
288    /// Creates column family options for block access list payloads.
289    fn block_access_lists_column_family_options(cache: &Cache) -> Options {
290        let mut cf_options = Self::default_column_family_options(cache);
291        cf_options.set_enable_blob_files(true);
292        cf_options.set_min_blob_size(DEFAULT_BAL_MIN_BLOB_SIZE);
293        cf_options.set_blob_file_size(DEFAULT_BAL_BLOB_FILE_SIZE);
294        cf_options.set_blob_compression_type(DBCompressionType::Lz4);
295        cf_options
296    }
297
298    /// Creates optimized column family options for `TransactionHashNumbers`.
299    ///
300    /// This table stores `B256 -> TxNumber` mappings where:
301    /// - Keys are incompressible 32-byte hashes (compression wastes CPU for zero benefit)
302    /// - Values are varint-encoded `u64` (a few bytes - too small to benefit from compression)
303    /// - Every lookup expects a hit (bloom filters only help when checking non-existent keys)
304    fn tx_hash_numbers_column_family_options(cache: &Cache) -> Options {
305        let mut table_options = BlockBasedOptions::default();
306        table_options.set_block_size(DEFAULT_BLOCK_SIZE);
307        table_options.set_cache_index_and_filter_blocks(true);
308        table_options.set_pin_l0_filter_and_index_blocks_in_cache(true);
309        table_options.set_block_cache(cache);
310        // Disable bloom filter: every lookup expects a hit, so bloom filters provide no benefit
311        // and waste memory
312
313        let mut cf_options = Options::default();
314        cf_options.set_block_based_table_factory(&table_options);
315        cf_options.set_level_compaction_dynamic_level_bytes(true);
316        // Disable compression: B256 keys are incompressible hashes, TxNumber values are
317        // varint-encoded u64 (a few bytes). Compression wastes CPU cycles for zero space savings.
318        cf_options.set_compression_type(DBCompressionType::None);
319        cf_options.set_bottommost_compression_type(DBCompressionType::None);
320
321        cf_options
322    }
323
324    /// Adds a column family for a specific table type.
325    pub fn with_table<T: Table>(mut self) -> Self {
326        self.column_families.push(T::NAME.to_string());
327        self
328    }
329
330    /// Registers the default tables used by reth for `RocksDB` storage.
331    ///
332    /// This registers:
333    /// - [`tables::TransactionHashNumbers`] - Transaction hash to number mapping
334    /// - [`tables::AccountsHistory`] - Account history index
335    /// - [`tables::StoragesHistory`] - Storage history index
336    /// - [`tables::BlockAccessLists`] - Block access list payloads
337    /// - [`tables::BlockAccessListBlockNumbers`] - Block access list hash index
338    pub fn with_default_tables(self) -> Self {
339        self.with_table::<tables::TransactionHashNumbers>()
340            .with_table::<tables::AccountsHistory>()
341            .with_table::<tables::StoragesHistory>()
342            .with_table::<tables::BlockAccessLists>()
343            .with_table::<tables::BlockAccessListBlockNumbers>()
344    }
345
346    /// Enables metrics.
347    pub const fn with_metrics(mut self) -> Self {
348        self.enable_metrics = true;
349        self
350    }
351
352    /// Enables `RocksDB` internal statistics collection.
353    pub const fn with_statistics(mut self) -> Self {
354        self.enable_statistics = true;
355        self
356    }
357
358    /// Sets the log level from `DatabaseArgs` configuration.
359    pub const fn with_database_log_level(mut self, log_level: Option<LogLevel>) -> Self {
360        if let Some(level) = log_level {
361            self.log_level = convert_log_level(level);
362        }
363        self
364    }
365
366    /// Sets a custom block cache size.
367    pub fn with_block_cache_size(mut self, capacity_bytes: usize) -> Self {
368        self.block_cache = Cache::new_lru_cache(capacity_bytes);
369        self
370    }
371
372    /// Sets a custom block cache size if provided, otherwise keeps the current cache.
373    pub fn with_block_cache_size_opt(self, capacity_bytes: Option<usize>) -> Self {
374        if let Some(capacity_bytes) = capacity_bytes {
375            self.with_block_cache_size(capacity_bytes)
376        } else {
377            self
378        }
379    }
380
381    /// Sets read-only mode.
382    ///
383    /// Opens the database as a secondary instance, which supports catching up
384    /// with the primary via [`RocksDBProvider::try_catch_up_with_primary`].
385    /// A temporary directory is created automatically for the secondary's LOG files.
386    ///
387    /// Note: Write operations on a read-only provider will panic at runtime.
388    pub const fn with_read_only(mut self, read_only: bool) -> Self {
389        self.read_only = read_only;
390        self
391    }
392
393    /// Builds the [`RocksDBProvider`].
394    pub fn build(self) -> ProviderResult<RocksDBProvider> {
395        let options =
396            Self::default_options(self.log_level, &self.block_cache, self.enable_statistics);
397
398        let mut cf_descriptors: Vec<ColumnFamilyDescriptor> = self
399            .column_families
400            .iter()
401            .map(|name| {
402                let cf_options = if name == tables::TransactionHashNumbers::NAME {
403                    Self::tx_hash_numbers_column_family_options(&self.block_cache)
404                } else if name == tables::BlockAccessLists::NAME {
405                    Self::block_access_lists_column_family_options(&self.block_cache)
406                } else {
407                    Self::default_column_family_options(&self.block_cache)
408                };
409                ColumnFamilyDescriptor::new(name.clone(), cf_options)
410            })
411            .collect();
412
413        // RocksDB requires every existing column family to be opened. Preserve column families
414        // unknown to this configuration so databases remain openable after a downgrade.
415        if RocksDBProvider::exists(&self.path) {
416            let existing_column_families = DB::list_cf(&options, &self.path).map_err(|e| {
417                ProviderError::Database(DatabaseError::Open(DatabaseErrorInfo {
418                    message: e.to_string().into(),
419                    code: -1,
420                }))
421            })?;
422            if self.read_only {
423                // Legacy databases have no BAL tables; secondary opens cannot create them.
424                // Keep existing BAL tables; only omit those absent from disk.
425                cf_descriptors.retain(|cf| {
426                    !matches!(
427                        cf.name(),
428                        tables::BlockAccessLists::NAME | tables::BlockAccessListBlockNumbers::NAME
429                    ) || existing_column_families.iter().any(|name| name == cf.name())
430                });
431            }
432            let unknown_column_families: Vec<String> = existing_column_families
433                .into_iter()
434                .filter(|name| {
435                    name != DEFAULT_COLUMN_FAMILY_NAME && !self.column_families.contains(name)
436                })
437                .collect();
438            if !unknown_column_families.is_empty() {
439                tracing::debug!(
440                    target: "providers::rocksdb",
441                    column_families = ?unknown_column_families,
442                    "Preserving unknown column families"
443                );
444                cf_descriptors.extend(unknown_column_families.into_iter().map(|name| {
445                    ColumnFamilyDescriptor::new(
446                        name,
447                        Self::default_column_family_options(&self.block_cache),
448                    )
449                }));
450            }
451        }
452
453        let metrics = self.enable_metrics.then(RocksDBMetrics::default);
454
455        if self.read_only {
456            // Open as secondary instance for catch-up capability.
457            // Secondary needs max_open_files = -1 to keep all FDs open.
458            let mut options = options;
459            options.set_max_open_files(KEEP_ALL_FILES_OPEN);
460
461            let secondary_path = self
462                .path
463                .parent()
464                .unwrap_or(&self.path)
465                .join(format!("rocksdb-secondary-tmp-{}", std::process::id()));
466            reth_fs_util::create_dir_all(&secondary_path).map_err(ProviderError::other)?;
467
468            let db = DB::open_cf_descriptors_as_secondary(
469                &options,
470                &self.path,
471                &secondary_path,
472                cf_descriptors,
473            )
474            .map_err(|e| {
475                ProviderError::Database(DatabaseError::Open(DatabaseErrorInfo {
476                    message: e.to_string().into(),
477                    code: -1,
478                }))
479            })?;
480            Ok(RocksDBProvider(Arc::new(RocksDBProviderInner::Secondary {
481                db,
482                metrics,
483                secondary_path,
484            })))
485        } else {
486            // Use OptimisticTransactionDB for MDBX-like transaction semantics (read-your-writes,
487            // rollback) OptimisticTransactionDB uses optimistic concurrency control (conflict
488            // detection at commit) and is backed by DBCommon, giving us access to
489            // cancel_all_background_work for clean shutdown.
490            let db =
491                OptimisticTransactionDB::open_cf_descriptors(&options, &self.path, cf_descriptors)
492                    .map_err(|e| {
493                        ProviderError::Database(DatabaseError::Open(DatabaseErrorInfo {
494                            message: e.to_string().into(),
495                            code: -1,
496                        }))
497                    })?;
498            Ok(RocksDBProvider(Arc::new(RocksDBProviderInner::ReadWrite { db, metrics })))
499        }
500    }
501}
502
503/// Some types don't support compression (eg. B256), and we don't want to be copying them to the
504/// allocated buffer when we can just use their reference.
505macro_rules! compress_to_buf_or_ref {
506    ($buf:expr, $value:expr) => {
507        if let Some(value) = $value.uncompressable_ref() {
508            Some(value)
509        } else {
510            $buf.clear();
511            $value.compress_to_buf(&mut $buf);
512            None
513        }
514    };
515}
516
517/// `RocksDB` provider for auxiliary storage layer beside main database MDBX.
518#[derive(Debug)]
519pub struct RocksDBProvider(Arc<RocksDBProviderInner>);
520
521/// Inner state for `RocksDB` provider.
522enum RocksDBProviderInner {
523    /// Read-write mode using `OptimisticTransactionDB`.
524    ReadWrite {
525        /// `RocksDB` database instance with optimistic transaction support.
526        db: OptimisticTransactionDB,
527        /// Metrics latency & operations.
528        metrics: Option<RocksDBMetrics>,
529    },
530    /// Secondary mode using `DB` opened with `open_cf_descriptors_as_secondary`.
531    /// Supports catching up with the primary via `try_catch_up_with_primary`.
532    /// Does not support snapshots; consistency is guaranteed externally.
533    Secondary {
534        /// Secondary `RocksDB` database instance.
535        db: DB,
536        /// Metrics latency & operations.
537        metrics: Option<RocksDBMetrics>,
538        /// Temporary directory for secondary LOG files, removed on drop.
539        secondary_path: PathBuf,
540    },
541}
542
543impl RocksDBProviderInner {
544    /// Returns the metrics for this provider.
545    const fn metrics(&self) -> Option<&RocksDBMetrics> {
546        match self {
547            Self::ReadWrite { metrics, .. } | Self::Secondary { metrics, .. } => metrics.as_ref(),
548        }
549    }
550
551    /// Returns the read-write database, panicking if in read-only mode.
552    fn db_rw(&self) -> &OptimisticTransactionDB {
553        match self {
554            Self::ReadWrite { db, .. } => db,
555            Self::Secondary { .. } => {
556                panic!("Cannot perform write operation on secondary RocksDB provider")
557            }
558        }
559    }
560
561    /// Gets the column family handle for a table.
562    fn cf_handle<T: Table>(&self) -> Result<&rocksdb::ColumnFamily, DatabaseError> {
563        let cf = match self {
564            Self::ReadWrite { db, .. } => db.cf_handle(T::NAME),
565            Self::Secondary { db, .. } => db.cf_handle(T::NAME),
566        };
567        cf.ok_or_else(|| DatabaseError::Other(format!("Column family '{}' not found", T::NAME)))
568    }
569
570    /// Gets a value from a column family.
571    fn get_cf(
572        &self,
573        cf: &rocksdb::ColumnFamily,
574        key: impl AsRef<[u8]>,
575    ) -> Result<Option<Vec<u8>>, rocksdb::Error> {
576        match self {
577            Self::ReadWrite { db, .. } => db.get_cf(cf, key),
578            Self::Secondary { db, .. } => db.get_cf(cf, key),
579        }
580    }
581
582    /// Puts a value into a column family.
583    fn put_cf(
584        &self,
585        cf: &rocksdb::ColumnFamily,
586        key: impl AsRef<[u8]>,
587        value: impl AsRef<[u8]>,
588    ) -> Result<(), rocksdb::Error> {
589        self.db_rw().put_cf(cf, key, value)
590    }
591
592    /// Deletes a value from a column family.
593    fn delete_cf(
594        &self,
595        cf: &rocksdb::ColumnFamily,
596        key: impl AsRef<[u8]>,
597    ) -> Result<(), rocksdb::Error> {
598        self.db_rw().delete_cf(cf, key)
599    }
600
601    /// Deletes a range of values from a column family.
602    fn delete_range_cf<K: AsRef<[u8]>>(
603        &self,
604        cf: &rocksdb::ColumnFamily,
605        from: K,
606        to: K,
607    ) -> Result<(), rocksdb::Error> {
608        self.db_rw().delete_range_cf(cf, from, to)
609    }
610
611    /// Returns an iterator over a column family.
612    fn iterator_cf(
613        &self,
614        cf: &rocksdb::ColumnFamily,
615        mode: IteratorMode<'_>,
616    ) -> RocksDBIterEnum<'_> {
617        match self {
618            Self::ReadWrite { db, .. } => RocksDBIterEnum::ReadWrite(db.iterator_cf(cf, mode)),
619            Self::Secondary { db, .. } => RocksDBIterEnum::ReadOnly(db.iterator_cf(cf, mode)),
620        }
621    }
622
623    /// Returns a raw iterator over a column family.
624    ///
625    /// Unlike [`Self::iterator_cf`], raw iterators support `seek()` for efficient
626    /// repositioning without creating a new iterator.
627    fn raw_iterator_cf(&self, cf: &rocksdb::ColumnFamily) -> RocksDBRawIterEnum<'_> {
628        match self {
629            Self::ReadWrite { db, .. } => RocksDBRawIterEnum::ReadWrite(db.raw_iterator_cf(cf)),
630            Self::Secondary { db, .. } => RocksDBRawIterEnum::ReadOnly(db.raw_iterator_cf(cf)),
631        }
632    }
633
634    /// Returns a read-only, point-in-time snapshot of the database.
635    fn snapshot(&self) -> RocksReadSnapshotInner<'_> {
636        match self {
637            Self::ReadWrite { db, .. } => RocksReadSnapshotInner::ReadWrite(db.snapshot()),
638            Self::Secondary { db, .. } => RocksReadSnapshotInner::Secondary(db),
639        }
640    }
641
642    /// Returns the path to the database directory.
643    fn path(&self) -> &Path {
644        match self {
645            Self::ReadWrite { db, .. } => db.path(),
646            Self::Secondary { db, .. } => db.path(),
647        }
648    }
649
650    /// Returns the total size of WAL (Write-Ahead Log) files in bytes.
651    ///
652    /// WAL files have a `.log` extension in the `RocksDB` directory.
653    fn wal_size_bytes(&self) -> u64 {
654        let path = self.path();
655
656        match std::fs::read_dir(path) {
657            Ok(entries) => entries
658                .filter_map(|e| e.ok())
659                .filter(|e| e.path().extension().is_some_and(|ext| ext == "log"))
660                .filter_map(|e| e.metadata().ok())
661                .map(|m| m.len())
662                .sum(),
663            Err(_) => 0,
664        }
665    }
666
667    /// Returns statistics for all column families in the database.
668    fn table_stats(&self) -> Vec<RocksDBTableStats> {
669        let mut stats = Vec::new();
670
671        macro_rules! collect_stats {
672            ($db:expr) => {
673                for cf_name in ROCKSDB_TABLES {
674                    if let Some(cf) = $db.cf_handle(cf_name) {
675                        let estimated_num_keys = $db
676                            .property_int_value_cf(cf, rocksdb::properties::ESTIMATE_NUM_KEYS)
677                            .ok()
678                            .flatten()
679                            .unwrap_or(0);
680
681                        // SST files size (on-disk) + memtable size (in-memory)
682                        let sst_size = $db
683                            .property_int_value_cf(cf, rocksdb::properties::LIVE_SST_FILES_SIZE)
684                            .ok()
685                            .flatten()
686                            .unwrap_or(0);
687
688                        let memtable_size = $db
689                            .property_int_value_cf(cf, rocksdb::properties::SIZE_ALL_MEM_TABLES)
690                            .ok()
691                            .flatten()
692                            .unwrap_or(0);
693
694                        let estimated_size_bytes = sst_size + memtable_size;
695
696                        let pending_compaction_bytes = $db
697                            .property_int_value_cf(
698                                cf,
699                                rocksdb::properties::ESTIMATE_PENDING_COMPACTION_BYTES,
700                            )
701                            .ok()
702                            .flatten()
703                            .unwrap_or(0);
704
705                        stats.push(RocksDBTableStats {
706                            sst_size_bytes: sst_size,
707                            memtable_size_bytes: memtable_size,
708                            name: cf_name.to_string(),
709                            estimated_num_keys,
710                            estimated_size_bytes,
711                            pending_compaction_bytes,
712                        });
713                    }
714                }
715            };
716        }
717
718        match self {
719            Self::ReadWrite { db, .. } => collect_stats!(db),
720            Self::Secondary { db, .. } => collect_stats!(db),
721        }
722
723        stats
724    }
725
726    /// Returns database-level statistics including per-table stats and WAL size.
727    fn db_stats(&self) -> RocksDBStats {
728        RocksDBStats { tables: self.table_stats(), wal_size_bytes: self.wal_size_bytes() }
729    }
730}
731
732impl fmt::Debug for RocksDBProviderInner {
733    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
734        match self {
735            Self::ReadWrite { metrics, .. } => f
736                .debug_struct("RocksDBProviderInner::ReadWrite")
737                .field("db", &"<OptimisticTransactionDB>")
738                .field("metrics", metrics)
739                .finish(),
740            Self::Secondary { metrics, .. } => f
741                .debug_struct("RocksDBProviderInner::Secondary")
742                .field("db", &"<DB (secondary)>")
743                .field("metrics", metrics)
744                .finish(),
745        }
746    }
747}
748
749impl Drop for RocksDBProviderInner {
750    fn drop(&mut self) {
751        match self {
752            Self::ReadWrite { db, .. } => {
753                // Flush all memtables if possible. If not, they will be rebuilt from the WAL on
754                // restart
755                if let Err(e) = db.flush_wal(true) {
756                    tracing::warn!(target: "providers::rocksdb", ?e, "Failed to flush WAL on drop");
757                }
758                for cf_name in ROCKSDB_TABLES {
759                    if let Some(cf) = db.cf_handle(cf_name) &&
760                        let Err(e) = db.flush_cf(&cf)
761                    {
762                        tracing::warn!(target: "providers::rocksdb", cf = cf_name, ?e, "Failed to flush CF on drop");
763                    }
764                }
765                db.cancel_all_background_work(true);
766            }
767            Self::Secondary { db, secondary_path, .. } => {
768                db.cancel_all_background_work(true);
769                let _ = std::fs::remove_dir_all(secondary_path);
770            }
771        }
772    }
773}
774
775impl Clone for RocksDBProvider {
776    fn clone(&self) -> Self {
777        Self(self.0.clone())
778    }
779}
780
781impl DatabaseMetrics for RocksDBProvider {
782    fn gauge_metrics(&self) -> Vec<(&'static str, f64, Vec<Label>)> {
783        let mut metrics = Vec::new();
784
785        for stat in self.table_stats() {
786            metrics.push((
787                "rocksdb.table_size",
788                stat.estimated_size_bytes as f64,
789                vec![Label::new("table", stat.name.clone())],
790            ));
791            metrics.push((
792                "rocksdb.table_entries",
793                stat.estimated_num_keys as f64,
794                vec![Label::new("table", stat.name.clone())],
795            ));
796            metrics.push((
797                "rocksdb.pending_compaction_bytes",
798                stat.pending_compaction_bytes as f64,
799                vec![Label::new("table", stat.name.clone())],
800            ));
801            metrics.push((
802                "rocksdb.sst_size",
803                stat.sst_size_bytes as f64,
804                vec![Label::new("table", stat.name.clone())],
805            ));
806            metrics.push((
807                "rocksdb.memtable_size",
808                stat.memtable_size_bytes as f64,
809                vec![Label::new("table", stat.name)],
810            ));
811        }
812
813        // WAL size (DB-level, shared across all tables)
814        metrics.push(("rocksdb.wal_size", self.wal_size_bytes() as f64, vec![]));
815
816        metrics
817    }
818}
819
820impl RocksDBProvider {
821    /// Creates a new `RocksDB` provider.
822    pub fn new(path: impl AsRef<Path>) -> ProviderResult<Self> {
823        RocksDBBuilder::new(path).build()
824    }
825
826    /// Creates a new `RocksDB` provider builder.
827    pub fn builder(path: impl AsRef<Path>) -> RocksDBBuilder {
828        RocksDBBuilder::new(path)
829    }
830
831    /// Returns `true` if a `RocksDB` database exists at the given path.
832    ///
833    /// Checks for the presence of the `CURRENT` file, which `RocksDB` creates
834    /// when initializing a database.
835    pub fn exists(path: impl AsRef<Path>) -> bool {
836        path.as_ref().join("CURRENT").exists()
837    }
838
839    /// Returns `true` if this provider is in read-only mode.
840    pub fn is_read_only(&self) -> bool {
841        matches!(self.0.as_ref(), RocksDBProviderInner::Secondary { .. })
842    }
843
844    /// Tries to catch up with the primary instance by reading new WAL and MANIFEST entries.
845    ///
846    /// This is a no-op for read-write and read-only providers.
847    /// For secondary providers, this incrementally syncs with the primary's latest state.
848    pub fn try_catch_up_with_primary(&self) -> ProviderResult<()> {
849        match self.0.as_ref() {
850            RocksDBProviderInner::Secondary { db, .. } => {
851                db.try_catch_up_with_primary().map_err(|e| {
852                    ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
853                        message: e.to_string().into(),
854                        code: -1,
855                    }))
856                })
857            }
858            _ => Ok(()),
859        }
860    }
861
862    /// Returns a read-only, point-in-time snapshot of the database.
863    ///
864    /// Lighter weight than [`RocksTx`] — no write-conflict tracking, and `Send + Sync`.
865    pub fn snapshot(&self) -> RocksReadSnapshot<'_> {
866        RocksReadSnapshot {
867            accounts_history_iter: Mutex::new(None),
868            storages_history_iter: Mutex::new(None),
869            inner: self.0.snapshot(),
870            provider: &self.0,
871        }
872    }
873
874    /// Returns a read-only, point-in-time snapshot that owns the provider handle it reads through.
875    ///
876    /// Lets a reader keep one snapshot (and the iterators cached in it) alive for its whole
877    /// lifetime, instead of creating a snapshot per lookup.
878    pub fn owned_snapshot(&self) -> OwnedRocksReadSnapshot {
879        OwnedRocksReadSnapshot::new(self)
880    }
881
882    /// Creates a new transaction with MDBX-like semantics (read-your-writes, rollback).
883    ///
884    /// Note: With `OptimisticTransactionDB`, commits may fail if there are conflicts.
885    /// Conflict detection happens at commit time, not at write time.
886    ///
887    /// # Panics
888    /// Panics if the provider is in read-only mode.
889    pub fn tx(&self) -> RocksTx<'_> {
890        let write_options = synced_write_options();
891        let txn_options = OptimisticTransactionOptions::default();
892        let inner = self.0.db_rw().transaction_opt(&write_options, &txn_options);
893        RocksTx { inner, provider: self }
894    }
895
896    /// Creates a new batch for atomic writes.
897    ///
898    /// Use [`Self::write_batch`] for closure-based atomic writes.
899    /// Use this method when the batch needs to be held by [`crate::EitherWriter`].
900    ///
901    /// # Panics
902    /// Panics if the provider is in read-only mode when attempting to commit.
903    pub fn batch(&self) -> RocksDBBatch<'_> {
904        RocksDBBatch {
905            provider: self,
906            inner: WriteBatchWithTransaction::<true>::default(),
907            buf: Vec::with_capacity(DEFAULT_COMPRESS_BUF_CAPACITY),
908            auto_commit_threshold: None,
909        }
910    }
911
912    /// Creates a new batch with auto-commit enabled.
913    ///
914    /// When the batch size exceeds the threshold (4 GiB), the batch is automatically
915    /// committed and reset. This prevents OOM during large bulk writes while maintaining
916    /// crash-safety via the consistency check on startup.
917    pub fn batch_with_auto_commit(&self) -> RocksDBBatch<'_> {
918        RocksDBBatch {
919            provider: self,
920            inner: WriteBatchWithTransaction::<true>::default(),
921            buf: Vec::with_capacity(DEFAULT_COMPRESS_BUF_CAPACITY),
922            auto_commit_threshold: Some(DEFAULT_AUTO_COMMIT_THRESHOLD),
923        }
924    }
925
926    /// Gets the column family handle for a table.
927    fn get_cf_handle<T: Table>(&self) -> Result<&rocksdb::ColumnFamily, DatabaseError> {
928        self.0.cf_handle::<T>()
929    }
930
931    /// Returns whether this provider opened the given table.
932    pub(crate) fn has_table<T: Table>(&self) -> bool {
933        match self.0.as_ref() {
934            RocksDBProviderInner::ReadWrite { db, .. } => db.cf_handle(T::NAME).is_some(),
935            RocksDBProviderInner::Secondary { db, .. } => db.cf_handle(T::NAME).is_some(),
936        }
937    }
938
939    /// Executes a function and records metrics with the given operation and table name.
940    fn execute_with_operation_metric<R>(
941        &self,
942        operation: RocksDBOperation,
943        table: &'static str,
944        f: impl FnOnce(&Self) -> R,
945    ) -> R {
946        let start = self.0.metrics().map(|_| Instant::now());
947        let res = f(self);
948
949        if let (Some(start), Some(metrics)) = (start, self.0.metrics()) {
950            metrics.record_operation(operation, table, start.elapsed());
951        }
952
953        res
954    }
955
956    /// Gets a value from the specified table.
957    pub fn get<T: Table>(&self, key: T::Key) -> ProviderResult<Option<T::Value>> {
958        self.get_encoded::<T>(&key.encode())
959    }
960
961    /// Gets a value from the specified table using pre-encoded key.
962    pub fn get_encoded<T: Table>(
963        &self,
964        key: &<T::Key as Encode>::Encoded,
965    ) -> ProviderResult<Option<T::Value>> {
966        self.execute_with_operation_metric(RocksDBOperation::Get, T::NAME, |this| {
967            let result = this.0.get_cf(this.get_cf_handle::<T>()?, key.as_ref()).map_err(|e| {
968                ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
969                    message: e.to_string().into(),
970                    code: -1,
971                }))
972            })?;
973
974            Ok(result.and_then(|value| T::Value::decompress(&value).ok()))
975        })
976    }
977
978    /// Gets raw bytes from the specified table without decompressing.
979    pub fn get_raw<T: Table>(&self, key: T::Key) -> ProviderResult<Option<Vec<u8>>> {
980        let encoded = key.encode();
981        self.execute_with_operation_metric(RocksDBOperation::Get, T::NAME, |this| {
982            this.0.get_cf(this.get_cf_handle::<T>()?, encoded.as_ref()).map_err(|e| {
983                ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
984                    message: e.to_string().into(),
985                    code: -1,
986                }))
987            })
988        })
989    }
990
991    /// Puts upsert a value into the specified table with the given key.
992    ///
993    /// # Panics
994    /// Panics if the provider is in read-only mode.
995    pub fn put<T: Table>(&self, key: T::Key, value: &T::Value) -> ProviderResult<()> {
996        let encoded_key = key.encode();
997        self.put_encoded::<T>(&encoded_key, value)
998    }
999
1000    /// Puts a value into the specified table using pre-encoded key.
1001    ///
1002    /// # Panics
1003    /// Panics if the provider is in read-only mode.
1004    pub fn put_encoded<T: Table>(
1005        &self,
1006        key: &<T::Key as Encode>::Encoded,
1007        value: &T::Value,
1008    ) -> ProviderResult<()> {
1009        self.execute_with_operation_metric(RocksDBOperation::Put, T::NAME, |this| {
1010            // for simplify the code, we need allocate buf here each time because `RocksDBProvider`
1011            // is thread safe if user want to avoid allocate buf each time, they can use
1012            // write_batch api
1013            let mut buf = Vec::new();
1014            let value_bytes = compress_to_buf_or_ref!(buf, value).unwrap_or(&buf);
1015
1016            this.0.put_cf(this.get_cf_handle::<T>()?, key, value_bytes).map_err(|e| {
1017                ProviderError::Database(DatabaseError::Write(Box::new(DatabaseWriteError {
1018                    info: DatabaseErrorInfo { message: e.to_string().into(), code: -1 },
1019                    operation: DatabaseWriteOperation::PutUpsert,
1020                    table_name: T::NAME,
1021                    key: key.as_ref().to_vec(),
1022                })))
1023            })
1024        })
1025    }
1026
1027    /// Deletes a value from the specified table.
1028    ///
1029    /// # Panics
1030    /// Panics if the provider is in read-only mode.
1031    pub fn delete<T: Table>(&self, key: T::Key) -> ProviderResult<()> {
1032        self.execute_with_operation_metric(RocksDBOperation::Delete, T::NAME, |this| {
1033            this.0.delete_cf(this.get_cf_handle::<T>()?, key.encode().as_ref()).map_err(|e| {
1034                ProviderError::Database(DatabaseError::Delete(DatabaseErrorInfo {
1035                    message: e.to_string().into(),
1036                    code: -1,
1037                }))
1038            })
1039        })
1040    }
1041
1042    /// Clears all entries from the specified table.
1043    ///
1044    /// Uses `delete_range_cf` from empty key to a max key (256 bytes of 0xFF).
1045    /// This end key must exceed the maximum encoded key size for any table.
1046    /// Current max is ~60 bytes (`StorageShardedKey` = 20 + 32 + 8).
1047    pub fn clear<T: Table>(&self) -> ProviderResult<()> {
1048        let cf = self.get_cf_handle::<T>()?;
1049
1050        self.0.delete_range_cf(cf, &[] as &[u8], &[0xFF; 256]).map_err(|e| {
1051            ProviderError::Database(DatabaseError::Delete(DatabaseErrorInfo {
1052                message: e.to_string().into(),
1053                code: -1,
1054            }))
1055        })?;
1056
1057        Ok(())
1058    }
1059
1060    /// Retrieves the first or last entry from a table based on the iterator mode.
1061    fn get_boundary<T: Table>(
1062        &self,
1063        mode: IteratorMode<'_>,
1064    ) -> ProviderResult<Option<(T::Key, T::Value)>> {
1065        self.execute_with_operation_metric(RocksDBOperation::Get, T::NAME, |this| {
1066            let cf = this.get_cf_handle::<T>()?;
1067            let mut iter = this.0.iterator_cf(cf, mode);
1068
1069            match iter.next() {
1070                Some(Ok((key_bytes, value_bytes))) => {
1071                    let key = <T::Key as reth_db_api::table::Decode>::decode(&key_bytes)
1072                        .map_err(|_| ProviderError::Database(DatabaseError::Decode))?;
1073                    let value = T::Value::decompress(&value_bytes)
1074                        .map_err(|_| ProviderError::Database(DatabaseError::Decode))?;
1075                    Ok(Some((key, value)))
1076                }
1077                Some(Err(e)) => {
1078                    Err(ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
1079                        message: e.to_string().into(),
1080                        code: -1,
1081                    })))
1082                }
1083                None => Ok(None),
1084            }
1085        })
1086    }
1087
1088    /// Gets the first (smallest key) entry from the specified table.
1089    #[inline]
1090    pub fn first<T: Table>(&self) -> ProviderResult<Option<(T::Key, T::Value)>> {
1091        self.get_boundary::<T>(IteratorMode::Start)
1092    }
1093
1094    /// Gets the last (largest key) entry from the specified table.
1095    #[inline]
1096    pub fn last<T: Table>(&self) -> ProviderResult<Option<(T::Key, T::Value)>> {
1097        self.get_boundary::<T>(IteratorMode::End)
1098    }
1099
1100    /// Creates an iterator over all entries in the specified table.
1101    ///
1102    /// Returns decoded `(Key, Value)` pairs in key order.
1103    pub fn iter<T: Table>(&self) -> ProviderResult<RocksDBIter<'_, T>> {
1104        let cf = self.get_cf_handle::<T>()?;
1105        let iter = self.0.iterator_cf(cf, IteratorMode::Start);
1106        Ok(RocksDBIter { inner: iter, _marker: std::marker::PhantomData })
1107    }
1108
1109    /// Creates an iterator starting from the given key (inclusive, seek forward).
1110    ///
1111    /// Returns decoded `(Key, Value)` pairs starting from the first key >= `key`.
1112    pub fn iter_from<T: Table>(&self, key: T::Key) -> ProviderResult<RocksDBIter<'_, T>> {
1113        let cf = self.get_cf_handle::<T>()?;
1114        let encoded_key = key.encode();
1115        let iter = self
1116            .0
1117            .iterator_cf(cf, IteratorMode::From(encoded_key.as_ref(), rocksdb::Direction::Forward));
1118        Ok(RocksDBIter { inner: iter, _marker: std::marker::PhantomData })
1119    }
1120
1121    /// Returns statistics for all column families in the database.
1122    ///
1123    /// Returns a vector of (`table_name`, `estimated_keys`, `estimated_size_bytes`) tuples.
1124    pub fn table_stats(&self) -> Vec<RocksDBTableStats> {
1125        self.0.table_stats()
1126    }
1127
1128    /// Returns the total size of WAL (Write-Ahead Log) files in bytes.
1129    ///
1130    /// This scans the `RocksDB` directory for `.log` files and sums their sizes.
1131    /// WAL files can be significant (e.g., 2.7GB observed) and are not included
1132    /// in `table_size`, `sst_size`, or `memtable_size` metrics.
1133    pub fn wal_size_bytes(&self) -> u64 {
1134        self.0.wal_size_bytes()
1135    }
1136
1137    /// Returns database-level statistics including per-table stats and WAL size.
1138    ///
1139    /// This combines [`Self::table_stats`] and [`Self::wal_size_bytes`] into a single struct.
1140    pub fn db_stats(&self) -> RocksDBStats {
1141        self.0.db_stats()
1142    }
1143
1144    /// Flushes pending writes for the specified tables to disk.
1145    ///
1146    /// This performs a flush of:
1147    /// 1. The column family memtables for the specified table names to SST files
1148    /// 2. The Write-Ahead Log (WAL) with sync
1149    ///
1150    /// After this call completes, all data for the specified tables is durably persisted to disk.
1151    ///
1152    /// # Panics
1153    /// Panics if the provider is in read-only mode.
1154    #[instrument(level = "debug", target = "providers::rocksdb", skip_all, fields(tables = ?tables))]
1155    pub fn flush(&self, tables: &[&'static str]) -> ProviderResult<()> {
1156        let db = self.0.db_rw();
1157
1158        for cf_name in tables {
1159            if let Some(cf) = db.cf_handle(cf_name) {
1160                db.flush_cf(&cf).map_err(|e| {
1161                    ProviderError::Database(DatabaseError::Write(Box::new(DatabaseWriteError {
1162                        info: DatabaseErrorInfo { message: e.to_string().into(), code: -1 },
1163                        operation: DatabaseWriteOperation::Flush,
1164                        table_name: cf_name,
1165                        key: Vec::new(),
1166                    })))
1167                })?;
1168            }
1169        }
1170
1171        db.flush_wal(true).map_err(|e| {
1172            ProviderError::Database(DatabaseError::Write(Box::new(DatabaseWriteError {
1173                info: DatabaseErrorInfo { message: e.to_string().into(), code: -1 },
1174                operation: DatabaseWriteOperation::Flush,
1175                table_name: "WAL",
1176                key: Vec::new(),
1177            })))
1178        })?;
1179
1180        Ok(())
1181    }
1182
1183    /// Flushes and compacts all tables in `RocksDB`.
1184    ///
1185    /// This:
1186    /// 1. Flushes all column family memtables to SST files
1187    /// 2. Flushes the Write-Ahead Log (WAL) with sync
1188    /// 3. Triggers manual compaction on all column families to reclaim disk space
1189    ///
1190    /// Use this after large delete operations (like pruning) to reclaim disk space.
1191    ///
1192    /// # Panics
1193    /// Panics if the provider is in read-only mode.
1194    #[instrument(level = "debug", target = "providers::rocksdb", skip_all)]
1195    pub fn flush_and_compact(&self) -> ProviderResult<()> {
1196        self.flush(ROCKSDB_TABLES)?;
1197
1198        let db = self.0.db_rw();
1199
1200        for cf_name in ROCKSDB_TABLES {
1201            if let Some(cf) = db.cf_handle(cf_name) {
1202                db.compact_range_cf(&cf, None::<&[u8]>, None::<&[u8]>);
1203            }
1204        }
1205
1206        Ok(())
1207    }
1208
1209    /// Creates a raw iterator over all entries in the specified table.
1210    ///
1211    /// Returns raw `(key_bytes, value_bytes)` pairs without decoding.
1212    pub fn raw_iter<T: Table>(&self) -> ProviderResult<RocksDBRawIter<'_>> {
1213        let cf = self.get_cf_handle::<T>()?;
1214        let iter = self.0.iterator_cf(cf, IteratorMode::Start);
1215        Ok(RocksDBRawIter { inner: iter })
1216    }
1217
1218    /// Creates a raw key iterator positioned at `key`.
1219    pub(crate) fn raw_key_iter_from<T: Table>(
1220        &self,
1221        key: T::Key,
1222    ) -> ProviderResult<RocksDBRawKeyIter<'_>> {
1223        let cf = self.get_cf_handle::<T>()?;
1224        let encoded_key = key.encode();
1225        let mut iter = self.0.raw_iterator_cf(cf);
1226        iter.seek(encoded_key.as_ref());
1227        Ok(RocksDBRawKeyIter { inner: iter })
1228    }
1229
1230    /// Returns all account history shards for the given address in ascending key order.
1231    ///
1232    /// This is used for unwind operations where we need to scan all shards for an address
1233    /// and potentially delete or truncate them.
1234    pub fn account_history_shards(
1235        &self,
1236        address: Address,
1237    ) -> ProviderResult<Vec<(ShardedKey<Address>, BlockNumberList)>> {
1238        // Get the column family handle for the AccountsHistory table.
1239        let cf = self.get_cf_handle::<tables::AccountsHistory>()?;
1240
1241        // Build a seek key starting at the first shard (highest_block_number = 0) for this address.
1242        // ShardedKey is (address, highest_block_number) so this positions us at the beginning.
1243        let start_key = ShardedKey::new(address, 0u64);
1244        let start_bytes = start_key.encode();
1245
1246        // Create a forward iterator starting from our seek position.
1247        let iter = self
1248            .0
1249            .iterator_cf(cf, IteratorMode::From(start_bytes.as_ref(), rocksdb::Direction::Forward));
1250
1251        let mut result = Vec::new();
1252        for item in iter {
1253            match item {
1254                Ok((key_bytes, value_bytes)) => {
1255                    // Decode the sharded key to check if we're still on the same address.
1256                    let key = ShardedKey::<Address>::decode(&key_bytes)
1257                        .map_err(|_| ProviderError::Database(DatabaseError::Decode))?;
1258
1259                    // Stop when we reach a different address (keys are sorted by address first).
1260                    if key.key != address {
1261                        break;
1262                    }
1263
1264                    // Decompress the block number list stored in this shard.
1265                    let value = BlockNumberList::decompress(&value_bytes)
1266                        .map_err(|_| ProviderError::Database(DatabaseError::Decode))?;
1267
1268                    result.push((key, value));
1269                }
1270                Err(e) => {
1271                    return Err(ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
1272                        message: e.to_string().into(),
1273                        code: -1,
1274                    })));
1275                }
1276            }
1277        }
1278
1279        Ok(result)
1280    }
1281
1282    /// Returns all storage history shards for the given `(address, storage_key)` pair.
1283    ///
1284    /// Iterates through all shards in ascending `highest_block_number` order until
1285    /// a different `(address, storage_key)` is encountered.
1286    pub fn storage_history_shards(
1287        &self,
1288        address: Address,
1289        storage_key: B256,
1290    ) -> ProviderResult<Vec<(StorageShardedKey, BlockNumberList)>> {
1291        let cf = self.get_cf_handle::<tables::StoragesHistory>()?;
1292
1293        let start_key = StorageShardedKey::new(address, storage_key, 0u64);
1294        let start_bytes = start_key.encode();
1295
1296        let iter = self
1297            .0
1298            .iterator_cf(cf, IteratorMode::From(start_bytes.as_ref(), rocksdb::Direction::Forward));
1299
1300        let mut result = Vec::new();
1301        for item in iter {
1302            match item {
1303                Ok((key_bytes, value_bytes)) => {
1304                    let key = StorageShardedKey::decode(&key_bytes)
1305                        .map_err(|_| ProviderError::Database(DatabaseError::Decode))?;
1306
1307                    if key.address != address || key.sharded_key.key != storage_key {
1308                        break;
1309                    }
1310
1311                    let value = BlockNumberList::decompress(&value_bytes)
1312                        .map_err(|_| ProviderError::Database(DatabaseError::Decode))?;
1313
1314                    result.push((key, value));
1315                }
1316                Err(e) => {
1317                    return Err(ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
1318                        message: e.to_string().into(),
1319                        code: -1,
1320                    })));
1321                }
1322            }
1323        }
1324
1325        Ok(result)
1326    }
1327
1328    /// Unwinds account history indices for the given `(address, block_number)` pairs.
1329    ///
1330    /// Groups addresses by their minimum block number and calls the appropriate unwind
1331    /// operations. For each address, keeps only blocks less than the minimum block
1332    /// (i.e., removes the minimum block and all higher blocks).
1333    ///
1334    /// Returns a `WriteBatchWithTransaction` that can be committed later.
1335    #[instrument(level = "debug", target = "providers::rocksdb", skip_all)]
1336    pub fn unwind_account_history_indices(
1337        &self,
1338        last_indices: &[(Address, BlockNumber)],
1339    ) -> ProviderResult<WriteBatchWithTransaction<true>> {
1340        let mut address_min_block: AddressMap<BlockNumber> =
1341            AddressMap::with_capacity_and_hasher(last_indices.len(), Default::default());
1342        for &(address, block_number) in last_indices {
1343            address_min_block
1344                .entry(address)
1345                .and_modify(|min| *min = (*min).min(block_number))
1346                .or_insert(block_number);
1347        }
1348
1349        let mut batch = self.batch();
1350        for (address, min_block) in address_min_block {
1351            match min_block.checked_sub(1) {
1352                Some(keep_to) => batch.unwind_account_history_to(address, keep_to)?,
1353                None => batch.clear_account_history(address)?,
1354            }
1355        }
1356
1357        Ok(batch.into_inner())
1358    }
1359
1360    /// Unwinds storage history indices for the given `(address, storage_key, block_number)` tuples.
1361    ///
1362    /// Groups by `(address, storage_key)` and finds the minimum block number for each.
1363    /// For each key, keeps only blocks less than the minimum block
1364    /// (i.e., removes the minimum block and all higher blocks).
1365    ///
1366    /// Returns a `WriteBatchWithTransaction` that can be committed later.
1367    pub fn unwind_storage_history_indices(
1368        &self,
1369        storage_changesets: &[(Address, B256, BlockNumber)],
1370    ) -> ProviderResult<WriteBatchWithTransaction<true>> {
1371        let mut key_min_block: HashMap<(Address, B256), BlockNumber> =
1372            HashMap::with_capacity_and_hasher(storage_changesets.len(), Default::default());
1373        for &(address, storage_key, block_number) in storage_changesets {
1374            key_min_block
1375                .entry((address, storage_key))
1376                .and_modify(|min| *min = (*min).min(block_number))
1377                .or_insert(block_number);
1378        }
1379
1380        let mut batch = self.batch();
1381        for ((address, storage_key), min_block) in key_min_block {
1382            match min_block.checked_sub(1) {
1383                Some(keep_to) => batch.unwind_storage_history_to(address, storage_key, keep_to)?,
1384                None => batch.clear_storage_history(address, storage_key)?,
1385            }
1386        }
1387
1388        Ok(batch.into_inner())
1389    }
1390
1391    /// Writes a batch of operations atomically.
1392    #[instrument(level = "debug", target = "providers::rocksdb", skip_all)]
1393    pub fn write_batch<F>(&self, f: F) -> ProviderResult<()>
1394    where
1395        F: FnOnce(&mut RocksDBBatch<'_>) -> ProviderResult<()>,
1396    {
1397        self.execute_with_operation_metric(RocksDBOperation::BatchWrite, "Batch", |this| {
1398            let mut batch_handle = this.batch();
1399            f(&mut batch_handle)?;
1400            batch_handle.commit()
1401        })
1402    }
1403
1404    /// Commits a raw `WriteBatchWithTransaction` to `RocksDB`.
1405    ///
1406    /// This is used when the batch was extracted via [`RocksDBBatch::into_inner`]
1407    /// and needs to be committed at a later point (e.g., at provider commit time).
1408    ///
1409    /// # Panics
1410    /// Panics if the provider is in read-only mode.
1411    #[instrument(level = "debug", target = "providers::rocksdb", skip_all, fields(batch_len = batch.len(), batch_size = batch.size_in_bytes()))]
1412    pub fn commit_batch(&self, batch: WriteBatchWithTransaction<true>) -> ProviderResult<()> {
1413        self.0.db_rw().write_opt(batch, &synced_write_options()).map_err(|e| {
1414            ProviderError::Database(DatabaseError::Commit(DatabaseErrorInfo {
1415                message: e.to_string().into(),
1416                code: -1,
1417            }))
1418        })
1419    }
1420
1421    /// Writes all `RocksDB` data for multiple blocks in parallel.
1422    ///
1423    /// This handles transaction hash numbers, account history, and storage history based on
1424    /// the provided storage settings. Each operation runs in parallel with its own batch,
1425    /// pushing to `ctx.pending_batches` for later commit.
1426    #[instrument(level = "debug", target = "providers::rocksdb", skip_all, fields(num_blocks = blocks.len(), first_block = ctx.first_block_number))]
1427    pub(crate) fn write_blocks_data<N: reth_node_types::NodePrimitives>(
1428        &self,
1429        blocks: &[ExecutedBlock<N>],
1430        tx_nums: &[TxNumber],
1431        ctx: RocksDBWriteCtx,
1432        runtime: &reth_tasks::Runtime,
1433    ) -> ProviderResult<()> {
1434        if !ctx.storage_settings.storage_v2 {
1435            return Ok(());
1436        }
1437
1438        let mut r_tx_hash = None;
1439        let mut r_account_history = None;
1440        let mut r_storage_history = None;
1441
1442        let write_tx_hash =
1443            ctx.storage_settings.storage_v2 && ctx.prune_tx_lookup.is_none_or(|m| !m.is_full());
1444        let write_account_history = ctx.storage_settings.storage_v2;
1445        let write_storage_history = ctx.storage_settings.storage_v2;
1446
1447        // Propagate tracing context into rayon-spawned threads so that RocksDB
1448        // write spans appear as children of write_blocks_data in traces.
1449        let span = tracing::Span::current();
1450        runtime.storage_pool().in_place_scope(|s| {
1451            if write_tx_hash {
1452                s.spawn(|_| {
1453                    let _guard = span.enter();
1454                    r_tx_hash = Some(self.write_tx_hash_numbers(blocks, tx_nums, &ctx));
1455                });
1456            }
1457
1458            if write_account_history {
1459                s.spawn(|_| {
1460                    let _guard = span.enter();
1461                    r_account_history = Some(self.write_account_history(blocks, &ctx));
1462                });
1463            }
1464
1465            if write_storage_history {
1466                s.spawn(|_| {
1467                    let _guard = span.enter();
1468                    r_storage_history = Some(self.write_storage_history(blocks, &ctx));
1469                });
1470            }
1471        });
1472
1473        if write_tx_hash {
1474            r_tx_hash.ok_or_else(|| {
1475                ProviderError::Database(DatabaseError::Other(
1476                    "rocksdb tx-hash write thread panicked".into(),
1477                ))
1478            })??;
1479        }
1480        if write_account_history {
1481            r_account_history.ok_or_else(|| {
1482                ProviderError::Database(DatabaseError::Other(
1483                    "rocksdb account-history write thread panicked".into(),
1484                ))
1485            })??;
1486        }
1487        if write_storage_history {
1488            r_storage_history.ok_or_else(|| {
1489                ProviderError::Database(DatabaseError::Other(
1490                    "rocksdb storage-history write thread panicked".into(),
1491                ))
1492            })??;
1493        }
1494
1495        Ok(())
1496    }
1497
1498    /// Writes transaction hash to number mappings for the given blocks.
1499    #[instrument(level = "debug", target = "providers::rocksdb", skip_all)]
1500    fn write_tx_hash_numbers<N: reth_node_types::NodePrimitives>(
1501        &self,
1502        blocks: &[ExecutedBlock<N>],
1503        tx_nums: &[TxNumber],
1504        ctx: &RocksDBWriteCtx,
1505    ) -> ProviderResult<()> {
1506        let mut batch = self.batch();
1507        for (block, &first_tx_num) in blocks.iter().zip(tx_nums) {
1508            let body = block.recovered_block().body();
1509            for (tx_num, transaction) in (first_tx_num..).zip(body.transactions_iter()) {
1510                batch.put::<tables::TransactionHashNumbers>(*transaction.tx_hash(), &tx_num)?;
1511            }
1512        }
1513        ctx.pending_batches.lock().push(batch.into_inner());
1514        Ok(())
1515    }
1516
1517    /// Writes account history indices for the given blocks.
1518    ///
1519    /// Derives history indices from reverts (same source as changesets) to ensure consistency.
1520    #[instrument(level = "debug", target = "providers::rocksdb", skip_all)]
1521    fn write_account_history<N: reth_node_types::NodePrimitives>(
1522        &self,
1523        blocks: &[ExecutedBlock<N>],
1524        ctx: &RocksDBWriteCtx,
1525    ) -> ProviderResult<()> {
1526        let mut batch = self.batch();
1527        let mut account_history: BTreeMap<Address, Vec<u64>> = BTreeMap::new();
1528
1529        for (block_idx, block) in blocks.iter().enumerate() {
1530            let block_number = ctx.first_block_number + block_idx as u64;
1531            let reverts = block.execution_outcome().state.reverts.to_plain_state_reverts();
1532
1533            // Iterate through account reverts - these are exactly the accounts that have
1534            // changesets written, ensuring history indices match changeset entries.
1535            for account_block_reverts in reverts.accounts {
1536                for (address, _) in account_block_reverts {
1537                    account_history.entry(address).or_default().push(block_number);
1538                }
1539            }
1540        }
1541
1542        // Write account history using proper shard append logic
1543        for (address, indices) in account_history {
1544            batch.append_account_history_shard(address, indices)?;
1545        }
1546        ctx.pending_batches.lock().push(batch.into_inner());
1547        Ok(())
1548    }
1549
1550    /// Writes storage history indices for the given blocks.
1551    ///
1552    /// Derives history indices from reverts (same source as changesets) to ensure consistency.
1553    #[instrument(level = "debug", target = "providers::rocksdb", skip_all)]
1554    fn write_storage_history<N: reth_node_types::NodePrimitives>(
1555        &self,
1556        blocks: &[ExecutedBlock<N>],
1557        ctx: &RocksDBWriteCtx,
1558    ) -> ProviderResult<()> {
1559        let mut storage_history: BTreeMap<(Address, B256), Vec<u64>> = BTreeMap::new();
1560
1561        for (block_idx, block) in blocks.iter().enumerate() {
1562            let block_number = ctx.first_block_number + block_idx as u64;
1563            let reverts = block.execution_outcome().state.reverts.to_plain_state_reverts();
1564
1565            // Iterate through storage reverts - these are exactly the slots that have
1566            // changesets written, ensuring history indices match changeset entries.
1567            for storage_block_reverts in reverts.storage {
1568                for revert in storage_block_reverts {
1569                    for (slot, _) in revert.storage_revert {
1570                        let plain_key = B256::new(slot.to_be_bytes());
1571                        storage_history
1572                            .entry((revert.address, plain_key))
1573                            .or_default()
1574                            .push(block_number);
1575                    }
1576                }
1577            }
1578        }
1579
1580        let shard_puts = storage_history
1581            .into_par_iter()
1582            .map(|((address, slot), indices)| {
1583                self.storage_history_shards_to_put(address, slot, indices)
1584            })
1585            .collect::<ProviderResult<Vec<_>>>()?;
1586
1587        let mut batch = self.batch();
1588        for shards in shard_puts {
1589            for (key, shard) in shards {
1590                batch.put::<tables::StoragesHistory>(key, &shard)?;
1591            }
1592        }
1593        ctx.pending_batches.lock().push(batch.into_inner());
1594        Ok(())
1595    }
1596
1597    /// Prepares storage history shard writes by reading the current last shard and appending
1598    /// indices.
1599    fn storage_history_shards_to_put(
1600        &self,
1601        address: Address,
1602        storage_key: B256,
1603        indices: Vec<u64>,
1604    ) -> ProviderResult<Vec<(StorageShardedKey, BlockNumberList)>> {
1605        if indices.is_empty() {
1606            return Ok(Vec::new());
1607        }
1608
1609        debug_assert!(
1610            indices.windows(2).all(|w| w[0] < w[1]),
1611            "indices must be strictly increasing: {:?}",
1612            indices
1613        );
1614
1615        let last_key = StorageShardedKey::last(address, storage_key);
1616        let last_shard_opt = self.get::<tables::StoragesHistory>(last_key.clone())?;
1617        let mut last_shard = last_shard_opt.unwrap_or_else(BlockNumberList::empty);
1618
1619        last_shard.append(indices).map_err(ProviderError::other)?;
1620
1621        if last_shard.len() <= NUM_OF_INDICES_IN_SHARD as u64 {
1622            return Ok(vec![(last_key, last_shard)]);
1623        }
1624
1625        let chunks = last_shard.iter().chunks(NUM_OF_INDICES_IN_SHARD);
1626        let mut chunks_peekable = chunks.into_iter().peekable();
1627        let mut shards = Vec::new();
1628
1629        while let Some(chunk) = chunks_peekable.next() {
1630            let shard = BlockNumberList::new_pre_sorted(chunk);
1631            let highest_block_number = if chunks_peekable.peek().is_some() {
1632                shard.iter().next_back().expect("`chunks` does not return empty list")
1633            } else {
1634                u64::MAX
1635            };
1636
1637            shards
1638                .push((StorageShardedKey::new(address, storage_key, highest_block_number), shard));
1639        }
1640
1641        Ok(shards)
1642    }
1643}
1644
1645/// A point-in-time read snapshot of the `RocksDB` database.
1646///
1647/// All reads through this snapshot see a consistent view of the database at the point
1648/// the snapshot was created, regardless of concurrent writes. This is the primary reader
1649/// used by [`EitherReader::RocksDB`](crate::either_writer::EitherReader) for history lookups.
1650///
1651/// Lighter weight than [`RocksTx`] — no transaction overhead, no write support.
1652pub struct RocksReadSnapshot<'db> {
1653    /// Raw iterator reused across account history lookups.
1654    ///
1655    /// Declared before `inner` so it is dropped before the snapshot it reads from.
1656    accounts_history_iter: Mutex<Option<RocksDBRawIterEnum<'db>>>,
1657    /// Raw iterator reused across storage history lookups.
1658    ///
1659    /// Declared before `inner` so it is dropped before the snapshot it reads from.
1660    storages_history_iter: Mutex<Option<RocksDBRawIterEnum<'db>>>,
1661    inner: RocksReadSnapshotInner<'db>,
1662    provider: &'db RocksDBProviderInner,
1663}
1664
1665/// Inner enum to hold the snapshot for either read-write or secondary mode.
1666enum RocksReadSnapshotInner<'db> {
1667    /// Snapshot from read-write `OptimisticTransactionDB`.
1668    ReadWrite(SnapshotWithThreadMode<'db, OptimisticTransactionDB>),
1669    /// Direct reads from a secondary `DB` instance (no snapshot).
1670    Secondary(&'db DB),
1671}
1672
1673impl fmt::Debug for RocksReadSnapshot<'_> {
1674    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1675        f.debug_struct("RocksReadSnapshot")
1676            .field("provider", &self.provider)
1677            .finish_non_exhaustive()
1678    }
1679}
1680
1681impl<'db> RocksReadSnapshot<'db> {
1682    /// Gets the column family handle for a table.
1683    fn cf_handle<T: Table>(&self) -> Result<&'db rocksdb::ColumnFamily, DatabaseError> {
1684        self.provider.cf_handle::<T>()
1685    }
1686
1687    /// Creates a raw iterator over `cf` that observes this snapshot.
1688    ///
1689    /// The iterator is created from the database handle rather than from the snapshot value, so
1690    /// its lifetime is tied to the database and it can be cached in `self`. The snapshot is
1691    /// attached through [`ReadOptions`], giving the same point-in-time view as reading through the
1692    /// snapshot directly.
1693    fn new_raw_iterator_cf(&self, cf: &rocksdb::ColumnFamily) -> RocksDBRawIterEnum<'db> {
1694        match &self.inner {
1695            RocksReadSnapshotInner::ReadWrite(snap) => {
1696                let mut readopts = ReadOptions::default();
1697                readopts.set_snapshot(snap);
1698                RocksDBRawIterEnum::ReadWrite(
1699                    self.provider.db_rw().raw_iterator_cf_opt(cf, readopts),
1700                )
1701            }
1702            RocksReadSnapshotInner::Secondary(db) => {
1703                RocksDBRawIterEnum::ReadOnly((*db).raw_iterator_cf(cf))
1704            }
1705        }
1706    }
1707
1708    /// Gets a value from the specified table.
1709    pub fn get<T: Table>(&self, key: T::Key) -> ProviderResult<Option<T::Value>> {
1710        let encoded_key = key.encode();
1711        let cf = self.cf_handle::<T>()?;
1712        let result = match &self.inner {
1713            RocksReadSnapshotInner::ReadWrite(snap) => snap.get_cf(cf, encoded_key.as_ref()),
1714            RocksReadSnapshotInner::Secondary(db) => db.get_cf(cf, encoded_key.as_ref()),
1715        }
1716        .map_err(|e| {
1717            ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
1718                message: e.to_string().into(),
1719                code: -1,
1720            }))
1721        })?;
1722
1723        Ok(result.and_then(|value| T::Value::decompress(&value).ok()))
1724    }
1725
1726    /// Lookup account history and return [`HistoryInfo`] directly.
1727    ///
1728    /// `visible_tip` is the highest block considered visible from the companion MDBX snapshot.
1729    /// History entries above it are ignored even if they already exist in `RocksDB`.
1730    pub fn account_history_info(
1731        &self,
1732        address: Address,
1733        block_number: BlockNumber,
1734        lowest_available_block_number: Option<BlockNumber>,
1735        visible_tip: BlockNumber,
1736    ) -> ProviderResult<HistoryInfo> {
1737        let key = ShardedKey::new(address, block_number);
1738        self.history_info::<tables::AccountsHistory>(
1739            &self.accounts_history_iter,
1740            key.encode().as_ref(),
1741            block_number,
1742            lowest_available_block_number,
1743            visible_tip,
1744            |key_bytes| Ok(<ShardedKey<Address> as Decode>::decode(key_bytes)?.key == address),
1745            |prev_bytes| {
1746                <ShardedKey<Address> as Decode>::decode(prev_bytes)
1747                    .map(|k| k.key == address)
1748                    .unwrap_or(false)
1749            },
1750        )
1751    }
1752
1753    /// Lookup storage history and return [`HistoryInfo`] directly.
1754    ///
1755    /// `visible_tip` is the highest block considered visible from the companion MDBX snapshot.
1756    /// History entries above it are ignored even if they already exist in `RocksDB`.
1757    pub fn storage_history_info(
1758        &self,
1759        address: Address,
1760        storage_key: B256,
1761        block_number: BlockNumber,
1762        lowest_available_block_number: Option<BlockNumber>,
1763        visible_tip: BlockNumber,
1764    ) -> ProviderResult<HistoryInfo> {
1765        let key = StorageShardedKey::new(address, storage_key, block_number);
1766        self.history_info::<tables::StoragesHistory>(
1767            &self.storages_history_iter,
1768            key.encode().as_ref(),
1769            block_number,
1770            lowest_available_block_number,
1771            visible_tip,
1772            |key_bytes| {
1773                let k = <StorageShardedKey as Decode>::decode(key_bytes)?;
1774                Ok(k.address == address && k.sharded_key.key == storage_key)
1775            },
1776            |prev_bytes| {
1777                <StorageShardedKey as Decode>::decode(prev_bytes)
1778                    .map(|k| k.address == address && k.sharded_key.key == storage_key)
1779                    .unwrap_or(false)
1780            },
1781        )
1782    }
1783
1784    /// Generic history lookup using the snapshot's raw iterator.
1785    ///
1786    /// The result is derived from the history that is visible through `visible_tip`, not from the
1787    /// full contents of `RocksDB`. This lets a reader combine an older MDBX snapshot with a newer
1788    /// Rocks snapshot without routing through history entries that MDBX cannot see yet.
1789    ///
1790    /// `iter_cache` holds the raw iterator for `T`'s column family. Seeking an existing iterator
1791    /// is much cheaper than constructing one, so the iterator is created on the first lookup and
1792    /// reused by every later lookup through this snapshot. A lookup that finds the cache held by
1793    /// another thread uses a private iterator instead of waiting.
1794    #[expect(clippy::too_many_arguments)]
1795    fn history_info<T>(
1796        &self,
1797        iter_cache: &Mutex<Option<RocksDBRawIterEnum<'db>>>,
1798        encoded_key: &[u8],
1799        block_number: BlockNumber,
1800        lowest_available_block_number: Option<BlockNumber>,
1801        visible_tip: BlockNumber,
1802        key_matches: impl FnOnce(&[u8]) -> Result<bool, reth_db_api::DatabaseError>,
1803        prev_key_matches: impl Fn(&[u8]) -> bool,
1804    ) -> ProviderResult<HistoryInfo>
1805    where
1806        T: Table<Value = BlockNumberList>,
1807    {
1808        let is_maybe_pruned = lowest_available_block_number.is_some();
1809        let fallback = || {
1810            Ok(if is_maybe_pruned {
1811                HistoryInfo::MaybeInPlainState
1812            } else {
1813                HistoryInfo::NotYetWritten
1814            })
1815        };
1816
1817        let cf = self.cf_handle::<T>()?;
1818        // A lookup on another thread may be holding the cached iterator; a private iterator costs
1819        // what every lookup used to cost, whereas waiting would serialize the two lookups.
1820        let mut guard = iter_cache.try_lock();
1821        let mut private_iter;
1822        let iter = match guard.as_mut() {
1823            Some(cached) => cached.get_or_insert_with(|| self.new_raw_iterator_cf(cf)),
1824            None => {
1825                private_iter = self.new_raw_iterator_cf(cf);
1826                &mut private_iter
1827            }
1828        };
1829
1830        iter.seek(encoded_key);
1831        iter.status().map_err(|e| {
1832            ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
1833                message: e.to_string().into(),
1834                code: -1,
1835            }))
1836        })?;
1837
1838        if !iter.valid() {
1839            return fallback();
1840        }
1841
1842        let Some(key_bytes) = iter.key() else {
1843            return fallback();
1844        };
1845        if !key_matches(key_bytes)? {
1846            return fallback();
1847        }
1848
1849        let Some(value_bytes) = iter.value() else {
1850            return fallback();
1851        };
1852        let chunk = BlockNumberList::decompress(value_bytes)?;
1853
1854        let (rank, found_block) = compute_history_rank(&chunk, block_number);
1855        // Ignore later Rocks history that is ahead of the companion MDBX snapshot.
1856        let found_block = found_block.filter(|block| *block <= visible_tip);
1857
1858        let is_before_first_write = if needs_prev_shard_check(rank, found_block, block_number) {
1859            iter.prev();
1860            iter.status().map_err(|e| {
1861                ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
1862                    message: e.to_string().into(),
1863                    code: -1,
1864                }))
1865            })?;
1866            let has_prev = iter.valid() && iter.key().is_some_and(&prev_key_matches);
1867
1868            // If the current shard only contains history above `visible_tip`, there is no usable
1869            // later change. Without a previous shard for the same key, fall back to the existing
1870            // not-written / maybe-pruned result instead of routing into plain state.
1871            if found_block.is_none() && !has_prev {
1872                return fallback()
1873            }
1874
1875            !has_prev
1876        } else {
1877            false
1878        };
1879
1880        Ok(HistoryInfo::from_lookup(
1881            found_block,
1882            is_before_first_write,
1883            lowest_available_block_number,
1884        ))
1885    }
1886}
1887
1888/// A [`RocksReadSnapshot`] that owns the [`RocksDBProvider`] handle it reads through.
1889///
1890/// [`RocksDBProvider::snapshot`] borrows the provider, so a reader that wants to keep a single
1891/// snapshot alive has to keep the provider handle next to it. This type does that, which lets the
1892/// snapshot cache its history iterators across lookups.
1893pub struct OwnedRocksReadSnapshot {
1894    /// Borrows the allocation retained by `provider`. Declared first so the snapshot and its
1895    /// iterators are released before the owning database handle. Boxing keeps its internal
1896    /// references out of the movable owner, including when the owner is passed by value.
1897    snapshot: Box<RocksReadSnapshot<'static>>,
1898    provider: RocksDBProvider,
1899}
1900
1901impl fmt::Debug for OwnedRocksReadSnapshot {
1902    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1903        f.debug_struct("OwnedRocksReadSnapshot").field("provider", &self.provider).finish()
1904    }
1905}
1906
1907impl OwnedRocksReadSnapshot {
1908    fn new(provider: &RocksDBProvider) -> Self {
1909        let provider = provider.clone();
1910        let snapshot = provider.snapshot();
1911        // SAFETY: Every reference in `snapshot` points into the Arc allocation retained by
1912        // `provider`, not into the movable provider handle. The handle is private and never
1913        // replaced, and `snapshot` is dropped before it. `as_snapshot` restricts access to the
1914        // lifetime of a borrow of this owner.
1915        let snapshot = unsafe {
1916            std::mem::transmute::<RocksReadSnapshot<'_>, RocksReadSnapshot<'static>>(snapshot)
1917        };
1918        Self { snapshot: Box::new(snapshot), provider }
1919    }
1920
1921    /// Returns the borrowed snapshot.
1922    pub fn as_snapshot(&self) -> &RocksReadSnapshot<'_> {
1923        // SAFETY: shortening the snapshot's lifetime to this borrow of `self` is sound; the owned
1924        // provider handle keeps the database open for at least that long.
1925        unsafe {
1926            std::mem::transmute::<&RocksReadSnapshot<'static>, &RocksReadSnapshot<'_>>(
1927                &self.snapshot,
1928            )
1929        }
1930    }
1931}
1932
1933/// Outcome of pruning a history shard in `RocksDB`.
1934#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1935pub enum PruneShardOutcome {
1936    /// Shard was deleted entirely.
1937    Deleted,
1938    /// Shard was updated with filtered block numbers.
1939    Updated,
1940    /// Shard was unchanged (no blocks <= `to_block`).
1941    Unchanged,
1942}
1943
1944/// Tracks pruning outcomes for batch operations.
1945#[derive(Debug, Default, Clone, Copy)]
1946pub struct PrunedIndices {
1947    /// Number of shards completely deleted.
1948    pub deleted: usize,
1949    /// Number of shards that were updated (filtered but still have entries).
1950    pub updated: usize,
1951    /// Number of shards that were unchanged.
1952    pub unchanged: usize,
1953}
1954
1955/// Handle for building a batch of operations atomically.
1956///
1957/// Uses `WriteBatchWithTransaction` for atomic writes without full transaction overhead.
1958/// Unlike [`RocksTx`], this does NOT support read-your-writes. Use for write-only flows
1959/// where you don't need to read back uncommitted data within the same operation
1960/// (e.g., history index writes).
1961///
1962/// When `auto_commit_threshold` is set, the batch will automatically commit and reset
1963/// when the batch size exceeds the threshold. This prevents OOM during large bulk writes.
1964#[must_use = "batch must be committed"]
1965pub struct RocksDBBatch<'a> {
1966    provider: &'a RocksDBProvider,
1967    inner: WriteBatchWithTransaction<true>,
1968    buf: Vec<u8>,
1969    /// If set, batch auto-commits when size exceeds this threshold (in bytes).
1970    auto_commit_threshold: Option<usize>,
1971}
1972
1973impl fmt::Debug for RocksDBBatch<'_> {
1974    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1975        f.debug_struct("RocksDBBatch")
1976            .field("provider", &self.provider)
1977            .field("batch", &"<WriteBatchWithTransaction>")
1978            // Number of operations in this batch
1979            .field("length", &self.inner.len())
1980            // Total serialized size (encoded key + compressed value + metadata) of this batch
1981            // in bytes
1982            .field("size_in_bytes", &self.inner.size_in_bytes())
1983            .finish()
1984    }
1985}
1986
1987impl<'a> RocksDBBatch<'a> {
1988    /// Puts a value into the batch.
1989    ///
1990    /// If auto-commit is enabled and the batch exceeds the threshold, commits and resets.
1991    pub fn put<T: Table>(&mut self, key: T::Key, value: &T::Value) -> ProviderResult<()> {
1992        let encoded_key = key.encode();
1993        self.put_encoded::<T>(&encoded_key, value)
1994    }
1995
1996    /// Puts a value into the batch using pre-encoded key.
1997    ///
1998    /// If auto-commit is enabled and the batch exceeds the threshold, commits and resets.
1999    pub fn put_encoded<T: Table>(
2000        &mut self,
2001        key: &<T::Key as Encode>::Encoded,
2002        value: &T::Value,
2003    ) -> ProviderResult<()> {
2004        let value_bytes = compress_to_buf_or_ref!(self.buf, value).unwrap_or(&self.buf);
2005        self.inner.put_cf(self.provider.get_cf_handle::<T>()?, key, value_bytes);
2006        self.maybe_auto_commit()?;
2007        Ok(())
2008    }
2009
2010    /// Deletes a value from the batch.
2011    ///
2012    /// If auto-commit is enabled and the batch exceeds the threshold, commits and resets.
2013    pub fn delete<T: Table>(&mut self, key: T::Key) -> ProviderResult<()> {
2014        self.inner.delete_cf(self.provider.get_cf_handle::<T>()?, key.encode().as_ref());
2015        self.maybe_auto_commit()?;
2016        Ok(())
2017    }
2018
2019    /// Commits and resets the batch if it exceeds the auto-commit threshold.
2020    ///
2021    /// This is called after each `put` or `delete` operation to prevent unbounded memory growth.
2022    /// Returns immediately if auto-commit is disabled or threshold not reached.
2023    fn maybe_auto_commit(&mut self) -> ProviderResult<()> {
2024        if let Some(threshold) = self.auto_commit_threshold &&
2025            self.inner.size_in_bytes() >= threshold
2026        {
2027            tracing::debug!(
2028                target: "providers::rocksdb",
2029                batch_size = self.inner.size_in_bytes(),
2030                threshold,
2031                "Auto-committing RocksDB batch"
2032            );
2033            let old_batch = std::mem::take(&mut self.inner);
2034            self.provider.0.db_rw().write_opt(old_batch, &synced_write_options()).map_err(|e| {
2035                ProviderError::Database(DatabaseError::Commit(DatabaseErrorInfo {
2036                    message: e.to_string().into(),
2037                    code: -1,
2038                }))
2039            })?;
2040        }
2041        Ok(())
2042    }
2043
2044    /// Commits the batch to the database.
2045    ///
2046    /// This consumes the batch and writes all operations atomically to `RocksDB`.
2047    ///
2048    /// # Panics
2049    /// Panics if the provider is in read-only mode.
2050    #[instrument(level = "debug", target = "providers::rocksdb", skip_all, fields(batch_len = self.inner.len(), batch_size = self.inner.size_in_bytes()))]
2051    pub fn commit(self) -> ProviderResult<()> {
2052        self.provider.0.db_rw().write_opt(self.inner, &synced_write_options()).map_err(|e| {
2053            ProviderError::Database(DatabaseError::Commit(DatabaseErrorInfo {
2054                message: e.to_string().into(),
2055                code: -1,
2056            }))
2057        })
2058    }
2059
2060    /// Returns the number of write operations (puts + deletes) queued in this batch.
2061    pub fn len(&self) -> usize {
2062        self.inner.len()
2063    }
2064
2065    /// Returns `true` if the batch contains no operations.
2066    pub fn is_empty(&self) -> bool {
2067        self.inner.is_empty()
2068    }
2069
2070    /// Returns the size of the batch in bytes.
2071    pub fn size_in_bytes(&self) -> usize {
2072        self.inner.size_in_bytes()
2073    }
2074
2075    /// Returns a reference to the underlying `RocksDB` provider.
2076    pub const fn provider(&self) -> &RocksDBProvider {
2077        self.provider
2078    }
2079
2080    /// Consumes the batch and returns the underlying `WriteBatchWithTransaction`.
2081    ///
2082    /// This is used to defer commits to the provider level.
2083    pub fn into_inner(self) -> WriteBatchWithTransaction<true> {
2084        self.inner
2085    }
2086
2087    /// Gets a value from the database.
2088    ///
2089    /// **Important constraint:** This reads only committed state, not pending writes in this
2090    /// batch or other pending batches in `pending_rocksdb_batches`.
2091    pub fn get<T: Table>(&self, key: T::Key) -> ProviderResult<Option<T::Value>> {
2092        self.provider.get::<T>(key)
2093    }
2094
2095    /// Appends indices to an account history shard with proper shard management.
2096    ///
2097    /// Loads the existing shard (if any), appends new indices, and rechunks into
2098    /// multiple shards if needed (respecting `NUM_OF_INDICES_IN_SHARD` limit).
2099    ///
2100    /// # Requirements
2101    ///
2102    /// - The `indices` MUST be strictly increasing and contain no duplicates.
2103    /// - This method MUST only be called once per address per batch. The batch reads existing
2104    ///   shards from committed DB state, not from pending writes. Calling twice for the same
2105    ///   address will cause the second call to overwrite the first.
2106    pub fn append_account_history_shard(
2107        &mut self,
2108        address: Address,
2109        indices: impl IntoIterator<Item = u64>,
2110    ) -> ProviderResult<()> {
2111        let indices: Vec<u64> = indices.into_iter().collect();
2112
2113        if indices.is_empty() {
2114            return Ok(());
2115        }
2116
2117        debug_assert!(
2118            indices.windows(2).all(|w| w[0] < w[1]),
2119            "indices must be strictly increasing: {:?}",
2120            indices
2121        );
2122
2123        let last_key = ShardedKey::new(address, u64::MAX);
2124        let last_shard_opt = self.provider.get::<tables::AccountsHistory>(last_key.clone())?;
2125        let mut last_shard = last_shard_opt.unwrap_or_else(BlockNumberList::empty);
2126
2127        last_shard.append(indices).map_err(ProviderError::other)?;
2128
2129        // Fast path: all indices fit in one shard
2130        if last_shard.len() <= NUM_OF_INDICES_IN_SHARD as u64 {
2131            self.put::<tables::AccountsHistory>(last_key, &last_shard)?;
2132            return Ok(());
2133        }
2134
2135        // Slow path: rechunk into multiple shards
2136        let chunks = last_shard.iter().chunks(NUM_OF_INDICES_IN_SHARD);
2137        let mut chunks_peekable = chunks.into_iter().peekable();
2138
2139        while let Some(chunk) = chunks_peekable.next() {
2140            let shard = BlockNumberList::new_pre_sorted(chunk);
2141            let highest_block_number = if chunks_peekable.peek().is_some() {
2142                shard.iter().next_back().expect("`chunks` does not return empty list")
2143            } else {
2144                u64::MAX
2145            };
2146
2147            self.put::<tables::AccountsHistory>(
2148                ShardedKey::new(address, highest_block_number),
2149                &shard,
2150            )?;
2151        }
2152
2153        Ok(())
2154    }
2155
2156    /// Appends indices to a storage history shard with proper shard management.
2157    ///
2158    /// Loads the existing shard (if any), appends new indices, and rechunks into
2159    /// multiple shards if needed (respecting `NUM_OF_INDICES_IN_SHARD` limit).
2160    ///
2161    /// # Requirements
2162    ///
2163    /// - The `indices` MUST be strictly increasing and contain no duplicates.
2164    /// - This method MUST only be called once per (address, `storage_key`) pair per batch. The
2165    ///   batch reads existing shards from committed DB state, not from pending writes. Calling
2166    ///   twice for the same key will cause the second call to overwrite the first.
2167    pub fn append_storage_history_shard(
2168        &mut self,
2169        address: Address,
2170        storage_key: B256,
2171        indices: impl IntoIterator<Item = u64>,
2172    ) -> ProviderResult<()> {
2173        let indices: Vec<u64> = indices.into_iter().collect();
2174
2175        for (key, shard) in
2176            self.provider.storage_history_shards_to_put(address, storage_key, indices)?
2177        {
2178            self.put::<tables::StoragesHistory>(key, &shard)?;
2179        }
2180
2181        Ok(())
2182    }
2183
2184    /// Unwinds account history for the given address, keeping only blocks <= `keep_to`.
2185    ///
2186    /// Mirrors MDBX `unwind_history_shards` behavior:
2187    /// - Deletes shards entirely above `keep_to`
2188    /// - Truncates boundary shards and re-keys to `u64::MAX` sentinel
2189    /// - Preserves shards entirely below `keep_to`
2190    pub fn unwind_account_history_to(
2191        &mut self,
2192        address: Address,
2193        keep_to: BlockNumber,
2194    ) -> ProviderResult<()> {
2195        let shards = self.provider.account_history_shards(address)?;
2196        if shards.is_empty() {
2197            return Ok(());
2198        }
2199
2200        // Find the first shard that might contain blocks > keep_to.
2201        // A shard is affected if it's the sentinel (u64::MAX) or its highest_block_number > keep_to
2202        let boundary_idx = shards.iter().position(|(key, _)| {
2203            key.highest_block_number == u64::MAX || key.highest_block_number > keep_to
2204        });
2205
2206        // Repair path: no shards affected means all blocks <= keep_to, just ensure sentinel exists
2207        let Some(boundary_idx) = boundary_idx else {
2208            let (last_key, last_value) = shards.last().expect("shards is non-empty");
2209            if last_key.highest_block_number != u64::MAX {
2210                self.delete::<tables::AccountsHistory>(last_key.clone())?;
2211                self.put::<tables::AccountsHistory>(
2212                    ShardedKey::new(address, u64::MAX),
2213                    last_value,
2214                )?;
2215            }
2216            return Ok(());
2217        };
2218
2219        // Delete all shards strictly after the boundary (they are entirely > keep_to)
2220        for (key, _) in shards.iter().skip(boundary_idx + 1) {
2221            self.delete::<tables::AccountsHistory>(key.clone())?;
2222        }
2223
2224        // Process the boundary shard: filter out blocks > keep_to
2225        let (boundary_key, boundary_list) = &shards[boundary_idx];
2226
2227        // Delete the boundary shard (we'll either drop it or rewrite at u64::MAX)
2228        self.delete::<tables::AccountsHistory>(boundary_key.clone())?;
2229
2230        // Build truncated list once; check emptiness directly (avoids double iteration)
2231        let new_last =
2232            BlockNumberList::new_pre_sorted(boundary_list.iter().take_while(|&b| b <= keep_to));
2233
2234        if new_last.is_empty() {
2235            // Boundary shard is now empty. Previous shard becomes the last and must be keyed
2236            // u64::MAX.
2237            if boundary_idx == 0 {
2238                // Nothing left for this address
2239                return Ok(());
2240            }
2241
2242            let (prev_key, prev_value) = &shards[boundary_idx - 1];
2243            if prev_key.highest_block_number != u64::MAX {
2244                self.delete::<tables::AccountsHistory>(prev_key.clone())?;
2245                self.put::<tables::AccountsHistory>(
2246                    ShardedKey::new(address, u64::MAX),
2247                    prev_value,
2248                )?;
2249            }
2250            return Ok(());
2251        }
2252
2253        self.put::<tables::AccountsHistory>(ShardedKey::new(address, u64::MAX), &new_last)?;
2254
2255        Ok(())
2256    }
2257
2258    /// Prunes history shards, removing blocks <= `to_block`.
2259    ///
2260    /// Generic implementation for both account and storage history pruning.
2261    /// Mirrors MDBX `prune_shard` semantics. After pruning, the last remaining shard
2262    /// (if any) will have the sentinel key (`u64::MAX`).
2263    ///
2264    /// `shards_complete` must be `false` when `shards` stops short of the key's last shard, so the
2265    /// survivor is not re-keyed over a sentinel that is still on disk.
2266    #[expect(clippy::too_many_arguments)]
2267    fn prune_history_shards_inner<K>(
2268        &mut self,
2269        shards: Vec<(K, BlockNumberList)>,
2270        shards_complete: bool,
2271        to_block: BlockNumber,
2272        get_highest: impl Fn(&K) -> u64,
2273        is_sentinel: impl Fn(&K) -> bool,
2274        delete_shard: impl Fn(&mut Self, K) -> ProviderResult<()>,
2275        put_shard: impl Fn(&mut Self, K, &BlockNumberList) -> ProviderResult<()>,
2276        create_sentinel: impl Fn() -> K,
2277    ) -> ProviderResult<PruneShardOutcome>
2278    where
2279        K: Clone,
2280    {
2281        if shards.is_empty() {
2282            return Ok(PruneShardOutcome::Unchanged);
2283        }
2284
2285        let mut deleted = false;
2286        let mut updated = false;
2287        let mut last_remaining: Option<(K, BlockNumberList)> = None;
2288
2289        for (key, mut block_list) in shards {
2290            if !is_sentinel(&key) && get_highest(&key) <= to_block {
2291                delete_shard(self, key)?;
2292                deleted = true;
2293            } else {
2294                let removed = block_list.remove_range(0..=to_block);
2295
2296                if block_list.is_empty() {
2297                    delete_shard(self, key)?;
2298                    deleted = true;
2299                } else if removed > 0 {
2300                    put_shard(self, key.clone(), &block_list)?;
2301                    last_remaining = Some((key, block_list));
2302                    updated = true;
2303                } else {
2304                    last_remaining = Some((key, block_list));
2305                }
2306            }
2307        }
2308
2309        if shards_complete &&
2310            let Some((last_key, last_value)) = last_remaining &&
2311            !is_sentinel(&last_key)
2312        {
2313            delete_shard(self, last_key)?;
2314            put_shard(self, create_sentinel(), &last_value)?;
2315            updated = true;
2316        }
2317
2318        if deleted {
2319            Ok(PruneShardOutcome::Deleted)
2320        } else if updated {
2321            Ok(PruneShardOutcome::Updated)
2322        } else {
2323            Ok(PruneShardOutcome::Unchanged)
2324        }
2325    }
2326
2327    /// Prunes account history for the given address, removing blocks <= `to_block`.
2328    ///
2329    /// Mirrors MDBX `prune_shard` semantics. After pruning, the last remaining shard
2330    /// (if any) will have the sentinel key (`u64::MAX`).
2331    pub fn prune_account_history_to(
2332        &mut self,
2333        address: Address,
2334        to_block: BlockNumber,
2335    ) -> ProviderResult<PruneShardOutcome> {
2336        let shards = self.provider.account_history_shards(address)?;
2337        self.prune_history_shards_inner(
2338            shards,
2339            true,
2340            to_block,
2341            |key| key.highest_block_number,
2342            |key| key.highest_block_number == u64::MAX,
2343            |batch, key| batch.delete::<tables::AccountsHistory>(key),
2344            |batch, key, value| batch.put::<tables::AccountsHistory>(key, value),
2345            || ShardedKey::new(address, u64::MAX),
2346        )
2347    }
2348
2349    /// Prunes account history for multiple addresses in a single iterator pass.
2350    ///
2351    /// This is more efficient than calling [`Self::prune_account_history_to`] repeatedly
2352    /// because it reuses a single raw iterator and skips seeks when the iterator is already
2353    /// positioned correctly (which happens when targets are sorted and adjacent in key order).
2354    ///
2355    /// `targets` MUST be sorted by address and contain each address at most once, for
2356    /// correctness and optimal performance (matches on-disk key order).
2357    pub fn prune_account_history_batch(
2358        &mut self,
2359        targets: &[(Address, BlockNumber)],
2360    ) -> ProviderResult<PrunedIndices> {
2361        if targets.is_empty() {
2362            return Ok(PrunedIndices::default());
2363        }
2364
2365        debug_assert!(
2366            targets.windows(2).all(|w| w[0].0 < w[1].0),
2367            "prune_account_history_batch: targets must be sorted and unique"
2368        );
2369
2370        // ShardedKey<Address> layout: [address: 20][block: 8] = 28 bytes
2371        // The first 20 bytes are the "prefix" that identifies the address
2372        const PREFIX_LEN: usize = 20;
2373
2374        let cf = self.provider.get_cf_handle::<tables::AccountsHistory>()?;
2375        let mut iter = self.provider.0.raw_iterator_cf(cf);
2376        let mut outcomes = PrunedIndices::default();
2377
2378        for (address, to_block) in targets {
2379            // Build the target prefix (first 20 bytes = address)
2380            let start_key = ShardedKey::new(*address, 0u64).encode();
2381            let target_prefix = &start_key[..PREFIX_LEN];
2382
2383            // Check if we need to seek or if the iterator is already positioned correctly.
2384            // After processing the previous target, the iterator is either:
2385            // 1. Positioned at a key with a different prefix (we iterated past our shards)
2386            // 2. Positioned on a later shard of the previous target (we stopped early), whose
2387            //    prefix is below ours because targets are sorted and unique
2388            // 3. Invalid (no more keys)
2389            // If the current key's prefix >= our target prefix, we may be able to skip the seek.
2390            let needs_seek = if iter.valid() {
2391                if let Some(current_key) = iter.key() {
2392                    // If current key's prefix < target prefix, we need to seek forward
2393                    // If current key's prefix > target prefix, this target has no shards (skip)
2394                    // If current key's prefix == target prefix, we're already positioned
2395                    current_key.get(..PREFIX_LEN).is_none_or(|p| p < target_prefix)
2396                } else {
2397                    true
2398                }
2399            } else {
2400                true
2401            };
2402
2403            if needs_seek {
2404                iter.seek(start_key);
2405                iter.status().map_err(|e| {
2406                    ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
2407                        message: e.to_string().into(),
2408                        code: -1,
2409                    }))
2410                })?;
2411            }
2412
2413            // Collect the shards for this address that pruning can touch, using raw prefix
2414            // comparison
2415            let mut shards = Vec::new();
2416            let mut shards_complete = true;
2417            while iter.valid() {
2418                let Some(key_bytes) = iter.key() else { break };
2419
2420                // Use raw prefix comparison instead of full decode for the prefix check
2421                let current_prefix = key_bytes.get(..PREFIX_LEN);
2422                if current_prefix != Some(target_prefix) {
2423                    break;
2424                }
2425
2426                // Now decode the full key (we need the block number)
2427                let key = ShardedKey::<Address>::decode(key_bytes)
2428                    .map_err(|_| ProviderError::Database(DatabaseError::Decode))?;
2429
2430                let Some(value_bytes) = iter.value() else { break };
2431                let value = BlockNumberList::decompress(value_bytes)
2432                    .map_err(|_| ProviderError::Database(DatabaseError::Decode))?;
2433
2434                let highest = key.highest_block_number;
2435                shards.push((key, value));
2436
2437                iter.next();
2438
2439                // Shards are ordered by their highest block and their contents partition the
2440                // address's history, so this is the last shard holding anything at or below the
2441                // target. Peek past it only to tell whether it was the address's last shard,
2442                // which decides whether a survivor may be re-keyed to the sentinel.
2443                if highest > *to_block {
2444                    shards_complete = iter.key().and_then(|next_key| next_key.get(..PREFIX_LEN)) !=
2445                        Some(target_prefix);
2446                    break;
2447                }
2448            }
2449
2450            // The iterator also goes invalid on a read error, which would otherwise pass a
2451            // truncated shard list off as the key's complete one.
2452            if !iter.valid() {
2453                iter.status().map_err(|e| {
2454                    ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
2455                        message: e.to_string().into(),
2456                        code: -1,
2457                    }))
2458                })?;
2459            }
2460
2461            match self.prune_history_shards_inner(
2462                shards,
2463                shards_complete,
2464                *to_block,
2465                |key| key.highest_block_number,
2466                |key| key.highest_block_number == u64::MAX,
2467                |batch, key| batch.delete::<tables::AccountsHistory>(key),
2468                |batch, key, value| batch.put::<tables::AccountsHistory>(key, value),
2469                || ShardedKey::new(*address, u64::MAX),
2470            )? {
2471                PruneShardOutcome::Deleted => outcomes.deleted += 1,
2472                PruneShardOutcome::Updated => outcomes.updated += 1,
2473                PruneShardOutcome::Unchanged => outcomes.unchanged += 1,
2474            }
2475        }
2476
2477        Ok(outcomes)
2478    }
2479
2480    /// Prunes storage history for the given address and storage key, removing blocks <=
2481    /// `to_block`.
2482    ///
2483    /// Mirrors MDBX `prune_shard` semantics. After pruning, the last remaining shard
2484    /// (if any) will have the sentinel key (`u64::MAX`).
2485    pub fn prune_storage_history_to(
2486        &mut self,
2487        address: Address,
2488        storage_key: B256,
2489        to_block: BlockNumber,
2490    ) -> ProviderResult<PruneShardOutcome> {
2491        let shards = self.provider.storage_history_shards(address, storage_key)?;
2492        self.prune_history_shards_inner(
2493            shards,
2494            true,
2495            to_block,
2496            |key| key.sharded_key.highest_block_number,
2497            |key| key.sharded_key.highest_block_number == u64::MAX,
2498            |batch, key| batch.delete::<tables::StoragesHistory>(key),
2499            |batch, key, value| batch.put::<tables::StoragesHistory>(key, value),
2500            || StorageShardedKey::last(address, storage_key),
2501        )
2502    }
2503
2504    /// Prunes storage history for multiple (address, `storage_key`) pairs in a single iterator
2505    /// pass.
2506    ///
2507    /// This is more efficient than calling [`Self::prune_storage_history_to`] repeatedly
2508    /// because it reuses a single raw iterator and skips seeks when the iterator is already
2509    /// positioned correctly (which happens when targets are sorted and adjacent in key order).
2510    ///
2511    /// `targets` MUST be sorted by (address, `storage_key`) and contain each pair at most once,
2512    /// for correctness and optimal performance (matches on-disk key order).
2513    pub fn prune_storage_history_batch(
2514        &mut self,
2515        targets: &[((Address, B256), BlockNumber)],
2516    ) -> ProviderResult<PrunedIndices> {
2517        if targets.is_empty() {
2518            return Ok(PrunedIndices::default());
2519        }
2520
2521        debug_assert!(
2522            targets.windows(2).all(|w| w[0].0 < w[1].0),
2523            "prune_storage_history_batch: targets must be sorted and unique"
2524        );
2525
2526        // StorageShardedKey layout: [address: 20][storage_key: 32][block: 8] = 60 bytes
2527        // The first 52 bytes are the "prefix" that identifies (address, storage_key)
2528        const PREFIX_LEN: usize = 52;
2529
2530        let cf = self.provider.get_cf_handle::<tables::StoragesHistory>()?;
2531        let mut iter = self.provider.0.raw_iterator_cf(cf);
2532        let mut outcomes = PrunedIndices::default();
2533
2534        for ((address, storage_key), to_block) in targets {
2535            // Build the target prefix (first 52 bytes of encoded key)
2536            let start_key = StorageShardedKey::new(*address, *storage_key, 0u64).encode();
2537            let target_prefix = &start_key[..PREFIX_LEN];
2538
2539            // Check if we need to seek or if the iterator is already positioned correctly.
2540            // After processing the previous target, the iterator is either:
2541            // 1. Positioned at a key with a different prefix (we iterated past our shards)
2542            // 2. Positioned on a later shard of the previous target (we stopped early), whose
2543            //    prefix is below ours because targets are sorted and unique
2544            // 3. Invalid (no more keys)
2545            // If the current key's prefix >= our target prefix, we may be able to skip the seek.
2546            let needs_seek = if iter.valid() {
2547                if let Some(current_key) = iter.key() {
2548                    // If current key's prefix < target prefix, we need to seek forward
2549                    // If current key's prefix > target prefix, this target has no shards (skip)
2550                    // If current key's prefix == target prefix, we're already positioned
2551                    current_key.get(..PREFIX_LEN).is_none_or(|p| p < target_prefix)
2552                } else {
2553                    true
2554                }
2555            } else {
2556                true
2557            };
2558
2559            if needs_seek {
2560                iter.seek(start_key);
2561                iter.status().map_err(|e| {
2562                    ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
2563                        message: e.to_string().into(),
2564                        code: -1,
2565                    }))
2566                })?;
2567            }
2568
2569            // Collect the shards for this (address, storage_key) pair that pruning can touch,
2570            // using prefix comparison
2571            let mut shards = Vec::new();
2572            let mut shards_complete = true;
2573            while iter.valid() {
2574                let Some(key_bytes) = iter.key() else { break };
2575
2576                // Use raw prefix comparison instead of full decode for the prefix check
2577                let current_prefix = key_bytes.get(..PREFIX_LEN);
2578                if current_prefix != Some(target_prefix) {
2579                    break;
2580                }
2581
2582                // Now decode the full key (we need the block number)
2583                let key = StorageShardedKey::decode(key_bytes)
2584                    .map_err(|_| ProviderError::Database(DatabaseError::Decode))?;
2585
2586                let Some(value_bytes) = iter.value() else { break };
2587                let value = BlockNumberList::decompress(value_bytes)
2588                    .map_err(|_| ProviderError::Database(DatabaseError::Decode))?;
2589
2590                let highest = key.sharded_key.highest_block_number;
2591                shards.push((key, value));
2592
2593                iter.next();
2594
2595                // Shards are ordered by their highest block and their contents partition the
2596                // key's history, so this is the last shard holding anything at or below the
2597                // target. Peek past it only to tell whether it was the key's last shard, which
2598                // decides whether a survivor may be re-keyed to the sentinel.
2599                if highest > *to_block {
2600                    shards_complete = iter.key().and_then(|next_key| next_key.get(..PREFIX_LEN)) !=
2601                        Some(target_prefix);
2602                    break;
2603                }
2604            }
2605
2606            // The iterator also goes invalid on a read error, which would otherwise pass a
2607            // truncated shard list off as the key's complete one.
2608            if !iter.valid() {
2609                iter.status().map_err(|e| {
2610                    ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
2611                        message: e.to_string().into(),
2612                        code: -1,
2613                    }))
2614                })?;
2615            }
2616
2617            // Use existing prune_history_shards_inner logic
2618            match self.prune_history_shards_inner(
2619                shards,
2620                shards_complete,
2621                *to_block,
2622                |key| key.sharded_key.highest_block_number,
2623                |key| key.sharded_key.highest_block_number == u64::MAX,
2624                |batch, key| batch.delete::<tables::StoragesHistory>(key),
2625                |batch, key, value| batch.put::<tables::StoragesHistory>(key, value),
2626                || StorageShardedKey::last(*address, *storage_key),
2627            )? {
2628                PruneShardOutcome::Deleted => outcomes.deleted += 1,
2629                PruneShardOutcome::Updated => outcomes.updated += 1,
2630                PruneShardOutcome::Unchanged => outcomes.unchanged += 1,
2631            }
2632        }
2633
2634        Ok(outcomes)
2635    }
2636
2637    /// Unwinds storage history to keep only blocks `<= keep_to`.
2638    ///
2639    /// Handles multi-shard scenarios by:
2640    /// 1. Loading all shards for the `(address, storage_key)` pair
2641    /// 2. Finding the boundary shard containing `keep_to`
2642    /// 3. Deleting all shards after the boundary
2643    /// 4. Truncating the boundary shard to keep only indices `<= keep_to`
2644    /// 5. Ensuring the last shard is keyed with `u64::MAX`
2645    pub fn unwind_storage_history_to(
2646        &mut self,
2647        address: Address,
2648        storage_key: B256,
2649        keep_to: BlockNumber,
2650    ) -> ProviderResult<()> {
2651        let shards = self.provider.storage_history_shards(address, storage_key)?;
2652        if shards.is_empty() {
2653            return Ok(());
2654        }
2655
2656        // Find the first shard that might contain blocks > keep_to.
2657        // A shard is affected if it's the sentinel (u64::MAX) or its highest_block_number > keep_to
2658        let boundary_idx = shards.iter().position(|(key, _)| {
2659            key.sharded_key.highest_block_number == u64::MAX ||
2660                key.sharded_key.highest_block_number > keep_to
2661        });
2662
2663        // Repair path: no shards affected means all blocks <= keep_to, just ensure sentinel exists
2664        let Some(boundary_idx) = boundary_idx else {
2665            let (last_key, last_value) = shards.last().expect("shards is non-empty");
2666            if last_key.sharded_key.highest_block_number != u64::MAX {
2667                self.delete::<tables::StoragesHistory>(last_key.clone())?;
2668                self.put::<tables::StoragesHistory>(
2669                    StorageShardedKey::last(address, storage_key),
2670                    last_value,
2671                )?;
2672            }
2673            return Ok(());
2674        };
2675
2676        // Delete all shards strictly after the boundary (they are entirely > keep_to)
2677        for (key, _) in shards.iter().skip(boundary_idx + 1) {
2678            self.delete::<tables::StoragesHistory>(key.clone())?;
2679        }
2680
2681        // Process the boundary shard: filter out blocks > keep_to
2682        let (boundary_key, boundary_list) = &shards[boundary_idx];
2683
2684        // Delete the boundary shard (we'll either drop it or rewrite at u64::MAX)
2685        self.delete::<tables::StoragesHistory>(boundary_key.clone())?;
2686
2687        // Build truncated list once; check emptiness directly (avoids double iteration)
2688        let new_last =
2689            BlockNumberList::new_pre_sorted(boundary_list.iter().take_while(|&b| b <= keep_to));
2690
2691        if new_last.is_empty() {
2692            // Boundary shard is now empty. Previous shard becomes the last and must be keyed
2693            // u64::MAX.
2694            if boundary_idx == 0 {
2695                // Nothing left for this (address, storage_key) pair
2696                return Ok(());
2697            }
2698
2699            let (prev_key, prev_value) = &shards[boundary_idx - 1];
2700            if prev_key.sharded_key.highest_block_number != u64::MAX {
2701                self.delete::<tables::StoragesHistory>(prev_key.clone())?;
2702                self.put::<tables::StoragesHistory>(
2703                    StorageShardedKey::last(address, storage_key),
2704                    prev_value,
2705                )?;
2706            }
2707            return Ok(());
2708        }
2709
2710        self.put::<tables::StoragesHistory>(
2711            StorageShardedKey::last(address, storage_key),
2712            &new_last,
2713        )?;
2714
2715        Ok(())
2716    }
2717
2718    /// Clears all account history shards for the given address.
2719    ///
2720    /// Used when unwinding from block 0 (i.e., removing all history).
2721    pub fn clear_account_history(&mut self, address: Address) -> ProviderResult<()> {
2722        let shards = self.provider.account_history_shards(address)?;
2723        for (key, _) in shards {
2724            self.delete::<tables::AccountsHistory>(key)?;
2725        }
2726        Ok(())
2727    }
2728
2729    /// Clears all storage history shards for the given `(address, storage_key)` pair.
2730    ///
2731    /// Used when unwinding from block 0 (i.e., removing all history for this storage slot).
2732    pub fn clear_storage_history(
2733        &mut self,
2734        address: Address,
2735        storage_key: B256,
2736    ) -> ProviderResult<()> {
2737        let shards = self.provider.storage_history_shards(address, storage_key)?;
2738        for (key, _) in shards {
2739            self.delete::<tables::StoragesHistory>(key)?;
2740        }
2741        Ok(())
2742    }
2743}
2744
2745/// `RocksDB` transaction wrapper providing MDBX-like semantics.
2746///
2747/// Supports:
2748/// - Read-your-writes: reads see uncommitted writes within the same transaction
2749/// - Atomic commit/rollback
2750/// - Iteration over uncommitted data
2751///
2752/// Note: `Transaction` is `Send` but NOT `Sync`. This wrapper does not implement
2753/// `DbTx`/`DbTxMut` traits directly; use RocksDB-specific methods instead.
2754pub struct RocksTx<'db> {
2755    inner: Transaction<'db, OptimisticTransactionDB>,
2756    provider: &'db RocksDBProvider,
2757}
2758
2759impl fmt::Debug for RocksTx<'_> {
2760    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2761        f.debug_struct("RocksTx").field("provider", &self.provider).finish_non_exhaustive()
2762    }
2763}
2764
2765impl<'db> RocksTx<'db> {
2766    /// Gets a value from the specified table. Sees uncommitted writes in this transaction.
2767    pub fn get<T: Table>(&self, key: T::Key) -> ProviderResult<Option<T::Value>> {
2768        let encoded_key = key.encode();
2769        self.get_encoded::<T>(&encoded_key)
2770    }
2771
2772    /// Gets a value using pre-encoded key. Sees uncommitted writes in this transaction.
2773    pub fn get_encoded<T: Table>(
2774        &self,
2775        key: &<T::Key as Encode>::Encoded,
2776    ) -> ProviderResult<Option<T::Value>> {
2777        let cf = self.provider.get_cf_handle::<T>()?;
2778        let result = self.inner.get_cf(cf, key.as_ref()).map_err(|e| {
2779            ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
2780                message: e.to_string().into(),
2781                code: -1,
2782            }))
2783        })?;
2784
2785        Ok(result.and_then(|value| T::Value::decompress(&value).ok()))
2786    }
2787
2788    /// Puts a value into the specified table.
2789    pub fn put<T: Table>(&self, key: T::Key, value: &T::Value) -> ProviderResult<()> {
2790        let encoded_key = key.encode();
2791        self.put_encoded::<T>(&encoded_key, value)
2792    }
2793
2794    /// Puts a value using pre-encoded key.
2795    pub fn put_encoded<T: Table>(
2796        &self,
2797        key: &<T::Key as Encode>::Encoded,
2798        value: &T::Value,
2799    ) -> ProviderResult<()> {
2800        let cf = self.provider.get_cf_handle::<T>()?;
2801        let mut buf = Vec::new();
2802        let value_bytes = compress_to_buf_or_ref!(buf, value).unwrap_or(&buf);
2803
2804        self.inner.put_cf(cf, key.as_ref(), value_bytes).map_err(|e| {
2805            ProviderError::Database(DatabaseError::Write(Box::new(DatabaseWriteError {
2806                info: DatabaseErrorInfo { message: e.to_string().into(), code: -1 },
2807                operation: DatabaseWriteOperation::PutUpsert,
2808                table_name: T::NAME,
2809                key: key.as_ref().to_vec(),
2810            })))
2811        })
2812    }
2813
2814    /// Deletes a value from the specified table.
2815    pub fn delete<T: Table>(&self, key: T::Key) -> ProviderResult<()> {
2816        let cf = self.provider.get_cf_handle::<T>()?;
2817        self.inner.delete_cf(cf, key.encode().as_ref()).map_err(|e| {
2818            ProviderError::Database(DatabaseError::Delete(DatabaseErrorInfo {
2819                message: e.to_string().into(),
2820                code: -1,
2821            }))
2822        })
2823    }
2824
2825    /// Creates an iterator for the specified table. Sees uncommitted writes in this transaction.
2826    ///
2827    /// Returns an iterator that yields `(encoded_key, compressed_value)` pairs.
2828    pub fn iter<T: Table>(&self) -> ProviderResult<RocksTxIter<'_, T>> {
2829        let cf = self.provider.get_cf_handle::<T>()?;
2830        let iter = self.inner.iterator_cf(cf, IteratorMode::Start);
2831        Ok(RocksTxIter { inner: iter, _marker: std::marker::PhantomData })
2832    }
2833
2834    /// Creates an iterator starting from the given key (inclusive).
2835    pub fn iter_from<T: Table>(&self, key: T::Key) -> ProviderResult<RocksTxIter<'_, T>> {
2836        let cf = self.provider.get_cf_handle::<T>()?;
2837        let encoded_key = key.encode();
2838        let iter = self
2839            .inner
2840            .iterator_cf(cf, IteratorMode::From(encoded_key.as_ref(), rocksdb::Direction::Forward));
2841        Ok(RocksTxIter { inner: iter, _marker: std::marker::PhantomData })
2842    }
2843
2844    /// Commits the transaction, persisting all changes.
2845    #[instrument(level = "debug", target = "providers::rocksdb", skip_all)]
2846    pub fn commit(self) -> ProviderResult<()> {
2847        self.inner.commit().map_err(|e| {
2848            ProviderError::Database(DatabaseError::Commit(DatabaseErrorInfo {
2849                message: e.to_string().into(),
2850                code: -1,
2851            }))
2852        })
2853    }
2854
2855    /// Rolls back the transaction, discarding all changes.
2856    #[instrument(level = "debug", target = "providers::rocksdb", skip_all)]
2857    pub fn rollback(self) -> ProviderResult<()> {
2858        self.inner.rollback().map_err(|e| {
2859            ProviderError::Database(DatabaseError::Other(format!("rollback failed: {e}")))
2860        })
2861    }
2862}
2863
2864/// Wrapper enum for `RocksDB` iterators that works in both read-write and read-only modes.
2865enum RocksDBIterEnum<'db> {
2866    /// Iterator from read-write `OptimisticTransactionDB`.
2867    ReadWrite(rocksdb::DBIteratorWithThreadMode<'db, OptimisticTransactionDB>),
2868    /// Iterator from read-only `DB`.
2869    ReadOnly(rocksdb::DBIteratorWithThreadMode<'db, DB>),
2870}
2871
2872impl Iterator for RocksDBIterEnum<'_> {
2873    type Item = Result<(Box<[u8]>, Box<[u8]>), rocksdb::Error>;
2874
2875    fn next(&mut self) -> Option<Self::Item> {
2876        match self {
2877            Self::ReadWrite(iter) => iter.next(),
2878            Self::ReadOnly(iter) => iter.next(),
2879        }
2880    }
2881}
2882
2883/// Wrapper enum for raw `RocksDB` iterators that works in both read-write and read-only modes.
2884///
2885/// Unlike [`RocksDBIterEnum`], raw iterators expose `seek()` for efficient repositioning
2886/// without reinitializing the iterator.
2887enum RocksDBRawIterEnum<'db> {
2888    /// Raw iterator from read-write `OptimisticTransactionDB`.
2889    ReadWrite(DBRawIteratorWithThreadMode<'db, OptimisticTransactionDB>),
2890    /// Raw iterator from read-only `DB`.
2891    ReadOnly(DBRawIteratorWithThreadMode<'db, DB>),
2892}
2893
2894impl RocksDBRawIterEnum<'_> {
2895    /// Positions the iterator at the first key >= `key`.
2896    fn seek(&mut self, key: impl AsRef<[u8]>) {
2897        match self {
2898            Self::ReadWrite(iter) => iter.seek(key),
2899            Self::ReadOnly(iter) => iter.seek(key),
2900        }
2901    }
2902
2903    /// Returns true if the iterator is positioned at a valid key-value pair.
2904    fn valid(&self) -> bool {
2905        match self {
2906            Self::ReadWrite(iter) => iter.valid(),
2907            Self::ReadOnly(iter) => iter.valid(),
2908        }
2909    }
2910
2911    /// Returns the current key, if valid.
2912    fn key(&self) -> Option<&[u8]> {
2913        match self {
2914            Self::ReadWrite(iter) => iter.key(),
2915            Self::ReadOnly(iter) => iter.key(),
2916        }
2917    }
2918
2919    /// Returns the current value, if valid.
2920    fn value(&self) -> Option<&[u8]> {
2921        match self {
2922            Self::ReadWrite(iter) => iter.value(),
2923            Self::ReadOnly(iter) => iter.value(),
2924        }
2925    }
2926
2927    /// Advances the iterator to the next key.
2928    fn next(&mut self) {
2929        match self {
2930            Self::ReadWrite(iter) => iter.next(),
2931            Self::ReadOnly(iter) => iter.next(),
2932        }
2933    }
2934
2935    /// Moves the iterator to the previous key.
2936    fn prev(&mut self) {
2937        match self {
2938            Self::ReadWrite(iter) => iter.prev(),
2939            Self::ReadOnly(iter) => iter.prev(),
2940        }
2941    }
2942
2943    /// Returns the status of the iterator.
2944    fn status(&self) -> Result<(), rocksdb::Error> {
2945        match self {
2946            Self::ReadWrite(iter) => iter.status(),
2947            Self::ReadOnly(iter) => iter.status(),
2948        }
2949    }
2950}
2951
2952/// Iterator over a `RocksDB` table (non-transactional).
2953///
2954/// Yields decoded `(Key, Value)` pairs in key order.
2955pub struct RocksDBIter<'db, T: Table> {
2956    inner: RocksDBIterEnum<'db>,
2957    _marker: std::marker::PhantomData<T>,
2958}
2959
2960impl<T: Table> fmt::Debug for RocksDBIter<'_, T> {
2961    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2962        f.debug_struct("RocksDBIter").field("table", &T::NAME).finish_non_exhaustive()
2963    }
2964}
2965
2966impl<T: Table> Iterator for RocksDBIter<'_, T> {
2967    type Item = ProviderResult<(T::Key, T::Value)>;
2968
2969    fn next(&mut self) -> Option<Self::Item> {
2970        Some(decode_iter_item::<T>(self.inner.next()?))
2971    }
2972}
2973
2974/// Raw iterator over a `RocksDB` table (non-transactional).
2975///
2976/// Yields raw `(key_bytes, value_bytes)` pairs without decoding.
2977pub struct RocksDBRawIter<'db> {
2978    inner: RocksDBIterEnum<'db>,
2979}
2980
2981impl fmt::Debug for RocksDBRawIter<'_> {
2982    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2983        f.debug_struct("RocksDBRawIter").finish_non_exhaustive()
2984    }
2985}
2986
2987impl Iterator for RocksDBRawIter<'_> {
2988    type Item = ProviderResult<(Box<[u8]>, Box<[u8]>)>;
2989
2990    fn next(&mut self) -> Option<Self::Item> {
2991        match self.inner.next()? {
2992            Ok(kv) => Some(Ok(kv)),
2993            Err(e) => Some(Err(ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
2994                message: e.to_string().into(),
2995                code: -1,
2996            })))),
2997        }
2998    }
2999}
3000
3001/// Raw key iterator over a `RocksDB` table (non-transactional).
3002pub(crate) struct RocksDBRawKeyIter<'db> {
3003    inner: RocksDBRawIterEnum<'db>,
3004}
3005
3006impl fmt::Debug for RocksDBRawKeyIter<'_> {
3007    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3008        f.debug_struct("RocksDBRawKeyIter").finish_non_exhaustive()
3009    }
3010}
3011
3012impl Iterator for RocksDBRawKeyIter<'_> {
3013    type Item = ProviderResult<Box<[u8]>>;
3014
3015    fn next(&mut self) -> Option<Self::Item> {
3016        if !self.inner.valid() {
3017            return self.inner.status().err().map(|e| {
3018                Err(ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
3019                    message: e.to_string().into(),
3020                    code: -1,
3021                })))
3022            })
3023        }
3024
3025        let Some(key) = self.inner.key() else {
3026            return Some(Err(ProviderError::Database(DatabaseError::Decode)))
3027        };
3028        let key = Box::from(key);
3029        self.inner.next();
3030        Some(Ok(key))
3031    }
3032}
3033
3034/// Iterator over a `RocksDB` table within a transaction.
3035///
3036/// Yields decoded `(Key, Value)` pairs. Sees uncommitted writes.
3037pub struct RocksTxIter<'tx, T: Table> {
3038    inner: rocksdb::DBIteratorWithThreadMode<'tx, Transaction<'tx, OptimisticTransactionDB>>,
3039    _marker: std::marker::PhantomData<T>,
3040}
3041
3042impl<T: Table> fmt::Debug for RocksTxIter<'_, T> {
3043    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3044        f.debug_struct("RocksTxIter").field("table", &T::NAME).finish_non_exhaustive()
3045    }
3046}
3047
3048impl<T: Table> Iterator for RocksTxIter<'_, T> {
3049    type Item = ProviderResult<(T::Key, T::Value)>;
3050
3051    fn next(&mut self) -> Option<Self::Item> {
3052        Some(decode_iter_item::<T>(self.inner.next()?))
3053    }
3054}
3055
3056/// Decodes a raw key-value pair from a `RocksDB` iterator into typed table entries.
3057///
3058/// Handles both error propagation from the underlying iterator and
3059/// decoding/decompression of the key and value bytes.
3060fn decode_iter_item<T: Table>(result: RawKVResult) -> ProviderResult<(T::Key, T::Value)> {
3061    let (key_bytes, value_bytes) = result.map_err(|e| {
3062        ProviderError::Database(DatabaseError::Read(DatabaseErrorInfo {
3063            message: e.to_string().into(),
3064            code: -1,
3065        }))
3066    })?;
3067
3068    let key = <T::Key as reth_db_api::table::Decode>::decode(&key_bytes)
3069        .map_err(|_| ProviderError::Database(DatabaseError::Decode))?;
3070
3071    let value = T::Value::decompress(&value_bytes)
3072        .map_err(|_| ProviderError::Database(DatabaseError::Decode))?;
3073
3074    Ok((key, value))
3075}
3076
3077/// Converts Reth's [`LogLevel`] to `RocksDB`'s [`rocksdb::LogLevel`].
3078const fn convert_log_level(level: LogLevel) -> rocksdb::LogLevel {
3079    match level {
3080        LogLevel::Fatal => rocksdb::LogLevel::Fatal,
3081        LogLevel::Error => rocksdb::LogLevel::Error,
3082        LogLevel::Warn => rocksdb::LogLevel::Warn,
3083        LogLevel::Notice | LogLevel::Verbose => rocksdb::LogLevel::Info,
3084        LogLevel::Debug | LogLevel::Trace | LogLevel::Extra => rocksdb::LogLevel::Debug,
3085    }
3086}
3087
3088/// Selects the `RocksDB` `max_open_files` setting from the current file descriptor limit
3089/// balancing performance vs. compatibility.
3090///
3091/// A mature database can use tens of thousands of file descriptors. With a
3092/// finite `max_open_files` value, `RocksDB` performs additional table-cache checks
3093/// even when enough capacity is available. The value `-1` keeps all table files
3094/// open and avoids this overhead.
3095///
3096/// During normal startup, Reth raises the soft file descriptor limit to the
3097/// system hard limit before it opens `RocksDB`. Therefore, we expect the current
3098/// limit to be the hard limit. A survey of modern Linux distributions shows
3099/// common default limits of `1024:524288`.
3100///
3101/// We use `-1` only when the current limit is above a conservative threshold.
3102/// This keeps enough file descriptors for other parts of the process. When the
3103/// limit is lower or cannot be read, we use a stricter fixed limit.
3104fn select_max_open_files() -> i32 {
3105    let file_descriptor_limit = current_file_descriptor_limit();
3106    let max_open_files = max_open_files_for_limit(file_descriptor_limit);
3107
3108    if max_open_files == LIMITED_MAX_OPEN_FILES {
3109        tracing::warn!(
3110            target: "providers::rocksdb",
3111            ?file_descriptor_limit,
3112            threshold = HIGH_FILE_DESCRIPTOR_LIMIT,
3113            max_open_files,
3114            "RocksDB will not keep all files open; performance may be reduced"
3115        );
3116    }
3117
3118    max_open_files
3119}
3120
3121const fn max_open_files_for_limit(file_descriptor_limit: Option<u64>) -> i32 {
3122    match file_descriptor_limit {
3123        Some(limit) if limit >= HIGH_FILE_DESCRIPTOR_LIMIT => KEEP_ALL_FILES_OPEN,
3124        _ => LIMITED_MAX_OPEN_FILES,
3125    }
3126}
3127
3128#[cfg(unix)]
3129#[allow(clippy::useless_conversion)]
3130fn current_file_descriptor_limit() -> Option<u64> {
3131    let mut limit = libc::rlimit { rlim_cur: 0, rlim_max: 0 };
3132    // SAFETY: `limit` points to initialized writable memory for the kernel to populate.
3133    let result = unsafe { libc::getrlimit(libc::RLIMIT_NOFILE, &raw mut limit) };
3134    (result == 0).then(|| limit.rlim_cur.into())
3135}
3136
3137#[cfg(not(unix))]
3138const fn current_file_descriptor_limit() -> Option<u64> {
3139    None
3140}
3141
3142#[cfg(test)]
3143mod tests {
3144    use super::*;
3145    use crate::providers::HistoryInfo;
3146    use alloy_primitives::{Address, Bytes, TxHash, B256};
3147    use reth_db_api::{
3148        models::{
3149            sharded_key::{ShardedKey, NUM_OF_INDICES_IN_SHARD},
3150            storage_sharded_key::StorageShardedKey,
3151            IntegerList,
3152        },
3153        table::Table,
3154        tables,
3155    };
3156    use tempfile::TempDir;
3157
3158    #[test]
3159    fn max_open_files_adapts_to_file_descriptor_limit() {
3160        assert_eq!(max_open_files_for_limit(Some(HIGH_FILE_DESCRIPTOR_LIMIT)), KEEP_ALL_FILES_OPEN);
3161        assert_eq!(
3162            max_open_files_for_limit(Some(HIGH_FILE_DESCRIPTOR_LIMIT - 1)),
3163            LIMITED_MAX_OPEN_FILES
3164        );
3165        assert_eq!(max_open_files_for_limit(None), LIMITED_MAX_OPEN_FILES);
3166    }
3167
3168    #[test]
3169    fn test_with_default_tables_registers_required_column_families() {
3170        let temp_dir = TempDir::new().unwrap();
3171
3172        // Build with default tables
3173        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3174
3175        // Should be able to write/read TransactionHashNumbers
3176        let tx_hash = TxHash::from(B256::from([1u8; 32]));
3177        provider.put::<tables::TransactionHashNumbers>(tx_hash, &100).unwrap();
3178        assert_eq!(provider.get::<tables::TransactionHashNumbers>(tx_hash).unwrap(), Some(100));
3179
3180        // Should be able to write/read AccountsHistory
3181        let key = ShardedKey::new(Address::ZERO, 100);
3182        let value = IntegerList::default();
3183        provider.put::<tables::AccountsHistory>(key.clone(), &value).unwrap();
3184        assert!(provider.get::<tables::AccountsHistory>(key).unwrap().is_some());
3185
3186        // Should be able to write/read StoragesHistory
3187        let key = StorageShardedKey::new(Address::ZERO, B256::ZERO, 100);
3188        provider.put::<tables::StoragesHistory>(key.clone(), &value).unwrap();
3189        assert!(provider.get::<tables::StoragesHistory>(key).unwrap().is_some());
3190
3191        let bal_key =
3192            reth_db_api::models::StoredBlockAccessListKey::new(1, B256::with_last_byte(1));
3193        let bal_value =
3194            reth_db_api::models::StoredBlockAccessList::new(Bytes::from_static(&[0xc0]));
3195        provider.put::<tables::BlockAccessLists>(bal_key, &bal_value).unwrap();
3196        assert_eq!(provider.get::<tables::BlockAccessLists>(bal_key).unwrap(), Some(bal_value));
3197        provider
3198            .put::<tables::BlockAccessListBlockNumbers>(bal_key.hash(), &bal_key.number())
3199            .unwrap();
3200        assert_eq!(
3201            provider.get::<tables::BlockAccessListBlockNumbers>(bal_key.hash()).unwrap(),
3202            Some(bal_key.number())
3203        );
3204    }
3205
3206    #[test]
3207    fn block_access_lists_store_large_payloads_in_blob_files() {
3208        let temp_dir = TempDir::new().unwrap();
3209        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3210        let bal_key =
3211            reth_db_api::models::StoredBlockAccessListKey::new(1, B256::with_last_byte(1));
3212        let bal_value = reth_db_api::models::StoredBlockAccessList::new(Bytes::from(vec![
3213            0;
3214            DEFAULT_BAL_MIN_BLOB_SIZE as usize +
3215                1
3216        ]));
3217
3218        provider.put::<tables::BlockAccessLists>(bal_key, &bal_value).unwrap();
3219        provider.flush(&[tables::BlockAccessLists::NAME]).unwrap();
3220
3221        let has_blob_file = std::fs::read_dir(temp_dir.path()).unwrap().any(|entry| {
3222            entry.unwrap().path().extension().is_some_and(|extension| extension == "blob")
3223        });
3224        assert!(has_blob_file);
3225    }
3226
3227    #[derive(Debug)]
3228    struct TestTable;
3229
3230    impl Table for TestTable {
3231        const NAME: &'static str = "TestTable";
3232        const DUPSORT: bool = false;
3233        type Key = u64;
3234        type Value = Vec<u8>;
3235    }
3236
3237    #[test]
3238    fn test_reopens_with_unknown_column_family() {
3239        let temp_dir = TempDir::new().unwrap();
3240        let value = b"test_value".to_vec();
3241
3242        let provider = RocksDBBuilder::new(temp_dir.path())
3243            .with_default_tables()
3244            .with_table::<TestTable>()
3245            .build()
3246            .unwrap();
3247        provider.put::<TestTable>(42, &value).unwrap();
3248        drop(provider);
3249
3250        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3251        assert_eq!(provider.get::<TestTable>(42).unwrap(), Some(value));
3252    }
3253
3254    #[test]
3255    fn test_reopens_blob_column_family_with_legacy_table_set() {
3256        let temp_dir = TempDir::new().unwrap();
3257        let bal_key =
3258            reth_db_api::models::StoredBlockAccessListKey::new(1, B256::with_last_byte(1));
3259        let bal_value = reth_db_api::models::StoredBlockAccessList::new(Bytes::from(vec![
3260            0;
3261            DEFAULT_BAL_MIN_BLOB_SIZE as usize +
3262                1
3263        ]));
3264
3265        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3266        provider.put::<tables::BlockAccessLists>(bal_key, &bal_value).unwrap();
3267        provider
3268            .put::<tables::BlockAccessListBlockNumbers>(bal_key.hash(), &bal_key.number())
3269            .unwrap();
3270        provider.flush(&[tables::BlockAccessLists::NAME]).unwrap();
3271        drop(provider);
3272
3273        let provider = RocksDBBuilder::new(temp_dir.path())
3274            .with_table::<tables::TransactionHashNumbers>()
3275            .with_table::<tables::AccountsHistory>()
3276            .with_table::<tables::StoragesHistory>()
3277            .build()
3278            .unwrap();
3279        assert_eq!(provider.get::<tables::BlockAccessLists>(bal_key).unwrap(), Some(bal_value));
3280        assert_eq!(
3281            provider.get::<tables::BlockAccessListBlockNumbers>(bal_key.hash()).unwrap(),
3282            Some(bal_key.number())
3283        );
3284    }
3285
3286    #[test]
3287    fn test_basic_operations() {
3288        let temp_dir = TempDir::new().unwrap();
3289
3290        let provider = RocksDBBuilder::new(temp_dir.path())
3291            .with_table::<TestTable>() // Type-safe!
3292            .build()
3293            .unwrap();
3294
3295        let key = 42u64;
3296        let value = b"test_value".to_vec();
3297
3298        // Test write
3299        provider.put::<TestTable>(key, &value).unwrap();
3300
3301        // Test read
3302        let result = provider.get::<TestTable>(key).unwrap();
3303        assert_eq!(result, Some(value));
3304
3305        // Test delete
3306        provider.delete::<TestTable>(key).unwrap();
3307
3308        // Verify deletion
3309        assert_eq!(provider.get::<TestTable>(key).unwrap(), None);
3310    }
3311
3312    #[test]
3313    fn test_batch_operations() {
3314        let temp_dir = TempDir::new().unwrap();
3315        let provider =
3316            RocksDBBuilder::new(temp_dir.path()).with_table::<TestTable>().build().unwrap();
3317
3318        // Write multiple entries in a batch
3319        provider
3320            .write_batch(|batch| {
3321                for i in 0..10u64 {
3322                    let value = format!("value_{i}").into_bytes();
3323                    batch.put::<TestTable>(i, &value)?;
3324                }
3325                Ok(())
3326            })
3327            .unwrap();
3328
3329        // Read all entries
3330        for i in 0..10u64 {
3331            let value = format!("value_{i}").into_bytes();
3332            assert_eq!(provider.get::<TestTable>(i).unwrap(), Some(value));
3333        }
3334
3335        // Delete all entries in a batch
3336        provider
3337            .write_batch(|batch| {
3338                for i in 0..10u64 {
3339                    batch.delete::<TestTable>(i)?;
3340                }
3341                Ok(())
3342            })
3343            .unwrap();
3344
3345        // Verify all deleted
3346        for i in 0..10u64 {
3347            assert_eq!(provider.get::<TestTable>(i).unwrap(), None);
3348        }
3349    }
3350
3351    #[test]
3352    fn test_with_real_table() {
3353        let temp_dir = TempDir::new().unwrap();
3354        let provider = RocksDBBuilder::new(temp_dir.path())
3355            .with_table::<tables::TransactionHashNumbers>()
3356            .with_metrics()
3357            .build()
3358            .unwrap();
3359
3360        let tx_hash = TxHash::from(B256::from([1u8; 32]));
3361
3362        // Insert and retrieve
3363        provider.put::<tables::TransactionHashNumbers>(tx_hash, &100).unwrap();
3364        assert_eq!(provider.get::<tables::TransactionHashNumbers>(tx_hash).unwrap(), Some(100));
3365
3366        // Batch insert multiple transactions
3367        provider
3368            .write_batch(|batch| {
3369                for i in 0..10u64 {
3370                    let hash = TxHash::from(B256::from([i as u8; 32]));
3371                    let value = i * 100;
3372                    batch.put::<tables::TransactionHashNumbers>(hash, &value)?;
3373                }
3374                Ok(())
3375            })
3376            .unwrap();
3377
3378        // Verify batch insertions
3379        for i in 0..10u64 {
3380            let hash = TxHash::from(B256::from([i as u8; 32]));
3381            assert_eq!(
3382                provider.get::<tables::TransactionHashNumbers>(hash).unwrap(),
3383                Some(i * 100)
3384            );
3385        }
3386    }
3387    #[test]
3388    fn test_statistics_enabled() {
3389        let temp_dir = TempDir::new().unwrap();
3390        // Just verify that building with statistics doesn't panic
3391        let provider = RocksDBBuilder::new(temp_dir.path())
3392            .with_table::<TestTable>()
3393            .with_statistics()
3394            .build()
3395            .unwrap();
3396
3397        // Do operations - data should be immediately readable with OptimisticTransactionDB
3398        for i in 0..10 {
3399            let value = vec![i as u8];
3400            provider.put::<TestTable>(i, &value).unwrap();
3401            // Verify write is visible
3402            assert_eq!(provider.get::<TestTable>(i).unwrap(), Some(value));
3403        }
3404    }
3405
3406    /// Guards the statistics level: the counters must keep working, the per-operation timer
3407    /// histograms must stay empty. `rocksdb`'s `StatsLevel` discriminants do not line up with the
3408    /// levels the C API defines, so the variant name alone does not tell us what was selected.
3409    #[test]
3410    fn test_statistics_level_skips_timers() {
3411        use rocksdb::statistics::{Histogram, Ticker};
3412
3413        let temp_dir = TempDir::new().unwrap();
3414        let cache = Cache::new_lru_cache(1 << 20);
3415        let options = RocksDBBuilder::default_options(rocksdb::LogLevel::Info, &cache, true);
3416
3417        let db = DB::open(&options, temp_dir.path()).unwrap();
3418        for i in 0..10u8 {
3419            db.put([i], [i]).unwrap();
3420            assert_eq!(db.get([i]).unwrap(), Some(vec![i]));
3421        }
3422
3423        assert!(options.get_ticker_count(Ticker::NumberKeysRead) > 0);
3424        assert_eq!(options.get_histogram_data(Histogram::DbGet).count(), 0);
3425    }
3426
3427    #[test]
3428    fn test_data_persistence() {
3429        let temp_dir = TempDir::new().unwrap();
3430        let provider =
3431            RocksDBBuilder::new(temp_dir.path()).with_table::<TestTable>().build().unwrap();
3432
3433        // Insert data - OptimisticTransactionDB writes are immediately visible
3434        let value = vec![42u8; 1000];
3435        for i in 0..100 {
3436            provider.put::<TestTable>(i, &value).unwrap();
3437        }
3438
3439        // Verify data is readable
3440        for i in 0..100 {
3441            assert!(provider.get::<TestTable>(i).unwrap().is_some(), "Data should be readable");
3442        }
3443    }
3444
3445    #[test]
3446    fn test_transaction_read_your_writes() {
3447        let temp_dir = TempDir::new().unwrap();
3448        let provider =
3449            RocksDBBuilder::new(temp_dir.path()).with_table::<TestTable>().build().unwrap();
3450
3451        // Create a transaction
3452        let tx = provider.tx();
3453
3454        // Write data within the transaction
3455        let key = 42u64;
3456        let value = b"test_value".to_vec();
3457        tx.put::<TestTable>(key, &value).unwrap();
3458
3459        // Read-your-writes: should see uncommitted data in same transaction
3460        let result = tx.get::<TestTable>(key).unwrap();
3461        assert_eq!(
3462            result,
3463            Some(value.clone()),
3464            "Transaction should see its own uncommitted writes"
3465        );
3466
3467        // Data should NOT be visible via provider (outside transaction)
3468        let provider_result = provider.get::<TestTable>(key).unwrap();
3469        assert_eq!(provider_result, None, "Uncommitted data should not be visible outside tx");
3470
3471        // Commit the transaction
3472        tx.commit().unwrap();
3473
3474        // Now data should be visible via provider
3475        let committed_result = provider.get::<TestTable>(key).unwrap();
3476        assert_eq!(committed_result, Some(value), "Committed data should be visible");
3477    }
3478
3479    #[test]
3480    fn test_transaction_rollback() {
3481        let temp_dir = TempDir::new().unwrap();
3482        let provider =
3483            RocksDBBuilder::new(temp_dir.path()).with_table::<TestTable>().build().unwrap();
3484
3485        // First, put some initial data
3486        let key = 100u64;
3487        let initial_value = b"initial".to_vec();
3488        provider.put::<TestTable>(key, &initial_value).unwrap();
3489
3490        // Create a transaction and modify data
3491        let tx = provider.tx();
3492        let new_value = b"modified".to_vec();
3493        tx.put::<TestTable>(key, &new_value).unwrap();
3494
3495        // Verify modification is visible within transaction
3496        assert_eq!(tx.get::<TestTable>(key).unwrap(), Some(new_value));
3497
3498        // Rollback instead of commit
3499        tx.rollback().unwrap();
3500
3501        // Data should be unchanged (initial value)
3502        let result = provider.get::<TestTable>(key).unwrap();
3503        assert_eq!(result, Some(initial_value), "Rollback should preserve original data");
3504    }
3505
3506    #[test]
3507    fn test_transaction_iterator() {
3508        let temp_dir = TempDir::new().unwrap();
3509        let provider =
3510            RocksDBBuilder::new(temp_dir.path()).with_table::<TestTable>().build().unwrap();
3511
3512        // Create a transaction
3513        let tx = provider.tx();
3514
3515        // Write multiple entries
3516        for i in 0..5u64 {
3517            let value = format!("value_{i}").into_bytes();
3518            tx.put::<TestTable>(i, &value).unwrap();
3519        }
3520
3521        // Iterate - should see uncommitted writes
3522        let mut count = 0;
3523        for result in tx.iter::<TestTable>().unwrap() {
3524            let (key, value) = result.unwrap();
3525            assert_eq!(value, format!("value_{key}").into_bytes());
3526            count += 1;
3527        }
3528        assert_eq!(count, 5, "Iterator should see all uncommitted writes");
3529
3530        // Commit
3531        tx.commit().unwrap();
3532    }
3533
3534    #[test]
3535    fn test_batch_manual_commit() {
3536        let temp_dir = TempDir::new().unwrap();
3537        let provider =
3538            RocksDBBuilder::new(temp_dir.path()).with_table::<TestTable>().build().unwrap();
3539
3540        // Create a batch via provider.batch()
3541        let mut batch = provider.batch();
3542
3543        // Add entries
3544        for i in 0..10u64 {
3545            let value = format!("batch_value_{i}").into_bytes();
3546            batch.put::<TestTable>(i, &value).unwrap();
3547        }
3548
3549        // Verify len/is_empty
3550        assert_eq!(batch.len(), 10);
3551        assert!(!batch.is_empty());
3552
3553        // Data should NOT be visible before commit
3554        assert_eq!(provider.get::<TestTable>(0).unwrap(), None);
3555
3556        // Commit the batch
3557        batch.commit().unwrap();
3558
3559        // Now data should be visible
3560        for i in 0..10u64 {
3561            let value = format!("batch_value_{i}").into_bytes();
3562            assert_eq!(provider.get::<TestTable>(i).unwrap(), Some(value));
3563        }
3564    }
3565
3566    #[test]
3567    fn test_first_and_last_entry() {
3568        let temp_dir = TempDir::new().unwrap();
3569        let provider =
3570            RocksDBBuilder::new(temp_dir.path()).with_table::<TestTable>().build().unwrap();
3571
3572        // Empty table should return None for both
3573        assert_eq!(provider.first::<TestTable>().unwrap(), None);
3574        assert_eq!(provider.last::<TestTable>().unwrap(), None);
3575
3576        // Insert some entries
3577        provider.put::<TestTable>(10, &b"value_10".to_vec()).unwrap();
3578        provider.put::<TestTable>(20, &b"value_20".to_vec()).unwrap();
3579        provider.put::<TestTable>(5, &b"value_5".to_vec()).unwrap();
3580
3581        // First should return the smallest key
3582        let first = provider.first::<TestTable>().unwrap();
3583        assert_eq!(first, Some((5, b"value_5".to_vec())));
3584
3585        // Last should return the largest key
3586        let last = provider.last::<TestTable>().unwrap();
3587        assert_eq!(last, Some((20, b"value_20".to_vec())));
3588    }
3589
3590    #[test]
3591    fn test_owned_history_snapshot_outlives_provider() {
3592        let temp_dir = TempDir::new().unwrap();
3593        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3594        let address = Address::repeat_byte(0x42);
3595        let slot = B256::repeat_byte(0x43);
3596        let chunk = IntegerList::new([100, 200, 300]).unwrap();
3597        provider
3598            .put::<tables::AccountsHistory>(ShardedKey::new(address, u64::MAX), &chunk)
3599            .unwrap();
3600        provider
3601            .put::<tables::StoragesHistory>(StorageShardedKey::new(address, slot, u64::MAX), &chunk)
3602            .unwrap();
3603
3604        let owned = provider.owned_snapshot();
3605        let snapshot = owned.as_snapshot();
3606        assert_eq!(
3607            snapshot.account_history_info(address, 200, None, u64::MAX).unwrap(),
3608            HistoryInfo::InChangeset(200)
3609        );
3610        assert_eq!(
3611            snapshot.storage_history_info(address, slot, 200, None, u64::MAX).unwrap(),
3612            HistoryInfo::InChangeset(200)
3613        );
3614        drop(provider);
3615
3616        // Move the owner with populated iterator caches, then seek both forwards and backwards.
3617        std::thread::spawn(move || {
3618            let snapshot = owned.as_snapshot();
3619            for (block, expected) in [
3620                (400, HistoryInfo::InPlainState),
3621                (50, HistoryInfo::NotYetWritten),
3622                (200, HistoryInfo::InChangeset(200)),
3623            ] {
3624                assert_eq!(
3625                    snapshot.account_history_info(address, block, None, u64::MAX).unwrap(),
3626                    expected
3627                );
3628                assert_eq!(
3629                    snapshot.storage_history_info(address, slot, block, None, u64::MAX).unwrap(),
3630                    expected
3631                );
3632            }
3633        })
3634        .join()
3635        .unwrap();
3636    }
3637
3638    #[test]
3639    fn test_history_snapshot_cached_and_private_iterators_keep_same_view() {
3640        let temp_dir = TempDir::new().unwrap();
3641        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3642        let address = Address::repeat_byte(0x42);
3643        let slot = B256::repeat_byte(0x43);
3644        let account_key = ShardedKey::new(address, u64::MAX);
3645        let storage_key = StorageShardedKey::new(address, slot, u64::MAX);
3646        let chunk = IntegerList::new([100, 200]).unwrap();
3647        provider.put::<tables::AccountsHistory>(account_key.clone(), &chunk).unwrap();
3648        provider.put::<tables::StoragesHistory>(storage_key.clone(), &chunk).unwrap();
3649
3650        let owned = provider.owned_snapshot();
3651        let snapshot = owned.as_snapshot();
3652        let check = |snapshot: &RocksReadSnapshot<'_>, expected| {
3653            assert_eq!(
3654                snapshot.account_history_info(address, 125, None, u64::MAX).unwrap(),
3655                expected
3656            );
3657            assert_eq!(
3658                snapshot.storage_history_info(address, slot, 125, None, u64::MAX).unwrap(),
3659                expected
3660            );
3661        };
3662        check(snapshot, HistoryInfo::InChangeset(200));
3663
3664        let updated = IntegerList::new([100, 150, 200]).unwrap();
3665        provider.put::<tables::AccountsHistory>(account_key, &updated).unwrap();
3666        provider.put::<tables::StoragesHistory>(storage_key, &updated).unwrap();
3667        check(&provider.snapshot(), HistoryInfo::InChangeset(150));
3668        check(snapshot, HistoryInfo::InChangeset(200));
3669
3670        // Force another thread to use private iterators after the database has changed.
3671        let _accounts = snapshot.accounts_history_iter.lock();
3672        let _storages = snapshot.storages_history_iter.lock();
3673        std::thread::scope(|scope| {
3674            scope.spawn(|| check(snapshot, HistoryInfo::InChangeset(200))).join().unwrap();
3675        });
3676    }
3677
3678    /// Tests the edge case where block < `lowest_available_block_number`.
3679    ///
3680    /// State queries reject this before the `RocksDB` lookup, so this verifies the low-level
3681    /// behavior directly.
3682    #[test]
3683    fn test_account_history_info_pruned_before_first_entry() {
3684        let temp_dir = TempDir::new().unwrap();
3685        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3686
3687        let address = Address::from([0x42; 20]);
3688
3689        // Create a single shard starting at block 100
3690        let chunk = IntegerList::new([100, 200, 300]).unwrap();
3691        let shard_key = ShardedKey::new(address, u64::MAX);
3692        provider.put::<tables::AccountsHistory>(shard_key, &chunk).unwrap();
3693
3694        // Query for block 50 with lowest_available_block_number = 100
3695        // This simulates a pruned state where data before block 100 is not available.
3696        // Since we're before the first write AND pruning boundary is set, we need to
3697        // check the changeset at the first write block.
3698        let result =
3699            provider.snapshot().account_history_info(address, 50, Some(100), u64::MAX).unwrap();
3700        assert_eq!(result, HistoryInfo::InChangeset(100));
3701    }
3702
3703    /// Verifies that a read-only (secondary) provider can catch up with primary writes.
3704    #[test]
3705    fn test_account_history_info_read_only_and_catch_up() {
3706        let temp_dir = TempDir::new().unwrap();
3707        let address = Address::from([0x42; 20]);
3708        let chunk = IntegerList::new([100, 200, 300]).unwrap();
3709        let shard_key = ShardedKey::new(address, u64::MAX);
3710
3711        // Write data with a read-write provider
3712        let rw_provider =
3713            RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3714        rw_provider.put::<tables::AccountsHistory>(shard_key, &chunk).unwrap();
3715
3716        // Open read-only provider — it sees the initial data.
3717        let ro_provider = RocksDBBuilder::new(temp_dir.path())
3718            .with_default_tables()
3719            .with_read_only(true)
3720            .build()
3721            .unwrap();
3722
3723        let result =
3724            ro_provider.snapshot().account_history_info(address, 200, None, u64::MAX).unwrap();
3725        assert_eq!(result, HistoryInfo::InChangeset(200));
3726
3727        let result =
3728            ro_provider.snapshot().account_history_info(address, 50, None, u64::MAX).unwrap();
3729        assert_eq!(result, HistoryInfo::NotYetWritten);
3730
3731        let result =
3732            ro_provider.snapshot().account_history_info(address, 400, None, u64::MAX).unwrap();
3733        assert_eq!(result, HistoryInfo::InPlainState);
3734
3735        // Write new data via the primary.
3736        let address2 = Address::from([0x43; 20]);
3737        let chunk2 = IntegerList::new([500, 600]).unwrap();
3738        let shard_key2 = ShardedKey::new(address2, u64::MAX);
3739        rw_provider.put::<tables::AccountsHistory>(shard_key2, &chunk2).unwrap();
3740
3741        // Read-only doesn't see the new data yet.
3742        let result =
3743            ro_provider.snapshot().account_history_info(address2, 500, None, u64::MAX).unwrap();
3744        assert_eq!(result, HistoryInfo::NotYetWritten);
3745
3746        // Catch up — now it sees the new data.
3747        ro_provider.try_catch_up_with_primary().unwrap();
3748
3749        let result =
3750            ro_provider.snapshot().account_history_info(address2, 500, None, u64::MAX).unwrap();
3751        assert_eq!(result, HistoryInfo::InChangeset(500));
3752    }
3753
3754    #[test]
3755    fn test_account_history_info_ignores_blocks_above_visible_tip() {
3756        let temp_dir = TempDir::new().unwrap();
3757        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3758
3759        let address = Address::from([0x42; 20]);
3760
3761        provider
3762            .put::<tables::AccountsHistory>(
3763                ShardedKey::new(address, 110),
3764                &IntegerList::new([100, 110]).unwrap(),
3765            )
3766            .unwrap();
3767        provider
3768            .put::<tables::AccountsHistory>(
3769                ShardedKey::new(address, u64::MAX),
3770                &IntegerList::new([200, 210]).unwrap(),
3771            )
3772            .unwrap();
3773
3774        let result = provider.snapshot().account_history_info(address, 150, None, 150).unwrap();
3775        assert_eq!(result, HistoryInfo::InPlainState);
3776    }
3777
3778    #[test]
3779    fn test_account_history_info_mixed_shard_respects_visible_tip() {
3780        let temp_dir = TempDir::new().unwrap();
3781        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3782
3783        let address = Address::from([0x42; 20]);
3784        provider
3785            .put::<tables::AccountsHistory>(
3786                ShardedKey::new(address, u64::MAX),
3787                &IntegerList::new([100, 150, 300]).unwrap(),
3788            )
3789            .unwrap();
3790
3791        let result = provider.snapshot().account_history_info(address, 120, None, 200).unwrap();
3792        assert_eq!(result, HistoryInfo::InChangeset(150));
3793
3794        let result = provider.snapshot().account_history_info(address, 201, None, 200).unwrap();
3795        assert_eq!(result, HistoryInfo::InPlainState);
3796    }
3797
3798    #[test]
3799    fn test_account_history_info_only_stale_entries_use_fallback() {
3800        let temp_dir = TempDir::new().unwrap();
3801        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3802
3803        let address = Address::from([0x42; 20]);
3804        provider
3805            .put::<tables::AccountsHistory>(
3806                ShardedKey::new(address, u64::MAX),
3807                &IntegerList::new([200, 210]).unwrap(),
3808            )
3809            .unwrap();
3810
3811        let result = provider.snapshot().account_history_info(address, 150, None, 150).unwrap();
3812        assert_eq!(result, HistoryInfo::NotYetWritten);
3813
3814        let result =
3815            provider.snapshot().account_history_info(address, 150, Some(100), 150).unwrap();
3816        assert_eq!(result, HistoryInfo::MaybeInPlainState);
3817    }
3818
3819    #[test]
3820    fn test_account_history_shard_split_at_boundary() {
3821        let temp_dir = TempDir::new().unwrap();
3822        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3823
3824        let address = Address::from([0x42; 20]);
3825        let limit = NUM_OF_INDICES_IN_SHARD;
3826
3827        // Add exactly NUM_OF_INDICES_IN_SHARD + 1 indices to trigger a split
3828        let indices: Vec<u64> = (0..=(limit as u64)).collect();
3829        let mut batch = provider.batch();
3830        batch.append_account_history_shard(address, indices).unwrap();
3831        batch.commit().unwrap();
3832
3833        // Should have 2 shards: one completed shard and one sentinel shard
3834        let completed_key = ShardedKey::new(address, (limit - 1) as u64);
3835        let sentinel_key = ShardedKey::new(address, u64::MAX);
3836
3837        let completed_shard = provider.get::<tables::AccountsHistory>(completed_key).unwrap();
3838        let sentinel_shard = provider.get::<tables::AccountsHistory>(sentinel_key).unwrap();
3839
3840        assert!(completed_shard.is_some(), "completed shard should exist");
3841        assert!(sentinel_shard.is_some(), "sentinel shard should exist");
3842
3843        let completed_shard = completed_shard.unwrap();
3844        let sentinel_shard = sentinel_shard.unwrap();
3845
3846        assert_eq!(completed_shard.len(), limit as u64, "completed shard should be full");
3847        assert_eq!(sentinel_shard.len(), 1, "sentinel shard should have 1 element");
3848    }
3849
3850    #[test]
3851    fn test_account_history_multiple_shard_splits() {
3852        let temp_dir = TempDir::new().unwrap();
3853        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3854
3855        let address = Address::from([0x43; 20]);
3856        let limit = NUM_OF_INDICES_IN_SHARD;
3857
3858        // First batch: add NUM_OF_INDICES_IN_SHARD indices
3859        let first_batch_indices: Vec<u64> = (0..limit as u64).collect();
3860        let mut batch = provider.batch();
3861        batch.append_account_history_shard(address, first_batch_indices).unwrap();
3862        batch.commit().unwrap();
3863
3864        // Should have just a sentinel shard (exactly at limit, not over)
3865        let sentinel_key = ShardedKey::new(address, u64::MAX);
3866        let shard = provider.get::<tables::AccountsHistory>(sentinel_key.clone()).unwrap();
3867        assert!(shard.is_some());
3868        assert_eq!(shard.unwrap().len(), limit as u64);
3869
3870        // Second batch: add another NUM_OF_INDICES_IN_SHARD + 1 indices (causing 2 more shards)
3871        let second_batch_indices: Vec<u64> = (limit as u64..=(2 * limit) as u64).collect();
3872        let mut batch = provider.batch();
3873        batch.append_account_history_shard(address, second_batch_indices).unwrap();
3874        batch.commit().unwrap();
3875
3876        // Now we should have: 2 completed shards + 1 sentinel shard
3877        let first_completed = ShardedKey::new(address, (limit - 1) as u64);
3878        let second_completed = ShardedKey::new(address, (2 * limit - 1) as u64);
3879
3880        assert!(
3881            provider.get::<tables::AccountsHistory>(first_completed).unwrap().is_some(),
3882            "first completed shard should exist"
3883        );
3884        assert!(
3885            provider.get::<tables::AccountsHistory>(second_completed).unwrap().is_some(),
3886            "second completed shard should exist"
3887        );
3888        assert!(
3889            provider.get::<tables::AccountsHistory>(sentinel_key).unwrap().is_some(),
3890            "sentinel shard should exist"
3891        );
3892    }
3893
3894    #[test]
3895    fn test_storage_history_shard_split_at_boundary() {
3896        let temp_dir = TempDir::new().unwrap();
3897        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3898
3899        let address = Address::from([0x44; 20]);
3900        let slot = B256::from([0x55; 32]);
3901        let limit = NUM_OF_INDICES_IN_SHARD;
3902
3903        // Add exactly NUM_OF_INDICES_IN_SHARD + 1 indices to trigger a split
3904        let indices: Vec<u64> = (0..=(limit as u64)).collect();
3905        let mut batch = provider.batch();
3906        batch.append_storage_history_shard(address, slot, indices).unwrap();
3907        batch.commit().unwrap();
3908
3909        // Should have 2 shards: one completed shard and one sentinel shard
3910        let completed_key = StorageShardedKey::new(address, slot, (limit - 1) as u64);
3911        let sentinel_key = StorageShardedKey::new(address, slot, u64::MAX);
3912
3913        let completed_shard = provider.get::<tables::StoragesHistory>(completed_key).unwrap();
3914        let sentinel_shard = provider.get::<tables::StoragesHistory>(sentinel_key).unwrap();
3915
3916        assert!(completed_shard.is_some(), "completed shard should exist");
3917        assert!(sentinel_shard.is_some(), "sentinel shard should exist");
3918
3919        let completed_shard = completed_shard.unwrap();
3920        let sentinel_shard = sentinel_shard.unwrap();
3921
3922        assert_eq!(completed_shard.len(), limit as u64, "completed shard should be full");
3923        assert_eq!(sentinel_shard.len(), 1, "sentinel shard should have 1 element");
3924    }
3925
3926    #[test]
3927    fn test_storage_history_multiple_shard_splits() {
3928        let temp_dir = TempDir::new().unwrap();
3929        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3930
3931        let address = Address::from([0x46; 20]);
3932        let slot = B256::from([0x57; 32]);
3933        let limit = NUM_OF_INDICES_IN_SHARD;
3934
3935        // First batch: add NUM_OF_INDICES_IN_SHARD indices
3936        let first_batch_indices: Vec<u64> = (0..limit as u64).collect();
3937        let mut batch = provider.batch();
3938        batch.append_storage_history_shard(address, slot, first_batch_indices).unwrap();
3939        batch.commit().unwrap();
3940
3941        // Should have just a sentinel shard (exactly at limit, not over)
3942        let sentinel_key = StorageShardedKey::new(address, slot, u64::MAX);
3943        let shard = provider.get::<tables::StoragesHistory>(sentinel_key.clone()).unwrap();
3944        assert!(shard.is_some());
3945        assert_eq!(shard.unwrap().len(), limit as u64);
3946
3947        // Second batch: add another NUM_OF_INDICES_IN_SHARD + 1 indices (causing 2 more shards)
3948        let second_batch_indices: Vec<u64> = (limit as u64..=(2 * limit) as u64).collect();
3949        let mut batch = provider.batch();
3950        batch.append_storage_history_shard(address, slot, second_batch_indices).unwrap();
3951        batch.commit().unwrap();
3952
3953        // Now we should have: 2 completed shards + 1 sentinel shard
3954        let first_completed = StorageShardedKey::new(address, slot, (limit - 1) as u64);
3955        let second_completed = StorageShardedKey::new(address, slot, (2 * limit - 1) as u64);
3956
3957        assert!(
3958            provider.get::<tables::StoragesHistory>(first_completed).unwrap().is_some(),
3959            "first completed shard should exist"
3960        );
3961        assert!(
3962            provider.get::<tables::StoragesHistory>(second_completed).unwrap().is_some(),
3963            "second completed shard should exist"
3964        );
3965        assert!(
3966            provider.get::<tables::StoragesHistory>(sentinel_key).unwrap().is_some(),
3967            "sentinel shard should exist"
3968        );
3969    }
3970
3971    #[test]
3972    fn test_clear_table() {
3973        let temp_dir = TempDir::new().unwrap();
3974        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3975
3976        let address = Address::from([0x42; 20]);
3977        let key = ShardedKey::new(address, u64::MAX);
3978        let blocks = BlockNumberList::new_pre_sorted([1, 2, 3]);
3979
3980        provider.put::<tables::AccountsHistory>(key.clone(), &blocks).unwrap();
3981        assert!(provider.get::<tables::AccountsHistory>(key.clone()).unwrap().is_some());
3982
3983        provider.clear::<tables::AccountsHistory>().unwrap();
3984
3985        assert!(
3986            provider.get::<tables::AccountsHistory>(key).unwrap().is_none(),
3987            "table should be empty after clear"
3988        );
3989        assert!(
3990            provider.first::<tables::AccountsHistory>().unwrap().is_none(),
3991            "first() should return None after clear"
3992        );
3993    }
3994
3995    #[test]
3996    fn test_clear_empty_table() {
3997        let temp_dir = TempDir::new().unwrap();
3998        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3999
4000        assert!(provider.first::<tables::AccountsHistory>().unwrap().is_none());
4001
4002        provider.clear::<tables::AccountsHistory>().unwrap();
4003
4004        assert!(provider.first::<tables::AccountsHistory>().unwrap().is_none());
4005    }
4006
4007    #[test]
4008    fn test_unwind_account_history_to_basic() {
4009        let temp_dir = TempDir::new().unwrap();
4010        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4011
4012        let address = Address::from([0x42; 20]);
4013
4014        // Add blocks 0-10
4015        let mut batch = provider.batch();
4016        batch.append_account_history_shard(address, 0..=10).unwrap();
4017        batch.commit().unwrap();
4018
4019        // Verify we have blocks 0-10
4020        let key = ShardedKey::new(address, u64::MAX);
4021        let result = provider.get::<tables::AccountsHistory>(key.clone()).unwrap();
4022        assert!(result.is_some());
4023        let blocks: Vec<u64> = result.unwrap().iter().collect();
4024        assert_eq!(blocks, (0..=10).collect::<Vec<_>>());
4025
4026        // Unwind to block 5 (keep blocks 0-5, remove 6-10)
4027        let mut batch = provider.batch();
4028        batch.unwind_account_history_to(address, 5).unwrap();
4029        batch.commit().unwrap();
4030
4031        // Verify only blocks 0-5 remain
4032        let result = provider.get::<tables::AccountsHistory>(key).unwrap();
4033        assert!(result.is_some());
4034        let blocks: Vec<u64> = result.unwrap().iter().collect();
4035        assert_eq!(blocks, (0..=5).collect::<Vec<_>>());
4036    }
4037
4038    #[test]
4039    fn test_unwind_account_history_to_removes_all() {
4040        let temp_dir = TempDir::new().unwrap();
4041        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4042
4043        let address = Address::from([0x42; 20]);
4044
4045        // Add blocks 5-10
4046        let mut batch = provider.batch();
4047        batch.append_account_history_shard(address, 5..=10).unwrap();
4048        batch.commit().unwrap();
4049
4050        // Unwind to block 4 (removes all blocks since they're all > 4)
4051        let mut batch = provider.batch();
4052        batch.unwind_account_history_to(address, 4).unwrap();
4053        batch.commit().unwrap();
4054
4055        // Verify no data remains for this address
4056        let key = ShardedKey::new(address, u64::MAX);
4057        let result = provider.get::<tables::AccountsHistory>(key).unwrap();
4058        assert!(result.is_none(), "Should have no data after full unwind");
4059    }
4060
4061    #[test]
4062    fn test_unwind_account_history_to_no_op() {
4063        let temp_dir = TempDir::new().unwrap();
4064        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4065
4066        let address = Address::from([0x42; 20]);
4067
4068        // Add blocks 0-5
4069        let mut batch = provider.batch();
4070        batch.append_account_history_shard(address, 0..=5).unwrap();
4071        batch.commit().unwrap();
4072
4073        // Unwind to block 10 (no-op since all blocks are <= 10)
4074        let mut batch = provider.batch();
4075        batch.unwind_account_history_to(address, 10).unwrap();
4076        batch.commit().unwrap();
4077
4078        // Verify blocks 0-5 still remain
4079        let key = ShardedKey::new(address, u64::MAX);
4080        let result = provider.get::<tables::AccountsHistory>(key).unwrap();
4081        assert!(result.is_some());
4082        let blocks: Vec<u64> = result.unwrap().iter().collect();
4083        assert_eq!(blocks, (0..=5).collect::<Vec<_>>());
4084    }
4085
4086    #[test]
4087    fn test_unwind_account_history_to_block_zero() {
4088        let temp_dir = TempDir::new().unwrap();
4089        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4090
4091        let address = Address::from([0x42; 20]);
4092
4093        // Add blocks 0-5 (including block 0)
4094        let mut batch = provider.batch();
4095        batch.append_account_history_shard(address, 0..=5).unwrap();
4096        batch.commit().unwrap();
4097
4098        // Unwind to block 0 (keep only block 0, remove 1-5)
4099        // This simulates the caller doing: unwind_to = min_block.checked_sub(1) where min_block = 1
4100        let mut batch = provider.batch();
4101        batch.unwind_account_history_to(address, 0).unwrap();
4102        batch.commit().unwrap();
4103
4104        // Verify only block 0 remains
4105        let key = ShardedKey::new(address, u64::MAX);
4106        let result = provider.get::<tables::AccountsHistory>(key).unwrap();
4107        assert!(result.is_some());
4108        let blocks: Vec<u64> = result.unwrap().iter().collect();
4109        assert_eq!(blocks, vec![0]);
4110    }
4111
4112    #[test]
4113    fn test_unwind_account_history_to_multi_shard() {
4114        let temp_dir = TempDir::new().unwrap();
4115        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4116
4117        let address = Address::from([0x42; 20]);
4118
4119        // Create multiple shards by adding more than NUM_OF_INDICES_IN_SHARD entries
4120        // For testing, we'll manually create shards with specific keys
4121        let mut batch = provider.batch();
4122
4123        // First shard: blocks 1-50, keyed by 50
4124        let shard1 = BlockNumberList::new_pre_sorted(1..=50);
4125        batch.put::<tables::AccountsHistory>(ShardedKey::new(address, 50), &shard1).unwrap();
4126
4127        // Second shard: blocks 51-100, keyed by MAX (sentinel)
4128        let shard2 = BlockNumberList::new_pre_sorted(51..=100);
4129        batch.put::<tables::AccountsHistory>(ShardedKey::new(address, u64::MAX), &shard2).unwrap();
4130
4131        batch.commit().unwrap();
4132
4133        // Verify we have 2 shards
4134        let shards = provider.account_history_shards(address).unwrap();
4135        assert_eq!(shards.len(), 2);
4136
4137        // Unwind to block 75 (keep 1-75, remove 76-100)
4138        let mut batch = provider.batch();
4139        batch.unwind_account_history_to(address, 75).unwrap();
4140        batch.commit().unwrap();
4141
4142        // Verify: shard1 should be untouched, shard2 should be truncated
4143        let shards = provider.account_history_shards(address).unwrap();
4144        assert_eq!(shards.len(), 2);
4145
4146        // First shard unchanged
4147        assert_eq!(shards[0].0.highest_block_number, 50);
4148        assert_eq!(shards[0].1.iter().collect::<Vec<_>>(), (1..=50).collect::<Vec<_>>());
4149
4150        // Second shard truncated and re-keyed to MAX
4151        assert_eq!(shards[1].0.highest_block_number, u64::MAX);
4152        assert_eq!(shards[1].1.iter().collect::<Vec<_>>(), (51..=75).collect::<Vec<_>>());
4153    }
4154
4155    #[test]
4156    fn test_unwind_account_history_to_multi_shard_boundary_empty() {
4157        let temp_dir = TempDir::new().unwrap();
4158        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4159
4160        let address = Address::from([0x42; 20]);
4161
4162        // Create two shards
4163        let mut batch = provider.batch();
4164
4165        // First shard: blocks 1-50, keyed by 50
4166        let shard1 = BlockNumberList::new_pre_sorted(1..=50);
4167        batch.put::<tables::AccountsHistory>(ShardedKey::new(address, 50), &shard1).unwrap();
4168
4169        // Second shard: blocks 75-100, keyed by MAX
4170        let shard2 = BlockNumberList::new_pre_sorted(75..=100);
4171        batch.put::<tables::AccountsHistory>(ShardedKey::new(address, u64::MAX), &shard2).unwrap();
4172
4173        batch.commit().unwrap();
4174
4175        // Unwind to block 60 (removes all of shard2 since 75 > 60, promotes shard1 to MAX)
4176        let mut batch = provider.batch();
4177        batch.unwind_account_history_to(address, 60).unwrap();
4178        batch.commit().unwrap();
4179
4180        // Verify: only shard1 remains, now keyed as MAX
4181        let shards = provider.account_history_shards(address).unwrap();
4182        assert_eq!(shards.len(), 1);
4183        assert_eq!(shards[0].0.highest_block_number, u64::MAX);
4184        assert_eq!(shards[0].1.iter().collect::<Vec<_>>(), (1..=50).collect::<Vec<_>>());
4185    }
4186
4187    #[test]
4188    fn test_account_history_shards_iterator() {
4189        let temp_dir = TempDir::new().unwrap();
4190        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4191
4192        let address = Address::from([0x42; 20]);
4193        let other_address = Address::from([0x43; 20]);
4194
4195        // Add data for two addresses
4196        let mut batch = provider.batch();
4197        batch.append_account_history_shard(address, 0..=5).unwrap();
4198        batch.append_account_history_shard(other_address, 10..=15).unwrap();
4199        batch.commit().unwrap();
4200
4201        // Query shards for first address only
4202        let shards = provider.account_history_shards(address).unwrap();
4203        assert_eq!(shards.len(), 1);
4204        assert_eq!(shards[0].0.key, address);
4205
4206        // Query shards for second address only
4207        let shards = provider.account_history_shards(other_address).unwrap();
4208        assert_eq!(shards.len(), 1);
4209        assert_eq!(shards[0].0.key, other_address);
4210
4211        // Query shards for non-existent address
4212        let non_existent = Address::from([0x99; 20]);
4213        let shards = provider.account_history_shards(non_existent).unwrap();
4214        assert!(shards.is_empty());
4215    }
4216
4217    #[test]
4218    fn test_clear_account_history() {
4219        let temp_dir = TempDir::new().unwrap();
4220        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4221
4222        let address = Address::from([0x42; 20]);
4223
4224        // Add blocks 0-10
4225        let mut batch = provider.batch();
4226        batch.append_account_history_shard(address, 0..=10).unwrap();
4227        batch.commit().unwrap();
4228
4229        // Clear all history (simulates unwind from block 0)
4230        let mut batch = provider.batch();
4231        batch.clear_account_history(address).unwrap();
4232        batch.commit().unwrap();
4233
4234        // Verify no data remains
4235        let shards = provider.account_history_shards(address).unwrap();
4236        assert!(shards.is_empty(), "All shards should be deleted");
4237    }
4238
4239    #[test]
4240    fn test_unwind_non_sentinel_boundary() {
4241        let temp_dir = TempDir::new().unwrap();
4242        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4243
4244        let address = Address::from([0x42; 20]);
4245
4246        // Create three shards with non-sentinel boundary
4247        let mut batch = provider.batch();
4248
4249        // Shard 1: blocks 1-50, keyed by 50
4250        let shard1 = BlockNumberList::new_pre_sorted(1..=50);
4251        batch.put::<tables::AccountsHistory>(ShardedKey::new(address, 50), &shard1).unwrap();
4252
4253        // Shard 2: blocks 51-100, keyed by 100 (non-sentinel, will be boundary)
4254        let shard2 = BlockNumberList::new_pre_sorted(51..=100);
4255        batch.put::<tables::AccountsHistory>(ShardedKey::new(address, 100), &shard2).unwrap();
4256
4257        // Shard 3: blocks 101-150, keyed by MAX (will be deleted)
4258        let shard3 = BlockNumberList::new_pre_sorted(101..=150);
4259        batch.put::<tables::AccountsHistory>(ShardedKey::new(address, u64::MAX), &shard3).unwrap();
4260
4261        batch.commit().unwrap();
4262
4263        // Verify 3 shards
4264        let shards = provider.account_history_shards(address).unwrap();
4265        assert_eq!(shards.len(), 3);
4266
4267        // Unwind to block 75 (truncates shard2, deletes shard3)
4268        let mut batch = provider.batch();
4269        batch.unwind_account_history_to(address, 75).unwrap();
4270        batch.commit().unwrap();
4271
4272        // Verify: shard1 unchanged, shard2 truncated and re-keyed to MAX, shard3 deleted
4273        let shards = provider.account_history_shards(address).unwrap();
4274        assert_eq!(shards.len(), 2);
4275
4276        // First shard unchanged
4277        assert_eq!(shards[0].0.highest_block_number, 50);
4278        assert_eq!(shards[0].1.iter().collect::<Vec<_>>(), (1..=50).collect::<Vec<_>>());
4279
4280        // Second shard truncated and re-keyed to MAX
4281        assert_eq!(shards[1].0.highest_block_number, u64::MAX);
4282        assert_eq!(shards[1].1.iter().collect::<Vec<_>>(), (51..=75).collect::<Vec<_>>());
4283    }
4284
4285    #[test]
4286    fn test_batch_auto_commit_on_threshold() {
4287        let temp_dir = TempDir::new().unwrap();
4288        let provider =
4289            RocksDBBuilder::new(temp_dir.path()).with_table::<TestTable>().build().unwrap();
4290
4291        // Create batch with tiny threshold (1KB) to force auto-commits
4292        let mut batch = RocksDBBatch {
4293            provider: &provider,
4294            inner: WriteBatchWithTransaction::<true>::default(),
4295            buf: Vec::new(),
4296            auto_commit_threshold: Some(1024), // 1KB
4297        };
4298
4299        // Write entries until we exceed threshold multiple times
4300        // Each entry is ~20 bytes, so 100 entries = ~2KB = 2 auto-commits
4301        for i in 0..100u64 {
4302            let value = format!("value_{i:04}").into_bytes();
4303            batch.put::<TestTable>(i, &value).unwrap();
4304        }
4305
4306        // Data should already be visible (auto-committed) even before final commit
4307        // At least some entries should be readable
4308        let first_visible = provider.get::<TestTable>(0).unwrap();
4309        assert!(first_visible.is_some(), "Auto-committed data should be visible");
4310
4311        // Final commit for remaining batch
4312        batch.commit().unwrap();
4313
4314        // All entries should now be visible
4315        for i in 0..100u64 {
4316            let value = format!("value_{i:04}").into_bytes();
4317            assert_eq!(provider.get::<TestTable>(i).unwrap(), Some(value));
4318        }
4319    }
4320
4321    // ==================== PARAMETERIZED PRUNE TESTS ====================
4322
4323    /// Test case for account history pruning
4324    struct AccountPruneCase {
4325        name: &'static str,
4326        initial_shards: &'static [(u64, &'static [u64])],
4327        prune_to: u64,
4328        expected_outcome: PruneShardOutcome,
4329        expected_shards: &'static [(u64, &'static [u64])],
4330    }
4331
4332    /// Test case for storage history pruning
4333    struct StoragePruneCase {
4334        name: &'static str,
4335        initial_shards: &'static [(u64, &'static [u64])],
4336        prune_to: u64,
4337        expected_outcome: PruneShardOutcome,
4338        expected_shards: &'static [(u64, &'static [u64])],
4339    }
4340
4341    #[test]
4342    fn test_prune_account_history_cases() {
4343        const MAX: u64 = u64::MAX;
4344        const CASES: &[AccountPruneCase] = &[
4345            AccountPruneCase {
4346                name: "single_shard_truncate",
4347                initial_shards: &[(MAX, &[10, 20, 30, 40])],
4348                prune_to: 25,
4349                expected_outcome: PruneShardOutcome::Updated,
4350                expected_shards: &[(MAX, &[30, 40])],
4351            },
4352            AccountPruneCase {
4353                name: "single_shard_delete_all",
4354                initial_shards: &[(MAX, &[10, 20])],
4355                prune_to: 20,
4356                expected_outcome: PruneShardOutcome::Deleted,
4357                expected_shards: &[],
4358            },
4359            AccountPruneCase {
4360                name: "single_shard_noop",
4361                initial_shards: &[(MAX, &[10, 20])],
4362                prune_to: 5,
4363                expected_outcome: PruneShardOutcome::Unchanged,
4364                expected_shards: &[(MAX, &[10, 20])],
4365            },
4366            AccountPruneCase {
4367                name: "no_shards",
4368                initial_shards: &[],
4369                prune_to: 100,
4370                expected_outcome: PruneShardOutcome::Unchanged,
4371                expected_shards: &[],
4372            },
4373            AccountPruneCase {
4374                name: "multi_shard_truncate_first",
4375                initial_shards: &[(30, &[10, 20, 30]), (MAX, &[40, 50, 60])],
4376                prune_to: 25,
4377                expected_outcome: PruneShardOutcome::Updated,
4378                expected_shards: &[(30, &[30]), (MAX, &[40, 50, 60])],
4379            },
4380            AccountPruneCase {
4381                name: "delete_first_shard_sentinel_unchanged",
4382                initial_shards: &[(20, &[10, 20]), (MAX, &[30, 40])],
4383                prune_to: 20,
4384                expected_outcome: PruneShardOutcome::Deleted,
4385                expected_shards: &[(MAX, &[30, 40])],
4386            },
4387            AccountPruneCase {
4388                name: "multi_shard_delete_all_but_last",
4389                initial_shards: &[(10, &[5, 10]), (20, &[15, 20]), (MAX, &[25, 30])],
4390                prune_to: 22,
4391                expected_outcome: PruneShardOutcome::Deleted,
4392                expected_shards: &[(MAX, &[25, 30])],
4393            },
4394            AccountPruneCase {
4395                name: "mid_shard_preserves_key",
4396                initial_shards: &[(50, &[10, 20, 30, 40, 50]), (MAX, &[60, 70])],
4397                prune_to: 25,
4398                expected_outcome: PruneShardOutcome::Updated,
4399                expected_shards: &[(50, &[30, 40, 50]), (MAX, &[60, 70])],
4400            },
4401            // Equivalence tests
4402            AccountPruneCase {
4403                name: "equiv_delete_early_shards_keep_sentinel",
4404                initial_shards: &[(20, &[10, 15, 20]), (50, &[30, 40, 50]), (MAX, &[60, 70])],
4405                prune_to: 55,
4406                expected_outcome: PruneShardOutcome::Deleted,
4407                expected_shards: &[(MAX, &[60, 70])],
4408            },
4409            AccountPruneCase {
4410                name: "equiv_sentinel_becomes_empty_with_prev",
4411                initial_shards: &[(50, &[30, 40, 50]), (MAX, &[35])],
4412                prune_to: 40,
4413                expected_outcome: PruneShardOutcome::Deleted,
4414                expected_shards: &[(MAX, &[50])],
4415            },
4416            AccountPruneCase {
4417                name: "equiv_all_shards_become_empty",
4418                initial_shards: &[(50, &[30, 40, 50]), (MAX, &[51])],
4419                prune_to: 51,
4420                expected_outcome: PruneShardOutcome::Deleted,
4421                expected_shards: &[],
4422            },
4423            AccountPruneCase {
4424                name: "equiv_non_sentinel_last_shard_promoted",
4425                initial_shards: &[(100, &[50, 75, 100])],
4426                prune_to: 60,
4427                expected_outcome: PruneShardOutcome::Updated,
4428                expected_shards: &[(MAX, &[75, 100])],
4429            },
4430            AccountPruneCase {
4431                name: "equiv_filter_within_shard",
4432                initial_shards: &[(MAX, &[10, 20, 30, 40])],
4433                prune_to: 25,
4434                expected_outcome: PruneShardOutcome::Updated,
4435                expected_shards: &[(MAX, &[30, 40])],
4436            },
4437            AccountPruneCase {
4438                name: "equiv_multi_shard_partial_delete",
4439                initial_shards: &[(20, &[10, 20]), (50, &[30, 40, 50]), (MAX, &[60, 70])],
4440                prune_to: 35,
4441                expected_outcome: PruneShardOutcome::Deleted,
4442                expected_shards: &[(50, &[40, 50]), (MAX, &[60, 70])],
4443            },
4444        ];
4445
4446        let address = Address::from([0x42; 20]);
4447
4448        for case in CASES {
4449            let temp_dir = TempDir::new().unwrap();
4450            let provider =
4451                RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4452
4453            // Setup initial shards
4454            let mut batch = provider.batch();
4455            for (highest, blocks) in case.initial_shards {
4456                let shard = BlockNumberList::new_pre_sorted(blocks.iter().copied());
4457                batch
4458                    .put::<tables::AccountsHistory>(ShardedKey::new(address, *highest), &shard)
4459                    .unwrap();
4460            }
4461            batch.commit().unwrap();
4462
4463            // Prune
4464            let mut batch = provider.batch();
4465            let outcome = batch.prune_account_history_to(address, case.prune_to).unwrap();
4466            batch.commit().unwrap();
4467
4468            // Assert outcome
4469            assert_eq!(outcome, case.expected_outcome, "case '{}': wrong outcome", case.name);
4470
4471            // Assert final shards
4472            let shards = provider.account_history_shards(address).unwrap();
4473            assert_eq!(
4474                shards.len(),
4475                case.expected_shards.len(),
4476                "case '{}': wrong shard count",
4477                case.name
4478            );
4479            for (i, ((key, blocks), (exp_key, exp_blocks))) in
4480                shards.iter().zip(case.expected_shards.iter()).enumerate()
4481            {
4482                assert_eq!(
4483                    key.highest_block_number, *exp_key,
4484                    "case '{}': shard {} wrong key",
4485                    case.name, i
4486                );
4487                assert_eq!(
4488                    blocks.iter().collect::<Vec<_>>(),
4489                    *exp_blocks,
4490                    "case '{}': shard {} wrong blocks",
4491                    case.name,
4492                    i
4493                );
4494            }
4495        }
4496    }
4497
4498    #[test]
4499    fn test_prune_storage_history_cases() {
4500        const MAX: u64 = u64::MAX;
4501        const CASES: &[StoragePruneCase] = &[
4502            StoragePruneCase {
4503                name: "single_shard_truncate",
4504                initial_shards: &[(MAX, &[10, 20, 30, 40])],
4505                prune_to: 25,
4506                expected_outcome: PruneShardOutcome::Updated,
4507                expected_shards: &[(MAX, &[30, 40])],
4508            },
4509            StoragePruneCase {
4510                name: "single_shard_delete_all",
4511                initial_shards: &[(MAX, &[10, 20])],
4512                prune_to: 20,
4513                expected_outcome: PruneShardOutcome::Deleted,
4514                expected_shards: &[],
4515            },
4516            StoragePruneCase {
4517                name: "noop",
4518                initial_shards: &[(MAX, &[10, 20])],
4519                prune_to: 5,
4520                expected_outcome: PruneShardOutcome::Unchanged,
4521                expected_shards: &[(MAX, &[10, 20])],
4522            },
4523            StoragePruneCase {
4524                name: "no_shards",
4525                initial_shards: &[],
4526                prune_to: 100,
4527                expected_outcome: PruneShardOutcome::Unchanged,
4528                expected_shards: &[],
4529            },
4530            StoragePruneCase {
4531                name: "mid_shard_preserves_key",
4532                initial_shards: &[(50, &[10, 20, 30, 40, 50]), (MAX, &[60, 70])],
4533                prune_to: 25,
4534                expected_outcome: PruneShardOutcome::Updated,
4535                expected_shards: &[(50, &[30, 40, 50]), (MAX, &[60, 70])],
4536            },
4537            // Equivalence tests
4538            StoragePruneCase {
4539                name: "equiv_sentinel_promotion",
4540                initial_shards: &[(100, &[50, 75, 100])],
4541                prune_to: 60,
4542                expected_outcome: PruneShardOutcome::Updated,
4543                expected_shards: &[(MAX, &[75, 100])],
4544            },
4545            StoragePruneCase {
4546                name: "equiv_delete_early_shards_keep_sentinel",
4547                initial_shards: &[(20, &[10, 15, 20]), (50, &[30, 40, 50]), (MAX, &[60, 70])],
4548                prune_to: 55,
4549                expected_outcome: PruneShardOutcome::Deleted,
4550                expected_shards: &[(MAX, &[60, 70])],
4551            },
4552            StoragePruneCase {
4553                name: "equiv_sentinel_becomes_empty_with_prev",
4554                initial_shards: &[(50, &[30, 40, 50]), (MAX, &[35])],
4555                prune_to: 40,
4556                expected_outcome: PruneShardOutcome::Deleted,
4557                expected_shards: &[(MAX, &[50])],
4558            },
4559            StoragePruneCase {
4560                name: "equiv_all_shards_become_empty",
4561                initial_shards: &[(50, &[30, 40, 50]), (MAX, &[51])],
4562                prune_to: 51,
4563                expected_outcome: PruneShardOutcome::Deleted,
4564                expected_shards: &[],
4565            },
4566            StoragePruneCase {
4567                name: "equiv_filter_within_shard",
4568                initial_shards: &[(MAX, &[10, 20, 30, 40])],
4569                prune_to: 25,
4570                expected_outcome: PruneShardOutcome::Updated,
4571                expected_shards: &[(MAX, &[30, 40])],
4572            },
4573            StoragePruneCase {
4574                name: "equiv_multi_shard_partial_delete",
4575                initial_shards: &[(20, &[10, 20]), (50, &[30, 40, 50]), (MAX, &[60, 70])],
4576                prune_to: 35,
4577                expected_outcome: PruneShardOutcome::Deleted,
4578                expected_shards: &[(50, &[40, 50]), (MAX, &[60, 70])],
4579            },
4580        ];
4581
4582        let address = Address::from([0x42; 20]);
4583        let storage_key = B256::from([0x01; 32]);
4584
4585        for case in CASES {
4586            let temp_dir = TempDir::new().unwrap();
4587            let provider =
4588                RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4589
4590            // Setup initial shards
4591            let mut batch = provider.batch();
4592            for (highest, blocks) in case.initial_shards {
4593                let shard = BlockNumberList::new_pre_sorted(blocks.iter().copied());
4594                let key = if *highest == MAX {
4595                    StorageShardedKey::last(address, storage_key)
4596                } else {
4597                    StorageShardedKey::new(address, storage_key, *highest)
4598                };
4599                batch.put::<tables::StoragesHistory>(key, &shard).unwrap();
4600            }
4601            batch.commit().unwrap();
4602
4603            // Prune
4604            let mut batch = provider.batch();
4605            let outcome =
4606                batch.prune_storage_history_to(address, storage_key, case.prune_to).unwrap();
4607            batch.commit().unwrap();
4608
4609            // Assert outcome
4610            assert_eq!(outcome, case.expected_outcome, "case '{}': wrong outcome", case.name);
4611
4612            // Assert final shards
4613            let shards = provider.storage_history_shards(address, storage_key).unwrap();
4614            assert_eq!(
4615                shards.len(),
4616                case.expected_shards.len(),
4617                "case '{}': wrong shard count",
4618                case.name
4619            );
4620            for (i, ((key, blocks), (exp_key, exp_blocks))) in
4621                shards.iter().zip(case.expected_shards.iter()).enumerate()
4622            {
4623                assert_eq!(
4624                    key.sharded_key.highest_block_number, *exp_key,
4625                    "case '{}': shard {} wrong key",
4626                    case.name, i
4627                );
4628                assert_eq!(
4629                    blocks.iter().collect::<Vec<_>>(),
4630                    *exp_blocks,
4631                    "case '{}': shard {} wrong blocks",
4632                    case.name,
4633                    i
4634                );
4635            }
4636        }
4637    }
4638
4639    #[test]
4640    fn test_prune_storage_history_does_not_affect_other_slots() {
4641        let temp_dir = TempDir::new().unwrap();
4642        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4643
4644        let address = Address::from([0x42; 20]);
4645        let slot1 = B256::from([0x01; 32]);
4646        let slot2 = B256::from([0x02; 32]);
4647
4648        // Two different storage slots
4649        let mut batch = provider.batch();
4650        batch
4651            .put::<tables::StoragesHistory>(
4652                StorageShardedKey::last(address, slot1),
4653                &BlockNumberList::new_pre_sorted([10u64, 20]),
4654            )
4655            .unwrap();
4656        batch
4657            .put::<tables::StoragesHistory>(
4658                StorageShardedKey::last(address, slot2),
4659                &BlockNumberList::new_pre_sorted([30u64, 40]),
4660            )
4661            .unwrap();
4662        batch.commit().unwrap();
4663
4664        // Prune slot1 to block 20 (deletes all)
4665        let mut batch = provider.batch();
4666        let outcome = batch.prune_storage_history_to(address, slot1, 20).unwrap();
4667        batch.commit().unwrap();
4668
4669        assert_eq!(outcome, PruneShardOutcome::Deleted);
4670
4671        // slot1 should be empty
4672        let shards1 = provider.storage_history_shards(address, slot1).unwrap();
4673        assert!(shards1.is_empty());
4674
4675        // slot2 should be unchanged
4676        let shards2 = provider.storage_history_shards(address, slot2).unwrap();
4677        assert_eq!(shards2.len(), 1);
4678        assert_eq!(shards2[0].1.iter().collect::<Vec<_>>(), vec![30, 40]);
4679    }
4680
4681    #[test]
4682    fn test_prune_invariants() {
4683        // Test invariants: no empty shards, sentinel is always last
4684        let address = Address::from([0x42; 20]);
4685        let storage_key = B256::from([0x01; 32]);
4686
4687        // Test cases that exercise invariants
4688        #[expect(clippy::type_complexity)]
4689        let invariant_cases: &[(&[(u64, &[u64])], u64)] = &[
4690            // Account: shards where middle becomes empty
4691            (&[(10, &[5, 10]), (20, &[15, 20]), (u64::MAX, &[25, 30])], 20),
4692            // Account: non-sentinel shard only, partial prune -> must become sentinel
4693            (&[(100, &[50, 100])], 60),
4694        ];
4695
4696        for (initial_shards, prune_to) in invariant_cases {
4697            // Test account history invariants
4698            {
4699                let temp_dir = TempDir::new().unwrap();
4700                let provider =
4701                    RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4702
4703                let mut batch = provider.batch();
4704                for (highest, blocks) in *initial_shards {
4705                    let shard = BlockNumberList::new_pre_sorted(blocks.iter().copied());
4706                    batch
4707                        .put::<tables::AccountsHistory>(ShardedKey::new(address, *highest), &shard)
4708                        .unwrap();
4709                }
4710                batch.commit().unwrap();
4711
4712                let mut batch = provider.batch();
4713                batch.prune_account_history_to(address, *prune_to).unwrap();
4714                batch.commit().unwrap();
4715
4716                let shards = provider.account_history_shards(address).unwrap();
4717
4718                // Invariant 1: no empty shards
4719                for (key, blocks) in &shards {
4720                    assert!(
4721                        !blocks.is_empty(),
4722                        "Account: empty shard at key {}",
4723                        key.highest_block_number
4724                    );
4725                }
4726
4727                // Invariant 2: last shard is sentinel
4728                if !shards.is_empty() {
4729                    let last = shards.last().unwrap();
4730                    assert_eq!(
4731                        last.0.highest_block_number,
4732                        u64::MAX,
4733                        "Account: last shard must be sentinel"
4734                    );
4735                }
4736            }
4737
4738            // Test storage history invariants
4739            {
4740                let temp_dir = TempDir::new().unwrap();
4741                let provider =
4742                    RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4743
4744                let mut batch = provider.batch();
4745                for (highest, blocks) in *initial_shards {
4746                    let shard = BlockNumberList::new_pre_sorted(blocks.iter().copied());
4747                    let key = if *highest == u64::MAX {
4748                        StorageShardedKey::last(address, storage_key)
4749                    } else {
4750                        StorageShardedKey::new(address, storage_key, *highest)
4751                    };
4752                    batch.put::<tables::StoragesHistory>(key, &shard).unwrap();
4753                }
4754                batch.commit().unwrap();
4755
4756                let mut batch = provider.batch();
4757                batch.prune_storage_history_to(address, storage_key, *prune_to).unwrap();
4758                batch.commit().unwrap();
4759
4760                let shards = provider.storage_history_shards(address, storage_key).unwrap();
4761
4762                // Invariant 1: no empty shards
4763                for (key, blocks) in &shards {
4764                    assert!(
4765                        !blocks.is_empty(),
4766                        "Storage: empty shard at key {}",
4767                        key.sharded_key.highest_block_number
4768                    );
4769                }
4770
4771                // Invariant 2: last shard is sentinel
4772                if !shards.is_empty() {
4773                    let last = shards.last().unwrap();
4774                    assert_eq!(
4775                        last.0.sharded_key.highest_block_number,
4776                        u64::MAX,
4777                        "Storage: last shard must be sentinel"
4778                    );
4779                }
4780            }
4781        }
4782    }
4783
4784    #[test]
4785    fn test_prune_account_history_batch_multiple_sorted_targets() {
4786        let temp_dir = TempDir::new().unwrap();
4787        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4788
4789        let addr1 = Address::from([0x01; 20]);
4790        let addr2 = Address::from([0x02; 20]);
4791        let addr3 = Address::from([0x03; 20]);
4792
4793        // Setup shards for each address
4794        let mut batch = provider.batch();
4795        batch
4796            .put::<tables::AccountsHistory>(
4797                ShardedKey::new(addr1, u64::MAX),
4798                &BlockNumberList::new_pre_sorted([10, 20, 30]),
4799            )
4800            .unwrap();
4801        batch
4802            .put::<tables::AccountsHistory>(
4803                ShardedKey::new(addr2, u64::MAX),
4804                &BlockNumberList::new_pre_sorted([5, 10, 15]),
4805            )
4806            .unwrap();
4807        batch
4808            .put::<tables::AccountsHistory>(
4809                ShardedKey::new(addr3, u64::MAX),
4810                &BlockNumberList::new_pre_sorted([100, 200]),
4811            )
4812            .unwrap();
4813        batch.commit().unwrap();
4814
4815        // Prune all three (sorted by address)
4816        let mut targets = vec![(addr1, 15), (addr2, 10), (addr3, 50)];
4817        targets.sort_by_key(|(addr, _)| *addr);
4818
4819        let mut batch = provider.batch();
4820        let outcomes = batch.prune_account_history_batch(&targets).unwrap();
4821        batch.commit().unwrap();
4822
4823        // addr1: prune <=15, keep [20, 30] -> updated
4824        // addr2: prune <=10, keep [15] -> updated
4825        // addr3: prune <=50, keep [100, 200] -> unchanged
4826        assert_eq!(outcomes.updated, 2);
4827        assert_eq!(outcomes.unchanged, 1);
4828
4829        let shards1 = provider.account_history_shards(addr1).unwrap();
4830        assert_eq!(shards1[0].1.iter().collect::<Vec<_>>(), vec![20, 30]);
4831
4832        let shards2 = provider.account_history_shards(addr2).unwrap();
4833        assert_eq!(shards2[0].1.iter().collect::<Vec<_>>(), vec![15]);
4834
4835        let shards3 = provider.account_history_shards(addr3).unwrap();
4836        assert_eq!(shards3[0].1.iter().collect::<Vec<_>>(), vec![100, 200]);
4837    }
4838
4839    #[test]
4840    fn test_prune_account_history_batch_target_with_no_shards() {
4841        let temp_dir = TempDir::new().unwrap();
4842        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4843
4844        let addr1 = Address::from([0x01; 20]);
4845        let addr2 = Address::from([0x02; 20]); // No shards for this one
4846        let addr3 = Address::from([0x03; 20]);
4847
4848        // Only setup shards for addr1 and addr3
4849        let mut batch = provider.batch();
4850        batch
4851            .put::<tables::AccountsHistory>(
4852                ShardedKey::new(addr1, u64::MAX),
4853                &BlockNumberList::new_pre_sorted([10, 20]),
4854            )
4855            .unwrap();
4856        batch
4857            .put::<tables::AccountsHistory>(
4858                ShardedKey::new(addr3, u64::MAX),
4859                &BlockNumberList::new_pre_sorted([30, 40]),
4860            )
4861            .unwrap();
4862        batch.commit().unwrap();
4863
4864        // Prune all three (addr2 has no shards - tests p > target_prefix case)
4865        let mut targets = vec![(addr1, 15), (addr2, 100), (addr3, 35)];
4866        targets.sort_by_key(|(addr, _)| *addr);
4867
4868        let mut batch = provider.batch();
4869        let outcomes = batch.prune_account_history_batch(&targets).unwrap();
4870        batch.commit().unwrap();
4871
4872        // addr1: updated (keep [20])
4873        // addr2: unchanged (no shards)
4874        // addr3: updated (keep [40])
4875        assert_eq!(outcomes.updated, 2);
4876        assert_eq!(outcomes.unchanged, 1);
4877
4878        let shards1 = provider.account_history_shards(addr1).unwrap();
4879        assert_eq!(shards1[0].1.iter().collect::<Vec<_>>(), vec![20]);
4880
4881        let shards3 = provider.account_history_shards(addr3).unwrap();
4882        assert_eq!(shards3[0].1.iter().collect::<Vec<_>>(), vec![40]);
4883    }
4884
4885    #[test]
4886    fn test_prune_storage_history_batch_multiple_sorted_targets() {
4887        let temp_dir = TempDir::new().unwrap();
4888        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4889
4890        let addr = Address::from([0x42; 20]);
4891        let slot1 = B256::from([0x01; 32]);
4892        let slot2 = B256::from([0x02; 32]);
4893
4894        // Setup shards
4895        let mut batch = provider.batch();
4896        batch
4897            .put::<tables::StoragesHistory>(
4898                StorageShardedKey::new(addr, slot1, u64::MAX),
4899                &BlockNumberList::new_pre_sorted([10, 20, 30]),
4900            )
4901            .unwrap();
4902        batch
4903            .put::<tables::StoragesHistory>(
4904                StorageShardedKey::new(addr, slot2, u64::MAX),
4905                &BlockNumberList::new_pre_sorted([5, 15, 25]),
4906            )
4907            .unwrap();
4908        batch.commit().unwrap();
4909
4910        // Prune both (sorted)
4911        let mut targets = vec![((addr, slot1), 15), ((addr, slot2), 10)];
4912        targets.sort_by_key(|((a, s), _)| (*a, *s));
4913
4914        let mut batch = provider.batch();
4915        let outcomes = batch.prune_storage_history_batch(&targets).unwrap();
4916        batch.commit().unwrap();
4917
4918        assert_eq!(outcomes.updated, 2);
4919
4920        let shards1 = provider.storage_history_shards(addr, slot1).unwrap();
4921        assert_eq!(shards1[0].1.iter().collect::<Vec<_>>(), vec![20, 30]);
4922
4923        let shards2 = provider.storage_history_shards(addr, slot2).unwrap();
4924        assert_eq!(shards2[0].1.iter().collect::<Vec<_>>(), vec![15, 25]);
4925    }
4926
4927    /// Shards for one address, keyed by highest block, as `(highest, blocks)`.
4928    fn account_shard_layout(provider: &RocksDBProvider, address: Address) -> Vec<(u64, Vec<u64>)> {
4929        provider
4930            .account_history_shards(address)
4931            .unwrap()
4932            .into_iter()
4933            .map(|(key, list)| (key.highest_block_number, list.iter().collect::<Vec<_>>()))
4934            .collect()
4935    }
4936
4937    /// Shards for one storage slot, keyed by highest block, as `(highest, blocks)`.
4938    fn storage_shard_layout(
4939        provider: &RocksDBProvider,
4940        address: Address,
4941        storage_key: B256,
4942    ) -> Vec<(u64, Vec<u64>)> {
4943        provider
4944            .storage_history_shards(address, storage_key)
4945            .unwrap()
4946            .into_iter()
4947            .map(|(key, list)| {
4948                (key.sharded_key.highest_block_number, list.iter().collect::<Vec<_>>())
4949            })
4950            .collect()
4951    }
4952
4953    fn seed_three_storage_shards(provider: &RocksDBProvider, address: Address, storage_key: B256) {
4954        let mut batch = provider.batch();
4955        batch
4956            .put::<tables::StoragesHistory>(
4957                StorageShardedKey::new(address, storage_key, 100),
4958                &BlockNumberList::new_pre_sorted([10, 50, 100]),
4959            )
4960            .unwrap();
4961        batch
4962            .put::<tables::StoragesHistory>(
4963                StorageShardedKey::new(address, storage_key, 200),
4964                &BlockNumberList::new_pre_sorted([150, 200]),
4965            )
4966            .unwrap();
4967        batch
4968            .put::<tables::StoragesHistory>(
4969                StorageShardedKey::last(address, storage_key),
4970                &BlockNumberList::new_pre_sorted([250, 300]),
4971            )
4972            .unwrap();
4973        batch.commit().unwrap();
4974    }
4975
4976    #[test]
4977    fn test_prune_storage_history_batch_leaves_shards_above_target_untouched() {
4978        let temp_dir = TempDir::new().unwrap();
4979        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
4980
4981        let addr = Address::from([0x42; 20]);
4982        let slot = B256::from([0x01; 32]);
4983        seed_three_storage_shards(&provider, addr, slot);
4984
4985        // Only the oldest shard holds blocks at or below the target.
4986        let mut batch = provider.batch();
4987        let outcomes = batch.prune_storage_history_batch(&[((addr, slot), 50)]).unwrap();
4988        batch.commit().unwrap();
4989
4990        assert_eq!(outcomes.updated, 1);
4991        // The trimmed shard keeps its own key. Re-keying it to the sentinel here would overwrite
4992        // the sentinel's blocks.
4993        assert_eq!(
4994            storage_shard_layout(&provider, addr, slot),
4995            vec![(100, vec![100]), (200, vec![150, 200]), (u64::MAX, vec![250, 300])]
4996        );
4997    }
4998
4999    #[test]
5000    fn test_prune_storage_history_batch_trims_sentinel_once_earlier_shards_expire() {
5001        let temp_dir = TempDir::new().unwrap();
5002        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
5003
5004        let addr = Address::from([0x42; 20]);
5005        let slot = B256::from([0x01; 32]);
5006        seed_three_storage_shards(&provider, addr, slot);
5007
5008        // Every non-sentinel shard expires whole and the sentinel loses its lowest block.
5009        let mut batch = provider.batch();
5010        let outcomes = batch.prune_storage_history_batch(&[((addr, slot), 250)]).unwrap();
5011        batch.commit().unwrap();
5012
5013        assert_eq!(outcomes.deleted, 1);
5014        assert_eq!(storage_shard_layout(&provider, addr, slot), vec![(u64::MAX, vec![300])]);
5015    }
5016
5017    #[test]
5018    fn test_prune_account_history_batch_leaves_shards_above_target_untouched() {
5019        let temp_dir = TempDir::new().unwrap();
5020        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
5021
5022        let addr = Address::from([0x42; 20]);
5023
5024        let mut batch = provider.batch();
5025        batch
5026            .put::<tables::AccountsHistory>(
5027                ShardedKey::new(addr, 100),
5028                &BlockNumberList::new_pre_sorted([10, 50, 100]),
5029            )
5030            .unwrap();
5031        batch
5032            .put::<tables::AccountsHistory>(
5033                ShardedKey::new(addr, u64::MAX),
5034                &BlockNumberList::new_pre_sorted([250, 300]),
5035            )
5036            .unwrap();
5037        batch.commit().unwrap();
5038
5039        let mut batch = provider.batch();
5040        let outcomes = batch.prune_account_history_batch(&[(addr, 50)]).unwrap();
5041        batch.commit().unwrap();
5042
5043        assert_eq!(outcomes.updated, 1);
5044        assert_eq!(
5045            account_shard_layout(&provider, addr),
5046            vec![(100, vec![100]), (u64::MAX, vec![250, 300])]
5047        );
5048    }
5049
5050    #[test]
5051    fn test_prune_account_history_batch_seeks_after_stopping_early() {
5052        let temp_dir = TempDir::new().unwrap();
5053        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
5054
5055        let addr1 = Address::from([0x01; 20]);
5056        let addr2 = Address::from([0x02; 20]);
5057
5058        let mut batch = provider.batch();
5059        batch
5060            .put::<tables::AccountsHistory>(
5061                ShardedKey::new(addr1, 100),
5062                &BlockNumberList::new_pre_sorted([10, 50, 100]),
5063            )
5064            .unwrap();
5065        batch
5066            .put::<tables::AccountsHistory>(
5067                ShardedKey::new(addr1, u64::MAX),
5068                &BlockNumberList::new_pre_sorted([250, 300]),
5069            )
5070            .unwrap();
5071        batch
5072            .put::<tables::AccountsHistory>(
5073                ShardedKey::new(addr2, u64::MAX),
5074                &BlockNumberList::new_pre_sorted([5, 10, 15]),
5075            )
5076            .unwrap();
5077        batch.commit().unwrap();
5078
5079        // The first target stops on addr1's oldest shard, leaving the iterator on addr1's
5080        // sentinel. The second target must seek past it instead of skipping addr2.
5081        let mut batch = provider.batch();
5082        let outcomes = batch.prune_account_history_batch(&[(addr1, 50), (addr2, 10)]).unwrap();
5083        batch.commit().unwrap();
5084
5085        assert_eq!(outcomes.updated, 2);
5086        assert_eq!(
5087            account_shard_layout(&provider, addr1),
5088            vec![(100, vec![100]), (u64::MAX, vec![250, 300])]
5089        );
5090        assert_eq!(account_shard_layout(&provider, addr2), vec![(u64::MAX, vec![15])]);
5091    }
5092
5093    #[test]
5094    fn test_prune_storage_history_batch_seeks_after_stopping_early() {
5095        let temp_dir = TempDir::new().unwrap();
5096        let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
5097
5098        let addr = Address::from([0x42; 20]);
5099        let slot1 = B256::from([0x01; 32]);
5100        let slot2 = B256::from([0x02; 32]);
5101        seed_three_storage_shards(&provider, addr, slot1);
5102
5103        let mut batch = provider.batch();
5104        batch
5105            .put::<tables::StoragesHistory>(
5106                StorageShardedKey::last(addr, slot2),
5107                &BlockNumberList::new_pre_sorted([20, 40]),
5108            )
5109            .unwrap();
5110        batch.commit().unwrap();
5111
5112        // The first target stops on slot1's oldest shard, leaving the iterator on slot1's next
5113        // shard. The second target must seek past it instead of skipping slot2.
5114        let mut batch = provider.batch();
5115        let outcomes =
5116            batch.prune_storage_history_batch(&[((addr, slot1), 50), ((addr, slot2), 30)]).unwrap();
5117        batch.commit().unwrap();
5118
5119        assert_eq!(outcomes.updated, 2);
5120        assert_eq!(
5121            storage_shard_layout(&provider, addr, slot1),
5122            vec![(100, vec![100]), (200, vec![150, 200]), (u64::MAX, vec![250, 300])]
5123        );
5124        assert_eq!(storage_shard_layout(&provider, addr, slot2), vec![(u64::MAX, vec![40])]);
5125    }
5126}