Skip to main content

reth_provider/providers/database/
mod.rs

1use crate::{
2    providers::{
3        state::latest::LatestStateProvider, NodeTypesForProvider, RocksDBProvider,
4        StaticFileProvider, StaticFileProviderRWRefMut,
5    },
6    to_range,
7    traits::{BlockSource, ReceiptProvider},
8    BalProvider, BalStoreHandle, BlockHashReader, BlockNumReader, BlockReader, ChainSpecProvider,
9    DatabaseProviderFactory, EitherWriterDestination, HeaderProvider, HeaderSyncGapProvider,
10    InMemoryBalStore, MetadataProvider, ProviderError, PruneCheckpointReader,
11    RocksDBProviderFactory, StageCheckpointReader, StateProviderBox, StaticFileProviderFactory,
12    StaticFileWriter, TransactionVariant, TransactionsProvider,
13};
14use alloy_consensus::transaction::TransactionMeta;
15use alloy_eips::BlockHashOrNumber;
16use alloy_primitives::{Address, BlockHash, BlockNumber, TxHash, TxNumber, B256};
17use core::fmt;
18use notify::{RecommendedWatcher, RecursiveMode, Watcher};
19use parking_lot::RwLock;
20use reth_chainspec::ChainInfo;
21use reth_db::{init_db, mdbx::DatabaseArguments, DatabaseEnv};
22use reth_db_api::{database::Database, models::StoredBlockBodyIndices};
23use reth_errors::{RethError, RethResult};
24use reth_node_types::{
25    BlockTy, HeaderTy, NodeTypesWithDB, NodeTypesWithDBAdapter, ReceiptTy, TxTy,
26};
27use reth_primitives_traits::{RecoveredBlock, SealedHeader};
28use reth_prune_types::{PruneCheckpoint, PruneModes, PruneSegment, MINIMUM_UNWIND_SAFE_DISTANCE};
29use reth_stages_types::{PipelineTarget, StageCheckpoint, StageId};
30use reth_static_file_types::StaticFileSegment;
31use reth_storage_api::{
32    BlockBodyIndicesProvider, ChainStateBlockReader, ChainStateBlockWriter, DBProvider,
33    NodePrimitivesProvider, StorageSettings, StorageSettingsCache, TryIntoHistoricalStateProvider,
34};
35use reth_storage_errors::provider::ProviderResult;
36use reth_storage_overlay::OverlayManager;
37use std::{
38    ops::{RangeBounds, RangeInclusive},
39    path::Path,
40    sync::{
41        atomic::{AtomicU64, Ordering},
42        Arc, Mutex,
43    },
44};
45use tracing::{info, instrument, trace, warn};
46
47mod provider;
48pub use provider::{CommitOrder, DatabaseProvider, DatabaseProviderRO, DatabaseProviderRW};
49
50mod save_blocks;
51pub use save_blocks::SaveBlocksInput;
52
53use super::ProviderNodeTypes;
54mod builder;
55pub use builder::{ProviderFactoryBuilder, ReadOnlyConfig};
56
57mod metrics;
58pub use metrics::DatabaseProviderMetrics;
59
60mod chain;
61pub use chain::*;
62
63/// Sync state for read-only [`ProviderFactory`] instances.
64struct ReadOnlySyncState {
65    /// Last MDBX txn ID we synced `RocksDB` secondary / static file indexes to.
66    last_synced_txnid: AtomicU64,
67    /// Serializes the slow-path catch-up (`RocksDB` + static file re-init).
68    sync_lock: Mutex<()>,
69}
70
71/// A common provider that fetches data from a database or static file.
72///
73/// This provider implements most provider or provider factory traits.
74pub struct ProviderFactory<N: NodeTypesWithDB> {
75    /// Database instance
76    db: N::DB,
77    /// Chain spec
78    chain_spec: Arc<N::ChainSpec>,
79    /// Static File Provider
80    static_file_provider: StaticFileProvider<N::Primitives>,
81    /// Optional pruning configuration
82    prune_modes: PruneModes,
83    /// The node storage handler.
84    storage: Arc<N::Storage>,
85    /// Storage configuration settings for this node
86    storage_settings: Arc<RwLock<StorageSettings>>,
87    /// `RocksDB` provider
88    rocksdb_provider: RocksDBProvider,
89    /// Manager for state trie overlays and cached changesets.
90    overlay_manager: OverlayManager<N::Primitives>,
91    /// Store for block access lists.
92    bal_store: BalStoreHandle,
93    /// Task runtime for spawning parallel I/O work.
94    runtime: reth_tasks::Runtime,
95    /// Minimum distance from tip required before pruning can occur.
96    minimum_pruning_distance: u64,
97    /// Database provider metrics shared by providers created from this factory.
98    database_provider_metrics: Arc<DatabaseProviderMetrics>,
99    /// State for on-demand syncing of `RocksDB` secondary and static file indexes.
100    ///
101    /// Only set for read-only factories. Can be disabled if there is no concurrent read-write
102    /// factory writing to the database (e.g as part of a running reth node).
103    read_only_sync: Option<Arc<ReadOnlySyncState>>,
104}
105
106impl<N: NodeTypesForProvider> ProviderFactory<NodeTypesWithDBAdapter<N, DatabaseEnv>> {
107    /// Instantiates the builder for this type
108    pub fn builder() -> ProviderFactoryBuilder<N> {
109        ProviderFactoryBuilder::default()
110    }
111}
112
113impl<N: ProviderNodeTypes> ProviderFactory<N> {
114    /// Create new database provider factory.
115    ///
116    /// The storage backends used by the produced factory MAY be inconsistent.
117    /// It is recommended to call [`Self::check_consistency`] after
118    /// creation to ensure consistency between the database and static files.
119    /// If the function returns unwind targets, the caller MUST unwind the
120    /// inner database to the minimum of the two targets to ensure consistency.
121    pub fn new(
122        db: N::DB,
123        chain_spec: Arc<N::ChainSpec>,
124        static_file_provider: StaticFileProvider<N::Primitives>,
125        rocksdb_provider: RocksDBProvider,
126        runtime: reth_tasks::Runtime,
127    ) -> ProviderResult<Self> {
128        // Load storage settings from database at init time. Creates a temporary provider
129        // to read persisted settings, falling back to legacy defaults if none exist.
130        //
131        // Both factory and all providers it creates should share these cached settings.
132        let legacy_settings = StorageSettings::v1();
133        let database_provider_metrics = Arc::new(DatabaseProviderMetrics::default());
134        let overlay_manager = OverlayManager::default();
135        let storage_settings = DatabaseProvider::<_, N>::new(
136            db.tx()?,
137            chain_spec.clone(),
138            static_file_provider.clone(),
139            Default::default(),
140            Default::default(),
141            Arc::new(RwLock::new(legacy_settings)),
142            rocksdb_provider.clone(),
143            overlay_manager.clone(),
144            runtime.clone(),
145            db.path(),
146            database_provider_metrics.clone(),
147        )
148        .storage_settings()?
149        .unwrap_or(legacy_settings);
150
151        Ok(Self {
152            db,
153            chain_spec,
154            static_file_provider,
155            prune_modes: PruneModes::default(),
156            storage: Default::default(),
157            storage_settings: Arc::new(RwLock::new(storage_settings)),
158            rocksdb_provider,
159            overlay_manager,
160            bal_store: BalStoreHandle::new(InMemoryBalStore::default()),
161            runtime,
162            minimum_pruning_distance: MINIMUM_UNWIND_SAFE_DISTANCE,
163            database_provider_metrics,
164            read_only_sync: None,
165        })
166    }
167
168    /// Create new database provider factory and perform consistency checks.
169    ///
170    /// This will call [`Self::check_consistency`] internally and return
171    /// [`ProviderError::MustUnwind`] if inconsistencies are found. It may also
172    /// return any [`ProviderError`] that [`Self::new`] may return, or that are
173    /// encountered during consistency checks.
174    pub fn new_checked(
175        db: N::DB,
176        chain_spec: Arc<N::ChainSpec>,
177        static_file_provider: StaticFileProvider<N::Primitives>,
178        rocksdb_provider: RocksDBProvider,
179        runtime: reth_tasks::Runtime,
180    ) -> ProviderResult<Self> {
181        Self::new(db, chain_spec, static_file_provider, rocksdb_provider, runtime)
182            .and_then(Self::assert_consistent)
183    }
184}
185
186impl<N: NodeTypesWithDB> ProviderFactory<N> {
187    /// Sets the pruning configuration for an existing [`ProviderFactory`].
188    pub fn with_prune_modes(mut self, prune_modes: PruneModes) -> Self {
189        self.prune_modes = prune_modes;
190        self
191    }
192
193    /// Sets the BAL store for an existing [`ProviderFactory`].
194    pub fn with_bal_store(mut self, bal_store: BalStoreHandle) -> Self {
195        self.bal_store = bal_store;
196        self
197    }
198
199    /// Sets the overlay manager for an existing [`ProviderFactory`].
200    pub fn with_overlay_manager(mut self, overlay_manager: OverlayManager<N::Primitives>) -> Self {
201        self.overlay_manager = overlay_manager;
202        self
203    }
204
205    /// Returns the shared overlay manager.
206    pub(crate) const fn overlay_manager(&self) -> &OverlayManager<N::Primitives> {
207        &self.overlay_manager
208    }
209
210    /// Sets the minimum pruning distance for an existing [`ProviderFactory`].
211    ///
212    /// This controls the minimum distance from tip required before pruning can occur.
213    /// The default is [`MINIMUM_UNWIND_SAFE_DISTANCE`].
214    pub const fn with_minimum_pruning_distance(mut self, distance: u64) -> Self {
215        self.minimum_pruning_distance = distance;
216        self
217    }
218
219    /// Enables on-demand syncing of `RocksDB` secondary and static file indexes for read-only
220    /// factories. Initializes the tracker to the current MDBX txn ID.
221    ///
222    /// Should be used for read-only factories that are running concurrently to a reth node writing
223    /// new data to the database. Would effectively be a no-op if database directory is unchanged.
224    pub fn with_read_only_sync(mut self, watch: bool) -> Self
225    where
226        N::DB: Database,
227    {
228        // Initialize to 0 so the first `sync_providers_if_needed` call always
229        // triggers a RocksDB/static-file catch-up, regardless of what MDBX txnid
230        // the database was at when we opened it.
231        let state = Arc::new(ReadOnlySyncState {
232            last_synced_txnid: AtomicU64::new(0),
233            sync_lock: Mutex::new(()),
234        });
235        self.read_only_sync = Some(state);
236
237        if watch {
238            self.watch_db_directory();
239        }
240        self
241    }
242
243    /// Watches the MDBX data directory for changes and eagerly syncs `RocksDB` secondary and
244    /// static file indexes when modifications are detected.
245    fn watch_db_directory(&self)
246    where
247        N::DB: Database,
248    {
249        let factory = self.clone();
250        let db_path = self.db.path();
251        reth_tasks::spawn_os_thread("ro-sync", move || {
252            let (tx, rx) = std::sync::mpsc::channel();
253            let mut watcher = RecommendedWatcher::new(
254                move |res| {
255                    let _ = tx.send(res);
256                },
257                notify::Config::default(),
258            )
259            .expect("failed to create watcher");
260
261            watcher
262                .watch(&db_path, RecursiveMode::NonRecursive)
263                .expect("failed to watch MDBX path");
264
265            while let Ok(res) = rx.recv() {
266                match res {
267                    Ok(event) => {
268                        if !matches!(
269                            event.kind,
270                            notify::EventKind::Modify(_) | notify::EventKind::Create(_)
271                        ) {
272                            continue;
273                        }
274
275                        if let Err(err) = factory.sync_providers_if_needed() {
276                            warn!(target: "reth::provider", %err, "background ro-sync failed");
277                        }
278                    }
279                    Err(err) => {
280                        warn!(target: "reth::provider", ?err, "MDBX directory watcher error");
281                    }
282                }
283            }
284        });
285    }
286
287    /// For read-only factories, checks whether the MDBX committed txn ID has advanced since the
288    /// last sync and, if so, catches up the `RocksDB` secondary instance and re-initializes the
289    /// static file index.
290    ///
291    /// No-op for read-write factories.
292    pub fn sync_providers_if_needed(&self) -> ProviderResult<()> {
293        let Some(sync_state) = &self.read_only_sync else { return Ok(()) };
294        let current_txnid = self.db.last_txnid().unwrap_or(0);
295
296        // Fast path: no contention when nothing changed.
297        if current_txnid == sync_state.last_synced_txnid.load(Ordering::Relaxed) {
298            return Ok(());
299        }
300
301        // Slow path: serialize the actual catch-up I/O.
302        let _guard = sync_state.sync_lock.lock().unwrap_or_else(|e| e.into_inner());
303
304        // Double-check after acquiring the lock — another thread may have already synced.
305        if current_txnid == sync_state.last_synced_txnid.load(Ordering::Relaxed) {
306            return Ok(());
307        }
308
309        self.rocksdb_provider.try_catch_up_with_primary()?;
310        self.static_file_provider.initialize_index()?;
311        sync_state.last_synced_txnid.store(current_txnid, Ordering::Relaxed);
312        Ok(())
313    }
314
315    /// Returns reference to the underlying database.
316    pub const fn db_ref(&self) -> &N::DB {
317        &self.db
318    }
319
320    #[cfg(any(test, feature = "test-utils"))]
321    /// Consumes Self and returns DB
322    pub fn into_db(self) -> N::DB {
323        self.db
324    }
325}
326
327impl<N: NodeTypesWithDB> StorageSettingsCache for ProviderFactory<N> {
328    fn cached_storage_settings(&self) -> StorageSettings {
329        *self.storage_settings.read()
330    }
331
332    fn set_storage_settings_cache(&self, settings: StorageSettings) {
333        *self.storage_settings.write() = settings;
334    }
335}
336
337impl<N: NodeTypesWithDB> RocksDBProviderFactory for ProviderFactory<N> {
338    fn rocksdb_provider(&self) -> RocksDBProvider {
339        self.rocksdb_provider.clone()
340    }
341
342    fn set_pending_rocksdb_batch(&self, _batch: rocksdb::WriteBatchWithTransaction<true>) {
343        unimplemented!("ProviderFactory is a factory, not a provider - use DatabaseProvider::set_pending_rocksdb_batch instead")
344    }
345
346    fn commit_pending_rocksdb_batches(&self) -> ProviderResult<()> {
347        unimplemented!("ProviderFactory is a factory, not a provider - use DatabaseProvider::commit_pending_rocksdb_batches instead")
348    }
349}
350
351impl<N: ProviderNodeTypes<DB = DatabaseEnv>> ProviderFactory<N> {
352    /// Create new database provider by passing a path. [`ProviderFactory`] will own the database
353    /// instance.
354    pub fn new_with_database_path<P: AsRef<Path>>(
355        path: P,
356        chain_spec: Arc<N::ChainSpec>,
357        args: DatabaseArguments,
358        static_file_provider: StaticFileProvider<N::Primitives>,
359        rocksdb_provider: RocksDBProvider,
360        runtime: reth_tasks::Runtime,
361    ) -> RethResult<Self> {
362        Self::new(
363            init_db(path, args).map_err(RethError::msg)?,
364            chain_spec,
365            static_file_provider,
366            rocksdb_provider,
367            runtime,
368        )
369        .map_err(RethError::Provider)
370    }
371}
372
373impl<N: ProviderNodeTypes> ProviderFactory<N> {
374    /// Returns a provider with a created `DbTx` inside, which allows fetching data from the
375    /// database using different types of providers. Example: [`HeaderProvider`]
376    /// [`BlockHashReader`]. This may fail if the inner read database transaction fails to open.
377    ///
378    /// This sets the [`PruneModes`] to [`None`], because they should only be relevant for writing
379    /// data.
380    #[track_caller]
381    pub fn provider(&self) -> ProviderResult<DatabaseProviderRO<N::DB, N>> {
382        let db_tx = self.db.tx()?;
383
384        // Sync providers after opening the database transaction to make
385        // sure that no data is pruned from rocksdb or static files.
386        //
387        // Reorg logic ensures that no data is pruned from rocksdb or static files while there is an
388        // mdbx transaction open that might rely on this data.
389        self.sync_providers_if_needed()?;
390
391        Ok(DatabaseProvider::new(
392            db_tx,
393            self.chain_spec.clone(),
394            self.static_file_provider.clone(),
395            self.prune_modes.clone(),
396            self.storage.clone(),
397            self.storage_settings.clone(),
398            self.rocksdb_provider.clone(),
399            self.overlay_manager.clone(),
400            self.runtime.clone(),
401            self.db.path(),
402            self.database_provider_metrics.clone(),
403        )
404        .with_minimum_pruning_distance(self.minimum_pruning_distance))
405    }
406
407    /// Returns a provider with a created `DbTxMut` inside, which allows fetching and updating
408    /// data from the database using different types of providers. Example: [`HeaderProvider`]
409    /// [`BlockHashReader`].  This may fail if the inner read/write database transaction fails to
410    /// open.
411    #[track_caller]
412    pub fn provider_rw(&self) -> ProviderResult<DatabaseProviderRW<N::DB, N>> {
413        Ok(DatabaseProviderRW(
414            DatabaseProvider::new_rw(
415                self.db.tx_mut()?,
416                self.chain_spec.clone(),
417                self.static_file_provider.clone(),
418                self.prune_modes.clone(),
419                self.storage.clone(),
420                self.storage_settings.clone(),
421                self.rocksdb_provider.clone(),
422                self.overlay_manager.clone(),
423                self.runtime.clone(),
424                self.db.path(),
425                self.database_provider_metrics.clone(),
426            )
427            .with_reader_txn_tracker(self.db.clone())
428            .with_minimum_pruning_distance(self.minimum_pruning_distance),
429        ))
430    }
431
432    /// Returns a provider with a created `DbTxMut` inside, configured for unwind operations.
433    /// Uses unwind commit order (MDBX first, then `RocksDB`, then static files) to allow
434    /// recovery by truncating static files on restart if interrupted.
435    ///
436    /// Unwind commits may wait for pre-existing readers to drain before finishing later
437    /// cross-store steps. Drop any long-lived read providers before committing this provider.
438    #[track_caller]
439    pub fn unwind_provider_rw(
440        &self,
441    ) -> ProviderResult<DatabaseProvider<<N::DB as Database>::TXMut, N>> {
442        Ok(DatabaseProvider::new_unwind_rw(
443            self.db.tx_mut()?,
444            self.chain_spec.clone(),
445            self.static_file_provider.clone(),
446            self.prune_modes.clone(),
447            self.storage.clone(),
448            self.storage_settings.clone(),
449            self.rocksdb_provider.clone(),
450            self.overlay_manager.clone(),
451            self.runtime.clone(),
452            self.db.path(),
453            self.database_provider_metrics.clone(),
454        )
455        .with_reader_txn_tracker(self.db.clone())
456        .with_minimum_pruning_distance(self.minimum_pruning_distance))
457    }
458
459    /// State provider for latest block
460    #[track_caller]
461    pub fn latest(&self) -> ProviderResult<StateProviderBox> {
462        trace!(target: "providers::db", "Returning latest state provider");
463        Ok(Box::new(LatestStateProvider::new(self.database_provider_ro()?)))
464    }
465
466    /// Storage provider for state at that given block
467    pub fn history_by_block_number(
468        &self,
469        block_number: BlockNumber,
470    ) -> ProviderResult<StateProviderBox> {
471        let state_provider = self.provider()?.try_into_history_at_block(block_number)?;
472        trace!(target: "providers::db", ?block_number, "Returning historical state provider for block number");
473        Ok(state_provider)
474    }
475
476    /// Storage provider for state at that given block hash
477    pub fn history_by_block_hash(&self, block_hash: BlockHash) -> ProviderResult<StateProviderBox> {
478        let provider = self.provider()?;
479
480        let block_number = provider
481            .block_number(block_hash)?
482            .ok_or(ProviderError::BlockHashNotFound(block_hash))?;
483
484        let state_provider = provider.try_into_history_at_block(block_number)?;
485        trace!(target: "providers::db", ?block_number, %block_hash, "Returning historical state provider for block hash");
486        Ok(state_provider)
487    }
488
489    /// Asserts that the static files and database are consistent. If not,
490    /// returns [`ProviderError::MustUnwind`] with the appropriate unwind
491    /// target. May also return any [`ProviderError`] that
492    /// [`Self::check_consistency`] may return.
493    pub fn assert_consistent(self) -> ProviderResult<Self> {
494        let (rocksdb_unwind, static_file_unwind) = self.check_consistency()?;
495
496        let source = match (rocksdb_unwind, static_file_unwind) {
497            (None, None) => return Ok(self),
498            (Some(_), Some(_)) => "RocksDB and Static Files",
499            (Some(_), None) => "RocksDB",
500            (None, Some(_)) => "Static Files",
501        };
502
503        Err(ProviderError::MustUnwind {
504            data_source: source,
505            unwind_to: rocksdb_unwind
506                .into_iter()
507                .chain(static_file_unwind)
508                .min()
509                .expect("at least one unwind target must be Some"),
510        })
511    }
512
513    /// Checks the consistency between the static files and the database. This
514    /// may result in static files being pruned or otherwise healed to ensure
515    /// consistency. I.e. this MAY result in writes to the static files.
516    #[instrument(err, skip(self))]
517    pub fn check_consistency(&self) -> ProviderResult<(Option<u64>, Option<u64>)> {
518        let provider_ro = self
519            .database_provider_ro()?
520            // Healing can run long-lived read transactions (e.g., iterating changesets
521            // over millions of blocks). Disable the default timeout so MDBX doesn't
522            // kill the transaction mid-heal, which causes a crash loop on startup.
523            .disable_long_read_transaction_safety();
524
525        // Step 1: heal file-level inconsistencies (no pruning)
526        self.static_file_provider().check_file_consistency(&provider_ro)?;
527
528        // Step 2: RocksDB consistency check (needs static files tx data)
529        let rocksdb_unwind = self.rocksdb_provider().check_consistency(&provider_ro)?;
530
531        // Step 3: Static file checkpoint consistency (may prune)
532        let static_file_unwind = self.static_file_provider().check_consistency(&provider_ro)?.map(
533            |target| match target {
534                PipelineTarget::Unwind(block) => block,
535                PipelineTarget::Sync(_) => unreachable!("check_consistency returns Unwind"),
536            },
537        );
538
539        // Step 4: Heal finalized/safe block numbers that may be ahead of the
540        // highest header on nodes coming from <=1.10.2.
541        //
542        // Unwinds already set it to the target block.
543        if rocksdb_unwind.is_none() && static_file_unwind.is_none() {
544            self.heal_chain_state_block_numbers(&provider_ro)?;
545        }
546
547        Ok((rocksdb_unwind, static_file_unwind))
548    }
549
550    /// If the stored finalized or safe block number is ahead of the highest
551    /// header, resets it to the highest header.
552    fn heal_chain_state_block_numbers(
553        &self,
554        provider_ro: &DatabaseProvider<<N::DB as Database>::TX, N>,
555    ) -> ProviderResult<()> {
556        let highest_header = self.last_block_number()?;
557
558        let finalized = provider_ro.last_finalized_block_number()?;
559        let safe = provider_ro.last_safe_block_number()?;
560
561        if finalized.is_none_or(|f| f <= highest_header) && safe.is_none_or(|s| s <= highest_header)
562        {
563            return Ok(());
564        }
565
566        let provider_rw = self.provider_rw()?;
567
568        if let Some(finalized) = finalized.filter(|&f| f > highest_header) {
569            info!(
570                target: "providers::db",
571                finalized,
572                highest_header,
573                "Healing finalized block number",
574            );
575            provider_rw.save_finalized_block_number(highest_header)?;
576        }
577
578        if let Some(safe) = safe.filter(|&s| s > highest_header) {
579            info!(
580                target: "providers::db",
581                safe,
582                highest_header,
583                "Healing safe block number",
584            );
585            provider_rw.save_safe_block_number(highest_header)?;
586        }
587
588        provider_rw.commit()?;
589
590        Ok(())
591    }
592
593    /// Returns a static file provider. For read-only instances, this will also invoke
594    /// [`Self::sync_providers_if_needed`] to make sure that the static file provider is up to date.
595    pub fn caught_up_static_file_provider(
596        &self,
597    ) -> ProviderResult<StaticFileProvider<N::Primitives>> {
598        self.sync_providers_if_needed()?;
599        Ok(self.static_file_provider.clone())
600    }
601}
602
603impl<N: NodeTypesWithDB> NodePrimitivesProvider for ProviderFactory<N> {
604    type Primitives = N::Primitives;
605}
606
607impl<N: ProviderNodeTypes> BalProvider for ProviderFactory<N> {
608    fn bal_store(&self) -> &BalStoreHandle {
609        &self.bal_store
610    }
611}
612
613impl<N: ProviderNodeTypes> DatabaseProviderFactory for ProviderFactory<N> {
614    type DB = N::DB;
615    type Provider = DatabaseProvider<<N::DB as Database>::TX, N>;
616    type ProviderRW = DatabaseProvider<<N::DB as Database>::TXMut, N>;
617
618    fn database_provider_ro(&self) -> ProviderResult<Self::Provider> {
619        self.provider()
620    }
621
622    fn database_provider_rw(&self) -> ProviderResult<Self::ProviderRW> {
623        self.provider_rw().map(|provider| provider.0)
624    }
625}
626
627impl<N: NodeTypesWithDB> StaticFileProviderFactory for ProviderFactory<N> {
628    /// Returns static file provider
629    fn static_file_provider(&self) -> StaticFileProvider<Self::Primitives> {
630        self.static_file_provider.clone()
631    }
632
633    fn get_static_file_writer(
634        &self,
635        block: BlockNumber,
636        segment: StaticFileSegment,
637    ) -> ProviderResult<StaticFileProviderRWRefMut<'_, Self::Primitives>> {
638        self.static_file_provider.get_writer(block, segment)
639    }
640}
641
642impl<N: ProviderNodeTypes> HeaderSyncGapProvider for ProviderFactory<N> {
643    type Header = HeaderTy<N>;
644    fn local_tip_header(
645        &self,
646        highest_uninterrupted_block: BlockNumber,
647    ) -> ProviderResult<SealedHeader<Self::Header>> {
648        self.provider()?.local_tip_header(highest_uninterrupted_block)
649    }
650}
651
652impl<N: ProviderNodeTypes> HeaderProvider for ProviderFactory<N> {
653    type Header = HeaderTy<N>;
654
655    fn header(&self, block_hash: BlockHash) -> ProviderResult<Option<Self::Header>> {
656        self.provider()?.header(block_hash)
657    }
658
659    fn header_by_number(&self, num: BlockNumber) -> ProviderResult<Option<Self::Header>> {
660        self.caught_up_static_file_provider()?.header_by_number(num)
661    }
662
663    fn headers_range(
664        &self,
665        range: impl RangeBounds<BlockNumber>,
666    ) -> ProviderResult<Vec<Self::Header>> {
667        self.caught_up_static_file_provider()?.headers_range(range)
668    }
669
670    fn sealed_header(
671        &self,
672        number: BlockNumber,
673    ) -> ProviderResult<Option<SealedHeader<Self::Header>>> {
674        self.caught_up_static_file_provider()?.sealed_header(number)
675    }
676
677    fn sealed_headers_range(
678        &self,
679        range: impl RangeBounds<BlockNumber>,
680    ) -> ProviderResult<Vec<SealedHeader<Self::Header>>> {
681        self.caught_up_static_file_provider()?.sealed_headers_range(range)
682    }
683
684    fn sealed_headers_while(
685        &self,
686        range: impl RangeBounds<BlockNumber>,
687        predicate: impl FnMut(&SealedHeader<Self::Header>) -> bool,
688    ) -> ProviderResult<Vec<SealedHeader<Self::Header>>> {
689        self.caught_up_static_file_provider()?.sealed_headers_while(range, predicate)
690    }
691}
692
693impl<N: ProviderNodeTypes> BlockHashReader for ProviderFactory<N> {
694    fn block_hash(&self, number: u64) -> ProviderResult<Option<B256>> {
695        self.caught_up_static_file_provider()?.block_hash(number)
696    }
697
698    fn canonical_hashes_range(
699        &self,
700        start: BlockNumber,
701        end: BlockNumber,
702    ) -> ProviderResult<Vec<B256>> {
703        self.caught_up_static_file_provider()?.canonical_hashes_range(start, end)
704    }
705}
706
707impl<N: ProviderNodeTypes> BlockNumReader for ProviderFactory<N> {
708    fn chain_info(&self) -> ProviderResult<ChainInfo> {
709        self.provider()?.chain_info()
710    }
711
712    fn best_block_number(&self) -> ProviderResult<BlockNumber> {
713        self.provider()?.best_block_number()
714    }
715
716    fn last_block_number(&self) -> ProviderResult<BlockNumber> {
717        self.caught_up_static_file_provider()?.last_block_number()
718    }
719
720    fn earliest_block_number(&self) -> ProviderResult<BlockNumber> {
721        // earliest history height tracks the lowest block number that has __not__ been expired, in
722        // other words, the first/earliest available block.
723        Ok(self.caught_up_static_file_provider()?.earliest_history_height())
724    }
725
726    fn block_number(&self, hash: B256) -> ProviderResult<Option<BlockNumber>> {
727        self.provider()?.block_number(hash)
728    }
729}
730
731impl<N: ProviderNodeTypes> BlockReader for ProviderFactory<N> {
732    type Block = BlockTy<N>;
733
734    fn find_block_by_hash(
735        &self,
736        hash: B256,
737        source: BlockSource,
738    ) -> ProviderResult<Option<Self::Block>> {
739        self.provider()?.find_block_by_hash(hash, source)
740    }
741
742    fn block(&self, id: BlockHashOrNumber) -> ProviderResult<Option<Self::Block>> {
743        self.provider()?.block(id)
744    }
745
746    fn pending_block(&self) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
747        self.provider()?.pending_block()
748    }
749
750    fn pending_block_and_receipts(
751        &self,
752    ) -> ProviderResult<Option<(RecoveredBlock<Self::Block>, Vec<Self::Receipt>)>> {
753        self.provider()?.pending_block_and_receipts()
754    }
755
756    fn recovered_block(
757        &self,
758        id: BlockHashOrNumber,
759        transaction_kind: TransactionVariant,
760    ) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
761        self.provider()?.recovered_block(id, transaction_kind)
762    }
763
764    fn sealed_block_with_senders(
765        &self,
766        id: BlockHashOrNumber,
767        transaction_kind: TransactionVariant,
768    ) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
769        self.provider()?.sealed_block_with_senders(id, transaction_kind)
770    }
771
772    fn block_range(&self, range: RangeInclusive<BlockNumber>) -> ProviderResult<Vec<Self::Block>> {
773        self.provider()?.block_range(range)
774    }
775
776    fn block_with_senders_range(
777        &self,
778        range: RangeInclusive<BlockNumber>,
779    ) -> ProviderResult<Vec<RecoveredBlock<Self::Block>>> {
780        self.provider()?.block_with_senders_range(range)
781    }
782
783    fn recovered_block_range(
784        &self,
785        range: RangeInclusive<BlockNumber>,
786    ) -> ProviderResult<Vec<RecoveredBlock<Self::Block>>> {
787        self.provider()?.recovered_block_range(range)
788    }
789
790    fn block_by_transaction_id(&self, id: TxNumber) -> ProviderResult<Option<BlockNumber>> {
791        self.provider()?.block_by_transaction_id(id)
792    }
793}
794
795impl<N: ProviderNodeTypes> TransactionsProvider for ProviderFactory<N> {
796    type Transaction = TxTy<N>;
797
798    fn transaction_id(&self, tx_hash: TxHash) -> ProviderResult<Option<TxNumber>> {
799        self.provider()?.transaction_id(tx_hash)
800    }
801
802    fn transaction_by_id(&self, id: TxNumber) -> ProviderResult<Option<Self::Transaction>> {
803        self.caught_up_static_file_provider()?.transaction_by_id(id)
804    }
805
806    fn transaction_by_id_unhashed(
807        &self,
808        id: TxNumber,
809    ) -> ProviderResult<Option<Self::Transaction>> {
810        self.caught_up_static_file_provider()?.transaction_by_id_unhashed(id)
811    }
812
813    fn transaction_by_hash(&self, hash: TxHash) -> ProviderResult<Option<Self::Transaction>> {
814        self.provider()?.transaction_by_hash(hash)
815    }
816
817    fn transaction_by_hash_with_meta(
818        &self,
819        tx_hash: TxHash,
820    ) -> ProviderResult<Option<(Self::Transaction, TransactionMeta)>> {
821        self.provider()?.transaction_by_hash_with_meta(tx_hash)
822    }
823
824    fn transactions_by_block(
825        &self,
826        id: BlockHashOrNumber,
827    ) -> ProviderResult<Option<Vec<Self::Transaction>>> {
828        self.provider()?.transactions_by_block(id)
829    }
830
831    fn transactions_by_block_range(
832        &self,
833        range: impl RangeBounds<BlockNumber>,
834    ) -> ProviderResult<Vec<Vec<Self::Transaction>>> {
835        self.provider()?.transactions_by_block_range(range)
836    }
837
838    fn transactions_by_tx_range(
839        &self,
840        range: impl RangeBounds<TxNumber>,
841    ) -> ProviderResult<Vec<Self::Transaction>> {
842        self.caught_up_static_file_provider()?.transactions_by_tx_range(range)
843    }
844
845    fn senders_by_tx_range(
846        &self,
847        range: impl RangeBounds<TxNumber>,
848    ) -> ProviderResult<Vec<Address>> {
849        if EitherWriterDestination::senders(self).is_static_file() {
850            self.caught_up_static_file_provider()?.senders_by_tx_range(range)
851        } else {
852            self.provider()?.senders_by_tx_range(range)
853        }
854    }
855
856    fn transaction_sender(&self, id: TxNumber) -> ProviderResult<Option<Address>> {
857        if EitherWriterDestination::senders(self).is_static_file() {
858            self.caught_up_static_file_provider()?.transaction_sender(id)
859        } else {
860            self.provider()?.transaction_sender(id)
861        }
862    }
863}
864
865impl<N: ProviderNodeTypes> ReceiptProvider for ProviderFactory<N> {
866    type Receipt = ReceiptTy<N>;
867
868    fn receipt(&self, id: TxNumber) -> ProviderResult<Option<Self::Receipt>> {
869        self.caught_up_static_file_provider()?.get_with_static_file_or_database(
870            StaticFileSegment::Receipts,
871            id,
872            |static_file| static_file.receipt(id),
873            || self.provider()?.receipt(id),
874        )
875    }
876
877    fn receipt_by_hash(&self, hash: TxHash) -> ProviderResult<Option<Self::Receipt>> {
878        self.provider()?.receipt_by_hash(hash)
879    }
880
881    fn receipts_by_block(
882        &self,
883        block: BlockHashOrNumber,
884    ) -> ProviderResult<Option<Vec<Self::Receipt>>> {
885        self.provider()?.receipts_by_block(block)
886    }
887
888    fn receipts_by_tx_range(
889        &self,
890        range: impl RangeBounds<TxNumber>,
891    ) -> ProviderResult<Vec<Self::Receipt>> {
892        self.caught_up_static_file_provider()?.get_range_with_static_file_or_database(
893            StaticFileSegment::Receipts,
894            to_range(range),
895            |static_file, range, _| static_file.receipts_by_tx_range(range),
896            |range, _| self.provider()?.receipts_by_tx_range(range),
897            |_| true,
898        )
899    }
900
901    fn receipts_by_block_range(
902        &self,
903        block_range: RangeInclusive<BlockNumber>,
904    ) -> ProviderResult<Vec<Vec<Self::Receipt>>> {
905        self.provider()?.receipts_by_block_range(block_range)
906    }
907}
908
909impl<N: ProviderNodeTypes> BlockBodyIndicesProvider for ProviderFactory<N> {
910    fn block_body_indices(
911        &self,
912        number: BlockNumber,
913    ) -> ProviderResult<Option<StoredBlockBodyIndices>> {
914        self.provider()?.block_body_indices(number)
915    }
916
917    fn block_body_indices_range(
918        &self,
919        range: RangeInclusive<BlockNumber>,
920    ) -> ProviderResult<Vec<StoredBlockBodyIndices>> {
921        self.provider()?.block_body_indices_range(range)
922    }
923}
924
925impl<N: ProviderNodeTypes> StageCheckpointReader for ProviderFactory<N> {
926    fn get_stage_checkpoint(&self, id: StageId) -> ProviderResult<Option<StageCheckpoint>> {
927        self.provider()?.get_stage_checkpoint(id)
928    }
929
930    fn get_stage_checkpoint_progress(&self, id: StageId) -> ProviderResult<Option<Vec<u8>>> {
931        self.provider()?.get_stage_checkpoint_progress(id)
932    }
933    fn get_all_checkpoints(&self) -> ProviderResult<Vec<(String, StageCheckpoint)>> {
934        self.provider()?.get_all_checkpoints()
935    }
936}
937
938impl<N: NodeTypesWithDB> ChainSpecProvider for ProviderFactory<N> {
939    type ChainSpec = N::ChainSpec;
940
941    fn chain_spec(&self) -> Arc<N::ChainSpec> {
942        self.chain_spec.clone()
943    }
944}
945
946impl<N: ProviderNodeTypes> PruneCheckpointReader for ProviderFactory<N> {
947    fn get_prune_checkpoint(
948        &self,
949        segment: PruneSegment,
950    ) -> ProviderResult<Option<PruneCheckpoint>> {
951        self.provider()?.get_prune_checkpoint(segment)
952    }
953
954    fn get_prune_checkpoints(&self) -> ProviderResult<Vec<(PruneSegment, PruneCheckpoint)>> {
955        self.provider()?.get_prune_checkpoints()
956    }
957}
958
959impl<N: ProviderNodeTypes> MetadataProvider for ProviderFactory<N> {
960    fn get_metadata(&self, key: &str) -> ProviderResult<Option<Vec<u8>>> {
961        self.provider()?.get_metadata(key)
962    }
963}
964
965impl<N> fmt::Debug for ProviderFactory<N>
966where
967    N: NodeTypesWithDB<DB: fmt::Debug, ChainSpec: fmt::Debug, Storage: fmt::Debug>,
968{
969    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
970        let Self {
971            db,
972            chain_spec,
973            static_file_provider,
974            prune_modes,
975            storage,
976            storage_settings,
977            rocksdb_provider,
978            overlay_manager,
979            bal_store,
980            runtime,
981            minimum_pruning_distance,
982            database_provider_metrics: _,
983            read_only_sync,
984        } = self;
985        f.debug_struct("ProviderFactory")
986            .field("db", &db)
987            .field("chain_spec", &chain_spec)
988            .field("static_file_provider", &static_file_provider)
989            .field("prune_modes", &prune_modes)
990            .field("storage", &storage)
991            .field("storage_settings", &*storage_settings.read())
992            .field("rocksdb_provider", &rocksdb_provider)
993            .field("overlay_manager", &overlay_manager)
994            .field("bal_store", &bal_store)
995            .field("runtime", &runtime)
996            .field("minimum_pruning_distance", &minimum_pruning_distance)
997            .field(
998                "read_only_sync",
999                &read_only_sync.as_ref().map(|s| s.last_synced_txnid.load(Ordering::Relaxed)),
1000            )
1001            .finish()
1002    }
1003}
1004
1005impl<N: NodeTypesWithDB> Clone for ProviderFactory<N> {
1006    fn clone(&self) -> Self {
1007        Self {
1008            db: self.db.clone(),
1009            chain_spec: self.chain_spec.clone(),
1010            static_file_provider: self.static_file_provider.clone(),
1011            prune_modes: self.prune_modes.clone(),
1012            storage: self.storage.clone(),
1013            storage_settings: self.storage_settings.clone(),
1014            rocksdb_provider: self.rocksdb_provider.clone(),
1015            overlay_manager: self.overlay_manager.clone(),
1016            bal_store: self.bal_store.clone(),
1017            runtime: self.runtime.clone(),
1018            minimum_pruning_distance: self.minimum_pruning_distance,
1019            database_provider_metrics: self.database_provider_metrics.clone(),
1020            read_only_sync: self.read_only_sync.clone(),
1021        }
1022    }
1023}
1024
1025#[cfg(test)]
1026mod tests {
1027    use super::*;
1028    use crate::{
1029        providers::{StaticFileProvider, StaticFileWriter},
1030        test_utils::{blocks::TEST_BLOCK, create_test_provider_factory, MockNodeTypesWithDB},
1031        BlockHashReader, BlockNumReader, BlockWriter, DBProvider, HeaderSyncGapProvider,
1032        TransactionsProvider,
1033    };
1034    use alloy_primitives::{TxNumber, B256};
1035    use assert_matches::assert_matches;
1036    use reth_chainspec::ChainSpecBuilder;
1037    use reth_db::{
1038        mdbx::DatabaseArguments,
1039        test_utils::{create_test_rocksdb_dir, create_test_static_files_dir, ERROR_TEMPDIR},
1040    };
1041    use reth_db_api::tables;
1042    use reth_primitives_traits::SignerRecoverable;
1043    use reth_prune_types::{PruneMode, PruneModes};
1044    use reth_storage_errors::provider::ProviderError;
1045    use reth_testing_utils::generators::{self, random_block, random_header, BlockParams};
1046    use std::{ops::RangeInclusive, sync::Arc};
1047
1048    #[test]
1049    fn common_history_provider() {
1050        let factory = create_test_provider_factory();
1051        let _ = factory.latest();
1052    }
1053
1054    #[test]
1055    fn default_chain_info() {
1056        let factory = create_test_provider_factory();
1057        let provider = factory.provider().unwrap();
1058
1059        let chain_info = provider.chain_info().expect("should be ok");
1060        assert_eq!(chain_info.best_number, 0);
1061        assert_eq!(chain_info.best_hash, B256::ZERO);
1062    }
1063
1064    #[test]
1065    fn provider_flow() {
1066        let factory = create_test_provider_factory();
1067        let provider = factory.provider().unwrap();
1068        provider.block_hash(0).unwrap();
1069        let provider_rw = factory.provider_rw().unwrap();
1070        provider_rw.block_hash(0).unwrap();
1071        provider.block_hash(0).unwrap();
1072    }
1073
1074    #[test]
1075    fn provider_factory_with_database_path() {
1076        let chain_spec = ChainSpecBuilder::mainnet().build();
1077        let (_static_dir, static_dir_path) = create_test_static_files_dir();
1078        let (_rocksdb_dir, rocksdb_path) = create_test_rocksdb_dir();
1079        let _db_tempdir = tempfile::TempDir::new().expect(ERROR_TEMPDIR);
1080        let factory = ProviderFactory::<MockNodeTypesWithDB<DatabaseEnv>>::new_with_database_path(
1081            _db_tempdir.path(),
1082            Arc::new(chain_spec),
1083            DatabaseArguments::new(Default::default()),
1084            StaticFileProvider::read_write(static_dir_path).unwrap(),
1085            RocksDBProvider::builder(&rocksdb_path).build().unwrap(),
1086            reth_tasks::Runtime::test(),
1087        )
1088        .unwrap();
1089        let provider = factory.provider().unwrap();
1090        provider.block_hash(0).unwrap();
1091        let provider_rw = factory.provider_rw().unwrap();
1092        provider_rw.block_hash(0).unwrap();
1093        provider.block_hash(0).unwrap();
1094    }
1095
1096    #[test]
1097    fn insert_block_with_prune_modes() {
1098        let block = TEST_BLOCK.clone();
1099
1100        {
1101            let factory = create_test_provider_factory();
1102            let provider = factory.provider_rw().unwrap();
1103            assert_matches!(provider.insert_block(&block.clone().try_recover().unwrap()), Ok(_));
1104            assert_matches!(
1105                provider.transaction_sender(0), Ok(Some(sender))
1106                if sender == block.body().transactions[0].recover_signer().unwrap()
1107            );
1108            assert_matches!(
1109                provider.transaction_id(*block.body().transactions[0].tx_hash()),
1110                Ok(Some(0))
1111            );
1112        }
1113
1114        {
1115            let prune_modes = PruneModes {
1116                sender_recovery: Some(PruneMode::Full),
1117                transaction_lookup: Some(PruneMode::Full),
1118                ..PruneModes::default()
1119            };
1120            // Keep factory alive until provider is dropped to prevent TempDatabase cleanup
1121            let factory = create_test_provider_factory().with_prune_modes(prune_modes);
1122            let provider = factory.provider_rw().unwrap();
1123            assert_matches!(provider.insert_block(&block.clone().try_recover().unwrap()), Ok(_));
1124            assert_matches!(provider.transaction_sender(0), Ok(None));
1125            assert_matches!(
1126                provider.transaction_id(*block.body().transactions[0].tx_hash()),
1127                Ok(None)
1128            );
1129        }
1130    }
1131
1132    #[test]
1133    fn take_block_transaction_range_recover_senders() {
1134        let mut rng = generators::rng();
1135        let block =
1136            random_block(&mut rng, 0, BlockParams { tx_count: Some(3), ..Default::default() });
1137
1138        let tx_ranges: Vec<RangeInclusive<TxNumber>> = vec![0..=0, 1..=1, 2..=2, 0..=1, 1..=2];
1139        for range in tx_ranges {
1140            let factory = create_test_provider_factory();
1141            let provider = factory.provider_rw().unwrap();
1142
1143            assert_matches!(provider.insert_block(&block.clone().try_recover().unwrap()), Ok(_));
1144
1145            let senders = provider.take::<tables::TransactionSenders>(range.clone()).unwrap();
1146            assert_eq!(
1147                senders,
1148                range
1149                    .clone()
1150                    .map(|tx_number| (
1151                        tx_number,
1152                        block.body().transactions[tx_number as usize].recover_signer().unwrap()
1153                    ))
1154                    .collect::<Vec<_>>()
1155            );
1156
1157            let db_senders = provider.senders_by_tx_range(range);
1158            assert!(matches!(db_senders, Ok(ref v) if v.is_empty()));
1159        }
1160    }
1161
1162    #[test]
1163    fn header_sync_gap_lookup() {
1164        let factory = create_test_provider_factory();
1165        let provider = factory.provider_rw().unwrap();
1166
1167        let mut rng = generators::rng();
1168
1169        // Genesis
1170        let checkpoint = 0;
1171        let head = random_header(&mut rng, 0, None);
1172
1173        // Empty database
1174        assert_matches!(
1175            provider.local_tip_header(checkpoint),
1176            Err(ProviderError::HeaderNotFound(block_number))
1177                if block_number.as_number().unwrap() == checkpoint
1178        );
1179
1180        // Checkpoint and no gap
1181        let static_file_provider = provider.static_file_provider();
1182        let mut static_file_writer =
1183            static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
1184        static_file_writer.append_header(head.header(), &head.hash()).unwrap();
1185        static_file_writer.commit().unwrap();
1186        drop(static_file_writer);
1187
1188        let local_head = provider.local_tip_header(checkpoint).unwrap();
1189
1190        assert_eq!(local_head, head);
1191    }
1192}