1use super::{
2 metrics::StaticFileProviderMetrics, writer::StaticFileWriters, LoadedJar,
3 StaticFileJarProvider, StaticFileProviderRW, StaticFileProviderRWRefMut,
4};
5use crate::{
6 changeset_walker::{StaticFileAccountChangesetWalker, StaticFileStorageChangesetWalker},
7 to_range, BlockHashReader, BlockNumReader, BlockReader, BlockSource, EitherWriter,
8 EitherWriterDestination, HeaderProvider, ReceiptProvider, StageCheckpointReader, StatsReader,
9 TransactionVariant, TransactionsProvider, TransactionsProviderExt,
10};
11use alloy_consensus::{
12 transaction::{TransactionMeta, TxHashRef},
13 Header,
14};
15use alloy_eips::BlockHashOrNumber;
16use alloy_primitives::{b256, Address, BlockHash, BlockNumber, TxHash, TxNumber, B256};
17
18use parking_lot::RwLock;
19use reth_chain_state::ExecutedBlock;
20use reth_chainspec::{ChainInfo, ChainSpecProvider, EthChainSpec, NamedChain};
21use reth_db::{
22 lockfile::StorageLock,
23 static_file::{
24 iter_static_files, BlockHashMask, HeaderMask, HeaderWithHashMask, ReceiptMask,
25 StaticFileCursor, StorageChangesetMask, TransactionMask, TransactionSenderMask,
26 },
27};
28use reth_db_api::{
29 cursor::DbCursorRO,
30 models::{AccountBeforeTx, BlockNumberAddress, StorageBeforeTx, StoredBlockBodyIndices},
31 table::{Decompress, Table, Value},
32 tables,
33 transaction::DbTx,
34};
35use reth_ethereum_primitives::{Receipt, TransactionSigned};
36use reth_execution_types::RecoveredBlockAndExecutionOutput;
37use reth_nippy_jar::{NippyJar, NippyJarChecker};
38use reth_node_types::NodePrimitives;
39use reth_primitives_traits::{
40 dashmap::DashMap, AlloyBlockHeader as _, BlockBody as _, RecoveredBlock, SealedHeader,
41 SignedTransaction, StorageEntry,
42};
43use reth_prune_types::PruneSegment;
44use reth_stages_types::PipelineTarget;
45use reth_static_file_types::{
46 find_fixed_range, HighestStaticFiles, SegmentHeader, SegmentRangeInclusive, StaticFileMap,
47 StaticFileSegment, DEFAULT_BLOCKS_PER_STATIC_FILE,
48};
49use reth_storage_api::{
50 BlockBodyIndicesProvider, ChangeSetReader, DBProvider, PruneCheckpointReader,
51 StorageChangeSetReader, StorageSettingsCache,
52};
53use reth_storage_errors::provider::{ProviderError, ProviderResult, StaticFileWriterError};
54use std::{
55 collections::BTreeMap,
56 fmt::Debug,
57 ops::{Bound, Deref, Range, RangeBounds, RangeInclusive},
58 path::{Path, PathBuf},
59 sync::{
60 atomic::{AtomicU64, Ordering},
61 mpsc, Arc,
62 },
63};
64use tracing::{debug, info, info_span, instrument, trace, warn};
65
66type SegmentRanges = BTreeMap<u64, SegmentRangeInclusive>;
69
70#[derive(Debug, Default, PartialEq, Eq)]
72pub enum StaticFileAccess {
73 #[default]
75 RO,
76 RW,
78}
79
80impl StaticFileAccess {
81 pub const fn is_read_only(&self) -> bool {
83 matches!(self, Self::RO)
84 }
85
86 pub const fn is_read_write(&self) -> bool {
88 matches!(self, Self::RW)
89 }
90}
91
92#[derive(Debug, Clone, Copy, Default)]
96pub struct StaticFileWriteCtx {
97 pub write_senders: bool,
99 pub write_receipts: bool,
101 pub write_account_changesets: bool,
103 pub write_storage_changesets: bool,
105 pub tip: BlockNumber,
107 pub receipts_prune_mode: Option<reth_prune_types::PruneMode>,
109 pub receipts_prunable: bool,
111}
112
113#[derive(Debug)]
122pub struct StaticFileProvider<N>(pub(crate) Arc<StaticFileProviderInner<N>>);
123
124impl<N> Clone for StaticFileProvider<N> {
125 fn clone(&self) -> Self {
126 Self(self.0.clone())
127 }
128}
129
130#[derive(Debug)]
132pub struct StaticFileProviderBuilder<P> {
133 access: StaticFileAccess,
134 use_metrics: bool,
135 blocks_per_file: StaticFileMap<u64>,
136 path: P,
137 genesis_block_number: u64,
138}
139
140impl<P: AsRef<Path>> StaticFileProviderBuilder<P> {
141 pub fn read_write(path: P) -> Self {
143 Self {
144 path,
145 access: StaticFileAccess::RW,
146 blocks_per_file: Default::default(),
147 use_metrics: false,
148 genesis_block_number: 0,
149 }
150 }
151
152 pub fn read_only(path: P) -> Self {
154 Self {
155 path,
156 access: StaticFileAccess::RO,
157 blocks_per_file: Default::default(),
158 use_metrics: false,
159 genesis_block_number: 0,
160 }
161 }
162
163 pub fn with_blocks_per_file_for_segments(
175 mut self,
176 segments: &<StaticFileMap<u64> as Deref>::Target,
177 ) -> Self {
178 for (segment, &blocks_per_file) in segments {
179 self.blocks_per_file.insert(segment, blocks_per_file);
180 }
181 self
182 }
183
184 pub fn with_blocks_per_file(mut self, blocks_per_file: u64) -> Self {
186 for segment in StaticFileSegment::iter() {
187 self.blocks_per_file.insert(segment, blocks_per_file);
188 }
189 self
190 }
191
192 pub fn with_blocks_per_file_for_segment(
194 mut self,
195 segment: StaticFileSegment,
196 blocks_per_file: u64,
197 ) -> Self {
198 self.blocks_per_file.insert(segment, blocks_per_file);
199 self
200 }
201
202 pub const fn with_metrics(mut self) -> Self {
204 self.use_metrics = true;
205 self
206 }
207
208 pub const fn with_genesis_block_number(mut self, genesis_block_number: u64) -> Self {
221 self.genesis_block_number = genesis_block_number;
222 self
223 }
224
225 pub fn build<N: NodePrimitives>(self) -> ProviderResult<StaticFileProvider<N>> {
227 let mut provider = StaticFileProviderInner::new(self.path, self.access)?;
228 if self.use_metrics {
229 provider.metrics = Some(Arc::new(StaticFileProviderMetrics::default()));
230 }
231
232 for (segment, blocks_per_file) in *self.blocks_per_file {
233 provider.blocks_per_file.insert(segment, blocks_per_file);
234 }
235 provider.genesis_block_number = self.genesis_block_number;
236
237 let provider = StaticFileProvider(Arc::new(provider));
238 provider.initialize_index()?;
239 Ok(provider)
240 }
241}
242
243impl<N: NodePrimitives> StaticFileProvider<N> {
244 fn new(path: impl AsRef<Path>, access: StaticFileAccess) -> ProviderResult<Self> {
246 let provider = Self(Arc::new(StaticFileProviderInner::new(path, access)?));
247 provider.initialize_index()?;
248 Ok(provider)
249 }
250}
251
252impl<N: NodePrimitives> StaticFileProvider<N> {
253 pub fn read_only(path: impl AsRef<Path>) -> ProviderResult<Self> {
258 Self::new(path, StaticFileAccess::RO)
259 }
260
261 pub fn read_write(path: impl AsRef<Path>) -> ProviderResult<Self> {
263 Self::new(path, StaticFileAccess::RW)
264 }
265}
266
267impl<N: NodePrimitives> Deref for StaticFileProvider<N> {
268 type Target = StaticFileProviderInner<N>;
269
270 fn deref(&self) -> &Self::Target {
271 &self.0
272 }
273}
274
275#[cfg(test)]
276type BeforeCacheInsert = Box<dyn FnOnce(StaticFileSegment) + Send + 'static>;
277
278#[cfg(test)]
279#[derive(Default)]
280struct BeforeCacheInsertHook(parking_lot::Mutex<Option<BeforeCacheInsert>>);
281
282#[cfg(test)]
283impl Debug for BeforeCacheInsertHook {
284 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
285 f.debug_struct("BeforeCacheInsertHook").finish_non_exhaustive()
286 }
287}
288
289#[cfg(test)]
290impl BeforeCacheInsertHook {
291 fn set(&self, hook: impl FnOnce(StaticFileSegment) + Send + 'static) {
292 let previous = self.0.lock().replace(Box::new(hook));
293 assert!(previous.is_none(), "cache fill hook already installed");
294 }
295
296 fn fire(&self, segment: StaticFileSegment) {
297 if let Some(hook) = self.0.lock().take() {
298 hook(segment);
299 }
300 }
301}
302
303#[derive(Debug)]
305pub struct StaticFileProviderInner<N> {
306 map: DashMap<(BlockNumber, StaticFileSegment), LoadedJar>,
309 cache_generation: AtomicU64,
312 indexes: RwLock<StaticFileMap<StaticFileSegmentIndex>>,
314 earliest_history_height: AtomicU64,
325 path: PathBuf,
327 writers: StaticFileWriters<N>,
329 metrics: Option<Arc<StaticFileProviderMetrics>>,
331 access: StaticFileAccess,
333 blocks_per_file: StaticFileMap<u64>,
335 _lock_file: Option<StorageLock>,
337 genesis_block_number: u64,
339 #[cfg(test)]
340 before_cache_insert: BeforeCacheInsertHook,
341 #[cfg(test)]
342 after_cache_validation: BeforeCacheInsertHook,
343}
344
345impl<N: NodePrimitives> StaticFileProviderInner<N> {
346 fn new(path: impl AsRef<Path>, access: StaticFileAccess) -> ProviderResult<Self> {
348 let _lock_file = if access.is_read_write() {
349 StorageLock::try_acquire(path.as_ref()).map_err(ProviderError::other)?.into()
350 } else {
351 None
352 };
353
354 let mut blocks_per_file = StaticFileMap::default();
355 for segment in StaticFileSegment::iter() {
356 blocks_per_file.insert(segment, DEFAULT_BLOCKS_PER_STATIC_FILE);
357 }
358
359 let provider = Self {
360 map: Default::default(),
361 cache_generation: Default::default(),
362 indexes: Default::default(),
363 writers: Default::default(),
364 earliest_history_height: Default::default(),
365 path: path.as_ref().to_path_buf(),
366 metrics: None,
367 access,
368 blocks_per_file,
369 _lock_file,
370 genesis_block_number: 0,
371 #[cfg(test)]
372 before_cache_insert: Default::default(),
373 #[cfg(test)]
374 after_cache_validation: Default::default(),
375 };
376
377 Ok(provider)
378 }
379
380 pub const fn is_read_only(&self) -> bool {
381 self.access.is_read_only()
382 }
383
384 pub fn find_fixed_range_with_block_index(
393 &self,
394 segment: StaticFileSegment,
395 block_index: Option<&SegmentRanges>,
396 block: BlockNumber,
397 ) -> SegmentRangeInclusive {
398 let blocks_per_file =
399 self.blocks_per_file.get(segment).copied().unwrap_or(DEFAULT_BLOCKS_PER_STATIC_FILE);
400
401 if let Some(block_index) = block_index {
402 if let Some((_, range)) = block_index.range(block..).next() {
404 return *range;
406 } else if let Some((_, range)) = block_index.last_key_value() {
407 let blocks_after_last_range = block - range.end();
413 let segments_to_skip = (blocks_after_last_range - 1) / blocks_per_file;
414 let start = range.end() + 1 + segments_to_skip * blocks_per_file;
415 return SegmentRangeInclusive::new(start, start + blocks_per_file - 1);
416 }
417 }
418 find_fixed_range(block, blocks_per_file)
421 }
422
423 pub fn find_fixed_range(
436 &self,
437 segment: StaticFileSegment,
438 block: BlockNumber,
439 ) -> SegmentRangeInclusive {
440 self.find_fixed_range_with_block_index(
441 segment,
442 self.indexes.read().get(segment).map(|index| &index.expected_block_ranges_by_max_block),
443 block,
444 )
445 }
446
447 pub const fn genesis_block_number(&self) -> u64 {
449 self.genesis_block_number
450 }
451}
452
453impl<N: NodePrimitives> StaticFileProvider<N> {
454 pub fn report_metrics(&self) -> ProviderResult<()> {
459 let Some(metrics) = &self.metrics else { return Ok(()) };
460
461 let static_files = iter_static_files(&self.path).map_err(ProviderError::other)?;
462 for (segment, headers) in &*static_files {
463 let mut entries = 0;
464 let mut size = 0;
465
466 for (block_range, _) in headers {
467 let fixed_block_range = self.find_fixed_range(segment, block_range.start());
468 let jar_provider = self
469 .get_segment_provider_for_range(segment, || Some(fixed_block_range), None)?
470 .ok_or_else(|| {
471 ProviderError::MissingStaticFileBlock(segment, block_range.start())
472 })?;
473
474 entries += jar_provider.rows();
475 size += jar_provider.size() as u64;
476 }
477
478 metrics.record_segment(segment, size, headers.len(), entries);
479 }
480
481 Ok(())
482 }
483
484 #[instrument(level = "debug", target = "providers::static_file", skip_all)]
486 fn write_headers(
487 w: &mut StaticFileProviderRWRefMut<'_, N>,
488 blocks: &[ExecutedBlock<N>],
489 ) -> ProviderResult<()> {
490 for block in blocks {
491 let b = block.recovered_block();
492 w.append_header(b.header(), &b.hash())?;
493 }
494 Ok(())
495 }
496
497 #[instrument(level = "debug", target = "providers::static_file", skip_all)]
499 fn write_transactions(
500 w: &mut StaticFileProviderRWRefMut<'_, N>,
501 blocks: &[ExecutedBlock<N>],
502 tx_nums: &[TxNumber],
503 ) -> ProviderResult<()> {
504 for (block, &first_tx) in blocks.iter().zip(tx_nums) {
505 let b = block.recovered_block();
506 w.increment_block(b.number())?;
507 for (i, tx) in b.body().transactions().iter().enumerate() {
508 w.append_transaction(first_tx + i as u64, tx)?;
509 }
510 }
511 Ok(())
512 }
513
514 #[instrument(level = "debug", target = "providers::static_file", skip_all)]
516 fn write_transaction_senders(
517 w: &mut StaticFileProviderRWRefMut<'_, N>,
518 blocks: &[ExecutedBlock<N>],
519 tx_nums: &[TxNumber],
520 ) -> ProviderResult<()> {
521 for (block, &first_tx) in blocks.iter().zip(tx_nums) {
522 let b = block.recovered_block();
523 w.increment_block(b.number())?;
524 for (i, sender) in b.senders_iter().enumerate() {
525 w.append_transaction_sender(first_tx + i as u64, sender)?;
526 }
527 }
528 Ok(())
529 }
530
531 #[instrument(level = "debug", target = "providers::static_file", skip_all)]
533 fn write_receipts(
534 w: &mut StaticFileProviderRWRefMut<'_, N>,
535 blocks: &[ExecutedBlock<N>],
536 tx_nums: &[TxNumber],
537 ctx: &StaticFileWriteCtx,
538 ) -> ProviderResult<()> {
539 for (block, &first_tx) in blocks.iter().zip(tx_nums) {
540 let block_number = block.recovered_block().number();
541 w.increment_block(block_number)?;
542
543 if ctx.receipts_prunable &&
545 ctx.receipts_prune_mode
546 .is_some_and(|mode| mode.should_prune(block_number, ctx.tip))
547 {
548 continue
549 }
550
551 for (i, receipt) in block.execution_outcome().receipts.iter().enumerate() {
552 w.append_receipt(first_tx + i as u64, receipt)?;
553 }
554 }
555 Ok(())
556 }
557
558 #[instrument(level = "debug", target = "providers::static_file", skip_all)]
560 fn write_account_changesets(
561 w: &mut StaticFileProviderRWRefMut<'_, N>,
562 blocks: &[ExecutedBlock<N>],
563 ) -> ProviderResult<()> {
564 for block in blocks {
565 let block_number = block.recovered_block().number();
566 let reverts = block.execution_outcome().state.reverts.to_plain_state_reverts();
567
568 let changeset: Vec<_> = reverts
569 .accounts
570 .into_iter()
571 .flatten()
572 .map(|(address, info)| AccountBeforeTx { address, info: info.map(Into::into) })
573 .collect();
574 w.append_account_changeset(changeset, block_number)?;
575 }
576 Ok(())
577 }
578
579 #[instrument(level = "debug", target = "providers::db", skip_all)]
581 fn write_storage_changesets(
582 w: &mut StaticFileProviderRWRefMut<'_, N>,
583 blocks: &[ExecutedBlock<N>],
584 ) -> ProviderResult<()> {
585 for block in blocks {
586 let block_number = block.recovered_block().number();
587 let reverts = block.execution_outcome().state.reverts.to_plain_state_reverts();
588
589 let changeset: Vec<_> = reverts
590 .storage
591 .into_iter()
592 .flatten()
593 .flat_map(|revert| {
594 revert.storage_revert.into_iter().map(move |(key, revert_to_slot)| {
595 StorageBeforeTx {
596 address: revert.address,
597 key: B256::from(key.to_be_bytes()),
598 value: revert_to_slot.to_previous_value(),
599 }
600 })
601 })
602 .collect();
603 w.append_storage_changeset(changeset, block_number)?;
604 }
605 Ok(())
606 }
607
608 #[instrument(level = "debug", target = "providers::static_file", skip_all, fields(?segment))]
613 fn write_segment<F>(
614 &self,
615 segment: StaticFileSegment,
616 first_block_number: BlockNumber,
617 f: F,
618 ) -> ProviderResult<()>
619 where
620 F: FnOnce(&mut StaticFileProviderRWRefMut<'_, N>) -> ProviderResult<()>,
621 {
622 let mut w = self.get_writer(first_block_number, segment)?;
623 f(&mut w)?;
624 w.sync_all()
625 }
626
627 #[instrument(level = "debug", target = "providers::static_file", skip_all)]
632 pub fn write_blocks_data(
633 &self,
634 blocks: &[ExecutedBlock<N>],
635 tx_nums: &[TxNumber],
636 ctx: StaticFileWriteCtx,
637 runtime: &reth_tasks::Runtime,
638 ) -> ProviderResult<()> {
639 if blocks.is_empty() {
640 return Ok(());
641 }
642
643 let first_block_number = blocks[0].recovered_block().number();
644
645 let mut r_headers = None;
646 let mut r_txs = None;
647 let mut r_senders = None;
648 let mut r_receipts = None;
649 let mut r_account_changesets = None;
650 let mut r_storage_changesets = None;
651
652 let span = tracing::Span::current();
655 runtime.storage_pool().in_place_scope(|s| {
656 s.spawn(|_| {
657 let _guard = span.enter();
658 r_headers =
659 Some(self.write_segment(StaticFileSegment::Headers, first_block_number, |w| {
660 Self::write_headers(w, blocks)
661 }));
662 });
663
664 s.spawn(|_| {
665 let _guard = span.enter();
666 r_txs = Some(self.write_segment(
667 StaticFileSegment::Transactions,
668 first_block_number,
669 |w| Self::write_transactions(w, blocks, tx_nums),
670 ));
671 });
672
673 if ctx.write_senders {
674 s.spawn(|_| {
675 let _guard = span.enter();
676 r_senders = Some(self.write_segment(
677 StaticFileSegment::TransactionSenders,
678 first_block_number,
679 |w| Self::write_transaction_senders(w, blocks, tx_nums),
680 ));
681 });
682 }
683
684 if ctx.write_receipts {
685 s.spawn(|_| {
686 let _guard = span.enter();
687 r_receipts = Some(self.write_segment(
688 StaticFileSegment::Receipts,
689 first_block_number,
690 |w| Self::write_receipts(w, blocks, tx_nums, &ctx),
691 ));
692 });
693 }
694
695 if ctx.write_account_changesets {
696 s.spawn(|_| {
697 let _guard = span.enter();
698 r_account_changesets = Some(self.write_segment(
699 StaticFileSegment::AccountChangeSets,
700 first_block_number,
701 |w| Self::write_account_changesets(w, blocks),
702 ));
703 });
704 }
705
706 if ctx.write_storage_changesets {
707 s.spawn(|_| {
708 let _guard = span.enter();
709 r_storage_changesets = Some(self.write_segment(
710 StaticFileSegment::StorageChangeSets,
711 first_block_number,
712 |w| Self::write_storage_changesets(w, blocks),
713 ));
714 });
715 }
716 });
717
718 r_headers.ok_or(StaticFileWriterError::ThreadPanic("headers"))??;
719 r_txs.ok_or(StaticFileWriterError::ThreadPanic("transactions"))??;
720 if ctx.write_senders {
721 r_senders.ok_or(StaticFileWriterError::ThreadPanic("senders"))??;
722 }
723 if ctx.write_receipts {
724 r_receipts.ok_or(StaticFileWriterError::ThreadPanic("receipts"))??;
725 }
726 if ctx.write_account_changesets {
727 r_account_changesets
728 .ok_or(StaticFileWriterError::ThreadPanic("account_changesets"))??;
729 }
730 if ctx.write_storage_changesets {
731 r_storage_changesets
732 .ok_or(StaticFileWriterError::ThreadPanic("storage_changesets"))??;
733 }
734 Ok(())
735 }
736
737 pub fn get_segment_provider(
740 &self,
741 segment: StaticFileSegment,
742 number: u64,
743 ) -> ProviderResult<StaticFileJarProvider<'_, N>> {
744 if segment.is_block_or_change_based() {
745 self.get_segment_provider_for_block(segment, number, None)
746 } else {
747 self.get_segment_provider_for_transaction(segment, number, None)
748 }
749 }
750
751 pub fn get_maybe_segment_provider(
756 &self,
757 segment: StaticFileSegment,
758 number: u64,
759 ) -> ProviderResult<Option<StaticFileJarProvider<'_, N>>> {
760 let provider = if segment.is_block_or_change_based() {
761 self.get_segment_provider_for_block(segment, number, None)
762 } else {
763 self.get_segment_provider_for_transaction(segment, number, None)
764 };
765
766 match provider {
767 Ok(provider) => Ok(Some(provider)),
768 Err(
769 ProviderError::MissingStaticFileBlock(_, _) |
770 ProviderError::MissingStaticFileTx(_, _),
771 ) => Ok(None),
772 Err(err) => Err(err),
773 }
774 }
775
776 pub fn get_segment_provider_for_block(
778 &self,
779 segment: StaticFileSegment,
780 block: BlockNumber,
781 path: Option<&Path>,
782 ) -> ProviderResult<StaticFileJarProvider<'_, N>> {
783 self.get_segment_provider_for_range(
784 segment,
785 || self.get_segment_ranges_from_block(segment, block),
786 path,
787 )?
788 .ok_or(ProviderError::MissingStaticFileBlock(segment, block))
789 }
790
791 pub fn get_segment_provider_for_transaction(
793 &self,
794 segment: StaticFileSegment,
795 tx: TxNumber,
796 path: Option<&Path>,
797 ) -> ProviderResult<StaticFileJarProvider<'_, N>> {
798 self.get_segment_provider_for_range(
799 segment,
800 || self.get_segment_ranges_from_transaction(segment, tx),
801 path,
802 )?
803 .ok_or(ProviderError::MissingStaticFileTx(segment, tx))
804 }
805
806 pub fn get_segment_provider_for_range(
810 &self,
811 segment: StaticFileSegment,
812 fn_range: impl Fn() -> Option<SegmentRangeInclusive>,
813 path: Option<&Path>,
814 ) -> ProviderResult<Option<StaticFileJarProvider<'_, N>>> {
815 let block_range = match path {
818 Some(path) => StaticFileSegment::parse_filename(
819 &path
820 .file_name()
821 .ok_or_else(|| {
822 ProviderError::MissingStaticFileSegmentPath(segment, path.to_path_buf())
823 })?
824 .to_string_lossy(),
825 )
826 .and_then(|(parsed_segment, block_range)| {
827 if parsed_segment == segment {
828 return Some(block_range);
829 }
830 None
831 }),
832 None => fn_range(),
833 };
834
835 if let Some(block_range) = block_range {
837 return Ok(Some(self.get_or_create_jar_provider(segment, &block_range)?));
838 }
839
840 Ok(None)
841 }
842
843 pub fn get_segment_provider_for_path(
845 &self,
846 path: &Path,
847 ) -> ProviderResult<Option<StaticFileJarProvider<'_, N>>> {
848 StaticFileSegment::parse_filename(
849 &path
850 .file_name()
851 .ok_or_else(|| ProviderError::MissingStaticFilePath(path.to_path_buf()))?
852 .to_string_lossy(),
853 )
854 .map(|(segment, block_range)| self.get_or_create_jar_provider(segment, &block_range))
855 .transpose()
856 }
857
858 pub fn remove_cached_provider(
862 &self,
863 segment: StaticFileSegment,
864 fixed_block_range_end: BlockNumber,
865 ) {
866 self.map.remove(&(fixed_block_range_end, segment));
867 }
868
869 pub fn delete_segment_below_block(
886 &self,
887 segment: StaticFileSegment,
888 block: BlockNumber,
889 ) -> ProviderResult<Vec<SegmentHeader>> {
890 if block == 0 {
892 return Ok(Vec::new());
893 }
894
895 let highest_block = self.get_highest_static_file_block(segment);
896 let mut deleted_headers = Vec::new();
897
898 loop {
899 let Some(block_height) = self.get_lowest_range_end(segment) else {
900 return Ok(deleted_headers);
901 };
902
903 if block_height >= block || Some(block_height) == highest_block {
905 return Ok(deleted_headers);
906 }
907
908 debug!(
909 target: "providers::static_file",
910 ?segment,
911 ?block_height,
912 "Deleting static file below block"
913 );
914
915 let header = self.delete_jar(segment, block_height).inspect_err(|err| {
918 warn!( target: "providers::static_file", ?segment, %block_height, ?err, "Failed to delete static file below block")
919 })?;
920
921 deleted_headers.push(header);
922 }
923 }
924
925 pub fn delete_jar(
933 &self,
934 segment: StaticFileSegment,
935 block: BlockNumber,
936 ) -> ProviderResult<SegmentHeader> {
937 let fixed_block_range = self.find_fixed_range(segment, block);
938 let key = (fixed_block_range.end(), segment);
939 let file = self.path.join(segment.filename(&fixed_block_range));
940 let jar = if let Some((_, jar)) = self.map.remove(&key) {
941 jar.jar
942 } else {
943 debug!(
944 target: "providers::static_file",
945 ?file,
946 ?fixed_block_range,
947 ?block,
948 "Loading static file jar for deletion"
949 );
950 NippyJar::<SegmentHeader>::load(&file).map_err(ProviderError::other)?
951 };
952
953 let header = jar.user_header().clone();
954
955 if segment.is_change_based() {
957 let csoff_path = file.with_extension("csoff");
958 if csoff_path.exists() {
959 std::fs::remove_file(&csoff_path).map_err(ProviderError::other)?;
960 }
961 }
962
963 jar.delete().map_err(ProviderError::other)?;
964
965 self.initialize_index()?;
968
969 Ok(header)
970 }
971
972 pub fn delete_segment(&self, segment: StaticFileSegment) -> ProviderResult<Vec<SegmentHeader>> {
980 let mut deleted_headers = Vec::new();
981
982 self.writers.remove(segment);
983
984 while let Some(block_height) = self.get_highest_static_file_block(segment) {
985 debug!(
986 target: "providers::static_file",
987 ?segment,
988 ?block_height,
989 "Deleting static file jar"
990 );
991
992 let header = self.delete_jar(segment, block_height).inspect_err(|err| {
993 warn!(target: "providers::static_file", ?segment, %block_height, ?err, "Failed to delete static file jar")
994 })?;
995
996 deleted_headers.push(header);
997 }
998
999 Ok(deleted_headers)
1000 }
1001
1002 fn get_or_create_jar_provider(
1006 &self,
1007 segment: StaticFileSegment,
1008 fixed_block_range: &SegmentRangeInclusive,
1009 ) -> ProviderResult<StaticFileJarProvider<'_, N>> {
1010 let key = (fixed_block_range.end(), segment);
1011
1012 trace!(target: "providers::static_file", ?segment, ?fixed_block_range, "Getting provider");
1014 let mut provider: StaticFileJarProvider<'_, N> = if let Some(jar) = self.map.get(&key) {
1015 trace!(target: "providers::static_file", ?segment, ?fixed_block_range, "Jar found in cache");
1016 jar.into()
1017 } else {
1018 let generation = self.cache_generation.load(Ordering::Acquire);
1019 trace!(target: "providers::static_file", ?segment, ?fixed_block_range, generation, "Creating jar from scratch");
1020 let path = self.path.join(segment.filename(fixed_block_range));
1021 let jar = NippyJar::load(&path).map_err(ProviderError::other)?;
1022 let loaded = LoadedJar::new(jar)?;
1023 #[cfg(test)]
1024 self.before_cache_insert.fire(segment);
1025
1026 self.map
1030 .entry(key)
1031 .or_try_insert_with(|| -> ProviderResult<_> {
1032 let loaded = if self.cache_generation.load(Ordering::Acquire) == generation {
1036 loaded
1037 } else {
1038 trace!(target: "providers::static_file", ?segment, ?fixed_block_range, generation, "Reloading jar after cache invalidation");
1039 drop(loaded);
1040 LoadedJar::new(NippyJar::load(&path).map_err(ProviderError::other)?)?
1041 };
1042 #[cfg(test)]
1043 self.after_cache_validation.fire(segment);
1044 Ok(loaded)
1045 })?
1046 .downgrade()
1047 .into()
1048 };
1049
1050 if let Some(metrics) = &self.metrics {
1051 provider = provider.with_metrics(metrics.clone());
1052 }
1053 Ok(provider)
1054 }
1055
1056 fn get_segment_ranges_from_block(
1059 &self,
1060 segment: StaticFileSegment,
1061 block: u64,
1062 ) -> Option<SegmentRangeInclusive> {
1063 let indexes = self.indexes.read();
1064 let index = indexes.get(segment)?;
1065
1066 (index.max_block >= block).then(|| {
1067 self.find_fixed_range_with_block_index(
1068 segment,
1069 Some(&index.expected_block_ranges_by_max_block),
1070 block,
1071 )
1072 })
1073 }
1074
1075 fn get_segment_ranges_from_transaction(
1078 &self,
1079 segment: StaticFileSegment,
1080 tx: u64,
1081 ) -> Option<SegmentRangeInclusive> {
1082 let indexes = self.indexes.read();
1083 let index = indexes.get(segment)?;
1084 let available_block_ranges_by_max_tx = index.available_block_ranges_by_max_tx.as_ref()?;
1085
1086 let mut static_files_rev_iter = available_block_ranges_by_max_tx.iter().rev().peekable();
1089
1090 while let Some((tx_end, block_range)) = static_files_rev_iter.next() {
1091 if tx > *tx_end {
1092 return None;
1094 }
1095 let tx_start = static_files_rev_iter.peek().map(|(tx_end, _)| *tx_end + 1).unwrap_or(0);
1096 if tx_start <= tx {
1097 return Some(self.find_fixed_range_with_block_index(
1098 segment,
1099 Some(&index.expected_block_ranges_by_max_block),
1100 block_range.end(),
1101 ));
1102 }
1103 }
1104 None
1105 }
1106
1107 pub fn update_index(
1114 &self,
1115 segment: StaticFileSegment,
1116 segment_max_block: Option<BlockNumber>,
1117 ) -> ProviderResult<()> {
1118 trace!(
1119 target: "providers::static_file",
1120 ?segment,
1121 ?segment_max_block,
1122 "Updating provider index"
1123 );
1124 let mut indexes = self.indexes.write();
1125
1126 match segment_max_block {
1127 Some(segment_max_block) => {
1128 let fixed_range = self.find_fixed_range_with_block_index(
1129 segment,
1130 indexes.get(segment).map(|index| &index.expected_block_ranges_by_max_block),
1131 segment_max_block,
1132 );
1133
1134 let jar = NippyJar::<SegmentHeader>::load(
1135 &self.path.join(segment.filename(&fixed_range)),
1136 )
1137 .map_err(ProviderError::other)?;
1138
1139 let index = indexes
1140 .entry(segment)
1141 .and_modify(|index| {
1142 index.max_block = segment_max_block;
1144
1145 index
1149 .expected_block_ranges_by_max_block
1150 .retain(|_, block_range| block_range.start() < fixed_range.start());
1151 index
1153 .expected_block_ranges_by_max_block
1154 .insert(fixed_range.end(), fixed_range);
1155 })
1156 .or_insert_with(|| StaticFileSegmentIndex {
1157 min_block_range: None,
1158 max_block: segment_max_block,
1159 expected_block_ranges_by_max_block: BTreeMap::from([(
1160 fixed_range.end(),
1161 fixed_range,
1162 )]),
1163 available_block_ranges_by_max_tx: None,
1164 });
1165
1166 if let Some(current_block_range) = jar.user_header().block_range() {
1182 if let Some(min_block_range) = index.min_block_range.as_mut() {
1183 if current_block_range.start() == min_block_range.start() {
1186 *min_block_range = current_block_range;
1187 }
1188 } else {
1189 index.min_block_range = Some(current_block_range);
1190 }
1191 }
1192
1193 if let Some(tx_range) = jar.user_header().tx_range() {
1196 if let Some(current_block_range) = jar.user_header().block_range() {
1199 let tx_end = tx_range.end();
1200
1201 if let Some(index) = index.available_block_ranges_by_max_tx.as_mut() {
1210 index
1211 .retain(|_, block_range| block_range.start() < fixed_range.start());
1212 index.insert(tx_end, current_block_range);
1213 } else {
1214 index.available_block_ranges_by_max_tx =
1215 Some(BTreeMap::from([(tx_end, current_block_range)]));
1216 }
1217 }
1218 } else if segment.is_tx_based() {
1219 if let Some(index) = index.available_block_ranges_by_max_tx.as_mut() {
1223 index.retain(|_, block_range| block_range.start() < fixed_range.start());
1224 }
1225
1226 index.available_block_ranges_by_max_tx.take_if(|index| index.is_empty());
1228 }
1229
1230 trace!(target: "providers::static_file", ?segment, "Inserting updated jar into cache");
1232 self.map.insert((fixed_range.end(), segment), LoadedJar::new(jar)?);
1233
1234 trace!(target: "providers::static_file", ?segment, "Cleaning up jar map");
1236 self.map.retain(|(end, seg), _| !(*seg == segment && *end > fixed_range.end()));
1237 }
1238 None => {
1239 debug!(target: "providers::static_file", ?segment, "Removing segment from index");
1240 indexes.remove(segment);
1241 }
1242 };
1243
1244 trace!(target: "providers::static_file", ?segment, "Updated provider index");
1245 Ok(())
1246 }
1247
1248 pub fn initialize_index(&self) -> ProviderResult<()> {
1250 let mut indexes = self.indexes.write();
1251 indexes.clear();
1252
1253 for (segment, headers) in &*iter_static_files(&self.path).map_err(ProviderError::other)? {
1254 let min_block_range = Some(headers.first().expect("headers are not empty").0);
1259 let max_block = headers.last().expect("headers are not empty").0.end();
1260
1261 let mut expected_block_ranges_by_max_block = BTreeMap::default();
1262 let mut available_block_ranges_by_max_tx = None;
1263
1264 for (block_range, header) in headers {
1265 expected_block_ranges_by_max_block
1267 .insert(header.expected_block_end(), header.expected_block_range());
1268
1269 if let Some(tx_range) = header.tx_range() {
1271 let tx_end = tx_range.end();
1272
1273 available_block_ranges_by_max_tx
1274 .get_or_insert_with(BTreeMap::default)
1275 .insert(tx_end, *block_range);
1276 }
1277 }
1278
1279 indexes.insert(
1280 segment,
1281 StaticFileSegmentIndex {
1282 min_block_range,
1283 max_block,
1284 expected_block_ranges_by_max_block,
1285 available_block_ranges_by_max_tx,
1286 },
1287 );
1288 }
1289
1290 self.cache_generation.fetch_add(1, Ordering::AcqRel);
1294 self.map.clear();
1295
1296 if let Some(lowest_range) =
1298 indexes.get(StaticFileSegment::Transactions).and_then(|index| index.min_block_range)
1299 {
1300 self.earliest_history_height
1302 .store(lowest_range.start(), std::sync::atomic::Ordering::Relaxed);
1303 }
1304
1305 Ok(())
1306 }
1307
1308 #[instrument(skip(self, provider), fields(read_only = self.is_read_only()))]
1332 pub fn check_consistency<Provider>(
1333 &self,
1334 provider: &Provider,
1335 ) -> ProviderResult<Option<PipelineTarget>>
1336 where
1337 Provider: DBProvider
1338 + BlockReader
1339 + StageCheckpointReader
1340 + PruneCheckpointReader
1341 + ChainSpecProvider
1342 + StorageSettingsCache,
1343 N: NodePrimitives<Receipt: Value, BlockHeader: Value, SignedTx: Value>,
1344 {
1345 if provider.chain_spec().is_optimism() &&
1352 reth_chainspec::Chain::optimism_mainnet() == provider.chain_spec().chain_id()
1353 {
1354 const OVM_HEADER_1_HASH: B256 =
1356 b256!("0xbee7192e575af30420cae0c7776304ac196077ee72b048970549e4f08e875453");
1357 if provider.block_number(OVM_HEADER_1_HASH)?.is_some() {
1358 info!(target: "reth::cli",
1359 "Skipping storage verification for OP mainnet, expected inconsistency in OVM chain"
1360 );
1361 return Ok(None);
1362 }
1363 }
1364
1365 info!(target: "reth::cli", "Verifying storage consistency.");
1366
1367 let mut unwind_target: Option<BlockNumber> = None;
1368
1369 let mut update_unwind_target = |new_target| {
1370 unwind_target =
1371 unwind_target.map(|current| current.min(new_target)).or(Some(new_target));
1372 };
1373
1374 for segment in self.segments_to_check(provider) {
1375 let span = info_span!(
1376 "Checking consistency for segment",
1377 ?segment,
1378 initial_highest_block = tracing::field::Empty,
1379 highest_block = tracing::field::Empty,
1380 highest_tx = tracing::field::Empty,
1381 );
1382 let _guard = span.enter();
1383
1384 debug!(target: "reth::providers::static_file", "Checking consistency for segment");
1385
1386 let (initial_highest_block, mut highest_block) = self.maybe_heal_segment(segment)?;
1388 span.record("initial_highest_block", initial_highest_block);
1389 span.record("highest_block", highest_block);
1390
1391 if initial_highest_block != highest_block {
1396 info!(
1397 target: "reth::providers::static_file",
1398 unwind_target = highest_block,
1399 "Setting unwind target."
1400 );
1401 update_unwind_target(highest_block.unwrap_or_default());
1402 }
1403
1404 let highest_tx = self.get_highest_static_file_tx(segment);
1410 span.record("highest_tx", highest_tx);
1411 debug!(target: "reth::providers::static_file", "Checking tx index segment");
1412
1413 if let Some(highest_tx) = highest_tx {
1414 let mut last_block = highest_block.unwrap_or_default();
1415 debug!(target: "reth::providers::static_file", last_block, highest_tx, "Verifying last transaction matches last block indices");
1416 loop {
1417 let Some(indices) = provider.block_body_indices(last_block)? else {
1418 debug!(target: "reth::providers::static_file", last_block, "Block body indices not found, static files ahead of database");
1419 break
1423 };
1424
1425 debug!(target: "reth::providers::static_file", last_block, last_tx_num = indices.last_tx_num(), "Found block body indices");
1426
1427 if indices.last_tx_num() <= highest_tx {
1428 break
1429 }
1430
1431 if last_block == 0 {
1432 debug!(target: "reth::providers::static_file", "Reached block 0 in verification loop");
1433 break
1434 }
1435
1436 last_block -= 1;
1437
1438 info!(
1439 target: "reth::providers::static_file",
1440 highest_block = self.get_highest_static_file_block(segment),
1441 unwind_target = last_block,
1442 "Setting unwind target."
1443 );
1444 span.record("highest_block", last_block);
1445 highest_block = Some(last_block);
1446 update_unwind_target(last_block);
1447 }
1448 }
1449
1450 debug!(target: "reth::providers::static_file", "Ensuring invariants for segment");
1451
1452 match self.ensure_invariants_for(provider, segment, highest_tx, highest_block)? {
1453 Some(unwind) => {
1454 debug!(target: "reth::providers::static_file", unwind_target=unwind, "Invariants check returned unwind target");
1455 update_unwind_target(unwind);
1456 }
1457 None => {
1458 debug!(target: "reth::providers::static_file", "Invariants check completed, no unwind needed")
1459 }
1460 }
1461 }
1462
1463 Ok(unwind_target.map(PipelineTarget::Unwind))
1464 }
1465
1466 pub fn check_file_consistency<Provider>(&self, provider: &Provider) -> ProviderResult<()>
1475 where
1476 Provider: DBProvider + ChainSpecProvider + StorageSettingsCache + PruneCheckpointReader,
1477 {
1478 info!(target: "reth::cli", "Healing static file inconsistencies.");
1479
1480 for segment in self.segments_to_check(provider) {
1481 let _guard = info_span!("healing_static_file_segment", ?segment).entered();
1482 let _ = self.maybe_heal_segment(segment)?;
1483 }
1484
1485 Ok(())
1486 }
1487
1488 fn segments_to_check<'a, Provider>(
1490 &'a self,
1491 provider: &'a Provider,
1492 ) -> impl Iterator<Item = StaticFileSegment> + 'a
1493 where
1494 Provider: DBProvider + ChainSpecProvider + StorageSettingsCache + PruneCheckpointReader,
1495 {
1496 StaticFileSegment::iter()
1497 .filter(move |segment| self.should_check_segment(provider, *segment))
1498 }
1499
1500 fn should_check_segment<Provider>(
1502 &self,
1503 provider: &Provider,
1504 segment: StaticFileSegment,
1505 ) -> bool
1506 where
1507 Provider: DBProvider + ChainSpecProvider + StorageSettingsCache + PruneCheckpointReader,
1508 {
1509 match segment {
1510 StaticFileSegment::Headers | StaticFileSegment::Transactions => true,
1511 StaticFileSegment::Receipts => {
1512 if EitherWriter::receipts_destination(provider).is_database() {
1513 debug!(target: "reth::providers::static_file", ?segment, "Skipping receipts segment: receipts stored in database");
1516 return false;
1517 }
1518
1519 if NamedChain::Gnosis == provider.chain_spec().chain_id() ||
1520 NamedChain::Chiado == provider.chain_spec().chain_id()
1521 {
1522 debug!(target: "reth::providers::static_file", ?segment, "Skipping receipts segment: broken historical import for gnosis/chiado");
1526 return false;
1527 }
1528
1529 true
1530 }
1531 StaticFileSegment::TransactionSenders => {
1532 if EitherWriterDestination::senders(provider).is_database() {
1533 debug!(target: "reth::providers::static_file", ?segment, "Skipping senders segment: senders stored in database");
1534 return false;
1535 }
1536
1537 if Self::is_segment_fully_pruned(provider, PruneSegment::SenderRecovery) {
1538 debug!(target: "reth::providers::static_file", ?segment, "Skipping senders segment: fully pruned");
1539 return false;
1540 }
1541
1542 true
1543 }
1544 StaticFileSegment::AccountChangeSets => {
1545 if EitherWriter::account_changesets_destination(provider).is_database() {
1546 debug!(target: "reth::providers::static_file", ?segment, "Skipping account changesets segment: changesets stored in database");
1547 return false;
1548 }
1549 true
1550 }
1551 StaticFileSegment::StorageChangeSets => {
1552 if EitherWriter::storage_changesets_destination(provider).is_database() {
1553 debug!(target: "reth::providers::static_file", ?segment, "Skipping storage changesets segment: changesets stored in database");
1554 return false
1555 }
1556 true
1557 }
1558 }
1559 }
1560
1561 fn is_segment_fully_pruned<Provider>(provider: &Provider, segment: PruneSegment) -> bool
1565 where
1566 Provider: PruneCheckpointReader,
1567 {
1568 provider
1569 .get_prune_checkpoint(segment)
1570 .ok()
1571 .flatten()
1572 .is_some_and(|checkpoint| checkpoint.prune_mode.is_full())
1573 }
1574
1575 fn check_segment_consistency(&self, segment: StaticFileSegment) -> ProviderResult<()> {
1580 debug!(target: "reth::providers::static_file", "Checking segment consistency");
1581 if let Some(latest_block) = self.get_highest_static_file_block(segment) {
1582 let file_path = self
1583 .directory()
1584 .join(segment.filename(&self.find_fixed_range(segment, latest_block)));
1585 debug!(target: "reth::providers::static_file", ?file_path, latest_block, "Loading NippyJar for consistency check");
1586
1587 let jar = NippyJar::<SegmentHeader>::load(&file_path).map_err(ProviderError::other)?;
1588 debug!(target: "reth::providers::static_file", "NippyJar loaded, checking consistency");
1589
1590 NippyJarChecker::new(jar).check_consistency().map_err(ProviderError::other)?;
1591 debug!(target: "reth::providers::static_file", "NippyJar consistency check passed");
1592 } else {
1593 debug!(target: "reth::providers::static_file", "No static file block found, skipping consistency check");
1594 }
1595 Ok(())
1596 }
1597
1598 fn maybe_heal_segment(
1614 &self,
1615 segment: StaticFileSegment,
1616 ) -> ProviderResult<(Option<BlockNumber>, Option<BlockNumber>)> {
1617 let initial_highest_block = self.get_highest_static_file_block(segment);
1618 debug!(target: "reth::providers::static_file", ?initial_highest_block, "Initial highest block for segment");
1619
1620 if self.access.is_read_only() {
1621 debug!(target: "reth::providers::static_file", "Checking segment consistency (read-only)");
1624 self.check_segment_consistency(segment)?;
1625 } else {
1626 debug!(target: "reth::providers::static_file", "Fetching latest writer which might heal any potential inconsistency");
1629 self.latest_writer(segment)?;
1630 }
1631
1632 let highest_block = self.get_highest_static_file_block(segment);
1635
1636 Ok((initial_highest_block, highest_block))
1637 }
1638
1639 fn ensure_invariants_for<Provider>(
1641 &self,
1642 provider: &Provider,
1643 segment: StaticFileSegment,
1644 highest_tx: Option<u64>,
1645 highest_block: Option<BlockNumber>,
1646 ) -> ProviderResult<Option<BlockNumber>>
1647 where
1648 Provider: DBProvider + BlockReader + StageCheckpointReader + PruneCheckpointReader,
1649 N: NodePrimitives<Receipt: Value, BlockHeader: Value, SignedTx: Value>,
1650 {
1651 match segment {
1652 StaticFileSegment::Headers => self
1653 .ensure_invariants::<_, tables::Headers<N::BlockHeader>>(
1654 provider,
1655 segment,
1656 highest_block,
1657 highest_block,
1658 ),
1659 StaticFileSegment::Transactions => self
1660 .ensure_invariants::<_, tables::Transactions<N::SignedTx>>(
1661 provider,
1662 segment,
1663 highest_tx,
1664 highest_block,
1665 ),
1666 StaticFileSegment::Receipts => self
1667 .ensure_invariants::<_, tables::Receipts<N::Receipt>>(
1668 provider,
1669 segment,
1670 highest_tx,
1671 highest_block,
1672 ),
1673 StaticFileSegment::TransactionSenders => self
1674 .ensure_invariants::<_, tables::TransactionSenders>(
1675 provider,
1676 segment,
1677 highest_tx,
1678 highest_block,
1679 ),
1680 StaticFileSegment::AccountChangeSets => self
1681 .ensure_invariants::<_, tables::AccountChangeSets>(
1682 provider,
1683 segment,
1684 highest_tx,
1685 highest_block,
1686 ),
1687 StaticFileSegment::StorageChangeSets => self
1688 .ensure_changeset_invariants_by_block::<_, tables::StorageChangeSets, _>(
1689 provider,
1690 segment,
1691 highest_block,
1692 |key| key.block_number(),
1693 ),
1694 }
1695 }
1696
1697 #[instrument(skip(self, provider, segment), fields(table = T::NAME))]
1712 fn ensure_invariants<Provider, T: Table<Key = u64>>(
1713 &self,
1714 provider: &Provider,
1715 segment: StaticFileSegment,
1716 highest_static_file_entry: Option<u64>,
1717 highest_static_file_block: Option<BlockNumber>,
1718 ) -> ProviderResult<Option<BlockNumber>>
1719 where
1720 Provider: DBProvider + BlockReader + StageCheckpointReader + PruneCheckpointReader,
1721 {
1722 debug!(target: "reth::providers::static_file", "Ensuring invariants");
1723 let mut db_cursor = provider.tx_ref().cursor_read::<T>()?;
1724
1725 if let Some((db_first_entry, _)) = db_cursor.first()? {
1726 debug!(target: "reth::providers::static_file", db_first_entry, "Found first database entry");
1727 if let (Some(highest_entry), Some(highest_block)) =
1728 (highest_static_file_entry, highest_static_file_block)
1729 {
1730 if !(db_first_entry <= highest_entry || highest_entry + 1 == db_first_entry) {
1734 info!(
1735 target: "reth::providers::static_file",
1736 ?db_first_entry,
1737 ?highest_entry,
1738 unwind_target = highest_block,
1739 "Setting unwind target."
1740 );
1741 return Ok(Some(highest_block));
1742 }
1743 }
1744
1745 if let Some((db_last_entry, _)) = db_cursor.last()? &&
1746 highest_static_file_entry
1747 .is_none_or(|highest_entry| db_last_entry > highest_entry)
1748 {
1749 debug!(target: "reth::providers::static_file", db_last_entry, "Database has entries beyond static files, no unwind needed");
1750 return Ok(None)
1751 }
1752 } else {
1753 debug!(target: "reth::providers::static_file", "No database entries found");
1754 }
1755
1756 let highest_static_file_entry = highest_static_file_entry.unwrap_or_default();
1757 let highest_static_file_block = highest_static_file_block.unwrap_or_default();
1758
1759 let stage_id = segment.to_stage_id();
1762 let checkpoint_block_number =
1763 provider.get_stage_checkpoint(stage_id)?.unwrap_or_default().block_number;
1764 debug!(target: "reth::providers::static_file", ?stage_id, checkpoint_block_number, "Retrieved stage checkpoint");
1765
1766 let effective_coverage_block =
1767 Self::effective_coverage_block(provider, segment, highest_static_file_block)?;
1768
1769 if checkpoint_block_number > effective_coverage_block {
1771 info!(
1772 target: "reth::providers::static_file",
1773 checkpoint_block_number,
1774 unwind_target = effective_coverage_block,
1775 "Setting unwind target."
1776 );
1777 return Ok(Some(effective_coverage_block));
1778 }
1779
1780 if checkpoint_block_number >= highest_static_file_block {
1782 debug!(target: "reth::providers::static_file", "Invariants ensured, returning None");
1783 return Ok(None);
1784 }
1785
1786 info!(
1792 target: "reth::providers",
1793 from = highest_static_file_block,
1794 to = checkpoint_block_number,
1795 "Unwinding static file segment."
1796 );
1797 let mut writer = self.latest_writer(segment)?;
1798
1799 match segment {
1800 StaticFileSegment::Headers => {
1801 let prune_count = highest_static_file_block - checkpoint_block_number;
1802 debug!(target: "reth::providers::static_file", prune_count, "Pruning headers");
1803 writer.prune_headers(prune_count)?;
1805 }
1806 StaticFileSegment::Transactions |
1807 StaticFileSegment::Receipts |
1808 StaticFileSegment::TransactionSenders => {
1809 if let Some(block) = provider.block_body_indices(checkpoint_block_number)? {
1810 let number = highest_static_file_entry
1813 .saturating_add(1)
1814 .saturating_sub(block.next_tx_num());
1815 debug!(target: "reth::providers::static_file", prune_count = number, checkpoint_block_number, "Pruning transaction based segment");
1816
1817 match segment {
1818 StaticFileSegment::Transactions => {
1819 writer.prune_transactions(number, checkpoint_block_number)?
1820 }
1821 StaticFileSegment::Receipts => {
1822 writer.prune_receipts(number, checkpoint_block_number)?
1823 }
1824 StaticFileSegment::TransactionSenders => {
1825 writer.prune_transaction_senders(number, checkpoint_block_number)?
1826 }
1827 StaticFileSegment::Headers |
1828 StaticFileSegment::AccountChangeSets |
1829 StaticFileSegment::StorageChangeSets => {
1830 unreachable!()
1831 }
1832 }
1833 } else {
1834 debug!(target: "reth::providers::static_file", checkpoint_block_number, "No block body indices found for checkpoint block");
1835 }
1836 }
1837 StaticFileSegment::AccountChangeSets => {
1838 writer.prune_account_changesets(checkpoint_block_number)?;
1839 }
1840 StaticFileSegment::StorageChangeSets => {
1841 writer.prune_storage_changesets(checkpoint_block_number)?;
1842 }
1843 }
1844
1845 debug!(target: "reth::providers::static_file", "Committing writer after pruning");
1846 writer.commit()?;
1847 debug!(target: "reth::providers::static_file", "Writer committed successfully");
1848
1849 debug!(target: "reth::providers::static_file", "Invariants ensured, returning None");
1850 Ok(None)
1851 }
1852
1853 fn ensure_changeset_invariants_by_block<Provider, T, F>(
1854 &self,
1855 provider: &Provider,
1856 segment: StaticFileSegment,
1857 highest_static_file_block: Option<BlockNumber>,
1858 block_from_key: F,
1859 ) -> ProviderResult<Option<BlockNumber>>
1860 where
1861 Provider: DBProvider + BlockReader + StageCheckpointReader + PruneCheckpointReader,
1862 T: Table,
1863 F: Fn(&T::Key) -> BlockNumber,
1864 {
1865 debug!(
1866 target: "reth::providers::static_file",
1867 ?segment,
1868 ?highest_static_file_block,
1869 "Ensuring changeset invariants"
1870 );
1871 let mut db_cursor = provider.tx_ref().cursor_read::<T>()?;
1872
1873 if let Some((db_first_key, _)) = db_cursor.first()? {
1874 let db_first_block = block_from_key(&db_first_key);
1875 if let Some(highest_block) = highest_static_file_block &&
1876 !(db_first_block <= highest_block || highest_block + 1 == db_first_block)
1877 {
1878 info!(
1879 target: "reth::providers::static_file",
1880 ?db_first_block,
1881 ?highest_block,
1882 unwind_target = highest_block,
1883 ?segment,
1884 "Setting unwind target."
1885 );
1886 return Ok(Some(highest_block))
1887 }
1888
1889 if let Some((db_last_key, _)) = db_cursor.last()? &&
1890 highest_static_file_block
1891 .is_none_or(|highest_block| block_from_key(&db_last_key) > highest_block)
1892 {
1893 debug!(
1894 target: "reth::providers::static_file",
1895 ?segment,
1896 "Database has entries beyond static files, no unwind needed"
1897 );
1898 return Ok(None)
1899 }
1900 } else {
1901 debug!(target: "reth::providers::static_file", ?segment, "No database entries found");
1902 }
1903
1904 let highest_static_file_block = highest_static_file_block.unwrap_or_default();
1905
1906 let stage_id = segment.to_stage_id();
1907 let checkpoint_block_number =
1908 provider.get_stage_checkpoint(stage_id)?.unwrap_or_default().block_number;
1909
1910 let effective_coverage_block =
1911 Self::effective_coverage_block(provider, segment, highest_static_file_block)?;
1912
1913 if checkpoint_block_number > effective_coverage_block {
1914 info!(
1915 target: "reth::providers::static_file",
1916 checkpoint_block_number,
1917 unwind_target = effective_coverage_block,
1918 ?segment,
1919 "Setting unwind target."
1920 );
1921 return Ok(Some(effective_coverage_block))
1922 }
1923
1924 if checkpoint_block_number < highest_static_file_block {
1925 info!(
1926 target: "reth::providers",
1927 ?segment,
1928 from = highest_static_file_block,
1929 to = checkpoint_block_number,
1930 "Unwinding static file segment."
1931 );
1932 let mut writer = self.latest_writer(segment)?;
1933 match segment {
1934 StaticFileSegment::AccountChangeSets => {
1935 writer.prune_account_changesets(checkpoint_block_number)?;
1936 }
1937 StaticFileSegment::StorageChangeSets => {
1938 writer.prune_storage_changesets(checkpoint_block_number)?;
1939 }
1940 _ => unreachable!("invalid segment for changeset invariants"),
1941 }
1942 writer.commit()?;
1943 }
1944
1945 Ok(None)
1946 }
1947
1948 fn effective_coverage_block<Provider>(
1957 provider: &Provider,
1958 segment: StaticFileSegment,
1959 highest_static_file_block: BlockNumber,
1960 ) -> ProviderResult<BlockNumber>
1961 where
1962 Provider: PruneCheckpointReader,
1963 {
1964 let Some(prune_segment) = Self::prune_segment_for_static_file(segment) else {
1965 return Ok(highest_static_file_block)
1966 };
1967
1968 let prune_checkpoint_block = provider
1969 .get_prune_checkpoint(prune_segment)?
1970 .and_then(|checkpoint| checkpoint.block_number)
1971 .unwrap_or_default();
1972
1973 Ok(highest_static_file_block.max(prune_checkpoint_block))
1974 }
1975
1976 const fn prune_segment_for_static_file(segment: StaticFileSegment) -> Option<PruneSegment> {
1979 match segment {
1980 StaticFileSegment::Receipts => Some(PruneSegment::Receipts),
1981 StaticFileSegment::TransactionSenders => Some(PruneSegment::SenderRecovery),
1982 StaticFileSegment::AccountChangeSets => Some(PruneSegment::AccountHistory),
1983 StaticFileSegment::StorageChangeSets => Some(PruneSegment::StorageHistory),
1984 StaticFileSegment::Headers | StaticFileSegment::Transactions => None,
1985 }
1986 }
1987
1988 pub fn earliest_history_height(&self) -> BlockNumber {
1996 self.earliest_history_height.load(std::sync::atomic::Ordering::Relaxed)
1997 }
1998
1999 pub fn get_lowest_range(&self, segment: StaticFileSegment) -> Option<SegmentRangeInclusive> {
2003 self.indexes.read().get(segment).and_then(|index| index.min_block_range)
2004 }
2005
2006 pub fn get_lowest_range_start(&self, segment: StaticFileSegment) -> Option<BlockNumber> {
2012 self.get_lowest_range(segment).map(|range| range.start())
2013 }
2014
2015 pub fn get_lowest_range_end(&self, segment: StaticFileSegment) -> Option<BlockNumber> {
2021 self.get_lowest_range(segment).map(|range| range.end())
2022 }
2023
2024 pub fn get_highest_static_file_block(&self, segment: StaticFileSegment) -> Option<BlockNumber> {
2028 self.indexes.read().get(segment).map(|index| index.max_block)
2029 }
2030
2031 fn bound_range(
2038 &self,
2039 range: impl RangeBounds<BlockNumber>,
2040 segment: StaticFileSegment,
2041 ) -> RangeInclusive<BlockNumber> {
2042 let highest_block = self.get_highest_static_file_block(segment).unwrap_or(0);
2043
2044 let start = match range.start_bound() {
2045 Bound::Included(&n) => n,
2046 Bound::Excluded(&n) => n.saturating_add(1),
2047 Bound::Unbounded => 0,
2048 };
2049 let end = match range.end_bound() {
2050 Bound::Included(&n) => n.min(highest_block),
2051 Bound::Excluded(&n) => n.saturating_sub(1).min(highest_block),
2052 Bound::Unbounded => highest_block,
2053 };
2054
2055 start..=end
2056 }
2057
2058 pub fn get_highest_static_file_tx(&self, segment: StaticFileSegment) -> Option<TxNumber> {
2062 self.indexes
2063 .read()
2064 .get(segment)
2065 .and_then(|index| index.available_block_ranges_by_max_tx.as_ref())
2066 .and_then(|index| index.last_key_value().map(|(last_tx, _)| *last_tx))
2067 }
2068
2069 pub fn get_highest_static_files(&self) -> HighestStaticFiles {
2071 HighestStaticFiles {
2072 receipts: self.get_highest_static_file_block(StaticFileSegment::Receipts),
2073 }
2074 }
2075
2076 pub fn find_static_file<T>(
2079 &self,
2080 segment: StaticFileSegment,
2081 func: impl Fn(StaticFileJarProvider<'_, N>) -> ProviderResult<Option<T>>,
2082 ) -> ProviderResult<Option<T>> {
2083 if let Some(ranges) =
2084 self.indexes.read().get(segment).map(|index| &index.expected_block_ranges_by_max_block)
2085 {
2086 for range in ranges.values().rev() {
2088 if let Some(res) = func(self.get_or_create_jar_provider(segment, range)?)? {
2089 return Ok(Some(res));
2090 }
2091 }
2092 }
2093
2094 Ok(None)
2095 }
2096
2097 pub fn fetch_range_with_predicate<T, F, P>(
2103 &self,
2104 segment: StaticFileSegment,
2105 range: Range<u64>,
2106 mut get_fn: F,
2107 mut predicate: P,
2108 ) -> ProviderResult<Vec<T>>
2109 where
2110 F: FnMut(&mut StaticFileCursor<'_>, u64) -> ProviderResult<Option<T>>,
2111 P: FnMut(&T) -> bool,
2112 {
2113 let mut result = Vec::with_capacity((range.end - range.start).min(100) as usize);
2114
2115 macro_rules! get_provider {
2119 ($number:expr) => {{
2120 match self.get_segment_provider(segment, $number) {
2121 Ok(provider) => provider,
2122 Err(
2123 ProviderError::MissingStaticFileBlock(_, _) |
2124 ProviderError::MissingStaticFileTx(_, _),
2125 ) => return Ok(result),
2126 Err(err) => return Err(err),
2127 }
2128 }};
2129 }
2130
2131 let mut provider = get_provider!(range.start);
2132 let mut cursor = provider.cursor()?;
2133
2134 'outer: for number in range {
2136 let mut retrying = false;
2140
2141 'inner: loop {
2143 match get_fn(&mut cursor, number)? {
2144 Some(res) => {
2145 if !predicate(&res) {
2146 break 'outer;
2147 }
2148 result.push(res);
2149 break 'inner;
2150 }
2151 None => {
2152 if retrying {
2153 return Ok(result);
2154 }
2155 drop(cursor);
2160 drop(provider);
2161 provider = get_provider!(number);
2162 cursor = provider.cursor()?;
2163 retrying = true;
2164 }
2165 }
2166 }
2167 }
2168
2169 result.shrink_to_fit();
2170
2171 Ok(result)
2172 }
2173
2174 pub fn fetch_range_iter<'a, T, F>(
2179 &'a self,
2180 segment: StaticFileSegment,
2181 range: Range<u64>,
2182 get_fn: F,
2183 ) -> ProviderResult<impl Iterator<Item = ProviderResult<Option<T>>> + 'a>
2184 where
2185 F: Fn(&mut StaticFileCursor<'_>, u64) -> ProviderResult<Option<T>> + 'a,
2186 T: std::fmt::Debug,
2187 {
2188 let mut provider = self.get_maybe_segment_provider(segment, range.start)?;
2189 Ok(range.map(move |number| {
2190 match provider
2191 .as_ref()
2192 .map(|provider| get_fn(&mut provider.cursor()?, number))
2193 .and_then(|result| result.transpose())
2194 {
2195 Some(result) => result.map(Some),
2196 None => {
2197 provider.take();
2201 provider = self.get_maybe_segment_provider(segment, number)?;
2202 provider
2203 .as_ref()
2204 .map(|provider| get_fn(&mut provider.cursor()?, number))
2205 .and_then(|result| result.transpose())
2206 .transpose()
2207 }
2208 }
2209 }))
2210 }
2211
2212 pub fn directory(&self) -> &Path {
2214 &self.path
2215 }
2216
2217 pub fn get_with_static_file_or_database<T, FS, FD>(
2227 &self,
2228 segment: StaticFileSegment,
2229 number: u64,
2230 fetch_from_static_file: FS,
2231 fetch_from_database: FD,
2232 ) -> ProviderResult<Option<T>>
2233 where
2234 FS: Fn(&Self) -> ProviderResult<Option<T>>,
2235 FD: Fn() -> ProviderResult<Option<T>>,
2236 {
2237 let static_file_upper_bound = if segment.is_block_or_change_based() {
2239 self.get_highest_static_file_block(segment)
2240 } else {
2241 self.get_highest_static_file_tx(segment)
2242 };
2243
2244 if static_file_upper_bound
2245 .is_some_and(|static_file_upper_bound| static_file_upper_bound >= number)
2246 {
2247 return fetch_from_static_file(self);
2248 }
2249 fetch_from_database()
2250 }
2251
2252 pub fn get_range_with_static_file_or_database<T, P, FS, FD>(
2264 &self,
2265 segment: StaticFileSegment,
2266 mut block_or_tx_range: Range<u64>,
2267 fetch_from_static_file: FS,
2268 mut fetch_from_database: FD,
2269 mut predicate: P,
2270 ) -> ProviderResult<Vec<T>>
2271 where
2272 FS: Fn(&Self, Range<u64>, &mut P) -> ProviderResult<Vec<T>>,
2273 FD: FnMut(Range<u64>, P) -> ProviderResult<Vec<T>>,
2274 P: FnMut(&T) -> bool,
2275 {
2276 let mut data = Vec::new();
2277
2278 if let Some(static_file_upper_bound) = if segment.is_block_or_change_based() {
2280 self.get_highest_static_file_block(segment)
2281 } else {
2282 self.get_highest_static_file_tx(segment)
2283 } && block_or_tx_range.start <= static_file_upper_bound
2284 {
2285 let end = block_or_tx_range.end.min(static_file_upper_bound + 1);
2286 data.extend(fetch_from_static_file(
2287 self,
2288 block_or_tx_range.start..end,
2289 &mut predicate,
2290 )?);
2291 block_or_tx_range.start = end;
2292 }
2293
2294 if block_or_tx_range.end > block_or_tx_range.start {
2295 data.extend(fetch_from_database(block_or_tx_range, predicate)?)
2296 }
2297
2298 Ok(data)
2299 }
2300
2301 #[cfg(any(test, feature = "test-utils"))]
2303 pub fn path(&self) -> &Path {
2304 &self.path
2305 }
2306
2307 #[cfg(any(test, feature = "test-utils"))]
2309 pub fn tx_index(&self, segment: StaticFileSegment) -> Option<SegmentRanges> {
2310 self.indexes
2311 .read()
2312 .get(segment)
2313 .and_then(|index| index.available_block_ranges_by_max_tx.as_ref())
2314 .cloned()
2315 }
2316
2317 #[cfg(any(test, feature = "test-utils"))]
2319 pub fn expected_block_index(&self, segment: StaticFileSegment) -> Option<SegmentRanges> {
2320 self.indexes
2321 .read()
2322 .get(segment)
2323 .map(|index| &index.expected_block_ranges_by_max_block)
2324 .cloned()
2325 }
2326}
2327
2328#[derive(Debug)]
2329struct StaticFileSegmentIndex {
2330 min_block_range: Option<SegmentRangeInclusive>,
2342 max_block: u64,
2344 expected_block_ranges_by_max_block: SegmentRanges,
2350 available_block_ranges_by_max_tx: Option<SegmentRanges>,
2357}
2358
2359pub trait StaticFileWriter {
2361 type Primitives: Send + Sync + 'static;
2363
2364 fn get_writer(
2366 &self,
2367 block: BlockNumber,
2368 segment: StaticFileSegment,
2369 ) -> ProviderResult<StaticFileProviderRWRefMut<'_, Self::Primitives>>;
2370
2371 fn latest_writer(
2374 &self,
2375 segment: StaticFileSegment,
2376 ) -> ProviderResult<StaticFileProviderRWRefMut<'_, Self::Primitives>>;
2377
2378 fn commit(&self) -> ProviderResult<()>;
2380
2381 fn has_unwind_queued(&self) -> bool;
2383
2384 fn finalize(&self) -> ProviderResult<()>;
2388}
2389
2390impl<N: NodePrimitives> StaticFileWriter for StaticFileProvider<N> {
2391 type Primitives = N;
2392
2393 fn get_writer(
2394 &self,
2395 block: BlockNumber,
2396 segment: StaticFileSegment,
2397 ) -> ProviderResult<StaticFileProviderRWRefMut<'_, Self::Primitives>> {
2398 if self.access.is_read_only() {
2399 return Err(ProviderError::ReadOnlyStaticFileAccess);
2400 }
2401
2402 trace!(target: "providers::static_file", ?block, ?segment, "Getting static file writer.");
2403 self.writers.get_or_create(segment, || {
2404 StaticFileProviderRW::new(segment, block, Arc::downgrade(&self.0), self.metrics.clone())
2405 })
2406 }
2407
2408 fn latest_writer(
2409 &self,
2410 segment: StaticFileSegment,
2411 ) -> ProviderResult<StaticFileProviderRWRefMut<'_, Self::Primitives>> {
2412 let genesis_number = self.0.as_ref().genesis_block_number();
2413 self.get_writer(
2414 self.get_highest_static_file_block(segment).unwrap_or(genesis_number),
2415 segment,
2416 )
2417 }
2418
2419 fn commit(&self) -> ProviderResult<()> {
2420 self.writers.commit()
2421 }
2422
2423 fn has_unwind_queued(&self) -> bool {
2424 self.writers.has_unwind_queued()
2425 }
2426
2427 fn finalize(&self) -> ProviderResult<()> {
2428 self.writers.finalize()
2429 }
2430}
2431
2432impl<N: NodePrimitives> ChangeSetReader for StaticFileProvider<N> {
2433 fn account_block_changeset(
2434 &self,
2435 block_number: BlockNumber,
2436 ) -> ProviderResult<Vec<reth_db::models::AccountBeforeTx>> {
2437 let provider = match self.get_segment_provider_for_block(
2438 StaticFileSegment::AccountChangeSets,
2439 block_number,
2440 None,
2441 ) {
2442 Ok(provider) => provider,
2443 Err(ProviderError::MissingStaticFileBlock(_, _)) => return Ok(Vec::new()),
2444 Err(err) => return Err(err),
2445 };
2446
2447 if let Some(offset) = provider.read_changeset_offset(block_number)? {
2448 let mut cursor = provider.cursor()?;
2449 let mut changeset = Vec::with_capacity(offset.num_changes() as usize);
2450
2451 for i in offset.changeset_range() {
2452 if let Some(change) =
2453 cursor.get_one::<reth_db::static_file::AccountChangesetMask>(i.into())?
2454 {
2455 changeset.push(change)
2456 }
2457 }
2458 Ok(changeset)
2459 } else {
2460 Ok(Vec::new())
2461 }
2462 }
2463
2464 fn get_account_before_block(
2465 &self,
2466 block_number: BlockNumber,
2467 address: Address,
2468 ) -> ProviderResult<Option<reth_db::models::AccountBeforeTx>> {
2469 let provider = match self.get_segment_provider_for_block(
2470 StaticFileSegment::AccountChangeSets,
2471 block_number,
2472 None,
2473 ) {
2474 Ok(provider) => provider,
2475 Err(ProviderError::MissingStaticFileBlock(_, _)) => return Ok(None),
2476 Err(err) => return Err(err),
2477 };
2478
2479 let Some(offset) = provider.read_changeset_offset(block_number)? else {
2480 return Ok(None);
2481 };
2482
2483 let mut cursor = provider.cursor()?;
2484 let range = offset.changeset_range();
2485 let mut low = range.start;
2486 let mut high = range.end;
2487
2488 while low < high {
2489 let mid = low + (high - low) / 2;
2490 if let Some(change) =
2491 cursor.get_one::<reth_db::static_file::AccountChangesetMask>(mid.into())?
2492 {
2493 if change.address < address {
2494 low = mid + 1;
2495 } else {
2496 high = mid;
2497 }
2498 } else {
2499 debug!(
2502 target: "providers::static_file",
2503 ?low,
2504 ?mid,
2505 ?high,
2506 ?range,
2507 ?block_number,
2508 ?address,
2509 "Cannot continue binary search for account changeset fetch"
2510 );
2511 low = range.end;
2512 break;
2513 }
2514 }
2515
2516 if low < range.end &&
2517 let Some(change) = cursor
2518 .get_one::<reth_db::static_file::AccountChangesetMask>(low.into())?
2519 .filter(|change| change.address == address)
2520 {
2521 return Ok(Some(change));
2522 }
2523
2524 Ok(None)
2525 }
2526
2527 fn account_changesets_range(
2528 &self,
2529 range: impl core::ops::RangeBounds<BlockNumber>,
2530 ) -> ProviderResult<Vec<(BlockNumber, reth_db::models::AccountBeforeTx)>> {
2531 let range = self.bound_range(range, StaticFileSegment::AccountChangeSets);
2532 self.walk_account_changeset_range(range).collect()
2533 }
2534}
2535
2536impl<N: NodePrimitives> StorageChangeSetReader for StaticFileProvider<N> {
2537 fn storage_changeset(
2538 &self,
2539 block_number: BlockNumber,
2540 ) -> ProviderResult<Vec<(BlockNumberAddress, StorageEntry)>> {
2541 let provider = match self.get_segment_provider_for_block(
2542 StaticFileSegment::StorageChangeSets,
2543 block_number,
2544 None,
2545 ) {
2546 Ok(provider) => provider,
2547 Err(ProviderError::MissingStaticFileBlock(_, _)) => return Ok(Vec::new()),
2548 Err(err) => return Err(err),
2549 };
2550
2551 if let Some(offset) = provider.read_changeset_offset(block_number)? {
2552 let mut cursor = provider.cursor()?;
2553 let mut changeset = Vec::with_capacity(offset.num_changes() as usize);
2554
2555 for i in offset.changeset_range() {
2556 if let Some(change) = cursor.get_one::<StorageChangesetMask>(i.into())? {
2557 let block_address = BlockNumberAddress((block_number, change.address));
2558 let entry = StorageEntry { key: change.key, value: change.value };
2559 changeset.push((block_address, entry));
2560 }
2561 }
2562 Ok(changeset)
2563 } else {
2564 Ok(Vec::new())
2565 }
2566 }
2567
2568 fn get_storage_before_block(
2569 &self,
2570 block_number: BlockNumber,
2571 address: Address,
2572 storage_key: B256,
2573 ) -> ProviderResult<Option<StorageEntry>> {
2574 let provider = match self.get_segment_provider_for_block(
2575 StaticFileSegment::StorageChangeSets,
2576 block_number,
2577 None,
2578 ) {
2579 Ok(provider) => provider,
2580 Err(ProviderError::MissingStaticFileBlock(_, _)) => return Ok(None),
2581 Err(err) => return Err(err),
2582 };
2583
2584 let Some(offset) = provider.read_changeset_offset(block_number)? else {
2585 return Ok(None);
2586 };
2587
2588 let mut cursor = provider.cursor()?;
2589 let range = offset.changeset_range();
2590 let mut low = range.start;
2591 let mut high = range.end;
2592
2593 while low < high {
2594 let mid = low + (high - low) / 2;
2595 if let Some(change) = cursor.get_one::<StorageChangesetMask>(mid.into())? {
2596 match (change.address, change.key).cmp(&(address, storage_key)) {
2597 std::cmp::Ordering::Less => low = mid + 1,
2598 _ => high = mid,
2599 }
2600 } else {
2601 debug!(
2602 target: "providers::static_file",
2603 ?low,
2604 ?mid,
2605 ?high,
2606 ?range,
2607 ?block_number,
2608 ?address,
2609 ?storage_key,
2610 "Cannot continue binary search for storage changeset fetch"
2611 );
2612 low = range.end;
2613 break;
2614 }
2615 }
2616
2617 if low < range.end &&
2618 let Some(change) = cursor
2619 .get_one::<StorageChangesetMask>(low.into())?
2620 .filter(|change| change.address == address && change.key == storage_key)
2621 {
2622 return Ok(Some(StorageEntry { key: change.key, value: change.value }));
2623 }
2624
2625 Ok(None)
2626 }
2627
2628 fn storage_changesets_range(
2629 &self,
2630 range: impl RangeBounds<BlockNumber>,
2631 ) -> ProviderResult<Vec<(BlockNumberAddress, StorageEntry)>> {
2632 let range = self.bound_range(range, StaticFileSegment::StorageChangeSets);
2633 self.walk_storage_changeset_range(range).collect()
2634 }
2635}
2636
2637impl<N: NodePrimitives> StaticFileProvider<N> {
2638 pub fn walk_account_changeset_range(
2648 &self,
2649 range: impl RangeBounds<BlockNumber>,
2650 ) -> StaticFileAccountChangesetWalker<Self> {
2651 StaticFileAccountChangesetWalker::new(self.clone(), range)
2652 }
2653
2654 pub fn walk_storage_changeset_range(
2656 &self,
2657 range: impl RangeBounds<BlockNumber>,
2658 ) -> StaticFileStorageChangesetWalker<Self> {
2659 StaticFileStorageChangesetWalker::new(self.clone(), range)
2660 }
2661}
2662
2663impl<N: NodePrimitives<BlockHeader: Value>> HeaderProvider for StaticFileProvider<N> {
2664 type Header = N::BlockHeader;
2665
2666 fn header(&self, block_hash: BlockHash) -> ProviderResult<Option<Self::Header>> {
2667 self.find_static_file(StaticFileSegment::Headers, |jar_provider| {
2668 Ok(jar_provider
2669 .cursor()?
2670 .get_two::<HeaderWithHashMask<Self::Header>>((&block_hash).into())?
2671 .and_then(|(header, hash)| {
2672 if hash == block_hash {
2673 return Some(header);
2674 }
2675 None
2676 }))
2677 })
2678 }
2679
2680 fn header_by_number(&self, num: BlockNumber) -> ProviderResult<Option<Self::Header>> {
2681 self.get_segment_provider_for_block(StaticFileSegment::Headers, num, None)
2682 .and_then(|provider| provider.header_by_number(num))
2683 .or_else(|err| {
2684 if let ProviderError::MissingStaticFileBlock(_, _) = err {
2685 Ok(None)
2686 } else {
2687 Err(err)
2688 }
2689 })
2690 }
2691
2692 fn headers_range(
2693 &self,
2694 range: impl RangeBounds<BlockNumber>,
2695 ) -> ProviderResult<Vec<Self::Header>> {
2696 self.fetch_range_with_predicate(
2697 StaticFileSegment::Headers,
2698 to_range(range),
2699 |cursor, number| cursor.get_one::<HeaderMask<Self::Header>>(number.into()),
2700 |_| true,
2701 )
2702 }
2703
2704 fn sealed_header(
2705 &self,
2706 num: BlockNumber,
2707 ) -> ProviderResult<Option<SealedHeader<Self::Header>>> {
2708 self.get_segment_provider_for_block(StaticFileSegment::Headers, num, None)
2709 .and_then(|provider| provider.sealed_header(num))
2710 .or_else(|err| {
2711 if let ProviderError::MissingStaticFileBlock(_, _) = err {
2712 Ok(None)
2713 } else {
2714 Err(err)
2715 }
2716 })
2717 }
2718
2719 fn sealed_headers_while(
2720 &self,
2721 range: impl RangeBounds<BlockNumber>,
2722 predicate: impl FnMut(&SealedHeader<Self::Header>) -> bool,
2723 ) -> ProviderResult<Vec<SealedHeader<Self::Header>>> {
2724 self.fetch_range_with_predicate(
2725 StaticFileSegment::Headers,
2726 to_range(range),
2727 |cursor, number| {
2728 Ok(cursor
2729 .get_two::<HeaderWithHashMask<Self::Header>>(number.into())?
2730 .map(|(header, hash)| SealedHeader::new(header, hash)))
2731 },
2732 predicate,
2733 )
2734 }
2735}
2736
2737impl<N: NodePrimitives> BlockHashReader for StaticFileProvider<N> {
2738 fn block_hash(&self, num: u64) -> ProviderResult<Option<B256>> {
2739 self.get_segment_provider_for_block(StaticFileSegment::Headers, num, None)
2740 .and_then(|provider| provider.block_hash(num))
2741 .or_else(|err| {
2742 if let ProviderError::MissingStaticFileBlock(_, _) = err {
2743 Ok(None)
2744 } else {
2745 Err(err)
2746 }
2747 })
2748 }
2749
2750 fn canonical_hashes_range(
2751 &self,
2752 start: BlockNumber,
2753 end: BlockNumber,
2754 ) -> ProviderResult<Vec<B256>> {
2755 self.fetch_range_with_predicate(
2756 StaticFileSegment::Headers,
2757 start..end,
2758 |cursor, number| cursor.get_one::<BlockHashMask>(number.into()),
2759 |_| true,
2760 )
2761 }
2762}
2763
2764impl<N: NodePrimitives<SignedTx: Value + SignedTransaction, Receipt: Value>> ReceiptProvider
2765 for StaticFileProvider<N>
2766{
2767 type Receipt = N::Receipt;
2768
2769 fn receipt(&self, num: TxNumber) -> ProviderResult<Option<Self::Receipt>> {
2770 self.get_segment_provider_for_transaction(StaticFileSegment::Receipts, num, None)
2771 .and_then(|provider| provider.receipt(num))
2772 .or_else(|err| {
2773 if let ProviderError::MissingStaticFileTx(_, _) = err {
2774 Ok(None)
2775 } else {
2776 Err(err)
2777 }
2778 })
2779 }
2780
2781 fn receipt_by_hash(&self, hash: TxHash) -> ProviderResult<Option<Self::Receipt>> {
2782 if let Some(num) = self.transaction_id(hash)? {
2783 return self.receipt(num);
2784 }
2785 Ok(None)
2786 }
2787
2788 fn receipts_by_block(
2789 &self,
2790 _block: BlockHashOrNumber,
2791 ) -> ProviderResult<Option<Vec<Self::Receipt>>> {
2792 unreachable!()
2793 }
2794
2795 fn receipts_by_tx_range(
2796 &self,
2797 range: impl RangeBounds<TxNumber>,
2798 ) -> ProviderResult<Vec<Self::Receipt>> {
2799 self.fetch_range_with_predicate(
2800 StaticFileSegment::Receipts,
2801 to_range(range),
2802 |cursor, number| cursor.get_one::<ReceiptMask<Self::Receipt>>(number.into()),
2803 |_| true,
2804 )
2805 }
2806
2807 fn receipts_by_block_range(
2808 &self,
2809 _block_range: RangeInclusive<BlockNumber>,
2810 ) -> ProviderResult<Vec<Vec<Self::Receipt>>> {
2811 Err(ProviderError::UnsupportedProvider)
2812 }
2813}
2814
2815impl<N: NodePrimitives<SignedTx: Value, Receipt: Value, BlockHeader: Value>> TransactionsProviderExt
2816 for StaticFileProvider<N>
2817{
2818 fn transaction_hashes_by_range(
2819 &self,
2820 tx_range: Range<TxNumber>,
2821 ) -> ProviderResult<Vec<(TxHash, TxNumber)>> {
2822 let tx_range_size = (tx_range.end - tx_range.start) as usize;
2823
2824 let chunk_size = 100;
2828
2829 let chunks = tx_range
2831 .clone()
2832 .step_by(chunk_size)
2833 .map(|start| start..std::cmp::min(start + chunk_size as u64, tx_range.end));
2834 let mut channels = Vec::with_capacity(tx_range_size.div_ceil(chunk_size));
2835
2836 for chunk_range in chunks {
2837 let (channel_tx, channel_rx) = mpsc::channel();
2838 channels.push(channel_rx);
2839
2840 let manager = self.clone();
2841
2842 rayon::spawn(move || {
2845 let _ = manager.fetch_range_with_predicate(
2846 StaticFileSegment::Transactions,
2847 chunk_range,
2848 |cursor, number| {
2849 Ok(cursor
2850 .get_one::<TransactionMask<Self::Transaction>>(number.into())?
2851 .map(|transaction| {
2852 let _ = channel_tx.send(transaction_hash((number, transaction)));
2853 }))
2854 },
2855 |_| true,
2856 );
2857 });
2858 }
2859
2860 let mut tx_list = Vec::with_capacity(tx_range_size);
2861
2862 for channel in channels {
2864 while let Ok(tx) = channel.recv() {
2865 let (tx_hash, tx_id) = tx.map_err(|boxed| *boxed)?;
2866 tx_list.push((tx_hash, tx_id));
2867 }
2868 }
2869
2870 Ok(tx_list)
2871 }
2872}
2873
2874impl<N: NodePrimitives<SignedTx: Decompress + SignedTransaction>> TransactionsProvider
2875 for StaticFileProvider<N>
2876{
2877 type Transaction = N::SignedTx;
2878
2879 fn transaction_id(&self, tx_hash: TxHash) -> ProviderResult<Option<TxNumber>> {
2880 self.find_static_file(StaticFileSegment::Transactions, |jar_provider| {
2881 let mut cursor = jar_provider.cursor()?;
2882 if cursor
2883 .get_one::<TransactionMask<Self::Transaction>>((&tx_hash).into())?
2884 .and_then(|tx| (*tx.tx_hash() == tx_hash).then_some(tx))
2885 .is_some()
2886 {
2887 Ok(cursor.number())
2888 } else {
2889 Ok(None)
2890 }
2891 })
2892 }
2893
2894 fn transaction_by_id(&self, num: TxNumber) -> ProviderResult<Option<Self::Transaction>> {
2895 self.get_segment_provider_for_transaction(StaticFileSegment::Transactions, num, None)
2896 .and_then(|provider| provider.transaction_by_id(num))
2897 .or_else(|err| {
2898 if let ProviderError::MissingStaticFileTx(_, _) = err {
2899 Ok(None)
2900 } else {
2901 Err(err)
2902 }
2903 })
2904 }
2905
2906 fn transaction_by_id_unhashed(
2907 &self,
2908 num: TxNumber,
2909 ) -> ProviderResult<Option<Self::Transaction>> {
2910 self.get_segment_provider_for_transaction(StaticFileSegment::Transactions, num, None)
2911 .and_then(|provider| provider.transaction_by_id_unhashed(num))
2912 .or_else(|err| {
2913 if let ProviderError::MissingStaticFileTx(_, _) = err {
2914 Ok(None)
2915 } else {
2916 Err(err)
2917 }
2918 })
2919 }
2920
2921 fn transaction_by_hash(&self, hash: TxHash) -> ProviderResult<Option<Self::Transaction>> {
2922 self.find_static_file(StaticFileSegment::Transactions, |jar_provider| {
2923 Ok(jar_provider
2924 .cursor()?
2925 .get_one::<TransactionMask<Self::Transaction>>((&hash).into())?
2926 .and_then(|tx| (*tx.tx_hash() == hash).then_some(tx)))
2927 })
2928 }
2929
2930 fn transaction_by_hash_with_meta(
2931 &self,
2932 _hash: TxHash,
2933 ) -> ProviderResult<Option<(Self::Transaction, TransactionMeta)>> {
2934 Err(ProviderError::UnsupportedProvider)
2936 }
2937
2938 fn transactions_by_block(
2939 &self,
2940 _block_id: BlockHashOrNumber,
2941 ) -> ProviderResult<Option<Vec<Self::Transaction>>> {
2942 Err(ProviderError::UnsupportedProvider)
2944 }
2945
2946 fn transactions_by_block_range(
2947 &self,
2948 _range: impl RangeBounds<BlockNumber>,
2949 ) -> ProviderResult<Vec<Vec<Self::Transaction>>> {
2950 Err(ProviderError::UnsupportedProvider)
2952 }
2953
2954 fn transactions_by_tx_range(
2955 &self,
2956 range: impl RangeBounds<TxNumber>,
2957 ) -> ProviderResult<Vec<Self::Transaction>> {
2958 self.fetch_range_with_predicate(
2959 StaticFileSegment::Transactions,
2960 to_range(range),
2961 |cursor, number| cursor.get_one::<TransactionMask<Self::Transaction>>(number.into()),
2962 |_| true,
2963 )
2964 }
2965
2966 fn senders_by_tx_range(
2967 &self,
2968 range: impl RangeBounds<TxNumber>,
2969 ) -> ProviderResult<Vec<Address>> {
2970 self.fetch_range_with_predicate(
2971 StaticFileSegment::TransactionSenders,
2972 to_range(range),
2973 |cursor, number| cursor.get_one::<TransactionSenderMask>(number.into()),
2974 |_| true,
2975 )
2976 }
2977
2978 fn transaction_sender(&self, id: TxNumber) -> ProviderResult<Option<Address>> {
2979 self.get_segment_provider_for_transaction(StaticFileSegment::TransactionSenders, id, None)
2980 .and_then(|provider| provider.transaction_sender(id))
2981 .or_else(|err| {
2982 if let ProviderError::MissingStaticFileTx(_, _) = err {
2983 Ok(None)
2984 } else {
2985 Err(err)
2986 }
2987 })
2988 }
2989}
2990
2991impl<N: NodePrimitives> BlockNumReader for StaticFileProvider<N> {
2992 fn chain_info(&self) -> ProviderResult<ChainInfo> {
2993 Err(ProviderError::UnsupportedProvider)
2995 }
2996
2997 fn best_block_number(&self) -> ProviderResult<BlockNumber> {
2998 Err(ProviderError::UnsupportedProvider)
3000 }
3001
3002 fn last_block_number(&self) -> ProviderResult<BlockNumber> {
3003 Ok(self.get_highest_static_file_block(StaticFileSegment::Headers).unwrap_or_default())
3004 }
3005
3006 fn block_number(&self, _hash: B256) -> ProviderResult<Option<BlockNumber>> {
3007 Err(ProviderError::UnsupportedProvider)
3009 }
3010}
3011
3012impl<N: NodePrimitives<SignedTx: Value, Receipt: Value, BlockHeader: Value>> BlockReader
3015 for StaticFileProvider<N>
3016{
3017 type Block = N::Block;
3018
3019 fn find_block_by_hash(
3020 &self,
3021 _hash: B256,
3022 _source: BlockSource,
3023 ) -> ProviderResult<Option<Self::Block>> {
3024 Err(ProviderError::UnsupportedProvider)
3026 }
3027
3028 fn block(&self, _id: BlockHashOrNumber) -> ProviderResult<Option<Self::Block>> {
3029 Err(ProviderError::UnsupportedProvider)
3031 }
3032
3033 fn pending_block(&self) -> ProviderResult<Option<Arc<RecoveredBlock<Self::Block>>>> {
3034 Err(ProviderError::UnsupportedProvider)
3036 }
3037
3038 fn pending_block_and_receipts(
3039 &self,
3040 ) -> ProviderResult<Option<RecoveredBlockAndExecutionOutput<Self::Block, Self::Receipt>>> {
3041 Err(ProviderError::UnsupportedProvider)
3043 }
3044
3045 fn recovered_block(
3046 &self,
3047 _id: BlockHashOrNumber,
3048 _transaction_kind: TransactionVariant,
3049 ) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
3050 Err(ProviderError::UnsupportedProvider)
3052 }
3053
3054 fn sealed_block_with_senders(
3055 &self,
3056 _id: BlockHashOrNumber,
3057 _transaction_kind: TransactionVariant,
3058 ) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
3059 Err(ProviderError::UnsupportedProvider)
3061 }
3062
3063 fn block_range(&self, _range: RangeInclusive<BlockNumber>) -> ProviderResult<Vec<Self::Block>> {
3064 Err(ProviderError::UnsupportedProvider)
3066 }
3067
3068 fn block_with_senders_range(
3069 &self,
3070 _range: RangeInclusive<BlockNumber>,
3071 ) -> ProviderResult<Vec<RecoveredBlock<Self::Block>>> {
3072 Err(ProviderError::UnsupportedProvider)
3073 }
3074
3075 fn recovered_block_range(
3076 &self,
3077 _range: RangeInclusive<BlockNumber>,
3078 ) -> ProviderResult<Vec<RecoveredBlock<Self::Block>>> {
3079 Err(ProviderError::UnsupportedProvider)
3080 }
3081
3082 fn block_by_transaction_id(&self, _id: TxNumber) -> ProviderResult<Option<BlockNumber>> {
3083 Err(ProviderError::UnsupportedProvider)
3084 }
3085}
3086
3087impl<N: NodePrimitives> BlockBodyIndicesProvider for StaticFileProvider<N> {
3088 fn block_body_indices(&self, _num: u64) -> ProviderResult<Option<StoredBlockBodyIndices>> {
3089 Err(ProviderError::UnsupportedProvider)
3090 }
3091
3092 fn block_body_indices_range(
3093 &self,
3094 _range: RangeInclusive<BlockNumber>,
3095 ) -> ProviderResult<Vec<StoredBlockBodyIndices>> {
3096 Err(ProviderError::UnsupportedProvider)
3097 }
3098}
3099
3100impl<N: NodePrimitives> StatsReader for StaticFileProvider<N> {
3101 fn count_entries<T: Table>(&self) -> ProviderResult<usize> {
3102 match T::NAME {
3103 tables::CanonicalHeaders::NAME |
3104 tables::Headers::<Header>::NAME |
3105 tables::HeaderTerminalDifficulties::NAME => Ok(self
3106 .get_highest_static_file_block(StaticFileSegment::Headers)
3107 .map(|block| block + 1)
3108 .unwrap_or_default()
3109 as usize),
3110 tables::Receipts::<Receipt>::NAME => Ok(self
3111 .get_highest_static_file_tx(StaticFileSegment::Receipts)
3112 .map(|receipts| receipts + 1)
3113 .unwrap_or_default() as usize),
3114 tables::Transactions::<TransactionSigned>::NAME => Ok(self
3115 .get_highest_static_file_tx(StaticFileSegment::Transactions)
3116 .map(|txs| txs + 1)
3117 .unwrap_or_default()
3118 as usize),
3119 tables::TransactionSenders::NAME => Ok(self
3120 .get_highest_static_file_tx(StaticFileSegment::TransactionSenders)
3121 .map(|txs| txs + 1)
3122 .unwrap_or_default() as usize),
3123 _ => Err(ProviderError::UnsupportedProvider),
3124 }
3125 }
3126}
3127
3128#[inline]
3130fn transaction_hash<T>(entry: (TxNumber, T)) -> Result<(B256, TxNumber), Box<ProviderError>>
3131where
3132 T: TxHashRef,
3133{
3134 let (tx_id, tx) = entry;
3135 Ok((*tx.tx_hash(), tx_id))
3136}
3137
3138#[cfg(test)]
3139mod tests {
3140 use std::{
3141 collections::BTreeMap,
3142 sync::{atomic::Ordering, mpsc},
3143 thread,
3144 time::{Duration, Instant},
3145 };
3146
3147 use alloy_consensus::Header;
3148 use alloy_primitives::B256;
3149 use reth_chain_state::EthPrimitives;
3150 use reth_db::test_utils::create_test_static_files_dir;
3151 use reth_static_file_types::{SegmentRangeInclusive, StaticFileSegment};
3152
3153 use super::StaticFileWriter;
3154 use crate::{providers::StaticFileProvider, BlockHashReader, StaticFileProviderBuilder};
3155
3156 #[test]
3157 fn recovery_prunes_transaction_zero_at_empty_genesis() -> eyre::Result<()> {
3158 use crate::{
3159 test_utils::create_test_provider_factory, StaticFileProviderFactory,
3160 TransactionsProvider,
3161 };
3162 use alloy_consensus::{SignableTransaction, TxLegacy};
3163 use alloy_primitives::Signature;
3164 use reth_db::{tables, transaction::DbTxMut};
3165
3166 let factory = create_test_provider_factory();
3167 let provider = factory.provider_rw()?;
3168 provider.tx_ref().put::<tables::BlockBodyIndices>(0, Default::default())?;
3169 provider.commit()?;
3170
3171 let static_files = factory.static_file_provider();
3172 let segment = StaticFileSegment::Transactions;
3173 let tx = TxLegacy::default().into_signed(Signature::test_signature()).into();
3174 {
3175 let mut writer = static_files.latest_writer(segment)?;
3176 writer.increment_block(0)?;
3177 writer.increment_block(1)?;
3178 writer.append_transaction(0, &tx)?;
3179 writer.commit()?;
3180 }
3181 assert!(static_files.transaction_by_id(0)?.is_some());
3182
3183 assert_eq!(static_files.check_consistency(&factory.provider()?)?, None);
3184 assert_eq!(static_files.get_highest_static_file_block(segment), Some(0));
3185 assert!(static_files.transaction_by_id(0)?.is_none());
3186
3187 let mut writer = static_files.latest_writer(segment)?;
3188 writer.increment_block(1)?;
3189 writer.append_transaction(0, &tx)?;
3190 writer.commit()?;
3191 Ok(())
3192 }
3193
3194 #[test]
3195 fn stale_cache_fill_does_not_survive_index_reinitialization() -> eyre::Result<()> {
3196 let (static_dir, _) = create_test_static_files_dir();
3197 let static_files: StaticFileProvider<EthPrimitives> =
3198 StaticFileProviderBuilder::read_write(&static_dir)
3199 .with_blocks_per_file_for_segment(StaticFileSegment::Headers, 10)
3200 .with_blocks_per_file_for_segment(StaticFileSegment::Receipts, 2)
3201 .build()?;
3202
3203 let hash_0 = B256::from([0x10; 32]);
3204 let header_0 = Header { number: 0, ..Default::default() };
3205 {
3206 let mut writer = static_files.latest_writer(StaticFileSegment::Headers)?;
3207 writer.append_header(&header_0, &hash_0)?;
3208 writer.commit()?;
3209 }
3210
3211 {
3214 let mut writer = static_files.latest_writer(StaticFileSegment::Receipts)?;
3215 for block in 0..=3 {
3216 writer.increment_block(block)?;
3217 }
3218 writer.commit()?;
3219 }
3220 static_files.delete_jar(StaticFileSegment::Receipts, 0)?;
3221
3222 let (loaded_tx, loaded_rx) = mpsc::sync_channel(0);
3223 let (resume_tx, resume_rx) = mpsc::sync_channel(0);
3224 static_files.before_cache_insert.set(move |segment| {
3225 assert_eq!(segment, StaticFileSegment::Headers);
3226 loaded_tx.send(()).expect("test controller dropped");
3227 resume_rx
3228 .recv_timeout(Duration::from_secs(10))
3229 .expect("test controller did not resume");
3230 });
3231
3232 let reader_static_files = static_files.clone();
3234 let reader = thread::spawn(move || reader_static_files.block_hash(0));
3235 loaded_rx
3236 .recv_timeout(Duration::from_secs(10))
3237 .expect("reader did not reach cache-fill hook");
3238
3239 let hash_1 = B256::from([0x11; 32]);
3242 let header_1 = Header { number: 1, ..Default::default() };
3243 {
3244 let mut writer = static_files.latest_writer(StaticFileSegment::Headers)?;
3245 writer.append_header(&header_1, &hash_1)?;
3246 writer.commit()?;
3247 }
3248 static_files.delete_jar(StaticFileSegment::Receipts, 2)?;
3249
3250 resume_tx.send(()).expect("reader dropped before cache publication");
3251 assert_eq!(reader.join().expect("reader panicked")?, Some(hash_0));
3252 assert_eq!(
3253 static_files.block_hash(1)?,
3254 Some(hash_1),
3255 "cache fill started before reinitialization repopulated an invalidated jar"
3256 );
3257
3258 Ok(())
3259 }
3260
3261 #[test]
3262 fn cache_fill_validation_is_serialized_with_index_reinitialization() -> eyre::Result<()> {
3263 let (static_dir, _) = create_test_static_files_dir();
3264 let static_files: StaticFileProvider<EthPrimitives> =
3265 StaticFileProviderBuilder::read_write(&static_dir).with_blocks_per_file(10).build()?;
3266 let hash = B256::from([0x10; 32]);
3267 {
3268 let mut writer = static_files.latest_writer(StaticFileSegment::Headers)?;
3269 writer.append_header(&Header::default(), &hash)?;
3270 writer.commit()?;
3271 }
3272 static_files.remove_cached_provider(StaticFileSegment::Headers, 9);
3273
3274 let generation = static_files.cache_generation.load(Ordering::Acquire);
3275 let (validated_tx, validated_rx) = mpsc::sync_channel(0);
3276 let (resume_tx, resume_rx) = mpsc::sync_channel(0);
3277 static_files.after_cache_validation.set(move |segment| {
3278 assert_eq!(segment, StaticFileSegment::Headers);
3279 validated_tx.send(()).expect("test controller dropped");
3280 resume_rx
3281 .recv_timeout(Duration::from_secs(10))
3282 .expect("test controller did not resume");
3283 });
3284
3285 let reader_static_files = static_files.clone();
3286 let reader = thread::spawn(move || reader_static_files.block_hash(0));
3287 validated_rx
3288 .recv_timeout(Duration::from_secs(10))
3289 .expect("reader did not reach cache validation");
3290
3291 assert!(static_files.map.try_get(&(9, StaticFileSegment::Headers)).is_locked());
3294 let invalidator_static_files = static_files.clone();
3295 let invalidator = thread::spawn(move || invalidator_static_files.initialize_index());
3296 let deadline = Instant::now() + Duration::from_secs(10);
3297 while static_files.cache_generation.load(Ordering::Acquire) == generation {
3298 assert!(Instant::now() < deadline, "reinitialization did not invalidate the cache");
3299 thread::sleep(Duration::from_millis(1));
3300 }
3301
3302 resume_tx.send(()).expect("reader dropped before cache publication");
3305 assert_eq!(reader.join().expect("reader panicked")?, Some(hash));
3306 invalidator.join().expect("invalidator panicked")?;
3307 assert!(static_files.map.is_empty());
3308 assert_eq!(static_files.block_hash(0)?, Some(hash));
3309
3310 Ok(())
3311 }
3312
3313 #[test]
3314 fn test_find_fixed_range_with_block_index() -> eyre::Result<()> {
3315 let (static_dir, _) = create_test_static_files_dir();
3316 let sf_rw: StaticFileProvider<EthPrimitives> =
3317 StaticFileProviderBuilder::read_write(&static_dir).with_blocks_per_file(100).build()?;
3318
3319 let segment = StaticFileSegment::Headers;
3320
3321 assert_eq!(
3323 sf_rw.find_fixed_range_with_block_index(segment, None, 0),
3324 SegmentRangeInclusive::new(0, 99)
3325 );
3326 assert_eq!(
3327 sf_rw.find_fixed_range_with_block_index(segment, None, 250),
3328 SegmentRangeInclusive::new(200, 299)
3329 );
3330
3331 assert_eq!(
3333 sf_rw.find_fixed_range_with_block_index(segment, Some(&BTreeMap::new()), 150),
3334 SegmentRangeInclusive::new(100, 199)
3335 );
3336
3337 let block_index = BTreeMap::from_iter([
3339 (99, SegmentRangeInclusive::new(0, 99)),
3340 (199, SegmentRangeInclusive::new(100, 199)),
3341 (299, SegmentRangeInclusive::new(200, 299)),
3342 ]);
3343
3344 assert_eq!(
3346 sf_rw.find_fixed_range_with_block_index(segment, Some(&block_index), 0),
3347 SegmentRangeInclusive::new(0, 99)
3348 );
3349 assert_eq!(
3350 sf_rw.find_fixed_range_with_block_index(segment, Some(&block_index), 50),
3351 SegmentRangeInclusive::new(0, 99)
3352 );
3353 assert_eq!(
3354 sf_rw.find_fixed_range_with_block_index(segment, Some(&block_index), 99),
3355 SegmentRangeInclusive::new(0, 99)
3356 );
3357 assert_eq!(
3358 sf_rw.find_fixed_range_with_block_index(segment, Some(&block_index), 100),
3359 SegmentRangeInclusive::new(100, 199)
3360 );
3361 assert_eq!(
3362 sf_rw.find_fixed_range_with_block_index(segment, Some(&block_index), 150),
3363 SegmentRangeInclusive::new(100, 199)
3364 );
3365 assert_eq!(
3366 sf_rw.find_fixed_range_with_block_index(segment, Some(&block_index), 199),
3367 SegmentRangeInclusive::new(100, 199)
3368 );
3369
3370 assert_eq!(
3373 sf_rw.find_fixed_range_with_block_index(segment, Some(&block_index), 300),
3374 SegmentRangeInclusive::new(300, 399)
3375 );
3376 assert_eq!(
3377 sf_rw.find_fixed_range_with_block_index(segment, Some(&block_index), 350),
3378 SegmentRangeInclusive::new(300, 399)
3379 );
3380
3381 assert_eq!(
3383 sf_rw.find_fixed_range_with_block_index(segment, Some(&block_index), 500),
3384 SegmentRangeInclusive::new(500, 599)
3385 );
3386
3387 assert_eq!(
3389 sf_rw.find_fixed_range_with_block_index(segment, Some(&block_index), 1000),
3390 SegmentRangeInclusive::new(1000, 1099)
3391 );
3392
3393 let mixed_size_index = BTreeMap::from_iter([
3396 (49, SegmentRangeInclusive::new(0, 49)), (149, SegmentRangeInclusive::new(50, 149)), (349, SegmentRangeInclusive::new(150, 349)), ]);
3400
3401 assert_eq!(
3403 sf_rw.find_fixed_range_with_block_index(segment, Some(&mixed_size_index), 25),
3404 SegmentRangeInclusive::new(0, 49)
3405 );
3406 assert_eq!(
3407 sf_rw.find_fixed_range_with_block_index(segment, Some(&mixed_size_index), 100),
3408 SegmentRangeInclusive::new(50, 149)
3409 );
3410 assert_eq!(
3411 sf_rw.find_fixed_range_with_block_index(segment, Some(&mixed_size_index), 200),
3412 SegmentRangeInclusive::new(150, 349)
3413 );
3414
3415 assert_eq!(
3418 sf_rw.find_fixed_range_with_block_index(segment, Some(&mixed_size_index), 350),
3419 SegmentRangeInclusive::new(350, 449)
3420 );
3421 assert_eq!(
3422 sf_rw.find_fixed_range_with_block_index(segment, Some(&mixed_size_index), 450),
3423 SegmentRangeInclusive::new(450, 549)
3424 );
3425 assert_eq!(
3426 sf_rw.find_fixed_range_with_block_index(segment, Some(&mixed_size_index), 550),
3427 SegmentRangeInclusive::new(550, 649)
3428 );
3429
3430 Ok(())
3431 }
3432}