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
42fn synced_write_options() -> WriteOptions {
44 let mut opts = WriteOptions::default();
45 opts.set_sync(true);
46 opts
47}
48
49pub(crate) type PendingRocksDBBatches = Arc<Mutex<Vec<WriteBatchWithTransaction<true>>>>;
51
52type RawKVResult = Result<(Box<[u8]>, Box<[u8]>), rocksdb::Error>;
54
55#[derive(Debug, Clone)]
57pub struct RocksDBTableStats {
58 pub sst_size_bytes: u64,
60 pub memtable_size_bytes: u64,
62 pub name: String,
64 pub estimated_num_keys: u64,
66 pub estimated_size_bytes: u64,
68 pub pending_compaction_bytes: u64,
70}
71
72#[derive(Debug, Clone)]
76pub struct RocksDBStats {
77 pub tables: Vec<RocksDBTableStats>,
79 pub wal_size_bytes: u64,
83}
84
85#[derive(Clone)]
87pub(crate) struct RocksDBWriteCtx {
88 pub first_block_number: BlockNumber,
90 pub prune_tx_lookup: Option<PruneMode>,
92 pub storage_settings: StorageSettings,
94 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
109const DEFAULT_CACHE_SIZE: usize = 128 << 20;
111
112const DEFAULT_BLOCK_SIZE: usize = 16 * 1024;
114
115const DEFAULT_MAX_BACKGROUND_JOBS: i32 = 6;
117
118const LIMITED_MAX_OPEN_FILES: i32 = 512;
124
125const KEEP_ALL_FILES_OPEN: i32 = -1;
127
128const HIGH_FILE_DESCRIPTOR_LIMIT: u64 = 128 * 1024;
134
135const DEFAULT_BYTES_PER_SYNC: u64 = 1_048_576;
137
138const DEFAULT_WRITE_BUFFER_SIZE: usize = 128 << 20;
144
145const DEFAULT_WRITE_BUFFER_MANAGER_SIZE: usize = 4 * 1024 * 1024 * 1024;
150
151const DEFAULT_COMPRESS_BUF_CAPACITY: usize = 4096;
155
156const DEFAULT_AUTO_COMMIT_THRESHOLD: usize = 512 * 1024 * 1024;
163
164const DEFAULT_BAL_MIN_BLOB_SIZE: u64 = 4 * 1024;
168
169const DEFAULT_BAL_BLOB_FILE_SIZE: u64 = 256 * 1024 * 1024;
171
172pub 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 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 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 table_options.set_block_cache(cache);
216 table_options
217 }
218
219 fn default_options(
221 log_level: rocksdb::LogLevel,
222 cache: &Cache,
223 enable_statistics: bool,
224 ) -> Options {
225 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 options.set_wal_ttl_seconds(0);
250 options.set_wal_size_limit_mb(0);
251
252 if enable_statistics {
254 options.enable_statistics();
255 options.set_statistics_level(StatsLevel::ExceptHistogramOrTimers);
265 }
266
267 options
268 }
269
270 fn default_column_family_options(cache: &Cache) -> Options {
272 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 cf_options.set_compression_type(DBCompressionType::Lz4);
280 cf_options.set_bottommost_compression_type(DBCompressionType::Zstd);
281 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 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 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 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 cf_options.set_compression_type(DBCompressionType::None);
319 cf_options.set_bottommost_compression_type(DBCompressionType::None);
320
321 cf_options
322 }
323
324 pub fn with_table<T: Table>(mut self) -> Self {
326 self.column_families.push(T::NAME.to_string());
327 self
328 }
329
330 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 pub const fn with_metrics(mut self) -> Self {
348 self.enable_metrics = true;
349 self
350 }
351
352 pub const fn with_statistics(mut self) -> Self {
354 self.enable_statistics = true;
355 self
356 }
357
358 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 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 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 pub const fn with_read_only(mut self, read_only: bool) -> Self {
389 self.read_only = read_only;
390 self
391 }
392
393 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 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 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 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 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
503macro_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#[derive(Debug)]
519pub struct RocksDBProvider(Arc<RocksDBProviderInner>);
520
521enum RocksDBProviderInner {
523 ReadWrite {
525 db: OptimisticTransactionDB,
527 metrics: Option<RocksDBMetrics>,
529 },
530 Secondary {
534 db: DB,
536 metrics: Option<RocksDBMetrics>,
538 secondary_path: PathBuf,
540 },
541}
542
543impl RocksDBProviderInner {
544 const fn metrics(&self) -> Option<&RocksDBMetrics> {
546 match self {
547 Self::ReadWrite { metrics, .. } | Self::Secondary { metrics, .. } => metrics.as_ref(),
548 }
549 }
550
551 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 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 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 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 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 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 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 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 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 fn path(&self) -> &Path {
644 match self {
645 Self::ReadWrite { db, .. } => db.path(),
646 Self::Secondary { db, .. } => db.path(),
647 }
648 }
649
650 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 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 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 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 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 metrics.push(("rocksdb.wal_size", self.wal_size_bytes() as f64, vec![]));
815
816 metrics
817 }
818}
819
820impl RocksDBProvider {
821 pub fn new(path: impl AsRef<Path>) -> ProviderResult<Self> {
823 RocksDBBuilder::new(path).build()
824 }
825
826 pub fn builder(path: impl AsRef<Path>) -> RocksDBBuilder {
828 RocksDBBuilder::new(path)
829 }
830
831 pub fn exists(path: impl AsRef<Path>) -> bool {
836 path.as_ref().join("CURRENT").exists()
837 }
838
839 pub fn is_read_only(&self) -> bool {
841 matches!(self.0.as_ref(), RocksDBProviderInner::Secondary { .. })
842 }
843
844 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 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 pub fn owned_snapshot(&self) -> OwnedRocksReadSnapshot {
879 OwnedRocksReadSnapshot::new(self)
880 }
881
882 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 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 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 fn get_cf_handle<T: Table>(&self) -> Result<&rocksdb::ColumnFamily, DatabaseError> {
928 self.0.cf_handle::<T>()
929 }
930
931 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 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 pub fn get<T: Table>(&self, key: T::Key) -> ProviderResult<Option<T::Value>> {
958 self.get_encoded::<T>(&key.encode())
959 }
960
961 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 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 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 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 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 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 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 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 #[inline]
1090 pub fn first<T: Table>(&self) -> ProviderResult<Option<(T::Key, T::Value)>> {
1091 self.get_boundary::<T>(IteratorMode::Start)
1092 }
1093
1094 #[inline]
1096 pub fn last<T: Table>(&self) -> ProviderResult<Option<(T::Key, T::Value)>> {
1097 self.get_boundary::<T>(IteratorMode::End)
1098 }
1099
1100 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 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 pub fn table_stats(&self) -> Vec<RocksDBTableStats> {
1125 self.0.table_stats()
1126 }
1127
1128 pub fn wal_size_bytes(&self) -> u64 {
1134 self.0.wal_size_bytes()
1135 }
1136
1137 pub fn db_stats(&self) -> RocksDBStats {
1141 self.0.db_stats()
1142 }
1143
1144 #[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 #[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 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 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 pub fn account_history_shards(
1235 &self,
1236 address: Address,
1237 ) -> ProviderResult<Vec<(ShardedKey<Address>, BlockNumberList)>> {
1238 let cf = self.get_cf_handle::<tables::AccountsHistory>()?;
1240
1241 let start_key = ShardedKey::new(address, 0u64);
1244 let start_bytes = start_key.encode();
1245
1246 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 let key = ShardedKey::<Address>::decode(&key_bytes)
1257 .map_err(|_| ProviderError::Database(DatabaseError::Decode))?;
1258
1259 if key.key != address {
1261 break;
1262 }
1263
1264 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 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 #[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 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 #[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 #[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 #[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 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 #[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 #[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 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 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 #[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 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 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
1645pub struct RocksReadSnapshot<'db> {
1653 accounts_history_iter: Mutex<Option<RocksDBRawIterEnum<'db>>>,
1657 storages_history_iter: Mutex<Option<RocksDBRawIterEnum<'db>>>,
1661 inner: RocksReadSnapshotInner<'db>,
1662 provider: &'db RocksDBProviderInner,
1663}
1664
1665enum RocksReadSnapshotInner<'db> {
1667 ReadWrite(SnapshotWithThreadMode<'db, OptimisticTransactionDB>),
1669 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 fn cf_handle<T: Table>(&self) -> Result<&'db rocksdb::ColumnFamily, DatabaseError> {
1684 self.provider.cf_handle::<T>()
1685 }
1686
1687 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 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 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 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 #[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 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 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 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
1888pub struct OwnedRocksReadSnapshot {
1894 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 let snapshot = unsafe {
1916 std::mem::transmute::<RocksReadSnapshot<'_>, RocksReadSnapshot<'static>>(snapshot)
1917 };
1918 Self { snapshot: Box::new(snapshot), provider }
1919 }
1920
1921 pub fn as_snapshot(&self) -> &RocksReadSnapshot<'_> {
1923 unsafe {
1926 std::mem::transmute::<&RocksReadSnapshot<'static>, &RocksReadSnapshot<'_>>(
1927 &self.snapshot,
1928 )
1929 }
1930 }
1931}
1932
1933#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1935pub enum PruneShardOutcome {
1936 Deleted,
1938 Updated,
1940 Unchanged,
1942}
1943
1944#[derive(Debug, Default, Clone, Copy)]
1946pub struct PrunedIndices {
1947 pub deleted: usize,
1949 pub updated: usize,
1951 pub unchanged: usize,
1953}
1954
1955#[must_use = "batch must be committed"]
1965pub struct RocksDBBatch<'a> {
1966 provider: &'a RocksDBProvider,
1967 inner: WriteBatchWithTransaction<true>,
1968 buf: Vec<u8>,
1969 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 .field("length", &self.inner.len())
1980 .field("size_in_bytes", &self.inner.size_in_bytes())
1983 .finish()
1984 }
1985}
1986
1987impl<'a> RocksDBBatch<'a> {
1988 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 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 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 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 #[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 pub fn len(&self) -> usize {
2062 self.inner.len()
2063 }
2064
2065 pub fn is_empty(&self) -> bool {
2067 self.inner.is_empty()
2068 }
2069
2070 pub fn size_in_bytes(&self) -> usize {
2072 self.inner.size_in_bytes()
2073 }
2074
2075 pub const fn provider(&self) -> &RocksDBProvider {
2077 self.provider
2078 }
2079
2080 pub fn into_inner(self) -> WriteBatchWithTransaction<true> {
2084 self.inner
2085 }
2086
2087 pub fn get<T: Table>(&self, key: T::Key) -> ProviderResult<Option<T::Value>> {
2092 self.provider.get::<T>(key)
2093 }
2094
2095 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 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 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 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 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 let boundary_idx = shards.iter().position(|(key, _)| {
2203 key.highest_block_number == u64::MAX || key.highest_block_number > keep_to
2204 });
2205
2206 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 for (key, _) in shards.iter().skip(boundary_idx + 1) {
2221 self.delete::<tables::AccountsHistory>(key.clone())?;
2222 }
2223
2224 let (boundary_key, boundary_list) = &shards[boundary_idx];
2226
2227 self.delete::<tables::AccountsHistory>(boundary_key.clone())?;
2229
2230 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 if boundary_idx == 0 {
2238 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 #[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 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 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 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 let start_key = ShardedKey::new(*address, 0u64).encode();
2381 let target_prefix = &start_key[..PREFIX_LEN];
2382
2383 let needs_seek = if iter.valid() {
2391 if let Some(current_key) = iter.key() {
2392 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 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 let current_prefix = key_bytes.get(..PREFIX_LEN);
2422 if current_prefix != Some(target_prefix) {
2423 break;
2424 }
2425
2426 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 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 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 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 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 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 let start_key = StorageShardedKey::new(*address, *storage_key, 0u64).encode();
2537 let target_prefix = &start_key[..PREFIX_LEN];
2538
2539 let needs_seek = if iter.valid() {
2547 if let Some(current_key) = iter.key() {
2548 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 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 let current_prefix = key_bytes.get(..PREFIX_LEN);
2578 if current_prefix != Some(target_prefix) {
2579 break;
2580 }
2581
2582 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 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 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 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 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 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 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 for (key, _) in shards.iter().skip(boundary_idx + 1) {
2678 self.delete::<tables::StoragesHistory>(key.clone())?;
2679 }
2680
2681 let (boundary_key, boundary_list) = &shards[boundary_idx];
2683
2684 self.delete::<tables::StoragesHistory>(boundary_key.clone())?;
2686
2687 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 if boundary_idx == 0 {
2695 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 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 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
2745pub 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 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 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 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 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 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 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 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 #[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 #[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
2864enum RocksDBIterEnum<'db> {
2866 ReadWrite(rocksdb::DBIteratorWithThreadMode<'db, OptimisticTransactionDB>),
2868 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
2883enum RocksDBRawIterEnum<'db> {
2888 ReadWrite(DBRawIteratorWithThreadMode<'db, OptimisticTransactionDB>),
2890 ReadOnly(DBRawIteratorWithThreadMode<'db, DB>),
2892}
2893
2894impl RocksDBRawIterEnum<'_> {
2895 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 fn valid(&self) -> bool {
2905 match self {
2906 Self::ReadWrite(iter) => iter.valid(),
2907 Self::ReadOnly(iter) => iter.valid(),
2908 }
2909 }
2910
2911 fn key(&self) -> Option<&[u8]> {
2913 match self {
2914 Self::ReadWrite(iter) => iter.key(),
2915 Self::ReadOnly(iter) => iter.key(),
2916 }
2917 }
2918
2919 fn value(&self) -> Option<&[u8]> {
2921 match self {
2922 Self::ReadWrite(iter) => iter.value(),
2923 Self::ReadOnly(iter) => iter.value(),
2924 }
2925 }
2926
2927 fn next(&mut self) {
2929 match self {
2930 Self::ReadWrite(iter) => iter.next(),
2931 Self::ReadOnly(iter) => iter.next(),
2932 }
2933 }
2934
2935 fn prev(&mut self) {
2937 match self {
2938 Self::ReadWrite(iter) => iter.prev(),
2939 Self::ReadOnly(iter) => iter.prev(),
2940 }
2941 }
2942
2943 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
2952pub 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
2974pub 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
3001pub(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
3034pub 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
3056fn 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
3077const 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
3088fn 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 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 let provider = RocksDBBuilder::new(temp_dir.path()).with_default_tables().build().unwrap();
3174
3175 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 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 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>() .build()
3293 .unwrap();
3294
3295 let key = 42u64;
3296 let value = b"test_value".to_vec();
3297
3298 provider.put::<TestTable>(key, &value).unwrap();
3300
3301 let result = provider.get::<TestTable>(key).unwrap();
3303 assert_eq!(result, Some(value));
3304
3305 provider.delete::<TestTable>(key).unwrap();
3307
3308 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 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 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 provider
3337 .write_batch(|batch| {
3338 for i in 0..10u64 {
3339 batch.delete::<TestTable>(i)?;
3340 }
3341 Ok(())
3342 })
3343 .unwrap();
3344
3345 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 provider.put::<tables::TransactionHashNumbers>(tx_hash, &100).unwrap();
3364 assert_eq!(provider.get::<tables::TransactionHashNumbers>(tx_hash).unwrap(), Some(100));
3365
3366 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 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 let provider = RocksDBBuilder::new(temp_dir.path())
3392 .with_table::<TestTable>()
3393 .with_statistics()
3394 .build()
3395 .unwrap();
3396
3397 for i in 0..10 {
3399 let value = vec![i as u8];
3400 provider.put::<TestTable>(i, &value).unwrap();
3401 assert_eq!(provider.get::<TestTable>(i).unwrap(), Some(value));
3403 }
3404 }
3405
3406 #[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 let value = vec![42u8; 1000];
3435 for i in 0..100 {
3436 provider.put::<TestTable>(i, &value).unwrap();
3437 }
3438
3439 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 let tx = provider.tx();
3453
3454 let key = 42u64;
3456 let value = b"test_value".to_vec();
3457 tx.put::<TestTable>(key, &value).unwrap();
3458
3459 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 let provider_result = provider.get::<TestTable>(key).unwrap();
3469 assert_eq!(provider_result, None, "Uncommitted data should not be visible outside tx");
3470
3471 tx.commit().unwrap();
3473
3474 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 let key = 100u64;
3487 let initial_value = b"initial".to_vec();
3488 provider.put::<TestTable>(key, &initial_value).unwrap();
3489
3490 let tx = provider.tx();
3492 let new_value = b"modified".to_vec();
3493 tx.put::<TestTable>(key, &new_value).unwrap();
3494
3495 assert_eq!(tx.get::<TestTable>(key).unwrap(), Some(new_value));
3497
3498 tx.rollback().unwrap();
3500
3501 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 let tx = provider.tx();
3514
3515 for i in 0..5u64 {
3517 let value = format!("value_{i}").into_bytes();
3518 tx.put::<TestTable>(i, &value).unwrap();
3519 }
3520
3521 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 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 let mut batch = provider.batch();
3542
3543 for i in 0..10u64 {
3545 let value = format!("batch_value_{i}").into_bytes();
3546 batch.put::<TestTable>(i, &value).unwrap();
3547 }
3548
3549 assert_eq!(batch.len(), 10);
3551 assert!(!batch.is_empty());
3552
3553 assert_eq!(provider.get::<TestTable>(0).unwrap(), None);
3555
3556 batch.commit().unwrap();
3558
3559 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 assert_eq!(provider.first::<TestTable>().unwrap(), None);
3574 assert_eq!(provider.last::<TestTable>().unwrap(), None);
3575
3576 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 let first = provider.first::<TestTable>().unwrap();
3583 assert_eq!(first, Some((5, b"value_5".to_vec())));
3584
3585 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 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 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 #[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 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 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 #[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 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 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 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 let result =
3743 ro_provider.snapshot().account_history_info(address2, 500, None, u64::MAX).unwrap();
3744 assert_eq!(result, HistoryInfo::NotYetWritten);
3745
3746 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 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 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 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 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 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 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 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 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 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 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 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 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 let mut batch = provider.batch();
4016 batch.append_account_history_shard(address, 0..=10).unwrap();
4017 batch.commit().unwrap();
4018
4019 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 let mut batch = provider.batch();
4028 batch.unwind_account_history_to(address, 5).unwrap();
4029 batch.commit().unwrap();
4030
4031 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 let mut batch = provider.batch();
4047 batch.append_account_history_shard(address, 5..=10).unwrap();
4048 batch.commit().unwrap();
4049
4050 let mut batch = provider.batch();
4052 batch.unwind_account_history_to(address, 4).unwrap();
4053 batch.commit().unwrap();
4054
4055 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 let mut batch = provider.batch();
4070 batch.append_account_history_shard(address, 0..=5).unwrap();
4071 batch.commit().unwrap();
4072
4073 let mut batch = provider.batch();
4075 batch.unwind_account_history_to(address, 10).unwrap();
4076 batch.commit().unwrap();
4077
4078 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 let mut batch = provider.batch();
4095 batch.append_account_history_shard(address, 0..=5).unwrap();
4096 batch.commit().unwrap();
4097
4098 let mut batch = provider.batch();
4101 batch.unwind_account_history_to(address, 0).unwrap();
4102 batch.commit().unwrap();
4103
4104 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 let mut batch = provider.batch();
4122
4123 let shard1 = BlockNumberList::new_pre_sorted(1..=50);
4125 batch.put::<tables::AccountsHistory>(ShardedKey::new(address, 50), &shard1).unwrap();
4126
4127 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 let shards = provider.account_history_shards(address).unwrap();
4135 assert_eq!(shards.len(), 2);
4136
4137 let mut batch = provider.batch();
4139 batch.unwind_account_history_to(address, 75).unwrap();
4140 batch.commit().unwrap();
4141
4142 let shards = provider.account_history_shards(address).unwrap();
4144 assert_eq!(shards.len(), 2);
4145
4146 assert_eq!(shards[0].0.highest_block_number, 50);
4148 assert_eq!(shards[0].1.iter().collect::<Vec<_>>(), (1..=50).collect::<Vec<_>>());
4149
4150 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 let mut batch = provider.batch();
4164
4165 let shard1 = BlockNumberList::new_pre_sorted(1..=50);
4167 batch.put::<tables::AccountsHistory>(ShardedKey::new(address, 50), &shard1).unwrap();
4168
4169 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 let mut batch = provider.batch();
4177 batch.unwind_account_history_to(address, 60).unwrap();
4178 batch.commit().unwrap();
4179
4180 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 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 let shards = provider.account_history_shards(address).unwrap();
4203 assert_eq!(shards.len(), 1);
4204 assert_eq!(shards[0].0.key, address);
4205
4206 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 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 let mut batch = provider.batch();
4226 batch.append_account_history_shard(address, 0..=10).unwrap();
4227 batch.commit().unwrap();
4228
4229 let mut batch = provider.batch();
4231 batch.clear_account_history(address).unwrap();
4232 batch.commit().unwrap();
4233
4234 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 let mut batch = provider.batch();
4248
4249 let shard1 = BlockNumberList::new_pre_sorted(1..=50);
4251 batch.put::<tables::AccountsHistory>(ShardedKey::new(address, 50), &shard1).unwrap();
4252
4253 let shard2 = BlockNumberList::new_pre_sorted(51..=100);
4255 batch.put::<tables::AccountsHistory>(ShardedKey::new(address, 100), &shard2).unwrap();
4256
4257 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 let shards = provider.account_history_shards(address).unwrap();
4265 assert_eq!(shards.len(), 3);
4266
4267 let mut batch = provider.batch();
4269 batch.unwind_account_history_to(address, 75).unwrap();
4270 batch.commit().unwrap();
4271
4272 let shards = provider.account_history_shards(address).unwrap();
4274 assert_eq!(shards.len(), 2);
4275
4276 assert_eq!(shards[0].0.highest_block_number, 50);
4278 assert_eq!(shards[0].1.iter().collect::<Vec<_>>(), (1..=50).collect::<Vec<_>>());
4279
4280 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 let mut batch = RocksDBBatch {
4293 provider: &provider,
4294 inner: WriteBatchWithTransaction::<true>::default(),
4295 buf: Vec::new(),
4296 auto_commit_threshold: Some(1024), };
4298
4299 for i in 0..100u64 {
4302 let value = format!("value_{i:04}").into_bytes();
4303 batch.put::<TestTable>(i, &value).unwrap();
4304 }
4305
4306 let first_visible = provider.get::<TestTable>(0).unwrap();
4309 assert!(first_visible.is_some(), "Auto-committed data should be visible");
4310
4311 batch.commit().unwrap();
4313
4314 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 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 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 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 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 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_eq!(outcome, case.expected_outcome, "case '{}': wrong outcome", case.name);
4470
4471 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 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 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 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_eq!(outcome, case.expected_outcome, "case '{}': wrong outcome", case.name);
4611
4612 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 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 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 let shards1 = provider.storage_history_shards(address, slot1).unwrap();
4673 assert!(shards1.is_empty());
4674
4675 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 let address = Address::from([0x42; 20]);
4685 let storage_key = B256::from([0x01; 32]);
4686
4687 #[expect(clippy::type_complexity)]
4689 let invariant_cases: &[(&[(u64, &[u64])], u64)] = &[
4690 (&[(10, &[5, 10]), (20, &[15, 20]), (u64::MAX, &[25, 30])], 20),
4692 (&[(100, &[50, 100])], 60),
4694 ];
4695
4696 for (initial_shards, prune_to) in invariant_cases {
4697 {
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 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 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 {
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 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 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 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 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 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]); let addr3 = Address::from([0x03; 20]);
4847
4848 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 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 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 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 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 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 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 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 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 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 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 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}