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
63struct ReadOnlySyncState {
65 last_synced_txnid: AtomicU64,
67 sync_lock: Mutex<()>,
69}
70
71pub struct ProviderFactory<N: NodeTypesWithDB> {
75 db: N::DB,
77 chain_spec: Arc<N::ChainSpec>,
79 static_file_provider: StaticFileProvider<N::Primitives>,
81 prune_modes: PruneModes,
83 storage: Arc<N::Storage>,
85 storage_settings: Arc<RwLock<StorageSettings>>,
87 rocksdb_provider: RocksDBProvider,
89 overlay_manager: OverlayManager<N::Primitives>,
91 bal_store: BalStoreHandle,
93 runtime: reth_tasks::Runtime,
95 minimum_pruning_distance: u64,
97 database_provider_metrics: Arc<DatabaseProviderMetrics>,
99 read_only_sync: Option<Arc<ReadOnlySyncState>>,
104}
105
106impl<N: NodeTypesForProvider> ProviderFactory<NodeTypesWithDBAdapter<N, DatabaseEnv>> {
107 pub fn builder() -> ProviderFactoryBuilder<N> {
109 ProviderFactoryBuilder::default()
110 }
111}
112
113impl<N: ProviderNodeTypes> ProviderFactory<N> {
114 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 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 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 pub fn with_prune_modes(mut self, prune_modes: PruneModes) -> Self {
189 self.prune_modes = prune_modes;
190 self
191 }
192
193 pub fn with_bal_store(mut self, bal_store: BalStoreHandle) -> Self {
195 self.bal_store = bal_store;
196 self
197 }
198
199 pub fn with_overlay_manager(mut self, overlay_manager: OverlayManager<N::Primitives>) -> Self {
201 self.overlay_manager = overlay_manager;
202 self
203 }
204
205 pub(crate) const fn overlay_manager(&self) -> &OverlayManager<N::Primitives> {
207 &self.overlay_manager
208 }
209
210 pub const fn with_minimum_pruning_distance(mut self, distance: u64) -> Self {
215 self.minimum_pruning_distance = distance;
216 self
217 }
218
219 pub fn with_read_only_sync(mut self, watch: bool) -> Self
225 where
226 N::DB: Database,
227 {
228 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 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 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 if current_txnid == sync_state.last_synced_txnid.load(Ordering::Relaxed) {
298 return Ok(());
299 }
300
301 let _guard = sync_state.sync_lock.lock().unwrap_or_else(|e| e.into_inner());
303
304 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 pub const fn db_ref(&self) -> &N::DB {
317 &self.db
318 }
319
320 #[cfg(any(test, feature = "test-utils"))]
321 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 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 #[track_caller]
381 pub fn provider(&self) -> ProviderResult<DatabaseProviderRO<N::DB, N>> {
382 let db_tx = self.db.tx()?;
383
384 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 #[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 #[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 #[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 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 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 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 #[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 .disable_long_read_transaction_safety();
524
525 self.static_file_provider().check_file_consistency(&provider_ro)?;
527
528 let rocksdb_unwind = self.rocksdb_provider().check_consistency(&provider_ro)?;
530
531 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 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 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 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 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 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 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 let checkpoint = 0;
1171 let head = random_header(&mut rng, 0, None);
1172
1173 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 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}