1use super::{
2 manager::StaticFileProviderInner, metrics::StaticFileProviderMetrics, StaticFileProvider,
3};
4use crate::providers::static_file::metrics::StaticFileProviderOperation;
5use alloy_consensus::BlockHeader;
6use alloy_primitives::{BlockHash, BlockNumber, TxNumber, U256};
7use parking_lot::{lock_api::RwLockWriteGuard, RawRwLock, RwLock};
8use reth_codecs::Compact;
9use reth_db::models::{AccountBeforeTx, StorageBeforeTx};
10use reth_db_api::models::CompactU256;
11use reth_nippy_jar::{NippyJar, NippyJarError, NippyJarWriter};
12use reth_node_types::NodePrimitives;
13use reth_primitives_traits::FastInstant as Instant;
14use reth_static_file_types::{
15 ChangesetOffset, ChangesetOffsetReader, ChangesetOffsetWriter, SegmentHeader,
16 SegmentRangeInclusive, StaticFileSegment,
17};
18use reth_storage_errors::provider::{ProviderError, ProviderResult, StaticFileWriterError};
19use std::{
20 borrow::Borrow,
21 cmp::Ordering,
22 fmt::Debug,
23 path::{Path, PathBuf},
24 sync::{Arc, Weak},
25};
26use tracing::{debug, instrument};
27
28#[derive(Debug, Clone, Copy)]
30enum PruneStrategy {
31 Headers {
33 num_blocks: u64,
35 },
36 Transactions {
38 num_rows: u64,
40 last_block: BlockNumber,
42 },
43 Receipts {
45 num_rows: u64,
47 last_block: BlockNumber,
49 },
50 TransactionSenders {
52 num_rows: u64,
54 last_block: BlockNumber,
56 },
57 AccountChangeSets {
59 last_block: BlockNumber,
61 },
62 StorageChangeSets {
64 last_block: BlockNumber,
66 },
67}
68
69#[derive(Debug)]
74pub(crate) struct StaticFileWriters<N> {
75 headers: RwLock<Option<StaticFileProviderRW<N>>>,
76 transactions: RwLock<Option<StaticFileProviderRW<N>>>,
77 receipts: RwLock<Option<StaticFileProviderRW<N>>>,
78 transaction_senders: RwLock<Option<StaticFileProviderRW<N>>>,
79 account_change_sets: RwLock<Option<StaticFileProviderRW<N>>>,
80 storage_change_sets: RwLock<Option<StaticFileProviderRW<N>>>,
81}
82
83impl<N> Default for StaticFileWriters<N> {
84 fn default() -> Self {
85 Self {
86 headers: Default::default(),
87 transactions: Default::default(),
88 receipts: Default::default(),
89 transaction_senders: Default::default(),
90 account_change_sets: Default::default(),
91 storage_change_sets: Default::default(),
92 }
93 }
94}
95
96impl<N: NodePrimitives> StaticFileWriters<N> {
97 pub(crate) fn get_or_create(
98 &self,
99 segment: StaticFileSegment,
100 create_fn: impl FnOnce() -> ProviderResult<StaticFileProviderRW<N>>,
101 ) -> ProviderResult<StaticFileProviderRWRefMut<'_, N>> {
102 let mut write_guard = match segment {
103 StaticFileSegment::Headers => self.headers.write(),
104 StaticFileSegment::Transactions => self.transactions.write(),
105 StaticFileSegment::Receipts => self.receipts.write(),
106 StaticFileSegment::TransactionSenders => self.transaction_senders.write(),
107 StaticFileSegment::AccountChangeSets => self.account_change_sets.write(),
108 StaticFileSegment::StorageChangeSets => self.storage_change_sets.write(),
109 };
110
111 if write_guard.is_none() {
112 *write_guard = Some(create_fn()?);
113 }
114
115 Ok(StaticFileProviderRWRefMut(write_guard))
116 }
117
118 pub(crate) fn remove(&self, segment: StaticFileSegment) {
120 let mut write_guard = match segment {
121 StaticFileSegment::Headers => self.headers.write(),
122 StaticFileSegment::Transactions => self.transactions.write(),
123 StaticFileSegment::Receipts => self.receipts.write(),
124 StaticFileSegment::TransactionSenders => self.transaction_senders.write(),
125 StaticFileSegment::AccountChangeSets => self.account_change_sets.write(),
126 StaticFileSegment::StorageChangeSets => self.storage_change_sets.write(),
127 };
128
129 *write_guard = None;
130 }
131
132 #[instrument(
133 name = "StaticFileWriters::commit",
134 level = "debug",
135 target = "providers::static_file",
136 skip_all
137 )]
138 pub(crate) fn commit(&self) -> ProviderResult<()> {
139 debug!(target: "providers::static_file", "Committing all static file segments");
140
141 for writer_lock in [
142 &self.headers,
143 &self.transactions,
144 &self.receipts,
145 &self.transaction_senders,
146 &self.account_change_sets,
147 &self.storage_change_sets,
148 ] {
149 let mut writer = writer_lock.write();
150 if let Some(writer) = writer.as_mut() {
151 writer.commit()?;
152 }
153 }
154
155 debug!(target: "providers::static_file", "Committed all static file segments");
156 Ok(())
157 }
158
159 pub(crate) fn has_unwind_queued(&self) -> bool {
160 for writer_lock in [
161 &self.headers,
162 &self.transactions,
163 &self.receipts,
164 &self.transaction_senders,
165 &self.account_change_sets,
166 &self.storage_change_sets,
167 ] {
168 let writer = writer_lock.read();
169 if let Some(writer) = writer.as_ref() &&
170 writer.will_prune_on_commit()
171 {
172 return true
173 }
174 }
175 false
176 }
177
178 #[instrument(
183 name = "StaticFileWriters::finalize",
184 level = "debug",
185 target = "providers::static_file",
186 skip_all
187 )]
188 pub(crate) fn finalize(&self) -> ProviderResult<()> {
189 debug!(target: "providers::static_file", "Finalizing all static file segments into disk");
190
191 for writer_lock in [
192 &self.headers,
193 &self.transactions,
194 &self.receipts,
195 &self.transaction_senders,
196 &self.account_change_sets,
197 &self.storage_change_sets,
198 ] {
199 let mut writer = writer_lock.write();
200 if let Some(writer) = writer.as_mut() {
201 writer.finalize()?;
202 }
203 }
204
205 debug!(target: "providers::static_file", "Finalized all static file segments into disk");
206 Ok(())
207 }
208}
209
210#[derive(Debug)]
212pub struct StaticFileProviderRWRefMut<'a, N>(
213 pub(crate) RwLockWriteGuard<'a, RawRwLock, Option<StaticFileProviderRW<N>>>,
214);
215
216impl<N> std::ops::DerefMut for StaticFileProviderRWRefMut<'_, N> {
217 fn deref_mut(&mut self) -> &mut Self::Target {
218 self.0.as_mut().expect("static file writer provider should be init")
220 }
221}
222
223impl<N> std::ops::Deref for StaticFileProviderRWRefMut<'_, N> {
224 type Target = StaticFileProviderRW<N>;
225
226 fn deref(&self) -> &Self::Target {
227 self.0.as_ref().expect("static file writer provider should be init")
229 }
230}
231
232#[derive(Debug)]
233pub struct StaticFileProviderRW<N> {
235 reader: Weak<StaticFileProviderInner<N>>,
240 writer: NippyJarWriter<SegmentHeader>,
242 data_path: PathBuf,
244 buf: Vec<u8>,
246 metrics: Option<Arc<StaticFileProviderMetrics>>,
248 prune_on_commit: Option<PruneStrategy>,
250 synced: bool,
252 changeset_offsets: Option<ChangesetOffsetWriter>,
254 current_changeset_offset: Option<ChangesetOffset>,
256}
257
258impl<N: NodePrimitives> StaticFileProviderRW<N> {
259 pub fn new(
264 segment: StaticFileSegment,
265 block: BlockNumber,
266 reader: Weak<StaticFileProviderInner<N>>,
267 metrics: Option<Arc<StaticFileProviderMetrics>>,
268 ) -> ProviderResult<Self> {
269 let (writer, data_path) = Self::open(segment, block, reader.clone(), metrics.clone())?;
270
271 let mut writer = Self {
273 writer,
274 data_path,
275 buf: Vec::with_capacity(100),
276 reader,
277 metrics,
278 prune_on_commit: None,
279 synced: false,
280 changeset_offsets: None,
281 current_changeset_offset: None,
282 };
283
284 writer.ensure_end_range_consistency()?;
287
288 if segment.is_change_based() {
290 writer.heal_changeset_sidecar()?;
291 }
292
293 Ok(writer)
294 }
295
296 fn open(
297 segment: StaticFileSegment,
298 block: u64,
299 reader: Weak<StaticFileProviderInner<N>>,
300 metrics: Option<Arc<StaticFileProviderMetrics>>,
301 ) -> ProviderResult<(NippyJarWriter<SegmentHeader>, PathBuf)> {
302 let start = Instant::now();
303
304 let static_file_provider = Self::upgrade_provider_to_strong_reference(&reader);
305
306 let block_range = static_file_provider.find_fixed_range(segment, block);
307 let (jar, path) = match static_file_provider.get_segment_provider_for_block(
308 segment,
309 block_range.start(),
310 None,
311 ) {
312 Ok(provider) => (
313 NippyJar::load(provider.data_path()).map_err(ProviderError::other)?,
314 provider.data_path().into(),
315 ),
316 Err(ProviderError::MissingStaticFileBlock(_, _)) => {
317 let path = static_file_provider.directory().join(segment.filename(&block_range));
318 (create_jar(segment, &path, block_range), path)
319 }
320 Err(err) => return Err(err),
321 };
322
323 let result = match NippyJarWriter::new(jar) {
324 Ok(writer) => Ok((writer, path)),
325 Err(NippyJarError::FrozenJar) => {
326 Err(ProviderError::FinalizedStaticFile(segment, block))
328 }
329 Err(e) => Err(ProviderError::other(e)),
330 }?;
331
332 if let Some(metrics) = &metrics {
333 metrics.record_segment_operation(
334 segment,
335 StaticFileProviderOperation::OpenWriter,
336 Some(start.elapsed()),
337 );
338 }
339
340 Ok(result)
341 }
342
343 fn ensure_end_range_consistency(&mut self) -> ProviderResult<()> {
352 let expected_rows = if self.user_header().segment().is_headers() {
354 self.user_header().block_len().unwrap_or_default()
355 } else {
356 self.user_header().tx_len().unwrap_or_default()
357 };
358 let actual_rows = self.writer.rows() as u64;
359 let pruned_rows = expected_rows.saturating_sub(actual_rows);
360 if pruned_rows > 0 {
361 self.user_header_mut().prune(pruned_rows);
362 }
363
364 debug!(
365 target: "providers::static_file",
366 segment = ?self.writer.user_header().segment(),
367 path = ?self.data_path,
368 pruned_rows,
369 "Ensuring end range consistency"
370 );
371
372 self.writer.commit().map_err(ProviderError::other)?;
373
374 self.update_index()?;
376 Ok(())
377 }
378
379 pub const fn will_prune_on_commit(&self) -> bool {
381 self.prune_on_commit.is_some()
382 }
383
384 fn heal_changeset_sidecar(&mut self) -> ProviderResult<()> {
392 let csoff_path = self.data_path.with_extension("csoff");
393
394 let header_claims_blocks = self.writer.user_header().changeset_offsets_len();
396 let actual_nippy_rows = self.writer.rows() as u64;
397
398 let actual_sidecar_blocks = if csoff_path.exists() {
400 let file_len = reth_fs_util::metadata(&csoff_path).map_err(ProviderError::other)?.len();
401 let aligned_len = file_len - (file_len % 16);
403 aligned_len / 16
404 } else {
405 0
406 };
407
408 if header_claims_blocks == 0 && actual_sidecar_blocks == 0 {
410 self.changeset_offsets =
411 Some(ChangesetOffsetWriter::new(&csoff_path, 0).map_err(ProviderError::other)?);
412 return Ok(());
413 }
414
415 let valid_blocks = if actual_sidecar_blocks > 0 {
417 let reader = ChangesetOffsetReader::new(&csoff_path, actual_sidecar_blocks)
418 .map_err(ProviderError::other)?;
419
420 let mut valid = 0u64;
423 for i in 0..actual_sidecar_blocks {
424 if let Some(offset) = reader.get(i).map_err(ProviderError::other)? {
425 if offset.offset() + offset.num_changes() <= actual_nippy_rows {
426 valid = i + 1;
427 } else {
428 break;
430 }
431 }
432 }
433 valid
434 } else {
435 0
436 };
437
438 let correct_blocks = valid_blocks.min(header_claims_blocks);
441
442 let mut needs_header_commit = false;
444
445 if correct_blocks != header_claims_blocks || actual_sidecar_blocks != correct_blocks {
446 tracing::warn!(
447 target: "reth::static_file",
448 path = %csoff_path.display(),
449 header_claims = header_claims_blocks,
450 sidecar_has = actual_sidecar_blocks,
451 valid_blocks = correct_blocks,
452 actual_rows = actual_nippy_rows,
453 "Three-way healing: syncing header, sidecar, and NippyJar state"
454 );
455
456 if actual_sidecar_blocks > correct_blocks {
458 use std::fs::OpenOptions;
459 let file = OpenOptions::new()
460 .write(true)
461 .open(&csoff_path)
462 .map_err(ProviderError::other)?;
463 file.set_len(correct_blocks * 16).map_err(ProviderError::other)?;
464 file.sync_all().map_err(ProviderError::other)?;
465
466 tracing::debug!(
467 target: "reth::static_file",
468 "Truncated sidecar from {} to {} blocks",
469 actual_sidecar_blocks,
470 correct_blocks
471 );
472 }
473
474 if correct_blocks < header_claims_blocks {
476 let blocks_removed = header_claims_blocks - correct_blocks;
479 self.writer.user_header_mut().prune(blocks_removed);
480
481 tracing::debug!(
482 target: "reth::static_file",
483 "Updated header: removed {} blocks (changeset_offsets_len: {} -> {})",
484 blocks_removed,
485 header_claims_blocks,
486 correct_blocks
487 );
488
489 needs_header_commit = true;
490 }
491 } else {
492 tracing::debug!(
493 target: "reth::static_file",
494 path = %csoff_path.display(),
495 blocks = correct_blocks,
496 "Changeset sidecar consistent, no healing needed"
497 );
498 }
499
500 let csoff_writer = ChangesetOffsetWriter::new(&csoff_path, correct_blocks)
502 .map_err(ProviderError::other)?;
503
504 self.changeset_offsets = Some(csoff_writer);
505
506 if needs_header_commit {
508 self.writer.commit().map_err(ProviderError::other)?;
509
510 tracing::info!(
511 target: "reth::static_file",
512 path = %csoff_path.display(),
513 blocks = correct_blocks,
514 "Committed healed changeset offset header"
515 );
516 }
517
518 Ok(())
519 }
520
521 fn flush_current_changeset_offset(&mut self) -> ProviderResult<()> {
529 if !self.writer.user_header().segment().is_change_based() {
530 return Ok(());
531 }
532
533 if let Some(offset) = self.current_changeset_offset.take() &&
534 let Some(writer) = &mut self.changeset_offsets
535 {
536 writer.append(&offset).map_err(ProviderError::other)?;
537 }
538 Ok(())
539 }
540
541 pub fn sync_all(&mut self) -> ProviderResult<()> {
548 if self.prune_on_commit.is_some() {
549 return Err(StaticFileWriterError::FinalizeWithPruneQueued.into());
550 }
551
552 self.flush_current_changeset_offset()?;
554 if let Some(writer) = &mut self.changeset_offsets {
555 writer.sync().map_err(ProviderError::other)?;
556 self.writer.user_header_mut().set_changeset_offsets_len(writer.len());
558 }
559
560 if self.writer.is_dirty() {
561 self.writer.sync_all().map_err(ProviderError::other)?;
562 }
563 self.synced = true;
564 Ok(())
565 }
566
567 #[instrument(
573 name = "StaticFileProviderRW::finalize",
574 level = "debug",
575 target = "providers::static_file",
576 skip_all
577 )]
578 pub fn finalize(&mut self) -> ProviderResult<()> {
579 if self.prune_on_commit.is_some() {
580 return Err(StaticFileWriterError::FinalizeWithPruneQueued.into());
581 }
582 if self.writer.is_dirty() {
583 if !self.synced {
584 self.sync_all()?;
587 }
588
589 self.writer.finalize().map_err(ProviderError::other)?;
590 self.update_index()?;
591 }
592 self.synced = false;
593 Ok(())
594 }
595
596 #[instrument(
598 name = "StaticFileProviderRW::commit",
599 level = "debug",
600 target = "providers::static_file",
601 skip_all
602 )]
603 pub fn commit(&mut self) -> ProviderResult<()> {
604 let start = Instant::now();
605
606 if let Some(strategy) = self.prune_on_commit.take() {
608 debug!(
609 target: "providers::static_file",
610 segment = ?self.writer.user_header().segment(),
611 "Pruning data on commit"
612 );
613 match strategy {
614 PruneStrategy::Headers { num_blocks } => self.prune_header_data(num_blocks)?,
615 PruneStrategy::Transactions { num_rows, last_block } => {
616 self.prune_transaction_data(num_rows, last_block)?
617 }
618 PruneStrategy::Receipts { num_rows, last_block } => {
619 self.prune_receipt_data(num_rows, last_block)?
620 }
621 PruneStrategy::TransactionSenders { num_rows, last_block } => {
622 self.prune_transaction_sender_data(num_rows, last_block)?
623 }
624 PruneStrategy::AccountChangeSets { last_block } => {
625 self.prune_account_changeset_data(last_block)?
626 }
627 PruneStrategy::StorageChangeSets { last_block } => {
628 self.prune_storage_changeset_data(last_block)?
629 }
630 }
631 }
632
633 self.flush_current_changeset_offset()?;
636 if let Some(writer) = &mut self.changeset_offsets {
637 writer.sync().map_err(ProviderError::other)?;
638 self.writer.user_header_mut().set_changeset_offsets_len(writer.len());
640 }
641
642 if self.writer.is_dirty() {
643 debug!(
644 target: "providers::static_file",
645 segment = ?self.writer.user_header().segment(),
646 "Committing writer to disk"
647 );
648
649 self.writer.commit().map_err(ProviderError::other)?;
651
652 if let Some(metrics) = &self.metrics {
653 metrics.record_segment_operation(
654 self.writer.user_header().segment(),
655 StaticFileProviderOperation::CommitWriter,
656 Some(start.elapsed()),
657 );
658 }
659
660 debug!(
661 target: "providers::static_file",
662 segment = ?self.writer.user_header().segment(),
663 path = ?self.data_path,
664 duration = ?start.elapsed(),
665 "Committed writer to disk"
666 );
667
668 self.update_index()?;
669 }
670
671 Ok(())
672 }
673
674 #[cfg(feature = "test-utils")]
678 pub fn commit_without_sync_all(&mut self) -> ProviderResult<()> {
679 let start = Instant::now();
680
681 debug!(
682 target: "providers::static_file",
683 segment = ?self.writer.user_header().segment(),
684 "Committing writer to disk (without sync)"
685 );
686
687 self.writer.commit_without_sync_all().map_err(ProviderError::other)?;
689
690 if let Some(metrics) = &self.metrics {
691 metrics.record_segment_operation(
692 self.writer.user_header().segment(),
693 StaticFileProviderOperation::CommitWriter,
694 Some(start.elapsed()),
695 );
696 }
697
698 debug!(
699 target: "providers::static_file",
700 segment = ?self.writer.user_header().segment(),
701 path = ?self.data_path,
702 duration = ?start.elapsed(),
703 "Committed writer to disk (without sync)"
704 );
705
706 self.update_index()?;
707
708 Ok(())
709 }
710
711 fn update_index(&self) -> ProviderResult<()> {
713 let segment = self.writer.user_header().segment();
714
715 let segment_max_block = self
723 .writer
724 .user_header()
725 .block_range()
726 .as_ref()
727 .map(|block_range| block_range.end())
728 .or_else(|| {
729 let expected_start = self.writer.user_header().expected_block_start();
730 if expected_start <= self.reader().genesis_block_number() {
731 return None;
732 }
733
734 let prev_block = expected_start - 1;
735 let prev_range = self.reader().find_fixed_range(segment, prev_block);
736 let prev_path = self.reader().directory().join(segment.filename(&prev_range));
737 prev_path.exists().then_some(prev_block)
738 });
739
740 self.reader().update_index(segment, segment_max_block)
741 }
742
743 pub fn ensure_at_block(&mut self, advance_to: BlockNumber) -> ProviderResult<()> {
748 let current_block = if let Some(current_block_number) = self.current_block_number() {
749 current_block_number
750 } else {
751 let first_block = self.writer.user_header().expected_block_start();
755 self.increment_block(first_block)?;
756 first_block
757 };
758
759 match current_block.cmp(&advance_to) {
760 Ordering::Less => {
761 for block in current_block + 1..=advance_to {
762 self.increment_block(block)?;
763 }
764 }
765 Ordering::Equal => {}
766 Ordering::Greater => {
767 return Err(ProviderError::UnexpectedStaticFileBlockNumber(
768 self.writer.user_header().segment(),
769 current_block,
770 advance_to,
771 ));
772 }
773 }
774
775 Ok(())
776 }
777
778 pub fn increment_block(&mut self, expected_block_number: BlockNumber) -> ProviderResult<()> {
781 let segment = self.writer.user_header().segment();
782
783 self.check_next_block_number(expected_block_number)?;
784
785 let start = Instant::now();
786 if let Some(last_block) = self.writer.user_header().block_end() {
787 if last_block == self.writer.user_header().expected_block_end() {
789 self.commit()?;
791
792 let (writer, data_path) =
794 Self::open(segment, last_block + 1, self.reader.clone(), self.metrics.clone())?;
795 self.writer = writer;
796 self.data_path = data_path.clone();
797
798 if segment.is_change_based() {
800 let csoff_path = data_path.with_extension("csoff");
801 self.changeset_offsets = Some(
802 ChangesetOffsetWriter::new(&csoff_path, 0).map_err(ProviderError::other)?,
803 );
804 }
805
806 *self.writer.user_header_mut() = SegmentHeader::new(
807 self.reader().find_fixed_range(segment, last_block + 1),
808 None,
809 None,
810 segment,
811 );
812 }
813 }
814
815 self.writer.user_header_mut().increment_block();
816
817 if segment.is_change_based() {
819 if let Some(offset) = self.current_changeset_offset.take() &&
821 let Some(writer) = &mut self.changeset_offsets
822 {
823 writer.append(&offset).map_err(ProviderError::other)?;
824 }
825 let new_offset = self.writer.rows() as u64;
827 self.current_changeset_offset = Some(ChangesetOffset::new(new_offset, 0));
828 }
829
830 if let Some(metrics) = &self.metrics {
831 metrics.record_segment_operation(
832 segment,
833 StaticFileProviderOperation::IncrementBlock,
834 Some(start.elapsed()),
835 );
836 }
837
838 Ok(())
839 }
840
841 pub fn current_block_number(&self) -> Option<u64> {
843 self.writer.user_header().block_end()
844 }
845
846 pub fn next_block_number(&self) -> u64 {
848 self.writer
852 .user_header()
853 .block_end()
854 .map(|b| b + 1)
855 .unwrap_or_else(|| self.writer.user_header().expected_block_start())
856 }
857
858 fn check_next_block_number(&self, expected_block_number: u64) -> ProviderResult<()> {
861 let next_static_file_block = self.next_block_number();
862
863 if expected_block_number != next_static_file_block {
864 return Err(ProviderError::UnexpectedStaticFileBlockNumber(
865 self.writer.user_header().segment(),
866 expected_block_number,
867 next_static_file_block,
868 ))
869 }
870 Ok(())
871 }
872
873 fn truncate_changesets(&mut self, last_block: u64) -> ProviderResult<()> {
879 let segment = self.writer.user_header().segment();
880 debug_assert!(segment.is_change_based());
881
882 let current_block_end = self
884 .writer
885 .user_header()
886 .block_end()
887 .ok_or(ProviderError::MissingStaticFileBlock(segment, 0))?;
888
889 if current_block_end <= last_block {
891 return Ok(())
892 }
893
894 let mut expected_block_start = self.writer.user_header().expected_block_start();
896 while last_block < expected_block_start && expected_block_start > 0 {
897 self.delete_current_and_open_previous()?;
898 expected_block_start = self.writer.user_header().expected_block_start();
899 }
900
901 let blocks_to_keep = if last_block >= expected_block_start {
903 last_block - expected_block_start + 1
904 } else {
905 0
906 };
907
908 let csoff_path = self.data_path.with_extension("csoff");
910 let changeset_offsets_len = self.writer.user_header().changeset_offsets_len();
911
912 self.flush_current_changeset_offset()?;
914
915 let rows_to_keep = if blocks_to_keep == 0 {
916 0
917 } else if blocks_to_keep >= changeset_offsets_len {
918 self.writer.rows() as u64
920 } else {
921 let reader = ChangesetOffsetReader::new(&csoff_path, changeset_offsets_len)
925 .map_err(ProviderError::other)?;
926 if let Some(next_offset) = reader.get(blocks_to_keep).map_err(ProviderError::other)? {
927 next_offset.offset()
928 } else {
929 self.writer.rows() as u64
931 }
932 };
933
934 let total_rows = self.writer.rows() as u64;
935 let rows_to_delete = total_rows.saturating_sub(rows_to_keep);
936
937 if rows_to_delete > 0 {
938 let current_block_end = self
940 .writer
941 .user_header()
942 .block_end()
943 .ok_or(ProviderError::MissingStaticFileBlock(segment, 0))?;
944 let blocks_to_remove = current_block_end - last_block;
945
946 self.writer.user_header_mut().prune(blocks_to_remove);
948
949 self.writer.prune_rows(rows_to_delete as usize).map_err(ProviderError::other)?;
951 }
952
953 self.writer.user_header_mut().set_block_range(expected_block_start, last_block);
955
956 self.writer.user_header_mut().sync_changeset_offsets();
958
959 if let Some(writer) = &mut self.changeset_offsets {
961 writer.truncate(blocks_to_keep).map_err(ProviderError::other)?;
962 }
963
964 self.current_changeset_offset = None;
966
967 self.commit()?;
969
970 Ok(())
971 }
972
973 fn truncate(&mut self, num_rows: u64, last_block: Option<u64>) -> ProviderResult<()> {
981 let mut remaining_rows = num_rows;
982 let segment = self.writer.user_header().segment();
983 while remaining_rows > 0 {
984 let len = if segment.is_block_based() {
985 self.writer.user_header().block_len().unwrap_or_default()
986 } else {
987 self.writer.user_header().tx_len().unwrap_or_default()
988 };
989
990 if remaining_rows >= len {
991 let block_start = self.writer.user_header().expected_block_start();
994
995 if block_start != 0 &&
1001 (segment.is_headers() || last_block.is_some_and(|b| b < block_start))
1002 {
1003 self.delete_current_and_open_previous()?;
1004 } else {
1005 self.writer.user_header_mut().prune(len);
1007 self.writer.prune_rows(len as usize).map_err(ProviderError::other)?;
1008 break
1009 }
1010
1011 remaining_rows -= len;
1012 } else {
1013 self.writer.user_header_mut().prune(remaining_rows);
1015
1016 self.writer.prune_rows(remaining_rows as usize).map_err(ProviderError::other)?;
1018 remaining_rows = 0;
1019 }
1020 }
1021
1022 if let Some(last_block) = last_block {
1024 let mut expected_block_start = self.writer.user_header().expected_block_start();
1025
1026 if num_rows == 0 {
1027 while last_block < expected_block_start {
1031 self.delete_current_and_open_previous()?;
1032 expected_block_start = self.writer.user_header().expected_block_start();
1033 }
1034 }
1035 self.writer.user_header_mut().set_block_range(expected_block_start, last_block);
1036 }
1037
1038 self.commit()?;
1040
1041 Ok(())
1042 }
1043
1044 fn delete_current_and_open_previous(&mut self) -> Result<(), ProviderError> {
1047 let segment = self.user_header().segment();
1048 let current_path = self.data_path.clone();
1049 let (previous_writer, data_path) = Self::open(
1050 segment,
1051 self.writer.user_header().expected_block_start() - 1,
1052 self.reader.clone(),
1053 self.metrics.clone(),
1054 )?;
1055 self.writer = previous_writer;
1056 self.writer.set_dirty();
1057 self.data_path = data_path.clone();
1058
1059 if segment.is_change_based() {
1061 let csoff_path = current_path.with_extension("csoff");
1062 if csoff_path.exists() {
1063 std::fs::remove_file(&csoff_path).map_err(ProviderError::other)?;
1064 }
1065 let new_csoff_path = data_path.with_extension("csoff");
1067 let committed_len = self.writer.user_header().changeset_offsets_len();
1068 self.changeset_offsets = Some(
1069 ChangesetOffsetWriter::new(&new_csoff_path, committed_len)
1070 .map_err(ProviderError::other)?,
1071 );
1072 }
1073
1074 self.current_changeset_offset = None;
1076
1077 NippyJar::<SegmentHeader>::load(¤t_path)
1078 .map_err(ProviderError::other)?
1079 .delete()
1080 .map_err(ProviderError::other)?;
1081 Ok(())
1082 }
1083
1084 fn append_column<T: Compact>(&mut self, column: T) -> ProviderResult<()> {
1086 self.buf.clear();
1087 column.to_compact(&mut self.buf);
1088
1089 self.writer.append_column(Some(Ok(&self.buf))).map_err(ProviderError::other)?;
1090 Ok(())
1091 }
1092
1093 fn append_with_tx_number<V: Compact>(
1095 &mut self,
1096 tx_num: TxNumber,
1097 value: V,
1098 ) -> ProviderResult<()> {
1099 if let Some(range) = self.writer.user_header().tx_range() {
1100 let next_tx = range.end() + 1;
1101 if next_tx != tx_num {
1102 return Err(ProviderError::UnexpectedStaticFileTxNumber(
1103 self.writer.user_header().segment(),
1104 tx_num,
1105 next_tx,
1106 ))
1107 }
1108 self.writer.user_header_mut().increment_tx();
1109 } else {
1110 self.writer.user_header_mut().set_tx_range(tx_num, tx_num);
1111 }
1112
1113 self.append_column(value)?;
1114
1115 Ok(())
1116 }
1117
1118 fn append_change<V: Compact>(&mut self, change: &V) -> ProviderResult<()> {
1120 if let Some(ref mut offset) = self.current_changeset_offset {
1121 offset.increment_num_changes();
1122 }
1123 self.append_column(change)?;
1124 Ok(())
1125 }
1126
1127 pub fn append_header(&mut self, header: &N::BlockHeader, hash: &BlockHash) -> ProviderResult<()>
1132 where
1133 N::BlockHeader: Compact,
1134 {
1135 self.append_header_with_td(header, U256::ZERO, hash)
1136 }
1137
1138 pub fn append_header_with_td(
1143 &mut self,
1144 header: &N::BlockHeader,
1145 total_difficulty: U256,
1146 hash: &BlockHash,
1147 ) -> ProviderResult<()>
1148 where
1149 N::BlockHeader: Compact,
1150 {
1151 let start = Instant::now();
1152 self.ensure_no_queued_prune()?;
1153
1154 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::Headers);
1155
1156 self.increment_block(header.number())?;
1157
1158 self.append_column(header)?;
1159 self.append_column(CompactU256::from(total_difficulty))?;
1160 self.append_column(hash)?;
1161
1162 if let Some(metrics) = &self.metrics {
1163 metrics.record_segment_operation(
1164 StaticFileSegment::Headers,
1165 StaticFileProviderOperation::Append,
1166 Some(start.elapsed()),
1167 );
1168 }
1169
1170 Ok(())
1171 }
1172
1173 pub fn append_header_direct(
1176 &mut self,
1177 header: &N::BlockHeader,
1178 total_difficulty: U256,
1179 hash: &BlockHash,
1180 ) -> ProviderResult<()>
1181 where
1182 N::BlockHeader: Compact,
1183 {
1184 let start = Instant::now();
1185 self.ensure_no_queued_prune()?;
1186
1187 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::Headers);
1188
1189 self.append_column(header)?;
1190 self.append_column(CompactU256::from(total_difficulty))?;
1191 self.append_column(hash)?;
1192
1193 if let Some(metrics) = &self.metrics {
1194 metrics.record_segment_operation(
1195 StaticFileSegment::Headers,
1196 StaticFileProviderOperation::Append,
1197 Some(start.elapsed()),
1198 );
1199 }
1200
1201 Ok(())
1202 }
1203
1204 pub fn append_transaction(&mut self, tx_num: TxNumber, tx: &N::SignedTx) -> ProviderResult<()>
1209 where
1210 N::SignedTx: Compact,
1211 {
1212 let start = Instant::now();
1213 self.ensure_no_queued_prune()?;
1214
1215 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::Transactions);
1216 self.append_with_tx_number(tx_num, tx)?;
1217
1218 if let Some(metrics) = &self.metrics {
1219 metrics.record_segment_operation(
1220 StaticFileSegment::Transactions,
1221 StaticFileProviderOperation::Append,
1222 Some(start.elapsed()),
1223 );
1224 }
1225
1226 Ok(())
1227 }
1228
1229 pub fn append_receipt(&mut self, tx_num: TxNumber, receipt: &N::Receipt) -> ProviderResult<()>
1234 where
1235 N::Receipt: Compact,
1236 {
1237 let start = Instant::now();
1238 self.ensure_no_queued_prune()?;
1239
1240 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::Receipts);
1241 self.append_with_tx_number(tx_num, receipt)?;
1242
1243 if let Some(metrics) = &self.metrics {
1244 metrics.record_segment_operation(
1245 StaticFileSegment::Receipts,
1246 StaticFileProviderOperation::Append,
1247 Some(start.elapsed()),
1248 );
1249 }
1250
1251 Ok(())
1252 }
1253
1254 pub fn append_receipts<I, R>(&mut self, receipts: I) -> ProviderResult<()>
1256 where
1257 I: Iterator<Item = Result<(TxNumber, R), ProviderError>>,
1258 R: Borrow<N::Receipt>,
1259 N::Receipt: Compact,
1260 {
1261 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::Receipts);
1262
1263 let mut receipts_iter = receipts.into_iter().peekable();
1264 if receipts_iter.peek().is_none() {
1266 return Ok(());
1267 }
1268
1269 let start = Instant::now();
1270 self.ensure_no_queued_prune()?;
1271
1272 let mut count: u64 = 0;
1274
1275 for receipt_result in receipts_iter {
1276 let (tx_num, receipt) = receipt_result?;
1277 self.append_with_tx_number(tx_num, receipt.borrow())?;
1278 count += 1;
1279 }
1280
1281 if let Some(metrics) = &self.metrics {
1282 metrics.record_segment_operations(
1283 StaticFileSegment::Receipts,
1284 StaticFileProviderOperation::Append,
1285 count,
1286 Some(start.elapsed()),
1287 );
1288 }
1289
1290 Ok(())
1291 }
1292
1293 pub fn append_transaction_sender(
1298 &mut self,
1299 tx_num: TxNumber,
1300 sender: &alloy_primitives::Address,
1301 ) -> ProviderResult<()> {
1302 let start = Instant::now();
1303 self.ensure_no_queued_prune()?;
1304
1305 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::TransactionSenders);
1306 self.append_with_tx_number(tx_num, sender)?;
1307
1308 if let Some(metrics) = &self.metrics {
1309 metrics.record_segment_operation(
1310 StaticFileSegment::TransactionSenders,
1311 StaticFileProviderOperation::Append,
1312 Some(start.elapsed()),
1313 );
1314 }
1315
1316 Ok(())
1317 }
1318
1319 pub fn append_transaction_senders<I>(&mut self, senders: I) -> ProviderResult<()>
1321 where
1322 I: Iterator<Item = (TxNumber, alloy_primitives::Address)>,
1323 {
1324 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::TransactionSenders);
1325
1326 let mut senders_iter = senders.into_iter().peekable();
1327 if senders_iter.peek().is_none() {
1329 return Ok(());
1330 }
1331
1332 let start = Instant::now();
1333 self.ensure_no_queued_prune()?;
1334
1335 let mut count: u64 = 0;
1337 for (tx_num, sender) in senders_iter {
1338 self.append_with_tx_number(tx_num, sender)?;
1339 count += 1;
1340 }
1341
1342 if let Some(metrics) = &self.metrics {
1343 metrics.record_segment_operations(
1344 StaticFileSegment::TransactionSenders,
1345 StaticFileProviderOperation::Append,
1346 count,
1347 Some(start.elapsed()),
1348 );
1349 }
1350
1351 Ok(())
1352 }
1353
1354 pub fn append_account_changeset(
1360 &mut self,
1361 mut changeset: Vec<AccountBeforeTx>,
1362 block_number: u64,
1363 ) -> ProviderResult<()> {
1364 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::AccountChangeSets);
1365 let start = Instant::now();
1366
1367 self.increment_block(block_number)?;
1368 self.ensure_no_queued_prune()?;
1369
1370 changeset.sort_by_key(|change| change.address);
1372
1373 let mut count: u64 = 0;
1374
1375 for change in changeset {
1376 self.append_change(&change)?;
1377 count += 1;
1378 }
1379
1380 if let Some(metrics) = &self.metrics {
1381 metrics.record_segment_operations(
1382 StaticFileSegment::AccountChangeSets,
1383 StaticFileProviderOperation::Append,
1384 count,
1385 Some(start.elapsed()),
1386 );
1387 }
1388
1389 Ok(())
1390 }
1391
1392 pub fn begin_account_changeset(&mut self, block_number: u64) -> ProviderResult<()> {
1397 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::AccountChangeSets);
1398
1399 self.increment_block(block_number)?;
1400 self.ensure_no_queued_prune()
1401 }
1402
1403 pub fn append_account_changeset_entry(
1407 &mut self,
1408 change: AccountBeforeTx,
1409 ) -> ProviderResult<()> {
1410 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::AccountChangeSets);
1411 if self.current_changeset_offset.is_none() {
1412 return Err(ProviderError::other(StaticFileWriterError::new(
1413 "account changeset stream must be started before appending entries",
1414 )))
1415 }
1416
1417 self.append_change(&change)
1418 }
1419
1420 pub fn append_storage_changeset(
1424 &mut self,
1425 mut changeset: Vec<StorageBeforeTx>,
1426 block_number: u64,
1427 ) -> ProviderResult<()> {
1428 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::StorageChangeSets);
1429 let start = Instant::now();
1430
1431 self.increment_block(block_number)?;
1432 self.ensure_no_queued_prune()?;
1433
1434 changeset.sort_by_key(|change| (change.address, change.key));
1436
1437 let mut count: u64 = 0;
1438 for change in changeset {
1439 self.append_change(&change)?;
1440 count += 1;
1441 }
1442
1443 if let Some(metrics) = &self.metrics {
1444 metrics.record_segment_operations(
1445 StaticFileSegment::StorageChangeSets,
1446 StaticFileProviderOperation::Append,
1447 count,
1448 Some(start.elapsed()),
1449 );
1450 }
1451
1452 Ok(())
1453 }
1454
1455 pub fn begin_storage_changeset(&mut self, block_number: u64) -> ProviderResult<()> {
1461 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::StorageChangeSets);
1462
1463 self.increment_block(block_number)?;
1464 self.ensure_no_queued_prune()
1465 }
1466
1467 pub fn append_storage_changeset_entry(
1471 &mut self,
1472 change: StorageBeforeTx,
1473 ) -> ProviderResult<()> {
1474 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::StorageChangeSets);
1475 if self.current_changeset_offset.is_none() {
1476 return Err(ProviderError::other(StaticFileWriterError::new(
1477 "storage changeset stream must be started before appending entries",
1478 )))
1479 }
1480
1481 self.append_change(&change)
1482 }
1483
1484 pub fn prune_transactions(
1488 &mut self,
1489 to_delete: u64,
1490 last_block: BlockNumber,
1491 ) -> ProviderResult<()> {
1492 debug_assert_eq!(self.writer.user_header().segment(), StaticFileSegment::Transactions);
1493 self.queue_prune(PruneStrategy::Transactions { num_rows: to_delete, last_block })
1494 }
1495
1496 pub fn prune_receipts(
1500 &mut self,
1501 to_delete: u64,
1502 last_block: BlockNumber,
1503 ) -> ProviderResult<()> {
1504 debug_assert_eq!(self.writer.user_header().segment(), StaticFileSegment::Receipts);
1505 self.queue_prune(PruneStrategy::Receipts { num_rows: to_delete, last_block })
1506 }
1507
1508 pub fn prune_transaction_senders(
1512 &mut self,
1513 to_delete: u64,
1514 last_block: BlockNumber,
1515 ) -> ProviderResult<()> {
1516 debug_assert_eq!(
1517 self.writer.user_header().segment(),
1518 StaticFileSegment::TransactionSenders
1519 );
1520 self.queue_prune(PruneStrategy::TransactionSenders { num_rows: to_delete, last_block })
1521 }
1522
1523 pub fn prune_headers(&mut self, to_delete: u64) -> ProviderResult<()> {
1525 debug_assert_eq!(self.writer.user_header().segment(), StaticFileSegment::Headers);
1526 self.queue_prune(PruneStrategy::Headers { num_blocks: to_delete })
1527 }
1528
1529 pub fn prune_account_changesets(&mut self, last_block: u64) -> ProviderResult<()> {
1531 debug_assert_eq!(self.writer.user_header().segment(), StaticFileSegment::AccountChangeSets);
1532 self.queue_prune(PruneStrategy::AccountChangeSets { last_block })
1533 }
1534
1535 pub fn prune_storage_changesets(&mut self, last_block: u64) -> ProviderResult<()> {
1537 debug_assert_eq!(self.writer.user_header().segment(), StaticFileSegment::StorageChangeSets);
1538 self.queue_prune(PruneStrategy::StorageChangeSets { last_block })
1539 }
1540
1541 fn queue_prune(&mut self, strategy: PruneStrategy) -> ProviderResult<()> {
1543 self.ensure_no_queued_prune()?;
1544 self.prune_on_commit = Some(strategy);
1545 Ok(())
1546 }
1547
1548 fn ensure_no_queued_prune(&self) -> ProviderResult<()> {
1550 if self.prune_on_commit.is_some() {
1551 return Err(ProviderError::other(StaticFileWriterError::new(
1552 "Pruning should be committed before appending or pruning more data",
1553 )));
1554 }
1555 Ok(())
1556 }
1557
1558 fn prune_transaction_data(
1560 &mut self,
1561 to_delete: u64,
1562 last_block: BlockNumber,
1563 ) -> ProviderResult<()> {
1564 let start = Instant::now();
1565
1566 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::Transactions);
1567
1568 self.truncate(to_delete, Some(last_block))?;
1569
1570 if let Some(metrics) = &self.metrics {
1571 metrics.record_segment_operation(
1572 StaticFileSegment::Transactions,
1573 StaticFileProviderOperation::Prune,
1574 Some(start.elapsed()),
1575 );
1576 }
1577
1578 Ok(())
1579 }
1580
1581 fn prune_account_changeset_data(&mut self, last_block: BlockNumber) -> ProviderResult<()> {
1583 let start = Instant::now();
1584
1585 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::AccountChangeSets);
1586
1587 self.truncate_changesets(last_block)?;
1588
1589 if let Some(metrics) = &self.metrics {
1590 metrics.record_segment_operation(
1591 StaticFileSegment::AccountChangeSets,
1592 StaticFileProviderOperation::Prune,
1593 Some(start.elapsed()),
1594 );
1595 }
1596
1597 Ok(())
1598 }
1599
1600 fn prune_storage_changeset_data(&mut self, last_block: BlockNumber) -> ProviderResult<()> {
1602 let start = Instant::now();
1603
1604 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::StorageChangeSets);
1605
1606 self.truncate_changesets(last_block)?;
1607
1608 if let Some(metrics) = &self.metrics {
1609 metrics.record_segment_operation(
1610 StaticFileSegment::StorageChangeSets,
1611 StaticFileProviderOperation::Prune,
1612 Some(start.elapsed()),
1613 );
1614 }
1615
1616 Ok(())
1617 }
1618
1619 fn prune_receipt_data(
1621 &mut self,
1622 to_delete: u64,
1623 last_block: BlockNumber,
1624 ) -> ProviderResult<()> {
1625 let start = Instant::now();
1626
1627 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::Receipts);
1628
1629 self.truncate(to_delete, Some(last_block))?;
1630
1631 if let Some(metrics) = &self.metrics {
1632 metrics.record_segment_operation(
1633 StaticFileSegment::Receipts,
1634 StaticFileProviderOperation::Prune,
1635 Some(start.elapsed()),
1636 );
1637 }
1638
1639 Ok(())
1640 }
1641
1642 fn prune_transaction_sender_data(
1644 &mut self,
1645 to_delete: u64,
1646 last_block: BlockNumber,
1647 ) -> ProviderResult<()> {
1648 let start = Instant::now();
1649
1650 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::TransactionSenders);
1651
1652 self.truncate(to_delete, Some(last_block))?;
1653
1654 if let Some(metrics) = &self.metrics {
1655 metrics.record_segment_operation(
1656 StaticFileSegment::TransactionSenders,
1657 StaticFileProviderOperation::Prune,
1658 Some(start.elapsed()),
1659 );
1660 }
1661
1662 Ok(())
1663 }
1664
1665 fn prune_header_data(&mut self, to_delete: u64) -> ProviderResult<()> {
1667 let start = Instant::now();
1668
1669 debug_assert!(self.writer.user_header().segment() == StaticFileSegment::Headers);
1670
1671 self.truncate(to_delete, None)?;
1672
1673 if let Some(metrics) = &self.metrics {
1674 metrics.record_segment_operation(
1675 StaticFileSegment::Headers,
1676 StaticFileProviderOperation::Prune,
1677 Some(start.elapsed()),
1678 );
1679 }
1680
1681 Ok(())
1682 }
1683
1684 pub fn reader(&self) -> StaticFileProvider<N> {
1686 Self::upgrade_provider_to_strong_reference(&self.reader)
1687 }
1688
1689 fn upgrade_provider_to_strong_reference(
1698 provider: &Weak<StaticFileProviderInner<N>>,
1699 ) -> StaticFileProvider<N> {
1700 provider.upgrade().map(StaticFileProvider).expect("StaticFileProvider is dropped")
1701 }
1702
1703 pub const fn user_header(&self) -> &SegmentHeader {
1705 self.writer.user_header()
1706 }
1707
1708 pub const fn user_header_mut(&mut self) -> &mut SegmentHeader {
1710 self.writer.user_header_mut()
1711 }
1712
1713 #[cfg(any(test, feature = "test-utils"))]
1715 pub const fn set_block_range(&mut self, block_range: std::ops::RangeInclusive<BlockNumber>) {
1716 self.writer.user_header_mut().set_block_range(*block_range.start(), *block_range.end())
1717 }
1718
1719 #[cfg(any(test, feature = "test-utils"))]
1721 pub const fn inner(&mut self) -> &mut NippyJarWriter<SegmentHeader> {
1722 &mut self.writer
1723 }
1724}
1725
1726fn create_jar(
1727 segment: StaticFileSegment,
1728 path: &Path,
1729 expected_block_range: SegmentRangeInclusive,
1730) -> NippyJar<SegmentHeader> {
1731 let mut jar = NippyJar::new(
1732 segment.columns(),
1733 path,
1734 SegmentHeader::new(expected_block_range, None, None, segment),
1735 );
1736
1737 if segment.is_headers() {
1740 jar = jar.with_lz4();
1741 }
1742
1743 jar
1744}