Skip to main content

reth_provider/providers/static_file/
manager.rs

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
66/// Alias type for a map that can be queried for block or transaction ranges. It uses `u64` to
67/// represent either a block or a transaction number end of a static file range.
68type SegmentRanges = BTreeMap<u64, SegmentRangeInclusive>;
69
70/// Access mode on a static file provider. RO/RW.
71#[derive(Debug, Default, PartialEq, Eq)]
72pub enum StaticFileAccess {
73    /// Read-only access.
74    #[default]
75    RO,
76    /// Read-write access.
77    RW,
78}
79
80impl StaticFileAccess {
81    /// Returns `true` if read-only access.
82    pub const fn is_read_only(&self) -> bool {
83        matches!(self, Self::RO)
84    }
85
86    /// Returns `true` if read-write access.
87    pub const fn is_read_write(&self) -> bool {
88        matches!(self, Self::RW)
89    }
90}
91
92/// Context for static file block writes.
93///
94/// Contains target segments and pruning configuration.
95#[derive(Debug, Clone, Copy, Default)]
96pub struct StaticFileWriteCtx {
97    /// Whether transaction senders should be written to static files.
98    pub write_senders: bool,
99    /// Whether receipts should be written to static files.
100    pub write_receipts: bool,
101    /// Whether account changesets should be written to static files.
102    pub write_account_changesets: bool,
103    /// Whether storage changesets should be written to static files.
104    pub write_storage_changesets: bool,
105    /// The current chain tip block number (for pruning).
106    pub tip: BlockNumber,
107    /// The prune mode for receipts, if any.
108    pub receipts_prune_mode: Option<reth_prune_types::PruneMode>,
109    /// Whether receipts are prunable (based on storage settings and prune distance).
110    pub receipts_prunable: bool,
111}
112
113/// [`StaticFileProvider`] manages all existing [`StaticFileJarProvider`].
114///
115/// "Static files" contain immutable chain history data, such as:
116///  - transactions
117///  - headers
118///  - receipts
119///
120/// This provider type is responsible for reading and writing to static files.
121#[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/// Builder for [`StaticFileProvider`] that allows configuration before initialization.
131#[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    /// Creates a new builder with read-write access.
142    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    /// Creates a new builder with read-only access.
153    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    /// Set custom blocks per file for specific segments.
164    ///
165    /// Each static file segment is stored across multiple files, and each of these files contains
166    /// up to the specified number of blocks of data. When the file gets full, a new file is
167    /// created with the new block range.
168    ///
169    /// This setting affects the size of each static file, and can be set per segment.
170    ///
171    /// If it is changed for an existing node, existing static files will not be affected and will
172    /// be finished with the old blocks per file setting, but new static files will use the new
173    /// setting.
174    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    /// Set a custom number of blocks per file for all segments.
185    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    /// Set a custom number of blocks per file for a specific segment.
193    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    /// Enables metrics on the [`StaticFileProvider`].
203    pub const fn with_metrics(mut self) -> Self {
204        self.use_metrics = true;
205        self
206    }
207
208    /// Sets the genesis block number for the [`StaticFileProvider`].
209    ///
210    /// This configures the genesis block number, which is used to determine the starting point
211    /// for block indexing and querying operations.
212    ///
213    /// # Arguments
214    ///
215    /// * `genesis_block_number` - The block number of the genesis block.
216    ///
217    /// # Returns
218    ///
219    /// Returns `Self` to allow method chaining.
220    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    /// Builds the final [`StaticFileProvider`] and initializes the index.
226    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    /// Creates a new [`StaticFileProvider`] with the given [`StaticFileAccess`].
245    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    /// Creates a new [`StaticFileProvider`] with read-only access.
254    ///
255    /// The caller is responsible for calling [`StaticFileProvider::initialize_index`] when
256    /// underlying data changes.
257    pub fn read_only(path: impl AsRef<Path>) -> ProviderResult<Self> {
258        Self::new(path, StaticFileAccess::RO)
259    }
260
261    /// Creates a new [`StaticFileProvider`] with read-write access.
262    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/// [`StaticFileProviderInner`] manages all existing [`StaticFileJarProvider`].
304#[derive(Debug)]
305pub struct StaticFileProviderInner<N> {
306    /// Maintains a map which allows for concurrent access to different `NippyJars`, over different
307    /// segments and ranges.
308    map: DashMap<(BlockNumber, StaticFileSegment), LoadedJar>,
309    /// Incremented before clearing the jar cache so in-flight loads cannot repopulate it with a
310    /// snapshot from before the invalidation.
311    cache_generation: AtomicU64,
312    /// Indexes per segment.
313    indexes: RwLock<StaticFileMap<StaticFileSegmentIndex>>,
314    /// This is an additional index that tracks the expired height, this will track the highest
315    /// block number that has been expired (missing). The first, non expired block is
316    /// `expired_history_height + 1`.
317    ///
318    /// This is effectively the transaction range that has been expired:
319    /// [`StaticFileProvider::delete_segment_below_block`] and mirrors
320    /// `static_files_min_block[transactions] - blocks_per_file`.
321    ///
322    /// This additional tracker exists for more efficient lookups because the node must be aware of
323    /// the expired height.
324    earliest_history_height: AtomicU64,
325    /// Directory where `static_files` are located
326    path: PathBuf,
327    /// Maintains a writer set of [`StaticFileSegment`].
328    writers: StaticFileWriters<N>,
329    /// Metrics for the static files.
330    metrics: Option<Arc<StaticFileProviderMetrics>>,
331    /// Access rights of the provider.
332    access: StaticFileAccess,
333    /// Number of blocks per file, per segment.
334    blocks_per_file: StaticFileMap<u64>,
335    /// Write lock for when access is [`StaticFileAccess::RW`].
336    _lock_file: Option<StorageLock>,
337    /// Genesis block number, default is 0;
338    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    /// Creates a new [`StaticFileProviderInner`].
347    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    /// Each static file has a fixed number of blocks. This gives out the range where the requested
385    /// block is positioned.
386    ///
387    /// If the specified block falls into one of the ranges of already initialized static files,
388    /// this function will return that range.
389    ///
390    /// If no matching file exists, this function will derive a new range from the end of the last
391    /// existing file, if any.
392    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            // Find first block range that contains the requested block
403            if let Some((_, range)) = block_index.range(block..).next() {
404                // Found matching range for an existing file using block index
405                return *range;
406            } else if let Some((_, range)) = block_index.last_key_value() {
407                // Didn't find matching range for an existing file, derive a new range from the end
408                // of the last existing file range.
409                //
410                // `block` is always higher than `range.end()` here, because `block_index` holds no
411                // range with a `max_block` greater than or equal to `block`
412                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        // No block index is available, derive a new range using the fixed number of blocks,
419        // starting from the beginning.
420        find_fixed_range(block, blocks_per_file)
421    }
422
423    /// Each static file has a fixed number of blocks. This gives out the range where the requested
424    /// block is positioned.
425    ///
426    /// If the specified block falls into one of the ranges of already initialized static files,
427    /// this function will return that range.
428    ///
429    /// If no matching file exists, this function will derive a new range from the end of the last
430    /// existing file, if any.
431    ///
432    /// This function will block indefinitely if a write lock for
433    /// [`Self::indexes`] is already acquired. In that case, use
434    /// [`Self::find_fixed_range_with_block_index`].
435    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    /// Get genesis block number
448    pub const fn genesis_block_number(&self) -> u64 {
449        self.genesis_block_number
450    }
451}
452
453impl<N: NodePrimitives> StaticFileProvider<N> {
454    /// Reports metrics for the static files.
455    ///
456    /// This uses the in-memory index to get file sizes from mmap handles instead of reading
457    /// filesystem metadata.
458    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    /// Writes headers for all blocks to the static file segment.
485    #[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    /// Writes transactions for all blocks to the static file segment.
498    #[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    /// Writes transaction senders for all blocks to the static file segment.
515    #[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    /// Writes receipts for all blocks to the static file segment.
532    #[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            // skip writing receipts if pruning configuration requires us to.
544            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    /// Writes account changesets for all blocks to the static file segment.
559    #[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    /// Writes storage changesets for all blocks to the static file segment.
580    #[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    /// Writes to a static file segment using the provided closure.
609    ///
610    /// The closure receives a mutable reference to the segment writer. After the closure completes,
611    /// `sync_all()` is called to flush writes to disk.
612    #[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    /// Writes all static file data for multiple blocks in parallel per-segment.
628    ///
629    /// This spawns tasks on the storage thread pool for each segment type and each task calls
630    /// `sync_all()` on its writer when done.
631    #[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        // Propagate tracing context into rayon-spawned threads so that per-segment
653        // write spans appear as children of write_blocks_data in traces.
654        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    /// Gets the [`StaticFileJarProvider`] of the requested segment and start index that can be
738    /// either block or transaction.
739    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    /// Gets the [`StaticFileJarProvider`] of the requested segment and start index that can be
752    /// either block or transaction.
753    ///
754    /// If the segment is not found, returns [`None`].
755    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    /// Gets the [`StaticFileJarProvider`] of the requested segment and block.
777    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    /// Gets the [`StaticFileJarProvider`] of the requested segment and transaction.
792    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    /// Gets the [`StaticFileJarProvider`] of the requested segment and block or transaction.
807    ///
808    /// `fn_range` should make sure the range goes through `find_fixed_range`.
809    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        // If we have a path, then get the block range from its name.
816        // Otherwise, check `self.available_static_files`
817        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        // Return cached `LoadedJar` or insert it for the first time, and then, return it.
836        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    /// Gets the [`StaticFileJarProvider`] of the requested path.
844    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    /// Given a segment and block range it removes the cached provider from the map.
859    ///
860    /// CAUTION: cached provider should be dropped before calling this or IT WILL deadlock.
861    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    /// This handles history expiry by deleting all static files for the given segment below the
870    /// given block.
871    ///
872    /// For example if block is 1M and the blocks per file are 500K this will delete all individual
873    /// files below 1M, so 0-499K and 500K-999K.
874    ///
875    /// This will not delete the file that contains the block itself, because files can only be
876    /// removed entirely.
877    ///
878    /// # Safety
879    ///
880    /// This method will never delete the highest static file for the segment, even if the
881    /// requested block is higher than the highest block in static files. This ensures we always
882    /// maintain at least one static file if any exist.
883    ///
884    /// Returns a list of `SegmentHeader`s from the deleted jars.
885    pub fn delete_segment_below_block(
886        &self,
887        segment: StaticFileSegment,
888        block: BlockNumber,
889    ) -> ProviderResult<Vec<SegmentHeader>> {
890        // Nothing to delete if block is 0.
891        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            // Stop if we've reached the target block or the highest static file
904            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            // now we need to wipe the static file, this will take care of updating the index and
916            // advance the lowest tracked block height for the segment.
917            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    /// Given a segment and block, it deletes the jar and all files from the respective block range.
926    ///
927    /// CAUTION: destructive. Deletes files on disk.
928    ///
929    /// This will re-initialize the index after deletion, so all files are tracked.
930    ///
931    /// Returns the `SegmentHeader` of the deleted jar.
932    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        // Delete the sidecar file for changeset segments before deleting the main jar
956        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        // SAFETY: this is currently necessary to ensure that certain indexes like
966        // `static_files_min_block` have the correct values after pruning.
967        self.initialize_index()?;
968
969        Ok(header)
970    }
971
972    /// Deletes ALL static file jars for the given segment, including the highest one.
973    ///
974    /// CAUTION: destructive. Deletes all files on disk for this segment.
975    ///
976    /// This is used for `PruneMode::Full` where all data should be removed.
977    ///
978    /// Returns a list of `SegmentHeader`s from the deleted jars.
979    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    /// Given a segment and block range it returns a cached
1003    /// [`StaticFileJarProvider`]. TODO(joshie): we should check the size and pop N if there's too
1004    /// many.
1005    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        // Avoid using `entry` directly to avoid a write lock in the common case.
1013        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            // The cache may have been populated since the initial miss, including by
1027            // `update_index` publishing a newer snapshot while we loaded this jar without a lock.
1028            // Preserve that entry instead of overwriting it with our potentially stale snapshot.
1029            self.map
1030                .entry(key)
1031                .or_try_insert_with(|| -> ProviderResult<_> {
1032                    // Validate while holding the entry's shard lock so clearing the cache cannot
1033                    // slip between validation and insertion. If invalidated, reload once under
1034                    // this lock rather than retrying unboundedly during frequent pruning.
1035                    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    /// Gets a static file segment's block range from the provider inner block
1057    /// index.
1058    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    /// Gets a static file segment's fixed block range from the provider inner
1076    /// transaction index.
1077    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        // It's more probable that the request comes from a newer tx height, so we iterate
1087        // the static_files in reverse.
1088        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                // request tx is higher than highest static file tx
1093                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    /// Updates the inner transaction and block indexes alongside the internal cached providers in
1108    /// `self.map`.
1109    ///
1110    /// Any entry higher than `segment_max_block` will be deleted from the previous structures.
1111    ///
1112    /// If `segment_max_block` is None it means there's no static file for this segment.
1113    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                        // Update max block
1143                        index.max_block = segment_max_block;
1144
1145                        // Update expected block range index
1146
1147                        // Remove all expected block ranges that are less than the new max block
1148                        index
1149                            .expected_block_ranges_by_max_block
1150                            .retain(|_, block_range| block_range.start() < fixed_range.start());
1151                        // Insert new expected block range
1152                        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                // Update min_block to track the lowest block range of the segment.
1167                // This is initially set by initialize_index() on node startup, but must be updated
1168                // as the file grows to prevent stale values.
1169                //
1170                // Without this update, min_block can remain at genesis (e.g. Some([0..=0]) or None)
1171                // even after syncing to higher blocks (e.g. [0..=100]). A stale
1172                // min_block causes get_lowest_static_file_block() to return the
1173                // wrong end value, which breaks pruning logic that relies on it for
1174                // safety checks.
1175                //
1176                // Example progression:
1177                // 1. Node starts, initialize_index() sets min_block = [0..=0]
1178                // 2. Sync to block 100, this update sets min_block = [0..=100]
1179                // 3. Pruner calls get_lowest_static_file_block() -> returns 100 (correct). Without
1180                //    this update, it would incorrectly return 0 (stale)
1181                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                        // delete_jar WILL ALWAYS re-initialize all indexes, so we are always
1184                        // sure that current_min is always the lowest.
1185                        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                // Updates the tx index by first removing all entries which have a higher
1194                // block_start than our current static file.
1195                if let Some(tx_range) = jar.user_header().tx_range() {
1196                    // Current block range has the same block start as `fixed_range``, but block end
1197                    // might be different if we are still filling this static file.
1198                    if let Some(current_block_range) = jar.user_header().block_range() {
1199                        let tx_end = tx_range.end();
1200
1201                        // Considering that `update_index` is called when we either append/truncate,
1202                        // we are sure that we are handling the latest data
1203                        // points.
1204                        //
1205                        // Here we remove every entry of the index that has a block start higher or
1206                        // equal than our current one. This is important in the case
1207                        // that we prune a lot of rows resulting in a file (and thus
1208                        // a higher block range) deletion.
1209                        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                    // The unwinded file has no more transactions/receipts. However, the highest
1220                    // block is within this files' block range. We only retain
1221                    // entries with block ranges before the current one.
1222                    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                    // If the index is empty, just remove it.
1227                    index.available_block_ranges_by_max_tx.take_if(|index| index.is_empty());
1228                }
1229
1230                // Update the cached provider.
1231                trace!(target: "providers::static_file", ?segment, "Inserting updated jar into cache");
1232                self.map.insert((fixed_range.end(), segment), LoadedJar::new(jar)?);
1233
1234                // Delete any cached provider that no longer has an associated jar.
1235                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    /// Initializes the inner transaction and block index
1249    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            // Update first and last block for each segment
1255            //
1256            // It's safe to call `expect` here, because every segment has at least one header
1257            // associated with it.
1258            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                // Update max expected block -> expected_block_range index
1266                expected_block_ranges_by_max_block
1267                    .insert(header.expected_block_end(), header.expected_block_range());
1268
1269                // Update max tx -> block_range index
1270                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        // If this is a re-initialization, invalidate in-flight cache fills before clearing the
1291        // cache. A fill that started before this point either reloads under its shard lock or
1292        // publishes before the clear reaches that shard and is removed by it.
1293        self.cache_generation.fetch_add(1, Ordering::AcqRel);
1294        self.map.clear();
1295
1296        // initialize the expired history height to the lowest static file block
1297        if let Some(lowest_range) =
1298            indexes.get(StaticFileSegment::Transactions).and_then(|index| index.min_block_range)
1299        {
1300            // the earliest height is the lowest available block number
1301            self.earliest_history_height
1302                .store(lowest_range.start(), std::sync::atomic::Ordering::Relaxed);
1303        }
1304
1305        Ok(())
1306    }
1307
1308    /// Ensures that any broken invariants which cannot be healed on the spot return a pipeline
1309    /// target to unwind to.
1310    ///
1311    /// Two types of consistency checks are done for:
1312    ///
1313    /// 1) When a static file fails to commit but the underlying data was changed.
1314    /// 2) When a static file was committed, but the required database transaction was not.
1315    ///
1316    /// For 1) it can self-heal if `self.access.is_read_only()` is set to `false`. Otherwise, it
1317    /// will return an error.
1318    /// For 2) the invariants below are checked, and if broken, might require a pipeline unwind
1319    /// to heal.
1320    ///
1321    /// For each static file segment:
1322    /// * the corresponding database table should overlap or have continuity in their keys
1323    ///   ([`TxNumber`] or [`BlockNumber`]).
1324    /// * its highest block should match the stage checkpoint block number if it's equal or higher
1325    ///   than the corresponding database table last entry.
1326    ///
1327    /// Returns a [`Option`] of [`PipelineTarget::Unwind`] if any healing is further required.
1328    ///
1329    /// WARNING: No static file writer should be held before calling this function, otherwise it
1330    /// will deadlock.
1331    #[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        // OVM historical import is broken and does not work with this check. It's importing
1346        // duplicated receipts resulting in having more receipts than the expected transaction
1347        // range.
1348        //
1349        // If we detect an OVM import was done (block #1 <https://optimistic.etherscan.io/block/1>), skip it.
1350        // More on [#11099](https://github.com/paradigmxyz/reth/pull/11099).
1351        if provider.chain_spec().is_optimism() &&
1352            reth_chainspec::Chain::optimism_mainnet() == provider.chain_spec().chain_id()
1353        {
1354            // check whether we have the first OVM block: <https://optimistic.etherscan.io/block/0xbee7192e575af30420cae0c7776304ac196077ee72b048970549e4f08e875453>
1355            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            // Heal file-level inconsistencies and get before/after highest block
1387            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            // Only applies to block-based static files. (Headers)
1392            //
1393            // The updated `highest_block` may have decreased if we healed from a pruning
1394            // interruption.
1395            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            // Only applies to transaction-based static files. (Receipts & Transactions)
1405            //
1406            // Make sure the last transaction matches the last block from its indices, since a heal
1407            // from a pruning interruption might have decreased the number of transactions without
1408            // being able to update the last block of the static file segment.
1409            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                        // If the block body indices can not be found, then it means that static
1420                        // files is ahead of database, and the `ensure_invariants` check will fix
1421                        // it by comparing with stage checkpoints.
1422                        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    /// Heals file-level (`NippyJar`) inconsistencies for eligible static file
1467    /// segments.
1468    ///
1469    /// Call before [`Self::check_consistency`] so files are internally
1470    /// consistent.
1471    ///
1472    /// Uses the same segment-skip logic as [`Self::check_consistency`], but
1473    /// does not compare with database checkpoints or prune against them.
1474    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    /// Returns the static file segments that should be checked/healed for this provider.
1489    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    /// True if the given segment should be checked/healed for this provider.
1501    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                    // Old pruned nodes (including full node) do not store receipts as static
1514                    // files.
1515                    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                    // Gnosis and Chiado's historical import is broken and does not work with
1523                    // this check. They are importing receipts along
1524                    // with importing headers/bodies.
1525                    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    /// Returns `true` if the given prune segment has a checkpoint with
1562    /// [`reth_prune_types::PruneMode::Full`], indicating all data for this segment has been
1563    /// intentionally deleted.
1564    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    /// Checks consistency of the latest static file segment and throws an
1576    /// error if at fault.
1577    ///
1578    /// Read-only.
1579    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    /// Attempts to heal file-level (`NippyJar`) inconsistencies for a single static file segment.
1599    ///
1600    /// Returns the highest block before and after healing, which can be used to detect
1601    /// if healing from a pruning interruption decreased the highest block.
1602    ///
1603    /// File consistency is broken if:
1604    ///
1605    /// * appending data was interrupted before a config commit, then data file will be truncated
1606    ///   according to the config.
1607    ///
1608    /// * pruning data was interrupted before a config commit, then we have deleted data that we are
1609    ///   expected to still have. We need to check the Database and unwind everything accordingly.
1610    ///
1611    /// **Note:** In read-only mode, this will return an error if a consistency issue is detected,
1612    /// since healing requires write access.
1613    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            // Read-only mode: cannot modify files, so just validate consistency and error if
1622            // broken.
1623            debug!(target: "reth::providers::static_file", "Checking segment consistency (read-only)");
1624            self.check_segment_consistency(segment)?;
1625        } else {
1626            // Writable mode: fetching the writer will automatically heal any file-level
1627            // inconsistency by truncating data to match the last committed config.
1628            debug!(target: "reth::providers::static_file", "Fetching latest writer which might heal any potential inconsistency");
1629            self.latest_writer(segment)?;
1630        }
1631
1632        // The updated `highest_block` may have decreased if we healed from a
1633        // pruning interruption.
1634        let highest_block = self.get_highest_static_file_block(segment);
1635
1636        Ok((initial_highest_block, highest_block))
1637    }
1638
1639    /// Ensure invariants for each corresponding table and static file segment.
1640    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    /// Check invariants for each corresponding table and static file segment:
1698    ///
1699    /// * the corresponding database table should overlap or have continuity in their keys
1700    ///   ([`TxNumber`] or [`BlockNumber`]).
1701    /// * its highest block should match the stage checkpoint block number if it's equal or higher
1702    ///   than the corresponding database table last entry.
1703    ///   * If the checkpoint block is higher, then request a pipeline unwind to the static file
1704    ///     block. This is expressed by returning [`Some`] with the requested pipeline unwind
1705    ///     target.
1706    ///   * If the checkpoint block is lower, then heal by removing rows from the static file. In
1707    ///     this case, the rows will be removed and [`None`] will be returned.
1708    ///
1709    /// * If the database tables overlap with static files and have contiguous keys, or the
1710    ///   checkpoint block matches the highest static files block, then [`None`] will be returned.
1711    #[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 there is a gap between the entry found in static file and
1731                // database, then we have most likely lost static file data and need to unwind so we
1732                // can load it again
1733                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        // If static file entry is ahead of the database entries, then ensure the checkpoint block
1760        // number matches.
1761        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 the checkpoint is ahead, then we lost static file data. May be data corruption.
1770        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 the checkpoint is ahead, or matches, then nothing to do.
1781        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        // If the checkpoint is behind, then we failed to do a database commit
1787        // **but committed** to static files on executing a stage, or the
1788        // reverse on unwinding a stage.
1789        //
1790        // All we need to do is to prune the extra static file rows.
1791        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                // TODO(joshie): is_block_meta
1804                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                    // `last_tx_num()` saturates to zero for an empty genesis block, but row zero
1811                    // belongs to the first non-empty block and must be removed as well.
1812                    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    /// Returns the highest block accounted for in this segment, either through data
1949    /// present in static files or through data intentionally removed by pruning.
1950    ///
1951    /// Data below a segment's prune checkpoint has been intentionally deleted, so its
1952    /// absence from static files is not an inconsistency. Without this, a pruned segment
1953    /// whose stage checkpoint is ahead of its (empty) static files is treated as data
1954    /// corruption and triggers an unwind to block 0, which aborts the node on startup.
1955    /// See <https://github.com/paradigmxyz/reth/issues/23463>.
1956    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    /// Returns the prune segment that governs data availability for a static file segment,
1977    /// or `None` if the segment is never pruned.
1978    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    /// Returns the earliest available block number that has not been expired and is still
1989    /// available.
1990    ///
1991    /// This means that the highest expired block (or expired block height) is
1992    /// `earliest_history_height.saturating_sub(1)`.
1993    ///
1994    /// Returns `0` if no history has been expired.
1995    pub fn earliest_history_height(&self) -> BlockNumber {
1996        self.earliest_history_height.load(std::sync::atomic::Ordering::Relaxed)
1997    }
1998
1999    /// Gets the lowest static file's block range if it exists for a static file segment.
2000    ///
2001    /// If there is nothing on disk for the given segment, this will return [`None`].
2002    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    /// Gets the lowest static file's block range start if it exists for a static file segment.
2007    ///
2008    /// For example if the lowest static file has blocks 0-499, this will return 0.
2009    ///
2010    /// If there is nothing on disk for the given segment, this will return [`None`].
2011    pub fn get_lowest_range_start(&self, segment: StaticFileSegment) -> Option<BlockNumber> {
2012        self.get_lowest_range(segment).map(|range| range.start())
2013    }
2014
2015    /// Gets the lowest static file's block range end if it exists for a static file segment.
2016    ///
2017    /// For example if the static file has blocks 0-499, this will return 499.
2018    ///
2019    /// If there is nothing on disk for the given segment, this will return [`None`].
2020    pub fn get_lowest_range_end(&self, segment: StaticFileSegment) -> Option<BlockNumber> {
2021        self.get_lowest_range(segment).map(|range| range.end())
2022    }
2023
2024    /// Gets the highest static file's block height if it exists for a static file segment.
2025    ///
2026    /// If there is nothing on disk for the given segment, this will return [`None`].
2027    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    /// Converts a range to a bounded `RangeInclusive` capped to the highest static file block.
2032    ///
2033    /// This is necessary because static file iteration beyond the tip would loop forever:
2034    /// blocks beyond the static file tip return `Ok(empty)` which is indistinguishable from
2035    /// blocks with no changes. We cap the end to the highest available block regardless of
2036    /// whether the input was unbounded or an explicit large value like `BlockNumber::MAX`.
2037    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    /// Gets the highest static file transaction.
2059    ///
2060    /// If there is nothing on disk for the given segment, this will return [`None`].
2061    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    /// Gets the highest static file block for all segments.
2070    pub fn get_highest_static_files(&self) -> HighestStaticFiles {
2071        HighestStaticFiles {
2072            receipts: self.get_highest_static_file_block(StaticFileSegment::Receipts),
2073        }
2074    }
2075
2076    /// Iterates through segment `static_files` in reverse order, executing a function until it
2077    /// returns some object. Useful for finding objects by [`TxHash`] or [`BlockHash`].
2078    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            // Iterate through all ranges in reverse order (highest to lowest)
2087            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    /// Fetches data within a specified range across multiple static files.
2098    ///
2099    /// This function iteratively retrieves data using `get_fn` for each item in the given range.
2100    /// It continues fetching until the end of the range is reached or the provided `predicate`
2101    /// returns false.
2102    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        /// Resolves to the provider for the given block or transaction number.
2116        ///
2117        /// If the static file is missing, the `result` is returned.
2118        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        // advances number in range
2135        'outer: for number in range {
2136            // The `retrying` flag ensures a single retry attempt per `number`. If `get_fn` fails to
2137            // access data in two different static files, it halts further attempts by returning
2138            // an error, effectively preventing infinite retry loops.
2139            let mut retrying = false;
2140
2141            // advances static files if `get_fn` returns None
2142            '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                        // There is a very small chance of hitting a deadlock if two consecutive
2156                        // static files share the same bucket in the
2157                        // internal dashmap and we don't drop the current provider
2158                        // before requesting the next one.
2159                        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    /// Fetches data within a specified range across multiple static files.
2175    ///
2176    /// Returns an iterator over the data. Yields [`None`] if the data for the specified number is
2177    /// not found.
2178    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                    // There is a very small chance of hitting a deadlock if two consecutive
2198                    // static files share the same bucket in the internal dashmap and we don't drop
2199                    // the current provider before requesting the next one.
2200                    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    /// Returns directory where `static_files` are located.
2213    pub fn directory(&self) -> &Path {
2214        &self.path
2215    }
2216
2217    /// Retrieves data from the database or static file, wherever it's available.
2218    ///
2219    /// # Arguments
2220    /// * `segment` - The segment of the static file to check against.
2221    /// * `index_key` - Requested index key, usually a block or transaction number.
2222    /// * `fetch_from_static_file` - A closure that defines how to fetch the data from the static
2223    ///   file provider.
2224    /// * `fetch_from_database` - A closure that defines how to fetch the data from the database
2225    ///   when the static file doesn't contain the required data or is not available.
2226    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        // If there is, check the maximum block or transaction number of the segment.
2238        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    /// Gets data within a specified range, potentially spanning different `static_files` and
2253    /// database.
2254    ///
2255    /// # Arguments
2256    /// * `segment` - The segment of the static file to query.
2257    /// * `block_or_tx_range` - The range of data to fetch.
2258    /// * `fetch_from_static_file` - A function to fetch data from the `static_file`.
2259    /// * `fetch_from_database` - A function to fetch data from the database.
2260    /// * `predicate` - A function used to evaluate each item in the fetched data. Fetching is
2261    ///   terminated when this function returns false, thereby filtering the data based on the
2262    ///   provided condition.
2263    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 there is, check the maximum block or transaction number of the segment.
2279        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    /// Returns static files directory
2302    #[cfg(any(test, feature = "test-utils"))]
2303    pub fn path(&self) -> &Path {
2304        &self.path
2305    }
2306
2307    /// Returns transaction index
2308    #[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    /// Returns expected block index
2318    #[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 static file block range.
2331    ///
2332    /// This index is initialized on launch to keep track of the lowest, non-expired static file
2333    /// per segment and gets updated on [`StaticFileProvider::update_index`].
2334    ///
2335    /// This tracks the lowest static file per segment together with the block range in that
2336    /// file. E.g. static file is batched in 500k block intervals then the lowest static file
2337    /// is [0..499K], and the block range is start = 0, end = 499K.
2338    ///
2339    /// This index is mainly used for history expiry, which targets transactions, e.g. pre-merge
2340    /// history expiry would lead to removing all static files below the merge height.
2341    min_block_range: Option<SegmentRangeInclusive>,
2342    /// Max static file block.
2343    max_block: u64,
2344    /// Expected static file block ranges indexed by max expected blocks.
2345    ///
2346    /// For example, a static file for expected block range `0..=499_000` may have only block range
2347    /// `0..=1000` contained in it, as it's not fully filled yet. This index maps the max expected
2348    /// block to the expected range, i.e. block `499_000` to block range `0..=499_000`.
2349    expected_block_ranges_by_max_block: SegmentRanges,
2350    /// Available on disk static file block ranges indexed by max transactions.
2351    ///
2352    /// For example, a static file for block range `0..=499_000` may only have block range
2353    /// `0..=1000` and transaction range `0..=2000` contained in it. This index maps the max
2354    /// available transaction to the available block range, i.e. transaction `2000` to block range
2355    /// `0..=1000`.
2356    available_block_ranges_by_max_tx: Option<SegmentRanges>,
2357}
2358
2359/// Helper trait to manage different [`StaticFileProviderRW`] of an `Arc<StaticFileProvider`
2360pub trait StaticFileWriter {
2361    /// The primitives type used by the static file provider.
2362    type Primitives: Send + Sync + 'static;
2363
2364    /// Returns a mutable reference to a [`StaticFileProviderRW`] of a [`StaticFileSegment`].
2365    fn get_writer(
2366        &self,
2367        block: BlockNumber,
2368        segment: StaticFileSegment,
2369    ) -> ProviderResult<StaticFileProviderRWRefMut<'_, Self::Primitives>>;
2370
2371    /// Returns a mutable reference to a [`StaticFileProviderRW`] of the latest
2372    /// [`StaticFileSegment`].
2373    fn latest_writer(
2374        &self,
2375        segment: StaticFileSegment,
2376    ) -> ProviderResult<StaticFileProviderRWRefMut<'_, Self::Primitives>>;
2377
2378    /// Commits all changes of all [`StaticFileProviderRW`] of all [`StaticFileSegment`].
2379    fn commit(&self) -> ProviderResult<()>;
2380
2381    /// Returns `true` if the static file provider has unwind queued.
2382    fn has_unwind_queued(&self) -> bool;
2383
2384    /// Finalizes all static file writers by committing their configuration to disk.
2385    ///
2386    /// Returns an error if prune is queued (use [`Self::commit`] instead).
2387    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                // This is not expected but means we are out of the range / file somehow, and can't
2500                // continue
2501                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    /// Creates an iterator for walking through account changesets in the specified block range.
2639    ///
2640    /// This returns a lazy iterator that fetches changesets block by block to avoid loading
2641    /// everything into memory at once.
2642    ///
2643    /// Accepts any range type that implements `RangeBounds<BlockNumber>`, including:
2644    /// - `Range<BlockNumber>` (e.g., `0..100`)
2645    /// - `RangeInclusive<BlockNumber>` (e.g., `0..=99`)
2646    /// - `RangeFrom<BlockNumber>` (e.g., `0..`) - iterates until exhausted
2647    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    /// Creates an iterator for walking through storage changesets in the specified block range.
2655    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        // Transactions are different size, so chunks will not all take the same processing time. If
2825        // chunks are too big, there will be idle threads waiting for work. Choosing an
2826        // arbitrary smaller value to make sure it doesn't happen.
2827        let chunk_size = 100;
2828
2829        // iterator over the chunks
2830        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            // Spawn the task onto the global rayon pool
2843            // This task will send the cached transaction hash through the channel.
2844            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        // Iterate over channels and append the tx hashes unsorted
2863        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        // Required data not present in static_files
2935        Err(ProviderError::UnsupportedProvider)
2936    }
2937
2938    fn transactions_by_block(
2939        &self,
2940        _block_id: BlockHashOrNumber,
2941    ) -> ProviderResult<Option<Vec<Self::Transaction>>> {
2942        // Required data not present in static_files
2943        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        // Required data not present in static_files
2951        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        // Required data not present in static_files
2994        Err(ProviderError::UnsupportedProvider)
2995    }
2996
2997    fn best_block_number(&self) -> ProviderResult<BlockNumber> {
2998        // Required data not present in static_files
2999        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        // Required data not present in static_files
3008        Err(ProviderError::UnsupportedProvider)
3009    }
3010}
3011
3012/* Cannot be successfully implemented but must exist for trait requirements */
3013
3014impl<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        // Required data not present in static_files
3025        Err(ProviderError::UnsupportedProvider)
3026    }
3027
3028    fn block(&self, _id: BlockHashOrNumber) -> ProviderResult<Option<Self::Block>> {
3029        // Required data not present in static_files
3030        Err(ProviderError::UnsupportedProvider)
3031    }
3032
3033    fn pending_block(&self) -> ProviderResult<Option<Arc<RecoveredBlock<Self::Block>>>> {
3034        // Required data not present in static_files
3035        Err(ProviderError::UnsupportedProvider)
3036    }
3037
3038    fn pending_block_and_receipts(
3039        &self,
3040    ) -> ProviderResult<Option<RecoveredBlockAndExecutionOutput<Self::Block, Self::Receipt>>> {
3041        // Required data not present in static_files
3042        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        // Required data not present in static_files
3051        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        // Required data not present in static_files
3060        Err(ProviderError::UnsupportedProvider)
3061    }
3062
3063    fn block_range(&self, _range: RangeInclusive<BlockNumber>) -> ProviderResult<Vec<Self::Block>> {
3064        // Required data not present in static_files
3065        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/// Returns the tx hash for the given transaction and its id.
3129#[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        // Create two receipt jars. Deleting either one reinitializes every segment index and
3212        // invalidates the shared jar cache, matching the pruning path from the incidents.
3213        {
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        // Pause a cache fill after it loads the header jar containing only block 0.
3233        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        // Publish a newer snapshot of the same header jar, then simulate unrelated segment
3240        // pruning invalidating the whole cache while the old load remains in flight.
3241        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        // The shard must remain locked between checking the generation and inserting the jar.
3292        // Otherwise a clear can finish here and the reader can repopulate the cache afterward.
3293        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        // Invalidation has started, but clearing this shard must wait for publication and the
3303        // reader's guard to be released. The published snapshot must not survive that clear.
3304        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        // Test with None - should use default behavior
3322        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        // Test with empty index - should fall back to default behavior
3332        assert_eq!(
3333            sf_rw.find_fixed_range_with_block_index(segment, Some(&BTreeMap::new()), 150),
3334            SegmentRangeInclusive::new(100, 199)
3335        );
3336
3337        // Create block index with existing ranges
3338        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        // Test blocks within existing ranges - should return the matching range
3345        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        // Test blocks beyond existing ranges - should derive new ranges from the last range
3371        // Block 300 is exactly one segment after the last range
3372        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        // Block 500 skips one segment (300-399)
3382        assert_eq!(
3383            sf_rw.find_fixed_range_with_block_index(segment, Some(&block_index), 500),
3384            SegmentRangeInclusive::new(500, 599)
3385        );
3386
3387        // Block 1000 skips many segments
3388        assert_eq!(
3389            sf_rw.find_fixed_range_with_block_index(segment, Some(&block_index), 1000),
3390            SegmentRangeInclusive::new(1000, 1099)
3391        );
3392
3393        // Test with block index having different sizes than blocks_per_file setting
3394        // This simulates the scenario where blocks_per_file was changed between runs
3395        let mixed_size_index = BTreeMap::from_iter([
3396            (49, SegmentRangeInclusive::new(0, 49)),     // 50 blocks
3397            (149, SegmentRangeInclusive::new(50, 149)),  // 100 blocks
3398            (349, SegmentRangeInclusive::new(150, 349)), // 200 blocks
3399        ]);
3400
3401        // Blocks within existing ranges should return those ranges regardless of size
3402        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        // Block after the last range should derive using current blocks_per_file (100)
3416        // from the end of the last range (349)
3417        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}