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,
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    /// Asserts that the static files and database are consistent. If not,
467    /// returns [`ProviderError::MustUnwind`] with the appropriate unwind
468    /// target. May also return any [`ProviderError`] that
469    /// [`Self::check_consistency`] may return.
470    pub fn assert_consistent(self) -> ProviderResult<Self> {
471        let (rocksdb_unwind, static_file_unwind) = self.check_consistency()?;
472
473        let source = match (rocksdb_unwind, static_file_unwind) {
474            (None, None) => return Ok(self),
475            (Some(_), Some(_)) => "RocksDB and Static Files",
476            (Some(_), None) => "RocksDB",
477            (None, Some(_)) => "Static Files",
478        };
479
480        Err(ProviderError::MustUnwind {
481            data_source: source,
482            unwind_to: rocksdb_unwind
483                .into_iter()
484                .chain(static_file_unwind)
485                .min()
486                .expect("at least one unwind target must be Some"),
487        })
488    }
489
490    /// Checks the consistency between the static files and the database. This
491    /// may result in static files being pruned or otherwise healed to ensure
492    /// consistency. I.e. this MAY result in writes to the static files.
493    #[instrument(err, skip(self))]
494    pub fn check_consistency(&self) -> ProviderResult<(Option<u64>, Option<u64>)> {
495        let provider_ro = self
496            .database_provider_ro()?
497            // Healing can run long-lived read transactions (e.g., iterating changesets
498            // over millions of blocks). Disable the default timeout so MDBX doesn't
499            // kill the transaction mid-heal, which causes a crash loop on startup.
500            .disable_long_read_transaction_safety();
501
502        // Step 1: heal file-level inconsistencies (no pruning)
503        self.static_file_provider().check_file_consistency(&provider_ro)?;
504
505        // Step 2: RocksDB consistency check (needs static files tx data)
506        let rocksdb_unwind = self.rocksdb_provider().check_consistency(&provider_ro)?;
507
508        // Step 3: Static file checkpoint consistency (may prune)
509        let static_file_unwind = self.static_file_provider().check_consistency(&provider_ro)?.map(
510            |target| match target {
511                PipelineTarget::Unwind(block) => block,
512                PipelineTarget::Sync(_) => unreachable!("check_consistency returns Unwind"),
513            },
514        );
515
516        // Step 4: Heal finalized/safe block numbers that may be ahead of the
517        // highest header on nodes coming from <=1.10.2.
518        //
519        // Unwinds already set it to the target block.
520        if rocksdb_unwind.is_none() && static_file_unwind.is_none() {
521            self.heal_chain_state_block_numbers(&provider_ro)?;
522        }
523
524        Ok((rocksdb_unwind, static_file_unwind))
525    }
526
527    /// If the stored finalized or safe block number is ahead of the highest
528    /// header, resets it to the highest header.
529    fn heal_chain_state_block_numbers(
530        &self,
531        provider_ro: &DatabaseProvider<<N::DB as Database>::TX, N>,
532    ) -> ProviderResult<()> {
533        let highest_header = self.last_block_number()?;
534
535        let finalized = provider_ro.last_finalized_block_number()?;
536        let safe = provider_ro.last_safe_block_number()?;
537
538        if finalized.is_none_or(|f| f <= highest_header) && safe.is_none_or(|s| s <= highest_header)
539        {
540            return Ok(());
541        }
542
543        let provider_rw = self.provider_rw()?;
544
545        if let Some(finalized) = finalized.filter(|&f| f > highest_header) {
546            info!(
547                target: "providers::db",
548                finalized,
549                highest_header,
550                "Healing finalized block number",
551            );
552            provider_rw.save_finalized_block_number(highest_header)?;
553        }
554
555        if let Some(safe) = safe.filter(|&s| s > highest_header) {
556            info!(
557                target: "providers::db",
558                safe,
559                highest_header,
560                "Healing safe block number",
561            );
562            provider_rw.save_safe_block_number(highest_header)?;
563        }
564
565        provider_rw.commit()?;
566
567        Ok(())
568    }
569
570    /// Returns a static file provider. For read-only instances, this will also invoke
571    /// [`Self::sync_providers_if_needed`] to make sure that the static file provider is up to date.
572    pub fn caught_up_static_file_provider(
573        &self,
574    ) -> ProviderResult<StaticFileProvider<N::Primitives>> {
575        self.sync_providers_if_needed()?;
576        Ok(self.static_file_provider.clone())
577    }
578}
579
580impl<N: NodeTypesWithDB> NodePrimitivesProvider for ProviderFactory<N> {
581    type Primitives = N::Primitives;
582}
583
584impl<N: ProviderNodeTypes> BalProvider for ProviderFactory<N> {
585    fn bal_store(&self) -> &BalStoreHandle {
586        &self.bal_store
587    }
588}
589
590impl<N: ProviderNodeTypes> DatabaseProviderFactory for ProviderFactory<N> {
591    type DB = N::DB;
592    type Provider = DatabaseProvider<<N::DB as Database>::TX, N>;
593    type ProviderRW = DatabaseProvider<<N::DB as Database>::TXMut, N>;
594
595    fn database_provider_ro(&self) -> ProviderResult<Self::Provider> {
596        self.provider()
597    }
598
599    fn database_provider_rw(&self) -> ProviderResult<Self::ProviderRW> {
600        self.provider_rw().map(|provider| provider.0)
601    }
602}
603
604impl<N: NodeTypesWithDB> StaticFileProviderFactory for ProviderFactory<N> {
605    /// Returns static file provider
606    fn static_file_provider(&self) -> StaticFileProvider<Self::Primitives> {
607        self.static_file_provider.clone()
608    }
609
610    fn get_static_file_writer(
611        &self,
612        block: BlockNumber,
613        segment: StaticFileSegment,
614    ) -> ProviderResult<StaticFileProviderRWRefMut<'_, Self::Primitives>> {
615        self.static_file_provider.get_writer(block, segment)
616    }
617}
618
619impl<N: ProviderNodeTypes> HeaderSyncGapProvider for ProviderFactory<N> {
620    type Header = HeaderTy<N>;
621    fn local_tip_header(
622        &self,
623        highest_uninterrupted_block: BlockNumber,
624    ) -> ProviderResult<SealedHeader<Self::Header>> {
625        self.provider()?.local_tip_header(highest_uninterrupted_block)
626    }
627}
628
629impl<N: ProviderNodeTypes> HeaderProvider for ProviderFactory<N> {
630    type Header = HeaderTy<N>;
631
632    fn header(&self, block_hash: BlockHash) -> ProviderResult<Option<Self::Header>> {
633        self.provider()?.header(block_hash)
634    }
635
636    fn header_by_number(&self, num: BlockNumber) -> ProviderResult<Option<Self::Header>> {
637        self.caught_up_static_file_provider()?.header_by_number(num)
638    }
639
640    fn headers_range(
641        &self,
642        range: impl RangeBounds<BlockNumber>,
643    ) -> ProviderResult<Vec<Self::Header>> {
644        self.caught_up_static_file_provider()?.headers_range(range)
645    }
646
647    fn sealed_header(
648        &self,
649        number: BlockNumber,
650    ) -> ProviderResult<Option<SealedHeader<Self::Header>>> {
651        self.caught_up_static_file_provider()?.sealed_header(number)
652    }
653
654    fn sealed_headers_range(
655        &self,
656        range: impl RangeBounds<BlockNumber>,
657    ) -> ProviderResult<Vec<SealedHeader<Self::Header>>> {
658        self.caught_up_static_file_provider()?.sealed_headers_range(range)
659    }
660
661    fn sealed_headers_while(
662        &self,
663        range: impl RangeBounds<BlockNumber>,
664        predicate: impl FnMut(&SealedHeader<Self::Header>) -> bool,
665    ) -> ProviderResult<Vec<SealedHeader<Self::Header>>> {
666        self.caught_up_static_file_provider()?.sealed_headers_while(range, predicate)
667    }
668}
669
670impl<N: ProviderNodeTypes> BlockHashReader for ProviderFactory<N> {
671    fn block_hash(&self, number: u64) -> ProviderResult<Option<B256>> {
672        self.caught_up_static_file_provider()?.block_hash(number)
673    }
674
675    fn canonical_hashes_range(
676        &self,
677        start: BlockNumber,
678        end: BlockNumber,
679    ) -> ProviderResult<Vec<B256>> {
680        self.caught_up_static_file_provider()?.canonical_hashes_range(start, end)
681    }
682}
683
684impl<N: ProviderNodeTypes> BlockNumReader for ProviderFactory<N> {
685    fn chain_info(&self) -> ProviderResult<ChainInfo> {
686        self.provider()?.chain_info()
687    }
688
689    fn best_block_number(&self) -> ProviderResult<BlockNumber> {
690        self.provider()?.best_block_number()
691    }
692
693    fn last_block_number(&self) -> ProviderResult<BlockNumber> {
694        self.caught_up_static_file_provider()?.last_block_number()
695    }
696
697    fn earliest_block_number(&self) -> ProviderResult<BlockNumber> {
698        // earliest history height tracks the lowest block number that has __not__ been expired, in
699        // other words, the first/earliest available block.
700        Ok(self.caught_up_static_file_provider()?.earliest_history_height())
701    }
702
703    fn block_number(&self, hash: B256) -> ProviderResult<Option<BlockNumber>> {
704        self.provider()?.block_number(hash)
705    }
706}
707
708impl<N: ProviderNodeTypes> BlockReader for ProviderFactory<N> {
709    type Block = BlockTy<N>;
710
711    fn find_block_by_hash(
712        &self,
713        hash: B256,
714        source: BlockSource,
715    ) -> ProviderResult<Option<Self::Block>> {
716        self.provider()?.find_block_by_hash(hash, source)
717    }
718
719    fn block(&self, id: BlockHashOrNumber) -> ProviderResult<Option<Self::Block>> {
720        self.provider()?.block(id)
721    }
722
723    fn pending_block(&self) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
724        self.provider()?.pending_block()
725    }
726
727    fn pending_block_and_receipts(
728        &self,
729    ) -> ProviderResult<Option<(RecoveredBlock<Self::Block>, Vec<Self::Receipt>)>> {
730        self.provider()?.pending_block_and_receipts()
731    }
732
733    fn recovered_block(
734        &self,
735        id: BlockHashOrNumber,
736        transaction_kind: TransactionVariant,
737    ) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
738        self.provider()?.recovered_block(id, transaction_kind)
739    }
740
741    fn sealed_block_with_senders(
742        &self,
743        id: BlockHashOrNumber,
744        transaction_kind: TransactionVariant,
745    ) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
746        self.provider()?.sealed_block_with_senders(id, transaction_kind)
747    }
748
749    fn block_range(&self, range: RangeInclusive<BlockNumber>) -> ProviderResult<Vec<Self::Block>> {
750        self.provider()?.block_range(range)
751    }
752
753    fn block_with_senders_range(
754        &self,
755        range: RangeInclusive<BlockNumber>,
756    ) -> ProviderResult<Vec<RecoveredBlock<Self::Block>>> {
757        self.provider()?.block_with_senders_range(range)
758    }
759
760    fn recovered_block_range(
761        &self,
762        range: RangeInclusive<BlockNumber>,
763    ) -> ProviderResult<Vec<RecoveredBlock<Self::Block>>> {
764        self.provider()?.recovered_block_range(range)
765    }
766
767    fn block_by_transaction_id(&self, id: TxNumber) -> ProviderResult<Option<BlockNumber>> {
768        self.provider()?.block_by_transaction_id(id)
769    }
770}
771
772impl<N: ProviderNodeTypes> TransactionsProvider for ProviderFactory<N> {
773    type Transaction = TxTy<N>;
774
775    fn transaction_id(&self, tx_hash: TxHash) -> ProviderResult<Option<TxNumber>> {
776        self.provider()?.transaction_id(tx_hash)
777    }
778
779    fn transaction_by_id(&self, id: TxNumber) -> ProviderResult<Option<Self::Transaction>> {
780        self.caught_up_static_file_provider()?.transaction_by_id(id)
781    }
782
783    fn transaction_by_id_unhashed(
784        &self,
785        id: TxNumber,
786    ) -> ProviderResult<Option<Self::Transaction>> {
787        self.caught_up_static_file_provider()?.transaction_by_id_unhashed(id)
788    }
789
790    fn transaction_by_hash(&self, hash: TxHash) -> ProviderResult<Option<Self::Transaction>> {
791        self.provider()?.transaction_by_hash(hash)
792    }
793
794    fn transaction_by_hash_with_meta(
795        &self,
796        tx_hash: TxHash,
797    ) -> ProviderResult<Option<(Self::Transaction, TransactionMeta)>> {
798        self.provider()?.transaction_by_hash_with_meta(tx_hash)
799    }
800
801    fn transactions_by_block(
802        &self,
803        id: BlockHashOrNumber,
804    ) -> ProviderResult<Option<Vec<Self::Transaction>>> {
805        self.provider()?.transactions_by_block(id)
806    }
807
808    fn transactions_by_block_range(
809        &self,
810        range: impl RangeBounds<BlockNumber>,
811    ) -> ProviderResult<Vec<Vec<Self::Transaction>>> {
812        self.provider()?.transactions_by_block_range(range)
813    }
814
815    fn transactions_by_tx_range(
816        &self,
817        range: impl RangeBounds<TxNumber>,
818    ) -> ProviderResult<Vec<Self::Transaction>> {
819        self.caught_up_static_file_provider()?.transactions_by_tx_range(range)
820    }
821
822    fn senders_by_tx_range(
823        &self,
824        range: impl RangeBounds<TxNumber>,
825    ) -> ProviderResult<Vec<Address>> {
826        if EitherWriterDestination::senders(self).is_static_file() {
827            self.caught_up_static_file_provider()?.senders_by_tx_range(range)
828        } else {
829            self.provider()?.senders_by_tx_range(range)
830        }
831    }
832
833    fn transaction_sender(&self, id: TxNumber) -> ProviderResult<Option<Address>> {
834        if EitherWriterDestination::senders(self).is_static_file() {
835            self.caught_up_static_file_provider()?.transaction_sender(id)
836        } else {
837            self.provider()?.transaction_sender(id)
838        }
839    }
840}
841
842impl<N: ProviderNodeTypes> ReceiptProvider for ProviderFactory<N> {
843    type Receipt = ReceiptTy<N>;
844
845    fn receipt(&self, id: TxNumber) -> ProviderResult<Option<Self::Receipt>> {
846        self.caught_up_static_file_provider()?.get_with_static_file_or_database(
847            StaticFileSegment::Receipts,
848            id,
849            |static_file| static_file.receipt(id),
850            || self.provider()?.receipt(id),
851        )
852    }
853
854    fn receipt_by_hash(&self, hash: TxHash) -> ProviderResult<Option<Self::Receipt>> {
855        self.provider()?.receipt_by_hash(hash)
856    }
857
858    fn receipts_by_block(
859        &self,
860        block: BlockHashOrNumber,
861    ) -> ProviderResult<Option<Vec<Self::Receipt>>> {
862        self.provider()?.receipts_by_block(block)
863    }
864
865    fn receipts_by_tx_range(
866        &self,
867        range: impl RangeBounds<TxNumber>,
868    ) -> ProviderResult<Vec<Self::Receipt>> {
869        self.caught_up_static_file_provider()?.get_range_with_static_file_or_database(
870            StaticFileSegment::Receipts,
871            to_range(range),
872            |static_file, range, _| static_file.receipts_by_tx_range(range),
873            |range, _| self.provider()?.receipts_by_tx_range(range),
874            |_| true,
875        )
876    }
877
878    fn receipts_by_block_range(
879        &self,
880        block_range: RangeInclusive<BlockNumber>,
881    ) -> ProviderResult<Vec<Vec<Self::Receipt>>> {
882        self.provider()?.receipts_by_block_range(block_range)
883    }
884}
885
886impl<N: ProviderNodeTypes> BlockBodyIndicesProvider for ProviderFactory<N> {
887    fn block_body_indices(
888        &self,
889        number: BlockNumber,
890    ) -> ProviderResult<Option<StoredBlockBodyIndices>> {
891        self.provider()?.block_body_indices(number)
892    }
893
894    fn block_body_indices_range(
895        &self,
896        range: RangeInclusive<BlockNumber>,
897    ) -> ProviderResult<Vec<StoredBlockBodyIndices>> {
898        self.provider()?.block_body_indices_range(range)
899    }
900}
901
902impl<N: ProviderNodeTypes> StageCheckpointReader for ProviderFactory<N> {
903    fn get_stage_checkpoint(&self, id: StageId) -> ProviderResult<Option<StageCheckpoint>> {
904        self.provider()?.get_stage_checkpoint(id)
905    }
906
907    fn get_stage_checkpoint_progress(&self, id: StageId) -> ProviderResult<Option<Vec<u8>>> {
908        self.provider()?.get_stage_checkpoint_progress(id)
909    }
910    fn get_all_checkpoints(&self) -> ProviderResult<Vec<(String, StageCheckpoint)>> {
911        self.provider()?.get_all_checkpoints()
912    }
913}
914
915impl<N: NodeTypesWithDB> ChainSpecProvider for ProviderFactory<N> {
916    type ChainSpec = N::ChainSpec;
917
918    fn chain_spec(&self) -> Arc<N::ChainSpec> {
919        self.chain_spec.clone()
920    }
921}
922
923impl<N: ProviderNodeTypes> PruneCheckpointReader for ProviderFactory<N> {
924    fn get_prune_checkpoint(
925        &self,
926        segment: PruneSegment,
927    ) -> ProviderResult<Option<PruneCheckpoint>> {
928        self.provider()?.get_prune_checkpoint(segment)
929    }
930
931    fn get_prune_checkpoints(&self) -> ProviderResult<Vec<(PruneSegment, PruneCheckpoint)>> {
932        self.provider()?.get_prune_checkpoints()
933    }
934}
935
936impl<N: ProviderNodeTypes> MetadataProvider for ProviderFactory<N> {
937    fn get_metadata(&self, key: &str) -> ProviderResult<Option<Vec<u8>>> {
938        self.provider()?.get_metadata(key)
939    }
940}
941
942impl<N> fmt::Debug for ProviderFactory<N>
943where
944    N: NodeTypesWithDB<DB: fmt::Debug, ChainSpec: fmt::Debug, Storage: fmt::Debug>,
945{
946    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
947        let Self {
948            db,
949            chain_spec,
950            static_file_provider,
951            prune_modes,
952            storage,
953            storage_settings,
954            rocksdb_provider,
955            overlay_manager,
956            bal_store,
957            runtime,
958            minimum_pruning_distance,
959            database_provider_metrics: _,
960            read_only_sync,
961        } = self;
962        f.debug_struct("ProviderFactory")
963            .field("db", &db)
964            .field("chain_spec", &chain_spec)
965            .field("static_file_provider", &static_file_provider)
966            .field("prune_modes", &prune_modes)
967            .field("storage", &storage)
968            .field("storage_settings", &*storage_settings.read())
969            .field("rocksdb_provider", &rocksdb_provider)
970            .field("overlay_manager", &overlay_manager)
971            .field("bal_store", &bal_store)
972            .field("runtime", &runtime)
973            .field("minimum_pruning_distance", &minimum_pruning_distance)
974            .field(
975                "read_only_sync",
976                &read_only_sync.as_ref().map(|s| s.last_synced_txnid.load(Ordering::Relaxed)),
977            )
978            .finish()
979    }
980}
981
982impl<N: NodeTypesWithDB> Clone for ProviderFactory<N> {
983    fn clone(&self) -> Self {
984        Self {
985            db: self.db.clone(),
986            chain_spec: self.chain_spec.clone(),
987            static_file_provider: self.static_file_provider.clone(),
988            prune_modes: self.prune_modes.clone(),
989            storage: self.storage.clone(),
990            storage_settings: self.storage_settings.clone(),
991            rocksdb_provider: self.rocksdb_provider.clone(),
992            overlay_manager: self.overlay_manager.clone(),
993            bal_store: self.bal_store.clone(),
994            runtime: self.runtime.clone(),
995            minimum_pruning_distance: self.minimum_pruning_distance,
996            database_provider_metrics: self.database_provider_metrics.clone(),
997            read_only_sync: self.read_only_sync.clone(),
998        }
999    }
1000}
1001
1002#[cfg(test)]
1003mod tests {
1004    use super::*;
1005    use crate::{
1006        providers::{StaticFileProvider, StaticFileWriter},
1007        test_utils::{blocks::TEST_BLOCK, create_test_provider_factory, MockNodeTypesWithDB},
1008        BlockHashReader, BlockNumReader, BlockWriter, DBProvider, HeaderSyncGapProvider,
1009        TransactionsProvider,
1010    };
1011    use alloy_primitives::{TxNumber, B256};
1012    use assert_matches::assert_matches;
1013    use reth_chainspec::ChainSpecBuilder;
1014    use reth_db::{
1015        mdbx::DatabaseArguments,
1016        test_utils::{create_test_rocksdb_dir, create_test_static_files_dir, ERROR_TEMPDIR},
1017    };
1018    use reth_db_api::tables;
1019    use reth_primitives_traits::SignerRecoverable;
1020    use reth_prune_types::{PruneMode, PruneModes};
1021    use reth_storage_errors::provider::ProviderError;
1022    use reth_testing_utils::generators::{self, random_block, random_header, BlockParams};
1023    use std::{ops::RangeInclusive, sync::Arc};
1024
1025    #[test]
1026    fn common_history_provider() {
1027        let factory = create_test_provider_factory();
1028        let _ = factory.latest();
1029    }
1030
1031    #[test]
1032    fn default_chain_info() {
1033        let factory = create_test_provider_factory();
1034        let provider = factory.provider().unwrap();
1035
1036        let chain_info = provider.chain_info().expect("should be ok");
1037        assert_eq!(chain_info.best_number, 0);
1038        assert_eq!(chain_info.best_hash, B256::ZERO);
1039    }
1040
1041    #[test]
1042    fn provider_flow() {
1043        let factory = create_test_provider_factory();
1044        let provider = factory.provider().unwrap();
1045        provider.block_hash(0).unwrap();
1046        let provider_rw = factory.provider_rw().unwrap();
1047        provider_rw.block_hash(0).unwrap();
1048        provider.block_hash(0).unwrap();
1049    }
1050
1051    #[test]
1052    fn provider_factory_with_database_path() {
1053        let chain_spec = ChainSpecBuilder::mainnet().build();
1054        let (_static_dir, static_dir_path) = create_test_static_files_dir();
1055        let (_rocksdb_dir, rocksdb_path) = create_test_rocksdb_dir();
1056        let _db_tempdir = tempfile::TempDir::new().expect(ERROR_TEMPDIR);
1057        let factory = ProviderFactory::<MockNodeTypesWithDB<DatabaseEnv>>::new_with_database_path(
1058            _db_tempdir.path(),
1059            Arc::new(chain_spec),
1060            DatabaseArguments::new(Default::default()),
1061            StaticFileProvider::read_write(static_dir_path).unwrap(),
1062            RocksDBProvider::builder(&rocksdb_path).build().unwrap(),
1063            reth_tasks::Runtime::test(),
1064        )
1065        .unwrap();
1066        let provider = factory.provider().unwrap();
1067        provider.block_hash(0).unwrap();
1068        let provider_rw = factory.provider_rw().unwrap();
1069        provider_rw.block_hash(0).unwrap();
1070        provider.block_hash(0).unwrap();
1071    }
1072
1073    #[test]
1074    fn insert_block_with_prune_modes() {
1075        let block = TEST_BLOCK.clone();
1076
1077        {
1078            let factory = create_test_provider_factory();
1079            let provider = factory.provider_rw().unwrap();
1080            assert_matches!(provider.insert_block(&block.clone().try_recover().unwrap()), Ok(_));
1081            assert_matches!(
1082                provider.transaction_sender(0), Ok(Some(sender))
1083                if sender == block.body().transactions[0].recover_signer().unwrap()
1084            );
1085            assert_matches!(
1086                provider.transaction_id(*block.body().transactions[0].tx_hash()),
1087                Ok(Some(0))
1088            );
1089        }
1090
1091        {
1092            let prune_modes = PruneModes {
1093                sender_recovery: Some(PruneMode::Full),
1094                transaction_lookup: Some(PruneMode::Full),
1095                ..PruneModes::default()
1096            };
1097            // Keep factory alive until provider is dropped to prevent TempDatabase cleanup
1098            let factory = create_test_provider_factory().with_prune_modes(prune_modes);
1099            let provider = factory.provider_rw().unwrap();
1100            assert_matches!(provider.insert_block(&block.clone().try_recover().unwrap()), Ok(_));
1101            assert_matches!(provider.transaction_sender(0), Ok(None));
1102            assert_matches!(
1103                provider.transaction_id(*block.body().transactions[0].tx_hash()),
1104                Ok(None)
1105            );
1106        }
1107    }
1108
1109    #[test]
1110    fn take_block_transaction_range_recover_senders() {
1111        let mut rng = generators::rng();
1112        let block =
1113            random_block(&mut rng, 0, BlockParams { tx_count: Some(3), ..Default::default() });
1114
1115        let tx_ranges: Vec<RangeInclusive<TxNumber>> = vec![0..=0, 1..=1, 2..=2, 0..=1, 1..=2];
1116        for range in tx_ranges {
1117            let factory = create_test_provider_factory();
1118            let provider = factory.provider_rw().unwrap();
1119
1120            assert_matches!(provider.insert_block(&block.clone().try_recover().unwrap()), Ok(_));
1121
1122            let senders = provider.take::<tables::TransactionSenders>(range.clone()).unwrap();
1123            assert_eq!(
1124                senders,
1125                range
1126                    .clone()
1127                    .map(|tx_number| (
1128                        tx_number,
1129                        block.body().transactions[tx_number as usize].recover_signer().unwrap()
1130                    ))
1131                    .collect::<Vec<_>>()
1132            );
1133
1134            let db_senders = provider.senders_by_tx_range(range);
1135            assert!(matches!(db_senders, Ok(ref v) if v.is_empty()));
1136        }
1137    }
1138
1139    #[test]
1140    fn header_sync_gap_lookup() {
1141        let factory = create_test_provider_factory();
1142        let provider = factory.provider_rw().unwrap();
1143
1144        let mut rng = generators::rng();
1145
1146        // Genesis
1147        let checkpoint = 0;
1148        let head = random_header(&mut rng, 0, None);
1149
1150        // Empty database
1151        assert_matches!(
1152            provider.local_tip_header(checkpoint),
1153            Err(ProviderError::HeaderNotFound(block_number))
1154                if block_number.as_number().unwrap() == checkpoint
1155        );
1156
1157        // Checkpoint and no gap
1158        let static_file_provider = provider.static_file_provider();
1159        let mut static_file_writer =
1160            static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
1161        static_file_writer.append_header(head.header(), &head.hash()).unwrap();
1162        static_file_writer.commit().unwrap();
1163        drop(static_file_writer);
1164
1165        let local_head = provider.local_tip_header(checkpoint).unwrap();
1166
1167        assert_eq!(local_head, head);
1168    }
1169}