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_execution_types::RecoveredBlockAndExecutionOutput;
25use reth_node_types::{
26 BlockTy, HeaderTy, NodeTypesWithDB, NodeTypesWithDBAdapter, ReceiptTy, TxTy,
27};
28use reth_primitives_traits::{RecoveredBlock, SealedHeader};
29use reth_prune_types::{PruneCheckpoint, PruneModes, PruneSegment, MINIMUM_UNWIND_SAFE_DISTANCE};
30use reth_stages_types::{PipelineTarget, StageCheckpoint, StageId};
31use reth_static_file_types::StaticFileSegment;
32use reth_storage_api::{
33 BlockBodyIndicesProvider, ChainStateBlockReader, ChainStateBlockWriter, DBProvider,
34 NodePrimitivesProvider, StorageSettings, StorageSettingsCache,
35};
36use reth_storage_errors::provider::ProviderResult;
37use reth_storage_overlay::OverlayManager;
38use std::{
39 ops::{RangeBounds, RangeInclusive},
40 path::Path,
41 sync::{
42 atomic::{AtomicU64, Ordering},
43 Arc, Mutex,
44 },
45};
46use tracing::{info, instrument, trace, warn};
47
48mod provider;
49pub use provider::{CommitOrder, DatabaseProvider, DatabaseProviderRO, DatabaseProviderRW};
50
51mod save_blocks;
52pub use save_blocks::SaveBlocksInput;
53
54use super::ProviderNodeTypes;
55mod builder;
56pub use builder::{ProviderFactoryBuilder, ReadOnlyConfig};
57
58mod metrics;
59pub use metrics::DatabaseProviderMetrics;
60
61mod chain;
62pub use chain::*;
63
64struct ReadOnlySyncState {
66 last_synced_txnid: AtomicU64,
68 sync_lock: Mutex<()>,
70}
71
72pub struct ProviderFactory<N: NodeTypesWithDB> {
76 db: N::DB,
78 chain_spec: Arc<N::ChainSpec>,
80 static_file_provider: StaticFileProvider<N::Primitives>,
82 prune_modes: PruneModes,
84 storage: Arc<N::Storage>,
86 storage_settings: Arc<RwLock<StorageSettings>>,
88 rocksdb_provider: RocksDBProvider,
90 overlay_manager: OverlayManager<N::Primitives>,
92 bal_store: BalStoreHandle,
94 runtime: reth_tasks::Runtime,
96 minimum_pruning_distance: u64,
98 database_provider_metrics: Arc<DatabaseProviderMetrics>,
100 read_only_sync: Option<Arc<ReadOnlySyncState>>,
105}
106
107impl<N: NodeTypesForProvider> ProviderFactory<NodeTypesWithDBAdapter<N, DatabaseEnv>> {
108 pub fn builder() -> ProviderFactoryBuilder<N> {
110 ProviderFactoryBuilder::default()
111 }
112}
113
114impl<N: ProviderNodeTypes> ProviderFactory<N> {
115 pub fn new(
123 db: N::DB,
124 chain_spec: Arc<N::ChainSpec>,
125 static_file_provider: StaticFileProvider<N::Primitives>,
126 rocksdb_provider: RocksDBProvider,
127 runtime: reth_tasks::Runtime,
128 ) -> ProviderResult<Self> {
129 let legacy_settings = StorageSettings::v1();
134 let database_provider_metrics = Arc::new(DatabaseProviderMetrics::default());
135 let overlay_manager = OverlayManager::default();
136 let storage_settings = DatabaseProvider::<_, N>::new(
137 db.tx()?,
138 chain_spec.clone(),
139 static_file_provider.clone(),
140 Default::default(),
141 Default::default(),
142 Arc::new(RwLock::new(legacy_settings)),
143 rocksdb_provider.clone(),
144 overlay_manager.clone(),
145 runtime.clone(),
146 db.path(),
147 database_provider_metrics.clone(),
148 )
149 .storage_settings()?
150 .unwrap_or(legacy_settings);
151
152 Ok(Self {
153 db,
154 chain_spec,
155 static_file_provider,
156 prune_modes: PruneModes::default(),
157 storage: Default::default(),
158 storage_settings: Arc::new(RwLock::new(storage_settings)),
159 rocksdb_provider,
160 overlay_manager,
161 bal_store: BalStoreHandle::new(InMemoryBalStore::default()),
162 runtime,
163 minimum_pruning_distance: MINIMUM_UNWIND_SAFE_DISTANCE,
164 database_provider_metrics,
165 read_only_sync: None,
166 })
167 }
168
169 pub fn new_checked(
176 db: N::DB,
177 chain_spec: Arc<N::ChainSpec>,
178 static_file_provider: StaticFileProvider<N::Primitives>,
179 rocksdb_provider: RocksDBProvider,
180 runtime: reth_tasks::Runtime,
181 ) -> ProviderResult<Self> {
182 Self::new(db, chain_spec, static_file_provider, rocksdb_provider, runtime)
183 .and_then(Self::assert_consistent)
184 }
185}
186
187impl<N: NodeTypesWithDB> ProviderFactory<N> {
188 pub fn with_prune_modes(mut self, prune_modes: PruneModes) -> Self {
190 self.prune_modes = prune_modes;
191 self
192 }
193
194 pub fn with_bal_store(mut self, bal_store: BalStoreHandle) -> Self {
196 self.bal_store = bal_store;
197 self
198 }
199
200 pub fn with_overlay_manager(mut self, overlay_manager: OverlayManager<N::Primitives>) -> Self {
202 self.overlay_manager = overlay_manager;
203 self
204 }
205
206 pub(crate) const fn overlay_manager(&self) -> &OverlayManager<N::Primitives> {
208 &self.overlay_manager
209 }
210
211 pub const fn with_minimum_pruning_distance(mut self, distance: u64) -> Self {
216 self.minimum_pruning_distance = distance;
217 self
218 }
219
220 pub fn with_read_only_sync(mut self, watch: bool) -> Self
226 where
227 N::DB: Database,
228 {
229 let state = Arc::new(ReadOnlySyncState {
233 last_synced_txnid: AtomicU64::new(0),
234 sync_lock: Mutex::new(()),
235 });
236 self.read_only_sync = Some(state);
237
238 if watch {
239 self.watch_db_directory();
240 }
241 self
242 }
243
244 fn watch_db_directory(&self)
247 where
248 N::DB: Database,
249 {
250 let factory = self.clone();
251 let db_path = self.db.path();
252 reth_tasks::spawn_os_thread("ro-sync", move || {
253 let (tx, rx) = std::sync::mpsc::channel();
254 let mut watcher = RecommendedWatcher::new(
255 move |res| {
256 let _ = tx.send(res);
257 },
258 notify::Config::default(),
259 )
260 .expect("failed to create watcher");
261
262 watcher
263 .watch(&db_path, RecursiveMode::NonRecursive)
264 .expect("failed to watch MDBX path");
265
266 while let Ok(res) = rx.recv() {
267 match res {
268 Ok(event) => {
269 if !matches!(
270 event.kind,
271 notify::EventKind::Modify(_) | notify::EventKind::Create(_)
272 ) {
273 continue;
274 }
275
276 if let Err(err) = factory.sync_providers_if_needed() {
277 warn!(target: "reth::provider", %err, "background ro-sync failed");
278 }
279 }
280 Err(err) => {
281 warn!(target: "reth::provider", ?err, "MDBX directory watcher error");
282 }
283 }
284 }
285 });
286 }
287
288 pub fn sync_providers_if_needed(&self) -> ProviderResult<()> {
294 let Some(sync_state) = &self.read_only_sync else { return Ok(()) };
295 let current_txnid = self.db.last_txnid().unwrap_or(0);
296
297 if current_txnid == sync_state.last_synced_txnid.load(Ordering::Relaxed) {
299 return Ok(());
300 }
301
302 let _guard = sync_state.sync_lock.lock().unwrap_or_else(|e| e.into_inner());
304
305 if current_txnid == sync_state.last_synced_txnid.load(Ordering::Relaxed) {
307 return Ok(());
308 }
309
310 self.rocksdb_provider.try_catch_up_with_primary()?;
311 self.static_file_provider.initialize_index()?;
312 sync_state.last_synced_txnid.store(current_txnid, Ordering::Relaxed);
313 Ok(())
314 }
315
316 pub const fn db_ref(&self) -> &N::DB {
318 &self.db
319 }
320
321 #[cfg(any(test, feature = "test-utils"))]
322 pub fn into_db(self) -> N::DB {
324 self.db
325 }
326}
327
328impl<N: NodeTypesWithDB> StorageSettingsCache for ProviderFactory<N> {
329 fn cached_storage_settings(&self) -> StorageSettings {
330 *self.storage_settings.read()
331 }
332
333 fn set_storage_settings_cache(&self, settings: StorageSettings) {
334 *self.storage_settings.write() = settings;
335 }
336}
337
338impl<N: NodeTypesWithDB> RocksDBProviderFactory for ProviderFactory<N> {
339 fn rocksdb_provider(&self) -> RocksDBProvider {
340 self.rocksdb_provider.clone()
341 }
342
343 fn set_pending_rocksdb_batch(&self, _batch: rocksdb::WriteBatchWithTransaction<true>) {
344 unimplemented!("ProviderFactory is a factory, not a provider - use DatabaseProvider::set_pending_rocksdb_batch instead")
345 }
346
347 fn commit_pending_rocksdb_batches(&self) -> ProviderResult<()> {
348 unimplemented!("ProviderFactory is a factory, not a provider - use DatabaseProvider::commit_pending_rocksdb_batches instead")
349 }
350}
351
352impl<N: ProviderNodeTypes<DB = DatabaseEnv>> ProviderFactory<N> {
353 pub fn new_with_database_path<P: AsRef<Path>>(
356 path: P,
357 chain_spec: Arc<N::ChainSpec>,
358 args: DatabaseArguments,
359 static_file_provider: StaticFileProvider<N::Primitives>,
360 rocksdb_provider: RocksDBProvider,
361 runtime: reth_tasks::Runtime,
362 ) -> RethResult<Self> {
363 Self::new(
364 init_db(path, args).map_err(RethError::msg)?,
365 chain_spec,
366 static_file_provider,
367 rocksdb_provider,
368 runtime,
369 )
370 .map_err(RethError::Provider)
371 }
372}
373
374impl<N: ProviderNodeTypes> ProviderFactory<N> {
375 #[track_caller]
382 pub fn provider(&self) -> ProviderResult<DatabaseProviderRO<N::DB, N>> {
383 let db_tx = self.db.tx()?;
384
385 self.sync_providers_if_needed()?;
391
392 Ok(DatabaseProvider::new(
393 db_tx,
394 self.chain_spec.clone(),
395 self.static_file_provider.clone(),
396 self.prune_modes.clone(),
397 self.storage.clone(),
398 self.storage_settings.clone(),
399 self.rocksdb_provider.clone(),
400 self.overlay_manager.clone(),
401 self.runtime.clone(),
402 self.db.path(),
403 self.database_provider_metrics.clone(),
404 )
405 .with_minimum_pruning_distance(self.minimum_pruning_distance))
406 }
407
408 #[track_caller]
413 pub fn provider_rw(&self) -> ProviderResult<DatabaseProviderRW<N::DB, N>> {
414 Ok(DatabaseProviderRW(
415 DatabaseProvider::new_rw(
416 self.db.tx_mut()?,
417 self.chain_spec.clone(),
418 self.static_file_provider.clone(),
419 self.prune_modes.clone(),
420 self.storage.clone(),
421 self.storage_settings.clone(),
422 self.rocksdb_provider.clone(),
423 self.overlay_manager.clone(),
424 self.runtime.clone(),
425 self.db.path(),
426 self.database_provider_metrics.clone(),
427 )
428 .with_reader_txn_tracker(self.db.clone())
429 .with_minimum_pruning_distance(self.minimum_pruning_distance),
430 ))
431 }
432
433 #[track_caller]
440 pub fn unwind_provider_rw(
441 &self,
442 ) -> ProviderResult<DatabaseProvider<<N::DB as Database>::TXMut, N>> {
443 Ok(DatabaseProvider::new_unwind_rw(
444 self.db.tx_mut()?,
445 self.chain_spec.clone(),
446 self.static_file_provider.clone(),
447 self.prune_modes.clone(),
448 self.storage.clone(),
449 self.storage_settings.clone(),
450 self.rocksdb_provider.clone(),
451 self.overlay_manager.clone(),
452 self.runtime.clone(),
453 self.db.path(),
454 self.database_provider_metrics.clone(),
455 )
456 .with_reader_txn_tracker(self.db.clone())
457 .with_minimum_pruning_distance(self.minimum_pruning_distance))
458 }
459
460 #[track_caller]
462 pub fn latest(&self) -> ProviderResult<StateProviderBox> {
463 trace!(target: "providers::db", "Returning latest state provider");
464 Ok(Box::new(LatestStateProvider::new(self.database_provider_ro()?)))
465 }
466
467 pub fn assert_consistent(self) -> ProviderResult<Self> {
472 let (rocksdb_unwind, static_file_unwind) = self.check_consistency()?;
473
474 let source = match (rocksdb_unwind, static_file_unwind) {
475 (None, None) => return Ok(self),
476 (Some(_), Some(_)) => "RocksDB and Static Files",
477 (Some(_), None) => "RocksDB",
478 (None, Some(_)) => "Static Files",
479 };
480
481 Err(ProviderError::MustUnwind {
482 data_source: source,
483 unwind_to: rocksdb_unwind
484 .into_iter()
485 .chain(static_file_unwind)
486 .min()
487 .expect("at least one unwind target must be Some"),
488 })
489 }
490
491 #[instrument(err, skip(self))]
495 pub fn check_consistency(&self) -> ProviderResult<(Option<u64>, Option<u64>)> {
496 let provider_ro = self
497 .database_provider_ro()?
498 .disable_long_read_transaction_safety();
502
503 self.static_file_provider().check_file_consistency(&provider_ro)?;
505
506 let rocksdb_unwind = self.rocksdb_provider().check_consistency(&provider_ro)?;
508
509 let static_file_unwind = self.static_file_provider().check_consistency(&provider_ro)?.map(
511 |target| match target {
512 PipelineTarget::Unwind(block) => block,
513 PipelineTarget::Sync(_) => unreachable!("check_consistency returns Unwind"),
514 },
515 );
516
517 if rocksdb_unwind.is_none() && static_file_unwind.is_none() {
522 self.heal_chain_state_block_numbers(&provider_ro)?;
523 }
524
525 Ok((rocksdb_unwind, static_file_unwind))
526 }
527
528 fn heal_chain_state_block_numbers(
531 &self,
532 provider_ro: &DatabaseProvider<<N::DB as Database>::TX, N>,
533 ) -> ProviderResult<()> {
534 let highest_header = self.last_block_number()?;
535
536 let finalized = provider_ro.last_finalized_block_number()?;
537 let safe = provider_ro.last_safe_block_number()?;
538
539 if finalized.is_none_or(|f| f <= highest_header) && safe.is_none_or(|s| s <= highest_header)
540 {
541 return Ok(());
542 }
543
544 let provider_rw = self.provider_rw()?;
545
546 if let Some(finalized) = finalized.filter(|&f| f > highest_header) {
547 info!(
548 target: "providers::db",
549 finalized,
550 highest_header,
551 "Healing finalized block number",
552 );
553 provider_rw.save_finalized_block_number(highest_header)?;
554 }
555
556 if let Some(safe) = safe.filter(|&s| s > highest_header) {
557 info!(
558 target: "providers::db",
559 safe,
560 highest_header,
561 "Healing safe block number",
562 );
563 provider_rw.save_safe_block_number(highest_header)?;
564 }
565
566 provider_rw.commit()?;
567
568 Ok(())
569 }
570
571 pub fn caught_up_static_file_provider(
574 &self,
575 ) -> ProviderResult<StaticFileProvider<N::Primitives>> {
576 self.sync_providers_if_needed()?;
577 Ok(self.static_file_provider.clone())
578 }
579}
580
581impl<N: NodeTypesWithDB> NodePrimitivesProvider for ProviderFactory<N> {
582 type Primitives = N::Primitives;
583}
584
585impl<N: ProviderNodeTypes> BalProvider for ProviderFactory<N> {
586 fn bal_store(&self) -> &BalStoreHandle {
587 &self.bal_store
588 }
589}
590
591impl<N: ProviderNodeTypes> DatabaseProviderFactory for ProviderFactory<N> {
592 type DB = N::DB;
593 type Provider = DatabaseProvider<<N::DB as Database>::TX, N>;
594 type ProviderRW = DatabaseProvider<<N::DB as Database>::TXMut, N>;
595
596 fn database_provider_ro(&self) -> ProviderResult<Self::Provider> {
597 self.provider()
598 }
599
600 fn database_provider_rw(&self) -> ProviderResult<Self::ProviderRW> {
601 self.provider_rw().map(|provider| provider.0)
602 }
603}
604
605impl<N: NodeTypesWithDB> StaticFileProviderFactory for ProviderFactory<N> {
606 fn static_file_provider(&self) -> StaticFileProvider<Self::Primitives> {
608 self.static_file_provider.clone()
609 }
610
611 fn get_static_file_writer(
612 &self,
613 block: BlockNumber,
614 segment: StaticFileSegment,
615 ) -> ProviderResult<StaticFileProviderRWRefMut<'_, Self::Primitives>> {
616 self.static_file_provider.get_writer(block, segment)
617 }
618}
619
620impl<N: ProviderNodeTypes> HeaderSyncGapProvider for ProviderFactory<N> {
621 type Header = HeaderTy<N>;
622 fn local_tip_header(
623 &self,
624 highest_uninterrupted_block: BlockNumber,
625 ) -> ProviderResult<SealedHeader<Self::Header>> {
626 self.provider()?.local_tip_header(highest_uninterrupted_block)
627 }
628}
629
630impl<N: ProviderNodeTypes> HeaderProvider for ProviderFactory<N> {
631 type Header = HeaderTy<N>;
632
633 fn header(&self, block_hash: BlockHash) -> ProviderResult<Option<Self::Header>> {
634 self.provider()?.header(block_hash)
635 }
636
637 fn header_by_number(&self, num: BlockNumber) -> ProviderResult<Option<Self::Header>> {
638 self.caught_up_static_file_provider()?.header_by_number(num)
639 }
640
641 fn headers_range(
642 &self,
643 range: impl RangeBounds<BlockNumber>,
644 ) -> ProviderResult<Vec<Self::Header>> {
645 self.caught_up_static_file_provider()?.headers_range(range)
646 }
647
648 fn sealed_header(
649 &self,
650 number: BlockNumber,
651 ) -> ProviderResult<Option<SealedHeader<Self::Header>>> {
652 self.caught_up_static_file_provider()?.sealed_header(number)
653 }
654
655 fn sealed_headers_range(
656 &self,
657 range: impl RangeBounds<BlockNumber>,
658 ) -> ProviderResult<Vec<SealedHeader<Self::Header>>> {
659 self.caught_up_static_file_provider()?.sealed_headers_range(range)
660 }
661
662 fn sealed_headers_while(
663 &self,
664 range: impl RangeBounds<BlockNumber>,
665 predicate: impl FnMut(&SealedHeader<Self::Header>) -> bool,
666 ) -> ProviderResult<Vec<SealedHeader<Self::Header>>> {
667 self.caught_up_static_file_provider()?.sealed_headers_while(range, predicate)
668 }
669}
670
671impl<N: ProviderNodeTypes> BlockHashReader for ProviderFactory<N> {
672 fn block_hash(&self, number: u64) -> ProviderResult<Option<B256>> {
673 self.caught_up_static_file_provider()?.block_hash(number)
674 }
675
676 fn canonical_hashes_range(
677 &self,
678 start: BlockNumber,
679 end: BlockNumber,
680 ) -> ProviderResult<Vec<B256>> {
681 self.caught_up_static_file_provider()?.canonical_hashes_range(start, end)
682 }
683}
684
685impl<N: ProviderNodeTypes> BlockNumReader for ProviderFactory<N> {
686 fn chain_info(&self) -> ProviderResult<ChainInfo> {
687 self.provider()?.chain_info()
688 }
689
690 fn best_block_number(&self) -> ProviderResult<BlockNumber> {
691 self.provider()?.best_block_number()
692 }
693
694 fn last_block_number(&self) -> ProviderResult<BlockNumber> {
695 self.caught_up_static_file_provider()?.last_block_number()
696 }
697
698 fn earliest_block_number(&self) -> ProviderResult<BlockNumber> {
699 Ok(self.caught_up_static_file_provider()?.earliest_history_height())
702 }
703
704 fn block_number(&self, hash: B256) -> ProviderResult<Option<BlockNumber>> {
705 self.provider()?.block_number(hash)
706 }
707}
708
709impl<N: ProviderNodeTypes> BlockReader for ProviderFactory<N> {
710 type Block = BlockTy<N>;
711
712 fn find_block_by_hash(
713 &self,
714 hash: B256,
715 source: BlockSource,
716 ) -> ProviderResult<Option<Self::Block>> {
717 self.provider()?.find_block_by_hash(hash, source)
718 }
719
720 fn block(&self, id: BlockHashOrNumber) -> ProviderResult<Option<Self::Block>> {
721 self.provider()?.block(id)
722 }
723
724 fn pending_block(&self) -> ProviderResult<Option<Arc<RecoveredBlock<Self::Block>>>> {
725 self.provider()?.pending_block()
726 }
727
728 fn pending_block_and_receipts(
729 &self,
730 ) -> ProviderResult<Option<RecoveredBlockAndExecutionOutput<Self::Block, Self::Receipt>>> {
731 self.provider()?.pending_block_and_receipts()
732 }
733
734 fn recovered_block(
735 &self,
736 id: BlockHashOrNumber,
737 transaction_kind: TransactionVariant,
738 ) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
739 self.provider()?.recovered_block(id, transaction_kind)
740 }
741
742 fn sealed_block_with_senders(
743 &self,
744 id: BlockHashOrNumber,
745 transaction_kind: TransactionVariant,
746 ) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
747 self.provider()?.sealed_block_with_senders(id, transaction_kind)
748 }
749
750 fn block_range(&self, range: RangeInclusive<BlockNumber>) -> ProviderResult<Vec<Self::Block>> {
751 self.provider()?.block_range(range)
752 }
753
754 fn block_with_senders_range(
755 &self,
756 range: RangeInclusive<BlockNumber>,
757 ) -> ProviderResult<Vec<RecoveredBlock<Self::Block>>> {
758 self.provider()?.block_with_senders_range(range)
759 }
760
761 fn recovered_block_range(
762 &self,
763 range: RangeInclusive<BlockNumber>,
764 ) -> ProviderResult<Vec<RecoveredBlock<Self::Block>>> {
765 self.provider()?.recovered_block_range(range)
766 }
767
768 fn block_by_transaction_id(&self, id: TxNumber) -> ProviderResult<Option<BlockNumber>> {
769 self.provider()?.block_by_transaction_id(id)
770 }
771}
772
773impl<N: ProviderNodeTypes> TransactionsProvider for ProviderFactory<N> {
774 type Transaction = TxTy<N>;
775
776 fn transaction_id(&self, tx_hash: TxHash) -> ProviderResult<Option<TxNumber>> {
777 self.provider()?.transaction_id(tx_hash)
778 }
779
780 fn transaction_by_id(&self, id: TxNumber) -> ProviderResult<Option<Self::Transaction>> {
781 self.caught_up_static_file_provider()?.transaction_by_id(id)
782 }
783
784 fn transaction_by_id_unhashed(
785 &self,
786 id: TxNumber,
787 ) -> ProviderResult<Option<Self::Transaction>> {
788 self.caught_up_static_file_provider()?.transaction_by_id_unhashed(id)
789 }
790
791 fn transaction_by_hash(&self, hash: TxHash) -> ProviderResult<Option<Self::Transaction>> {
792 self.provider()?.transaction_by_hash(hash)
793 }
794
795 fn transaction_by_hash_with_meta(
796 &self,
797 tx_hash: TxHash,
798 ) -> ProviderResult<Option<(Self::Transaction, TransactionMeta)>> {
799 self.provider()?.transaction_by_hash_with_meta(tx_hash)
800 }
801
802 fn transactions_by_block(
803 &self,
804 id: BlockHashOrNumber,
805 ) -> ProviderResult<Option<Vec<Self::Transaction>>> {
806 self.provider()?.transactions_by_block(id)
807 }
808
809 fn transactions_by_block_range(
810 &self,
811 range: impl RangeBounds<BlockNumber>,
812 ) -> ProviderResult<Vec<Vec<Self::Transaction>>> {
813 self.provider()?.transactions_by_block_range(range)
814 }
815
816 fn transactions_by_tx_range(
817 &self,
818 range: impl RangeBounds<TxNumber>,
819 ) -> ProviderResult<Vec<Self::Transaction>> {
820 self.caught_up_static_file_provider()?.transactions_by_tx_range(range)
821 }
822
823 fn senders_by_tx_range(
824 &self,
825 range: impl RangeBounds<TxNumber>,
826 ) -> ProviderResult<Vec<Address>> {
827 if EitherWriterDestination::senders(self).is_static_file() {
828 self.caught_up_static_file_provider()?.senders_by_tx_range(range)
829 } else {
830 self.provider()?.senders_by_tx_range(range)
831 }
832 }
833
834 fn transaction_sender(&self, id: TxNumber) -> ProviderResult<Option<Address>> {
835 if EitherWriterDestination::senders(self).is_static_file() {
836 self.caught_up_static_file_provider()?.transaction_sender(id)
837 } else {
838 self.provider()?.transaction_sender(id)
839 }
840 }
841}
842
843impl<N: ProviderNodeTypes> ReceiptProvider for ProviderFactory<N> {
844 type Receipt = ReceiptTy<N>;
845
846 fn receipt(&self, id: TxNumber) -> ProviderResult<Option<Self::Receipt>> {
847 self.caught_up_static_file_provider()?.get_with_static_file_or_database(
848 StaticFileSegment::Receipts,
849 id,
850 |static_file| static_file.receipt(id),
851 || self.provider()?.receipt(id),
852 )
853 }
854
855 fn receipt_by_hash(&self, hash: TxHash) -> ProviderResult<Option<Self::Receipt>> {
856 self.provider()?.receipt_by_hash(hash)
857 }
858
859 fn receipts_by_block(
860 &self,
861 block: BlockHashOrNumber,
862 ) -> ProviderResult<Option<Vec<Self::Receipt>>> {
863 self.provider()?.receipts_by_block(block)
864 }
865
866 fn receipts_by_tx_range(
867 &self,
868 range: impl RangeBounds<TxNumber>,
869 ) -> ProviderResult<Vec<Self::Receipt>> {
870 self.caught_up_static_file_provider()?.get_range_with_static_file_or_database(
871 StaticFileSegment::Receipts,
872 to_range(range),
873 |static_file, range, _| static_file.receipts_by_tx_range(range),
874 |range, _| self.provider()?.receipts_by_tx_range(range),
875 |_| true,
876 )
877 }
878
879 fn receipts_by_block_range(
880 &self,
881 block_range: RangeInclusive<BlockNumber>,
882 ) -> ProviderResult<Vec<Vec<Self::Receipt>>> {
883 self.provider()?.receipts_by_block_range(block_range)
884 }
885}
886
887impl<N: ProviderNodeTypes> BlockBodyIndicesProvider for ProviderFactory<N> {
888 fn block_body_indices(
889 &self,
890 number: BlockNumber,
891 ) -> ProviderResult<Option<StoredBlockBodyIndices>> {
892 self.provider()?.block_body_indices(number)
893 }
894
895 fn block_body_indices_range(
896 &self,
897 range: RangeInclusive<BlockNumber>,
898 ) -> ProviderResult<Vec<StoredBlockBodyIndices>> {
899 self.provider()?.block_body_indices_range(range)
900 }
901}
902
903impl<N: ProviderNodeTypes> StageCheckpointReader for ProviderFactory<N> {
904 fn get_stage_checkpoint(&self, id: StageId) -> ProviderResult<Option<StageCheckpoint>> {
905 self.provider()?.get_stage_checkpoint(id)
906 }
907
908 fn get_stage_checkpoint_progress(&self, id: StageId) -> ProviderResult<Option<Vec<u8>>> {
909 self.provider()?.get_stage_checkpoint_progress(id)
910 }
911 fn get_all_checkpoints(&self) -> ProviderResult<Vec<(String, StageCheckpoint)>> {
912 self.provider()?.get_all_checkpoints()
913 }
914}
915
916impl<N: NodeTypesWithDB> ChainSpecProvider for ProviderFactory<N> {
917 type ChainSpec = N::ChainSpec;
918
919 fn chain_spec(&self) -> Arc<N::ChainSpec> {
920 self.chain_spec.clone()
921 }
922}
923
924impl<N: ProviderNodeTypes> PruneCheckpointReader for ProviderFactory<N> {
925 fn get_prune_checkpoint(
926 &self,
927 segment: PruneSegment,
928 ) -> ProviderResult<Option<PruneCheckpoint>> {
929 self.provider()?.get_prune_checkpoint(segment)
930 }
931
932 fn get_prune_checkpoints(&self) -> ProviderResult<Vec<(PruneSegment, PruneCheckpoint)>> {
933 self.provider()?.get_prune_checkpoints()
934 }
935}
936
937impl<N: ProviderNodeTypes> MetadataProvider for ProviderFactory<N> {
938 fn get_metadata(&self, key: &str) -> ProviderResult<Option<Vec<u8>>> {
939 self.provider()?.get_metadata(key)
940 }
941}
942
943impl<N> fmt::Debug for ProviderFactory<N>
944where
945 N: NodeTypesWithDB<DB: fmt::Debug, ChainSpec: fmt::Debug, Storage: fmt::Debug>,
946{
947 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
948 let Self {
949 db,
950 chain_spec,
951 static_file_provider,
952 prune_modes,
953 storage,
954 storage_settings,
955 rocksdb_provider,
956 overlay_manager,
957 bal_store,
958 runtime,
959 minimum_pruning_distance,
960 database_provider_metrics: _,
961 read_only_sync,
962 } = self;
963 f.debug_struct("ProviderFactory")
964 .field("db", &db)
965 .field("chain_spec", &chain_spec)
966 .field("static_file_provider", &static_file_provider)
967 .field("prune_modes", &prune_modes)
968 .field("storage", &storage)
969 .field("storage_settings", &*storage_settings.read())
970 .field("rocksdb_provider", &rocksdb_provider)
971 .field("overlay_manager", &overlay_manager)
972 .field("bal_store", &bal_store)
973 .field("runtime", &runtime)
974 .field("minimum_pruning_distance", &minimum_pruning_distance)
975 .field(
976 "read_only_sync",
977 &read_only_sync.as_ref().map(|s| s.last_synced_txnid.load(Ordering::Relaxed)),
978 )
979 .finish()
980 }
981}
982
983impl<N: NodeTypesWithDB> Clone for ProviderFactory<N> {
984 fn clone(&self) -> Self {
985 Self {
986 db: self.db.clone(),
987 chain_spec: self.chain_spec.clone(),
988 static_file_provider: self.static_file_provider.clone(),
989 prune_modes: self.prune_modes.clone(),
990 storage: self.storage.clone(),
991 storage_settings: self.storage_settings.clone(),
992 rocksdb_provider: self.rocksdb_provider.clone(),
993 overlay_manager: self.overlay_manager.clone(),
994 bal_store: self.bal_store.clone(),
995 runtime: self.runtime.clone(),
996 minimum_pruning_distance: self.minimum_pruning_distance,
997 database_provider_metrics: self.database_provider_metrics.clone(),
998 read_only_sync: self.read_only_sync.clone(),
999 }
1000 }
1001}
1002
1003#[cfg(test)]
1004mod tests {
1005 use super::*;
1006 use crate::{
1007 providers::{StaticFileProvider, StaticFileWriter},
1008 test_utils::{blocks::TEST_BLOCK, create_test_provider_factory, MockNodeTypesWithDB},
1009 BlockHashReader, BlockNumReader, BlockWriter, DBProvider, HeaderSyncGapProvider,
1010 TransactionsProvider,
1011 };
1012 use alloy_primitives::{TxNumber, B256};
1013 use assert_matches::assert_matches;
1014 use reth_chainspec::ChainSpecBuilder;
1015 use reth_db::{
1016 mdbx::DatabaseArguments,
1017 test_utils::{create_test_rocksdb_dir, create_test_static_files_dir, ERROR_TEMPDIR},
1018 };
1019 use reth_db_api::tables;
1020 use reth_primitives_traits::SignerRecoverable;
1021 use reth_prune_types::{PruneMode, PruneModes};
1022 use reth_storage_errors::provider::ProviderError;
1023 use reth_testing_utils::generators::{self, random_block, random_header, BlockParams};
1024 use std::{ops::RangeInclusive, sync::Arc};
1025
1026 #[test]
1027 fn common_history_provider() {
1028 let factory = create_test_provider_factory();
1029 let _ = factory.latest();
1030 }
1031
1032 #[test]
1033 fn default_chain_info() {
1034 let factory = create_test_provider_factory();
1035 let provider = factory.provider().unwrap();
1036
1037 let chain_info = provider.chain_info().expect("should be ok");
1038 assert_eq!(chain_info.best_number, 0);
1039 assert_eq!(chain_info.best_hash, B256::ZERO);
1040 }
1041
1042 #[test]
1043 fn provider_flow() {
1044 let factory = create_test_provider_factory();
1045 let provider = factory.provider().unwrap();
1046 provider.block_hash(0).unwrap();
1047 let provider_rw = factory.provider_rw().unwrap();
1048 provider_rw.block_hash(0).unwrap();
1049 provider.block_hash(0).unwrap();
1050 }
1051
1052 #[test]
1053 fn provider_factory_with_database_path() {
1054 let chain_spec = ChainSpecBuilder::mainnet().build();
1055 let (_static_dir, static_dir_path) = create_test_static_files_dir();
1056 let (_rocksdb_dir, rocksdb_path) = create_test_rocksdb_dir();
1057 let _db_tempdir = tempfile::TempDir::new().expect(ERROR_TEMPDIR);
1058 let factory = ProviderFactory::<MockNodeTypesWithDB<DatabaseEnv>>::new_with_database_path(
1059 _db_tempdir.path(),
1060 Arc::new(chain_spec),
1061 DatabaseArguments::new(Default::default()),
1062 StaticFileProvider::read_write(static_dir_path).unwrap(),
1063 RocksDBProvider::builder(&rocksdb_path).build().unwrap(),
1064 reth_tasks::Runtime::test(),
1065 )
1066 .unwrap();
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 block_range_readers_reject_expired_history() {
1076 use crate::BlockReader;
1077 use reth_static_file_types::{SegmentHeader, SegmentRangeInclusive, StaticFileSegment};
1078
1079 let factory = create_test_provider_factory();
1080 let mut rng = generators::rng();
1081 let provider_rw = factory.provider_rw().unwrap();
1082 let mut parent = None;
1083 for number in 0..6 {
1084 let block = random_block(
1085 &mut rng,
1086 number,
1087 BlockParams { parent, tx_count: Some(1), ..Default::default() },
1088 );
1089 parent = Some(block.hash());
1090 provider_rw.insert_block(&block.try_recover().unwrap()).unwrap();
1091 }
1092 provider_rw.commit().unwrap();
1093
1094 let static_provider = factory.static_file_provider();
1096 {
1097 let mut writer =
1098 static_provider.latest_writer(StaticFileSegment::Transactions).unwrap();
1099 let header = writer.user_header().clone();
1100 *writer.user_header_mut() = SegmentHeader::new(
1101 header.expected_block_range(),
1102 Some(SegmentRangeInclusive::new(3, 5)),
1103 header.tx_range(),
1104 StaticFileSegment::Transactions,
1105 );
1106 writer.inner().set_dirty();
1107 writer.commit().unwrap();
1108 }
1109 static_provider.initialize_index().unwrap();
1110 assert_eq!(static_provider.earliest_history_height(), 3);
1111
1112 let provider = factory.provider().unwrap();
1113 assert_matches!(
1115 provider.block(1.into()),
1116 Err(ProviderError::BlockExpired { requested: 1, earliest_available: 3 })
1117 );
1118 assert_matches!(
1119 provider.recovered_block(1.into(), Default::default()),
1120 Err(ProviderError::BlockExpired { requested: 1, earliest_available: 3 })
1121 );
1122 assert_matches!(
1124 provider.block_range(1..=4),
1125 Err(ProviderError::BlockExpired { requested: 1, earliest_available: 3 })
1126 );
1127 assert_matches!(
1128 provider.block_with_senders_range(1..=4),
1129 Err(ProviderError::BlockExpired { requested: 1, earliest_available: 3 })
1130 );
1131 assert_matches!(
1132 provider.recovered_block_range(1..=4),
1133 Err(ProviderError::BlockExpired { requested: 1, earliest_available: 3 })
1134 );
1135 assert_eq!(provider.block_range(3..=5).unwrap().len(), 3);
1137 assert_eq!(provider.recovered_block_range(3..=5).unwrap().len(), 3);
1138 }
1139
1140 #[test]
1141 fn insert_block_with_prune_modes() {
1142 let block = TEST_BLOCK.clone();
1143
1144 {
1145 let factory = create_test_provider_factory();
1146 let provider = factory.provider_rw().unwrap();
1147 assert_matches!(provider.insert_block(&block.clone().try_recover().unwrap()), Ok(_));
1148 assert_matches!(
1149 provider.transaction_sender(0), Ok(Some(sender))
1150 if sender == block.body().transactions[0].recover_signer().unwrap()
1151 );
1152 assert_matches!(
1153 provider.transaction_id(*block.body().transactions[0].tx_hash()),
1154 Ok(Some(0))
1155 );
1156 }
1157
1158 {
1159 let prune_modes = PruneModes {
1160 sender_recovery: Some(PruneMode::Full),
1161 transaction_lookup: Some(PruneMode::Full),
1162 ..PruneModes::default()
1163 };
1164 let factory = create_test_provider_factory().with_prune_modes(prune_modes);
1166 let provider = factory.provider_rw().unwrap();
1167 assert_matches!(provider.insert_block(&block.clone().try_recover().unwrap()), Ok(_));
1168 assert_matches!(provider.transaction_sender(0), Ok(None));
1169 assert_matches!(
1170 provider.transaction_id(*block.body().transactions[0].tx_hash()),
1171 Ok(None)
1172 );
1173 }
1174 }
1175
1176 #[test]
1177 fn take_block_transaction_range_recover_senders() {
1178 let mut rng = generators::rng();
1179 let block =
1180 random_block(&mut rng, 0, BlockParams { tx_count: Some(3), ..Default::default() });
1181
1182 let tx_ranges: Vec<RangeInclusive<TxNumber>> = vec![0..=0, 1..=1, 2..=2, 0..=1, 1..=2];
1183 for range in tx_ranges {
1184 let factory = create_test_provider_factory();
1185 let provider = factory.provider_rw().unwrap();
1186
1187 assert_matches!(provider.insert_block(&block.clone().try_recover().unwrap()), Ok(_));
1188
1189 let senders = provider.take::<tables::TransactionSenders>(range.clone()).unwrap();
1190 assert_eq!(
1191 senders,
1192 range
1193 .clone()
1194 .map(|tx_number| (
1195 tx_number,
1196 block.body().transactions[tx_number as usize].recover_signer().unwrap()
1197 ))
1198 .collect::<Vec<_>>()
1199 );
1200
1201 let db_senders = provider.senders_by_tx_range(range);
1202 assert!(matches!(db_senders, Ok(ref v) if v.is_empty()));
1203 }
1204 }
1205
1206 #[test]
1207 fn header_sync_gap_lookup() {
1208 let factory = create_test_provider_factory();
1209 let provider = factory.provider_rw().unwrap();
1210
1211 let mut rng = generators::rng();
1212
1213 let checkpoint = 0;
1215 let head = random_header(&mut rng, 0, None);
1216
1217 assert_matches!(
1219 provider.local_tip_header(checkpoint),
1220 Err(ProviderError::HeaderNotFound(block_number))
1221 if block_number.as_number().unwrap() == checkpoint
1222 );
1223
1224 let static_file_provider = provider.static_file_provider();
1226 let mut static_file_writer =
1227 static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
1228 static_file_writer.append_header(head.header(), &head.hash()).unwrap();
1229 static_file_writer.commit().unwrap();
1230 drop(static_file_writer);
1231
1232 let local_head = provider.local_tip_header(checkpoint).unwrap();
1233
1234 assert_eq!(local_head, head);
1235 }
1236
1237 #[test]
1238 fn snap_sync_requires_the_hashed_state_layout() {
1239 let factory = create_test_provider_factory();
1240
1241 factory.set_storage_settings_cache(StorageSettings::v2());
1242 assert!(factory.database_provider_ro().unwrap().ensure_snap_sync_layout().is_ok());
1243
1244 factory.set_storage_settings_cache(StorageSettings::v1());
1245 assert_matches!(
1246 factory.database_provider_ro().unwrap().ensure_snap_sync_layout(),
1247 Err(ProviderError::SnapStorageLayoutUnsupported)
1248 );
1249 }
1250}