1use alloy_consensus::{
2 proofs::calculate_receipt_root, BlockHeader, ReceiptEnvelope, ReceiptWithBloom,
3 RlpDecodableReceipt, TxReceipt,
4};
5use alloy_primitives::{BlockHash, BlockNumber, Bloom, U256};
6use alloy_rlp::Decodable;
7use futures_util::{Stream, StreamExt};
8use reth_codecs::Compact;
9use reth_db_api::{
10 cursor::{DbCursorRO, DbCursorRW},
11 table::Value,
12 tables,
13 transaction::{DbTx, DbTxMut},
14 RawKey, RawTable, RawValue,
15};
16use reth_era::{
17 common::{decode::DecodeCompressedRlp, file_ops::StreamReader},
18 e2s::error::E2sError,
19 era::{file::EraReader, types::consensus::CompressedSignedBeaconBlock},
20 era1::{file::Era1Reader, types::execution::BlockTuple},
21 ere::{file::EreReader, types::execution::BlockTuple as EreBlockTuple},
22};
23use reth_era_downloader::EraMeta;
24use reth_etl::Collector;
25use reth_fs_util as fs;
26use reth_primitives_traits::{
27 Block, BlockBody, FullBlockBody, FullBlockHeader, NodePrimitives, Receipt,
28};
29use reth_provider::{
30 providers::StaticFileProviderRWRefMut, BlockReader, BlockWriter, EitherWriter,
31 EitherWriterDestination, StaticFileProviderFactory, StaticFileSegment, StaticFileWriter,
32};
33use reth_stages_types::{
34 CheckpointBlockRange, EntitiesCheckpoint, HeadersCheckpoint, StageCheckpoint, StageId,
35};
36use reth_storage_api::{
37 errors::{ProviderError, ProviderResult},
38 BlockBodyIndicesProvider, BlockHashReader, DBProvider, DatabaseProviderFactory,
39 NodePrimitivesProvider, StageCheckpointReader, StageCheckpointWriter, StorageSettingsCache,
40};
41use std::{collections::Bound, error::Error, ops::RangeBounds, sync::mpsc};
42use tracing::info;
43
44type EraBlock<BH, BB, R> = (BH, BB, Option<Vec<R>>);
46
47type ReceiptOf<P> = <<P as NodePrimitivesProvider>::Primitives as NodePrimitives>::Receipt;
49
50pub trait EraBlockReader<BH, BB, R> {
57 fn blocks<M: EraMeta + ?Sized>(
59 meta: &M,
60 decode_receipts: bool,
61 ) -> eyre::Result<impl Iterator<Item = eyre::Result<EraBlock<BH, BB, R>>>>;
62}
63
64#[derive(Debug)]
66pub struct Era1;
67
68impl<BH, BB, R> EraBlockReader<BH, BB, R> for Era1
69where
70 BH: FullBlockHeader + Value,
71 BB: FullBlockBody<OmmerHeader = BH>,
72 R: RlpDecodableReceipt,
73{
74 fn blocks<M: EraMeta + ?Sized>(
75 meta: &M,
76 decode_receipts: bool,
77 ) -> eyre::Result<impl Iterator<Item = eyre::Result<EraBlock<BH, BB, R>>>> {
78 let reader: Era1Reader<std::fs::File> = open(meta)?;
79 Ok(reader
80 .iter()
81 .map(move |block| decode_with_receipts::<BH, BB, R, E2sError>(block, decode_receipts)))
82 }
83}
84
85impl<BH, BB, R> EraBlockReader<BH, BB, R> for Ere
86where
87 BH: FullBlockHeader + Value,
88 BB: FullBlockBody<OmmerHeader = BH>,
89 R: RlpDecodableReceipt,
90{
91 fn blocks<M: EraMeta + ?Sized>(
92 meta: &M,
93 decode_receipts: bool,
94 ) -> eyre::Result<impl Iterator<Item = eyre::Result<EraBlock<BH, BB, R>>>> {
95 let reader: EreReader<std::fs::File> = open(meta)?;
96 Ok(reader.iter().map(move |block| Self::decode_with_receipts(block, decode_receipts)))
97 }
98}
99
100#[derive(Debug)]
102pub struct Ere;
103
104impl Ere {
105 pub fn decode<BH, BB, E>(block: Result<EreBlockTuple, E>) -> eyre::Result<(BH, BB)>
109 where
110 BH: FullBlockHeader + Value,
111 BB: FullBlockBody<OmmerHeader = BH>,
112 E: From<E2sError> + Error + Send + Sync + 'static,
113 {
114 let block = block?;
115 let header: BH = block.header.decode()?;
116 let body: BB = block.body.decode()?;
117 Ok((header, body))
118 }
119
120 pub fn decode_with_receipts<BH, BB, R, E>(
124 block: Result<EreBlockTuple, E>,
125 decode_receipts: bool,
126 ) -> eyre::Result<EraBlock<BH, BB, R>>
127 where
128 BH: FullBlockHeader + Value,
129 BB: FullBlockBody<OmmerHeader = BH>,
130 R: RlpDecodableReceipt,
131 E: From<E2sError> + Error + Send + Sync + 'static,
132 {
133 let block = block?;
134 let header: BH = block.header.decode()?;
135 let body: BB = block.body.decode()?;
136 let number = header.number();
137 let receipts = decode_receipts
138 .then(|| block.receipts.as_ref().map(|r| r.decode_receipts()))
139 .flatten()
140 .transpose()?
141 .map(|slim| receipts_from_envelopes(number, slim.into_iter().map(Into::into).collect()))
142 .transpose()?;
143
144 Ok((header, body, receipts))
145 }
146}
147
148#[derive(Debug)]
156pub struct Era;
157
158impl<BH, BB, R> EraBlockReader<BH, BB, R> for Era
159where
160 BH: FullBlockHeader,
161 BB: FullBlockBody,
162{
163 fn blocks<M: EraMeta + ?Sized>(
165 meta: &M,
166 _decode_receipts: bool,
167 ) -> eyre::Result<impl Iterator<Item = eyre::Result<EraBlock<BH, BB, R>>>> {
168 let reader: EraReader<std::fs::File> = open(meta)?;
169 let mut buf = Vec::new();
170 Ok(reader.iter().filter_map(move |block| {
171 Self::decode(block, &mut buf)
172 .map(|opt| opt.map(|(header, body)| (header, body, None)))
173 .transpose()
174 }))
175 }
176}
177
178impl Era {
179 pub fn decode<BH, BB>(
183 block: Result<CompressedSignedBeaconBlock, E2sError>,
184 buf: &mut Vec<u8>,
185 ) -> eyre::Result<Option<(BH, BB)>>
186 where
187 BH: FullBlockHeader,
188 BB: FullBlockBody,
189 {
190 let Some(alloy_consensus::Block { header, body }) =
191 block?.decode_execution_block::<<BB as BlockBody>::Transaction>()?
192 else {
193 return Ok(None);
194 };
195 Ok(Some((reencode_rlp(&header, buf)?, reencode_rlp(&body, buf)?)))
198 }
199}
200
201fn reencode_rlp<T, U>(value: &T, buf: &mut Vec<u8>) -> eyre::Result<U>
207where
208 T: alloy_rlp::Encodable,
209 U: alloy_rlp::Decodable,
210{
211 buf.clear();
212 alloy_rlp::Encodable::encode(value, buf);
213 Ok(<U as alloy_rlp::Decodable>::decode(&mut buf.as_slice())?)
214}
215
216pub fn open<Reader>(meta: &(impl EraMeta + ?Sized)) -> eyre::Result<Reader>
218where
219 Reader: StreamReader<std::fs::File>,
220{
221 Ok(Reader::new(fs::open(meta.path())?))
222}
223
224pub fn import<S, Downloader, Era, PF, B, BB, BH>(
239 mut downloader: Downloader,
240 provider_factory: &PF,
241 hash_collector: &mut Collector<BlockHash, BlockNumber>,
242 to_block: Option<BlockNumber>,
243 store_receipts: bool,
244 is_receipt_verifiable: &dyn Fn(BlockNumber) -> bool,
245) -> eyre::Result<BlockNumber>
246where
247 S: EraBlockReader<BH, BB, ReceiptOf<<PF as DatabaseProviderFactory>::ProviderRW>>,
248 B: Block<Header = BH, Body = BB>,
249 BH: FullBlockHeader + Value,
250 BB: FullBlockBody<
251 Transaction = <<<PF as DatabaseProviderFactory>::ProviderRW as NodePrimitivesProvider>::Primitives as NodePrimitives>::SignedTx,
252 OmmerHeader = BH,
253 >,
254 Downloader: Stream<Item = eyre::Result<Era>> + Send + 'static + Unpin,
255 Era: EraMeta + Send + 'static,
256 PF: DatabaseProviderFactory<
257 ProviderRW: BlockWriter<Block = B>
258 + DBProvider
259 + BlockBodyIndicesProvider
260 + BlockHashReader
261 + StaticFileProviderFactory<Primitives: NodePrimitives<Block = B, BlockHeader = BH, BlockBody = BB>>
262 + StageCheckpointReader
263 + StageCheckpointWriter
264 + StorageSettingsCache,
265 > + StaticFileProviderFactory<Primitives = <<PF as DatabaseProviderFactory>::ProviderRW as NodePrimitivesProvider>::Primitives>,
266 ReceiptOf<<PF as DatabaseProviderFactory>::ProviderRW>: Compact + Receipt,
267{
268 let (tx, rx) = mpsc::channel();
269
270 tokio::spawn(async move {
272 while let Some(file) = downloader.next().await {
273 tx.send(Some(file))?;
274 }
275 tx.send(None)
276 });
277
278 let static_file_provider = provider_factory.static_file_provider();
279
280 let genesis_block_number = static_file_provider.genesis_block_number();
282
283 let headers_tip = static_file_provider
286 .get_highest_static_file_block(StaticFileSegment::Headers)
287 .unwrap_or(genesis_block_number);
288
289 let receipts_tip = store_receipts
292 .then(|| static_file_provider.get_highest_static_file_block(StaticFileSegment::Receipts));
293 let mut height = match receipts_tip {
294 Some(tip) => headers_tip.min(tip.unwrap_or(genesis_block_number)),
295 None => headers_tip,
296 };
297
298 let receipts_target = store_receipts
301 .then(|| -> eyre::Result<_> {
302 let provider = provider_factory.database_provider_rw()?;
303
304 if !matches!(
307 EitherWriter::receipts_destination(&provider),
308 EitherWriterDestination::StaticFile
309 ) {
310 eyre::bail!(
311 "receipt import writes the Receipts static file segment, but this node's \
312 prune configuration keeps receipts elsewhere. Remove the receipt pruning \
313 configuration, or import without receipts"
314 );
315 }
316
317 let target = provider
318 .get_stage_checkpoint(StageId::Execution)?
319 .map(|checkpoint| checkpoint.block_number)
320 .unwrap_or(genesis_block_number);
321
322 if target <= genesis_block_number {
323 eyre::bail!(
324 "receipt import repairs the receipt static files of an already-executed \
325 range, but this database has executed no blocks. Sync the node first, or \
326 import without receipts"
327 );
328 }
329 if height >= target {
330 eyre::bail!(
331 "receipts already cover the executed range up to block {target}, nothing to \
332 repair"
333 );
334 }
335 if let Some(to_block) = to_block &&
336 to_block != target
337 {
338 eyre::bail!(
339 "--to-block {to_block} does not match the Execution checkpoint {target}. A \
340 receipt repair must cover the executed range exactly, so either drop \
341 --to-block or set it to {target}"
342 );
343 }
344 if !is_receipt_verifiable(height + 1) {
347 eyre::bail!(
348 "receipt repair would start at block {}, which predates Byzantium. Those \
349 receipts commit to a post-state root this node's receipt type cannot \
350 represent, so the Receipts segment has to already cover the pre-Byzantium \
351 range",
352 height + 1,
353 );
354 }
355
356 Ok(target)
357 })
358 .transpose()?;
359
360 if matches!(receipts_tip, Some(None)) {
365 let mut writer = static_file_provider.latest_writer(StaticFileSegment::Receipts)?;
366 writer.user_header_mut().set_block_range(genesis_block_number, genesis_block_number);
367 writer.commit()?;
368 }
369
370 let to_block = receipts_target.or(to_block);
371 let end = to_block.map_or(Bound::Unbounded, Bound::Included);
372 let policy = ImportPolicy { headers_tip, is_receipt_verifiable };
373
374 while let Some(meta) = rx.recv()? {
375 let meta = meta?;
376 let from = height;
377 let provider = provider_factory.database_provider_rw()?;
378
379 let mut receipts_writer = store_receipts
382 .then(|| static_file_provider.latest_writer(StaticFileSegment::Receipts))
383 .transpose()?;
384
385 height = process::<S, _, _, _, _>(
386 &meta,
387 &mut static_file_provider.latest_writer(StaticFileSegment::Headers)?,
388 receipts_writer.as_mut(),
389 &provider,
390 hash_collector,
391 (Bound::Included(height), end),
392 policy,
393 )?;
394
395 drop(receipts_writer);
398
399 let checkpoint = height.max(headers_tip);
402 save_stage_checkpoints(&provider, from, checkpoint, checkpoint, checkpoint)?;
403
404 provider.commit()?;
405
406 info!(target: "era::history::import", first = from, last = height, file = %meta.path().display(), "Imported ERA file");
407
408 if to_block.is_some_and(|to| height >= to) {
409 break;
410 }
411 }
412
413 let provider = provider_factory.database_provider_rw()?;
414
415 build_index(&provider, hash_collector)?;
416
417 provider.commit()?;
418
419 if let Some(target) = receipts_target &&
422 height < target
423 {
424 eyre::bail!(
425 "receipt repair reached block {height} but the executed range ends at {target}. The \
426 source ran out of files, re-run with the files covering blocks {}..={target}",
427 height + 1,
428 );
429 }
430
431 Ok(height)
432}
433
434pub fn save_stage_checkpoints<P>(
439 provider: P,
440 from: BlockNumber,
441 to: BlockNumber,
442 processed: u64,
443 total: u64,
444) -> ProviderResult<()>
445where
446 P: StageCheckpointWriter,
447{
448 provider.save_stage_checkpoint(
449 StageId::Headers,
450 StageCheckpoint::new(to).with_headers_stage_checkpoint(HeadersCheckpoint {
451 block_range: CheckpointBlockRange { from, to },
452 progress: EntitiesCheckpoint { processed, total },
453 }),
454 )?;
455 provider.save_stage_checkpoint(
456 StageId::Bodies,
457 StageCheckpoint::new(to)
458 .with_entities_stage_checkpoint(EntitiesCheckpoint { processed, total }),
459 )?;
460 Ok(())
461}
462
463pub fn process<S, P, B, BB, BH>(
466 meta: &(impl EraMeta + ?Sized),
467 writer: &mut StaticFileProviderRWRefMut<'_, <P as NodePrimitivesProvider>::Primitives>,
468 receipts_writer: Option<
469 &mut StaticFileProviderRWRefMut<'_, <P as NodePrimitivesProvider>::Primitives>,
470 >,
471 provider: &P,
472 hash_collector: &mut Collector<BlockHash, BlockNumber>,
473 block_numbers: impl RangeBounds<BlockNumber>,
474 policy: ImportPolicy<'_>,
475) -> eyre::Result<BlockNumber>
476where
477 S: EraBlockReader<BH, BB, ReceiptOf<P>>,
478 B: Block<Header = BH, Body = BB>,
479 BH: FullBlockHeader + Value,
480 BB: FullBlockBody<
481 Transaction = <<P as NodePrimitivesProvider>::Primitives as NodePrimitives>::SignedTx,
482 OmmerHeader = BH,
483 >,
484 P: DBProvider<Tx: DbTxMut>
485 + NodePrimitivesProvider
486 + BlockWriter<Block = B>
487 + BlockBodyIndicesProvider
488 + BlockHashReader,
489 <P as NodePrimitivesProvider>::Primitives: NodePrimitives<BlockHeader = BH, BlockBody = BB>,
490 ReceiptOf<P>: Compact + Receipt,
491{
492 let decode_receipts = receipts_writer.is_some();
493 let iter = S::blocks(meta, decode_receipts)?
494 .map(Some)
495 .chain(std::iter::once_with(|| match meta.mark_as_processed() {
496 Ok(()) => None,
497 Err(error) => Some(Err(error)),
498 }))
499 .flatten();
500
501 process_iter(iter, writer, receipts_writer, provider, hash_collector, block_numbers, policy)
502}
503
504#[derive(Clone, Copy)]
506pub struct ImportPolicy<'a> {
507 pub headers_tip: BlockNumber,
509 pub is_receipt_verifiable: &'a dyn Fn(BlockNumber) -> bool,
511}
512
513impl std::fmt::Debug for ImportPolicy<'_> {
514 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
515 f.debug_struct("ImportPolicy")
517 .field("headers_tip", &self.headers_tip)
518 .finish_non_exhaustive()
519 }
520}
521
522pub fn decode<BH, BB, E>(block: Result<BlockTuple, E>) -> eyre::Result<(BH, BB)>
524where
525 BH: FullBlockHeader + Value,
526 BB: FullBlockBody<OmmerHeader = BH>,
527 E: From<E2sError> + Error + Send + Sync + 'static,
528{
529 let block = block?;
530 let header: BH = block.header.decode()?;
531 let body: BB = block.body.decode()?;
532 Ok((header, body))
533}
534
535pub fn decode_with_receipts<BH, BB, R, E>(
539 block: Result<BlockTuple, E>,
540 decode_receipts: bool,
541) -> eyre::Result<EraBlock<BH, BB, R>>
542where
543 BH: FullBlockHeader + Value,
544 BB: FullBlockBody<OmmerHeader = BH>,
545 R: RlpDecodableReceipt,
546 E: From<E2sError> + Error + Send + Sync + 'static,
547{
548 let block = block?;
549 let header: BH = block.header.decode()?;
550 let body: BB = block.body.decode()?;
551 let number = header.number();
552 let receipts = decode_receipts
553 .then(|| -> eyre::Result<_> {
554 match block.receipts.decode::<Vec<ReceiptWithBloom<R>>>() {
555 Ok(receipts) => {
556 Ok(receipts.into_iter().map(|with_bloom| with_bloom.receipt).collect())
557 }
558 Err(err) => match block.receipts.decode::<Vec<ReceiptEnvelope>>() {
561 Ok(envelopes) => receipts_from_envelopes(number, envelopes),
562 Err(_) => Err(err.into()),
563 },
564 }
565 })
566 .transpose()?;
567
568 Ok((header, body, receipts))
569}
570
571fn receipts_from_envelopes<R: RlpDecodableReceipt>(
578 number: BlockNumber,
579 envelopes: Vec<ReceiptEnvelope>,
580) -> eyre::Result<Vec<R>> {
581 for envelope in &envelopes {
582 if envelope.status_or_post_state().is_post_state() {
583 eyre::bail!(
584 "block {number} has pre-Byzantium receipts, which commit to a post-state root \
585 rather than a success status and so cannot be represented by this node's receipt \
586 type. Receipt import is only supported from Byzantium onwards"
587 );
588 }
589 }
590
591 let encoded = alloy_rlp::encode(&envelopes);
592
593 Ok(Vec::<ReceiptWithBloom<R>>::decode(&mut encoded.as_slice())?
594 .into_iter()
595 .map(|with_bloom| with_bloom.receipt)
596 .collect())
597}
598
599pub fn process_iter<P, B, BB, BH>(
620 mut iter: impl Iterator<Item = eyre::Result<EraBlock<BH, BB, ReceiptOf<P>>>>,
621 writer: &mut StaticFileProviderRWRefMut<'_, <P as NodePrimitivesProvider>::Primitives>,
622 mut receipts_writer: Option<
623 &mut StaticFileProviderRWRefMut<'_, <P as NodePrimitivesProvider>::Primitives>,
624 >,
625 provider: &P,
626 hash_collector: &mut Collector<BlockHash, BlockNumber>,
627 block_numbers: impl RangeBounds<BlockNumber>,
628 policy: ImportPolicy<'_>,
629) -> eyre::Result<BlockNumber>
630where
631 B: Block<Header = BH, Body = BB>,
632 BH: FullBlockHeader + Value,
633 BB: FullBlockBody<
634 Transaction = <<P as NodePrimitivesProvider>::Primitives as NodePrimitives>::SignedTx,
635 OmmerHeader = BH,
636 >,
637 P: DBProvider<Tx: DbTxMut>
638 + NodePrimitivesProvider
639 + BlockWriter<Block = B>
640 + BlockBodyIndicesProvider
641 + BlockHashReader,
642 <P as NodePrimitivesProvider>::Primitives: NodePrimitives<BlockHeader = BH, BlockBody = BB>,
643 ReceiptOf<P>: Compact + Receipt,
644{
645 let mut last_header_number = match block_numbers.start_bound() {
646 Bound::Included(&number) => number,
647 Bound::Excluded(&number) => number.saturating_add(1),
648 Bound::Unbounded => 0,
649 };
650 let target = match block_numbers.end_bound() {
651 Bound::Included(&number) => Some(number),
652 Bound::Excluded(&number) => Some(number.saturating_sub(1)),
653 Bound::Unbounded => None,
654 };
655
656 for block in &mut iter {
657 let (header, body, receipts) = block?;
658 let number = header.number();
659
660 if number <= last_header_number {
661 continue;
662 }
663 if let Some(target) = target &&
664 number > target
665 {
666 break;
667 }
668
669 if number != last_header_number + 1 {
672 eyre::bail!(
673 "non-contiguous ERA import: expected block {}, got {number}; the execution \
674 database must be synced up to block {} before importing this file",
675 last_header_number + 1,
676 number - 1,
677 );
678 }
679
680 last_header_number = number;
681
682 if number > policy.headers_tip {
685 let hash = header.hash_slow();
686 writer.append_header(&header, &hash)?;
687 provider.append_block_bodies(vec![(header.number(), Some(&body))])?;
688 hash_collector.insert(hash, number)?;
689 } else if provider.block_hash(number)? != Some(header.hash_slow()) {
690 eyre::bail!(
693 "block {number} in this ERA file does not match the block already imported at \
694 that height"
695 );
696 }
697
698 if let Some(receipts_writer) = receipts_writer.as_deref_mut() {
699 if let Some(receipts) = receipts.as_deref() {
700 verify_receipts(&header, receipts, (policy.is_receipt_verifiable)(number))?;
701 }
702 provider.write_block_receipts(receipts_writer, number, receipts)?;
703 }
704 }
705
706 Ok(last_header_number)
707}
708
709fn verify_receipts<BH, R>(header: &BH, receipts: &[R], check_root: bool) -> eyre::Result<()>
714where
715 BH: FullBlockHeader,
716 R: Receipt,
717{
718 let with_bloom = receipts.iter().map(TxReceipt::with_bloom_ref).collect::<Vec<_>>();
719 let logs_bloom = with_bloom.iter().fold(Bloom::ZERO, |bloom, r| bloom | r.bloom_ref());
720
721 if logs_bloom != header.logs_bloom() {
722 eyre::bail!("logs bloom mismatch for block {}", header.number());
723 }
724
725 if check_root {
726 let receipts_root = calculate_receipt_root(&with_bloom);
727 if receipts_root != header.receipts_root() {
728 eyre::bail!(
729 "receipts root mismatch for block {}: computed {receipts_root}, header has {}",
730 header.number(),
731 header.receipts_root(),
732 );
733 }
734 }
735
736 Ok(())
737}
738
739trait BlockReceiptsWriterExt: BlockBodyIndicesProvider {
742 fn write_block_receipts<N: NodePrimitives>(
747 &self,
748 receipts_writer: &mut StaticFileProviderRWRefMut<'_, N>,
749 number: BlockNumber,
750 receipts: Option<Vec<N::Receipt>>,
751 ) -> eyre::Result<()>
752 where
753 N::Receipt: Compact;
754}
755
756impl<P: BlockBodyIndicesProvider> BlockReceiptsWriterExt for P {
757 fn write_block_receipts<N: NodePrimitives>(
758 &self,
759 receipts_writer: &mut StaticFileProviderRWRefMut<'_, N>,
760 number: BlockNumber,
761 receipts: Option<Vec<N::Receipt>>,
762 ) -> eyre::Result<()>
763 where
764 N::Receipt: Compact,
765 {
766 let Some(block_receipts) = receipts else {
767 eyre::bail!(
768 "block {number} has no receipts in the imported ERA file; drop --with-receipts, \
769 or import files that carry receipts for every block (`.era1` always includes \
770 them; `.ere` receipts are optional per spec)"
771 );
772 };
773
774 let indices = self
775 .block_body_indices(number)?
776 .ok_or_else(|| eyre::eyre!("missing block body indices for block {number}"))?;
777
778 if block_receipts.len() as u64 != indices.tx_count {
779 eyre::bail!(
780 "receipt count mismatch for block {number}: {} receipt(s) for {} transaction(s)",
781 block_receipts.len(),
782 indices.tx_count,
783 );
784 }
785
786 receipts_writer.increment_block(number)?;
787 receipts_writer.append_receipts(
788 (indices.first_tx_num..).zip(block_receipts.iter()).map(Ok::<_, ProviderError>),
789 )?;
790
791 Ok(())
792 }
793}
794
795pub fn build_index<P>(
797 provider: &P,
798 hash_collector: &mut Collector<BlockHash, BlockNumber>,
799) -> eyre::Result<()>
800where
801 P: DBProvider<Tx: DbTxMut>,
802{
803 let total_headers = hash_collector.len();
804 info!(target: "era::history::import", total = total_headers, "Writing headers hash index");
805
806 let mut cursor_header_numbers =
808 provider.tx_ref().cursor_write::<RawTable<tables::HeaderNumbers>>()?;
809 let first_sync = if provider.tx_ref().entries::<RawTable<tables::HeaderNumbers>>()? == 1 &&
812 let Some((hash, block_number)) = cursor_header_numbers.last()? &&
813 block_number.value()? == 0
814 {
815 hash_collector.insert(hash.key()?, 0)?;
816 cursor_header_numbers.delete_current()?;
817 true
818 } else {
819 false
820 };
821
822 let interval = (total_headers / 10).max(8192);
823
824 for (index, hash_to_number) in hash_collector.iter()?.enumerate() {
826 let (hash, number) = hash_to_number?;
827
828 if index != 0 && index.is_multiple_of(interval) {
829 info!(target: "era::history::import", progress = %format_args!("{:.2}%", (index as f64 / total_headers as f64) * 100.0), "Writing headers hash index");
830 }
831
832 let hash = RawKey::<BlockHash>::from_vec(hash);
833 let number = RawValue::<BlockNumber>::from_vec(number);
834
835 if first_sync {
836 cursor_header_numbers.append(hash, &number)?;
837 } else {
838 cursor_header_numbers.upsert(hash, &number)?;
839 }
840 }
841
842 Ok(())
843}
844
845pub fn calculate_td_by_number<P>(provider: &P, num: BlockNumber) -> eyre::Result<U256>
852where
853 P: BlockReader,
854{
855 let mut total_difficulty = U256::ZERO;
856 let mut start = 0;
857
858 while start <= num {
859 let end = (start + 1000 - 1).min(num);
860
861 total_difficulty +=
862 provider.headers_range(start..=end)?.iter().map(|h| h.difficulty()).sum::<U256>();
863
864 start = end + 1;
865 }
866
867 Ok(total_difficulty)
868}
869
870#[cfg(test)]
871mod tests {
872 use super::*;
873 use alloy_consensus::{
874 Eip658Value, Header, Receipt as RlpReceipt, ReceiptWithBloom, TxLegacy, TxType,
875 };
876 use alloy_primitives::{Address, Bytes, Log, Signature, B256};
877 use reth_db_common::init::init_genesis;
878 use reth_era::{
879 era1::types::execution::{
880 CompressedBody, CompressedHeader, CompressedReceipts, TotalDifficulty,
881 },
882 ere::types::execution::{
883 CompressedBody as EreCompressedBody, CompressedHeader as EreCompressedHeader,
884 CompressedSlimReceipts, SlimReceipt,
885 },
886 };
887 use reth_ethereum_primitives::{Block, BlockBody, Receipt, TransactionSigned};
888 use reth_provider::{
889 test_utils::{
890 create_test_provider_factory, create_test_provider_factory_with_genesis_block_number,
891 },
892 DatabaseProviderFactory, ReceiptProvider, StageCheckpointReader, StaticFileProviderFactory,
893 StaticFileSegment, StaticFileWriter,
894 };
895 use reth_prune_types::{PruneMode, PruneModes, ReceiptsLogPruneConfig};
896 use std::{cell::Cell, path::Path};
897 use tempfile::tempdir;
898
899 fn era1_block_tuple(number: u64, receipts: Vec<ReceiptEnvelope>) -> BlockTuple {
901 let header = Header { number, ..Default::default() };
902 BlockTuple::new(
903 CompressedHeader::from_rlp(&alloy_rlp::encode(&header)).unwrap(),
904 CompressedBody::from_rlp(&alloy_rlp::encode(BlockBody::default())).unwrap(),
905 CompressedReceipts::from_rlp(&alloy_rlp::encode(&receipts)).unwrap(),
906 TotalDifficulty::new(U256::ZERO),
907 )
908 }
909
910 fn block_with_one_receipt(number: u64) -> (Header, BlockBody, Receipt) {
912 let tx = TransactionSigned::new_unhashed(
913 TxLegacy::default().into(),
914 Signature::test_signature(),
915 );
916 let receipt = Receipt {
917 tx_type: TxType::Legacy,
918 success: true,
919 cumulative_gas_used: 21_000,
920 logs: vec![],
921 };
922
923 let with_bloom = vec![TxReceipt::with_bloom_ref(&receipt)];
924 let header = Header {
925 number,
926 receipts_root: calculate_receipt_root(&with_bloom),
927 logs_bloom: with_bloom.iter().fold(Bloom::ZERO, |bloom, r| bloom | r.bloom_ref()),
928 ..Default::default()
929 };
930
931 (header, BlockBody { transactions: vec![tx], ..Default::default() }, receipt)
932 }
933
934 fn era1_receipt(status: Eip658Value) -> ReceiptEnvelope {
936 ReceiptEnvelope::Legacy(ReceiptWithBloom::new(
937 RlpReceipt { status, cumulative_gas_used: 21_000, logs: vec![] },
938 Bloom::ZERO,
939 ))
940 }
941
942 fn ere_block_tuple(number: u64, receipts: &[SlimReceipt]) -> EreBlockTuple {
944 let header = Header { number, ..Default::default() };
945 EreBlockTuple::new(
946 EreCompressedHeader::from_rlp(&alloy_rlp::encode(&header)).unwrap(),
947 EreCompressedBody::from_rlp(&alloy_rlp::encode(BlockBody::default())).unwrap(),
948 )
949 .with_receipts(CompressedSlimReceipts::from_receipts(receipts).unwrap())
950 }
951
952 struct TestEra;
953
954 impl<R> EraBlockReader<Header, BlockBody, R> for TestEra {
955 fn blocks<M: EraMeta + ?Sized>(
956 _meta: &M,
957 _decode_receipts: bool,
958 ) -> eyre::Result<impl Iterator<Item = eyre::Result<(Header, BlockBody, Option<Vec<R>>)>>>
959 {
960 Ok([1, 2].into_iter().map(|number| {
961 Ok((Header { number, ..Default::default() }, BlockBody::default(), None))
962 }))
963 }
964 }
965
966 struct TestEraWithEmptyReceipts;
968
969 impl<R> EraBlockReader<Header, BlockBody, R> for TestEraWithEmptyReceipts {
970 fn blocks<M: EraMeta + ?Sized>(
971 _meta: &M,
972 _decode_receipts: bool,
973 ) -> eyre::Result<impl Iterator<Item = eyre::Result<(Header, BlockBody, Option<Vec<R>>)>>>
974 {
975 Ok([1, 2].into_iter().map(|number| {
976 Ok((Header { number, ..Default::default() }, BlockBody::default(), Some(vec![])))
977 }))
978 }
979 }
980
981 struct TestEraWithNonZeroGenesis;
983
984 impl<R> EraBlockReader<Header, BlockBody, R> for TestEraWithNonZeroGenesis {
985 fn blocks<M: EraMeta + ?Sized>(
986 _meta: &M,
987 _decode_receipts: bool,
988 ) -> eyre::Result<impl Iterator<Item = eyre::Result<(Header, BlockBody, Option<Vec<R>>)>>>
989 {
990 Ok([101, 102].into_iter().map(|number| {
991 Ok((Header { number, ..Default::default() }, BlockBody::default(), Some(vec![])))
992 }))
993 }
994 }
995
996 struct TestEraWithMismatchedReceipts;
999
1000 impl EraBlockReader<Header, BlockBody, reth_ethereum_primitives::Receipt>
1001 for TestEraWithMismatchedReceipts
1002 {
1003 fn blocks<M: EraMeta + ?Sized>(
1004 _meta: &M,
1005 _decode_receipts: bool,
1006 ) -> eyre::Result<
1007 impl Iterator<
1008 Item = eyre::Result<EraBlock<Header, BlockBody, reth_ethereum_primitives::Receipt>>,
1009 >,
1010 > {
1011 Ok(std::iter::once(Ok((
1012 Header { number: 1, ..Default::default() },
1013 BlockBody::default(),
1014 Some(vec![reth_ethereum_primitives::Receipt {
1015 tx_type: alloy_consensus::TxType::Legacy,
1016 success: true,
1017 cumulative_gas_used: 0,
1018 logs: vec![],
1019 }]),
1020 ))))
1021 }
1022 }
1023
1024 #[derive(Debug)]
1025 struct TestMeta {
1026 marked: Cell<bool>,
1027 }
1028
1029 impl EraMeta for TestMeta {
1030 fn mark_as_processed(&self) -> eyre::Result<()> {
1031 self.marked.set(true);
1032 Ok(())
1033 }
1034
1035 fn path(&self) -> &Path {
1036 Path::new("test.era1")
1037 }
1038 }
1039
1040 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1041 async fn import_stops_at_to_block() {
1042 let pf = create_test_provider_factory();
1043 init_genesis(&pf).unwrap();
1044
1045 let folder = tempdir().unwrap();
1046 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1047
1048 let stream = futures_util::stream::iter(vec![
1050 Ok(TestMeta { marked: Cell::new(false) }),
1051 Ok(TestMeta { marked: Cell::new(false) }),
1052 ]);
1053
1054 let height = import::<TestEra, _, _, _, Block, _, _>(
1055 stream,
1056 &pf,
1057 &mut hash_collector,
1058 Some(1),
1059 false,
1060 &|_| false,
1061 )
1062 .unwrap();
1063
1064 assert_eq!(height, 1);
1065 }
1066
1067 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1068 async fn backfill_does_not_move_header_checkpoints_backwards() {
1069 let pf = create_test_provider_factory();
1070 init_genesis(&pf).unwrap();
1071 let static_file_provider = pf.static_file_provider();
1072 let folder = tempdir().unwrap();
1073 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1074
1075 let provider = pf.database_provider_rw().unwrap();
1077 {
1078 let mut writer =
1079 static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
1080 let meta = TestMeta { marked: Cell::new(false) };
1081 process::<TestEra, _, Block, _, _>(
1082 &meta,
1083 &mut writer,
1084 None,
1085 &provider,
1086 &mut hash_collector,
1087 0..=2,
1088 ImportPolicy { headers_tip: 0, is_receipt_verifiable: &|_| false },
1089 )
1090 .unwrap();
1091 writer.commit().unwrap();
1092 }
1093 save_stage_checkpoints(&provider, 0, 2, 2, 2).unwrap();
1094 provider.save_stage_checkpoint(StageId::Execution, StageCheckpoint::new(1)).unwrap();
1096 provider.commit().unwrap();
1097
1098 let stream = futures_util::stream::iter(vec![Ok(TestMeta { marked: Cell::new(false) })]);
1099 let height = import::<TestEraWithEmptyReceipts, _, _, _, Block, _, _>(
1100 stream,
1101 &pf,
1102 &mut hash_collector,
1103 None,
1104 true,
1105 &|_| true,
1106 )
1107 .unwrap();
1108
1109 assert_eq!(height, 1);
1110 let provider = pf.database_provider_rw().unwrap();
1111 for stage in [StageId::Headers, StageId::Bodies] {
1112 assert_eq!(
1113 provider.get_stage_checkpoint(stage).unwrap().map(|c| c.block_number),
1114 Some(2),
1115 "{stage} checkpoint must not regress below the headers already in static files"
1116 );
1117 }
1118 }
1119
1120 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1121 async fn receipt_import_requires_an_executed_range() {
1122 let pf = create_test_provider_factory();
1123 init_genesis(&pf).unwrap();
1124 let folder = tempdir().unwrap();
1125 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1126
1127 let stream = futures_util::stream::iter(vec![Ok(TestMeta { marked: Cell::new(false) })]);
1129 let result = import::<TestEraWithEmptyReceipts, _, _, _, Block, _, _>(
1130 stream,
1131 &pf,
1132 &mut hash_collector,
1133 None,
1134 true,
1135 &|_| true,
1136 );
1137
1138 assert!(result.is_err());
1139 }
1140
1141 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1142 async fn receipt_import_rejects_a_to_block_off_the_execution_checkpoint() {
1143 let pf = create_test_provider_factory();
1144 init_genesis(&pf).unwrap();
1145 let folder = tempdir().unwrap();
1146 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1147
1148 let provider = pf.database_provider_rw().unwrap();
1149 provider.save_stage_checkpoint(StageId::Execution, StageCheckpoint::new(2)).unwrap();
1150 provider.commit().unwrap();
1151
1152 let stream = futures_util::stream::iter(vec![Ok(TestMeta { marked: Cell::new(false) })]);
1154 let result = import::<TestEraWithEmptyReceipts, _, _, _, Block, _, _>(
1155 stream,
1156 &pf,
1157 &mut hash_collector,
1158 Some(1),
1159 true,
1160 &|_| true,
1161 );
1162
1163 assert!(result.is_err());
1164 }
1165
1166 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1167 async fn receipt_import_backfills_a_completely_absent_segment() {
1168 let pf = create_test_provider_factory();
1170 let static_file_provider = pf.static_file_provider();
1171 let folder = tempdir().unwrap();
1172 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1173
1174 {
1177 let mut writer =
1178 static_file_provider.latest_writer(StaticFileSegment::Transactions).unwrap();
1179 writer.increment_block(0).unwrap();
1180 writer.commit().unwrap();
1181 }
1182
1183 let provider = pf.database_provider_rw().unwrap();
1184 {
1185 let mut writer =
1186 static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
1187 let genesis = Header::default();
1188 writer.append_header(&genesis, &genesis.hash_slow()).unwrap();
1189
1190 let meta = TestMeta { marked: Cell::new(false) };
1191 process::<TestEra, _, Block, _, _>(
1192 &meta,
1193 &mut writer,
1194 None,
1195 &provider,
1196 &mut hash_collector,
1197 0..=2,
1198 ImportPolicy { headers_tip: 0, is_receipt_verifiable: &|_| true },
1199 )
1200 .unwrap();
1201 writer.commit().unwrap();
1202 }
1203 provider.save_stage_checkpoint(StageId::Execution, StageCheckpoint::new(2)).unwrap();
1204 provider.commit().unwrap();
1205
1206 assert_eq!(
1207 static_file_provider.get_highest_static_file_block(StaticFileSegment::Receipts),
1208 None
1209 );
1210
1211 let stream = futures_util::stream::iter(vec![Ok(TestMeta { marked: Cell::new(false) })]);
1212 let height = import::<TestEraWithEmptyReceipts, _, _, _, Block, _, _>(
1213 stream,
1214 &pf,
1215 &mut hash_collector,
1216 None,
1217 true,
1218 &|_| true,
1219 )
1220 .unwrap();
1221
1222 assert_eq!(height, 2);
1223 assert_eq!(
1224 static_file_provider.get_highest_static_file_block(StaticFileSegment::Receipts),
1225 Some(2)
1226 );
1227 }
1228
1229 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1230 async fn receipt_import_backfills_an_absent_segment_after_non_zero_genesis() {
1231 const GENESIS: u64 = 100;
1232
1233 let pf = create_test_provider_factory_with_genesis_block_number(GENESIS);
1234 let static_file_provider = pf.static_file_provider();
1235 let folder = tempdir().unwrap();
1236 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1237
1238 {
1239 let mut writer =
1240 static_file_provider.get_writer(GENESIS, StaticFileSegment::Transactions).unwrap();
1241 writer.user_header_mut().set_block_range(GENESIS, GENESIS);
1242 writer.commit().unwrap();
1243 }
1244
1245 let provider = pf.database_provider_rw().unwrap();
1246 {
1247 let mut writer =
1248 static_file_provider.get_writer(GENESIS, StaticFileSegment::Headers).unwrap();
1249 let genesis = Header { number: GENESIS, ..Default::default() };
1250 writer.user_header_mut().set_block_range(GENESIS, GENESIS);
1251 writer
1252 .append_header_direct(&genesis, genesis.difficulty, &genesis.hash_slow())
1253 .unwrap();
1254
1255 let meta = TestMeta { marked: Cell::new(false) };
1256 process::<TestEraWithNonZeroGenesis, _, Block, _, _>(
1257 &meta,
1258 &mut writer,
1259 None,
1260 &provider,
1261 &mut hash_collector,
1262 GENESIS..=102,
1263 ImportPolicy { headers_tip: GENESIS, is_receipt_verifiable: &|_| true },
1264 )
1265 .unwrap();
1266 writer.commit().unwrap();
1267 }
1268 provider.save_stage_checkpoint(StageId::Execution, StageCheckpoint::new(102)).unwrap();
1269 provider.commit().unwrap();
1270
1271 assert_eq!(static_file_provider.genesis_block_number(), GENESIS);
1272 assert_eq!(
1273 static_file_provider.get_highest_static_file_block(StaticFileSegment::Receipts),
1274 None
1275 );
1276
1277 let stream = futures_util::stream::iter(vec![Ok(TestMeta { marked: Cell::new(false) })]);
1278 let height = import::<TestEraWithNonZeroGenesis, _, _, _, Block, _, _>(
1279 stream,
1280 &pf,
1281 &mut hash_collector,
1282 None,
1283 true,
1284 &|_| true,
1285 )
1286 .unwrap();
1287
1288 assert_eq!(height, 102);
1289 assert_eq!(
1290 static_file_provider.get_highest_static_file_block(StaticFileSegment::Receipts),
1291 Some(102)
1292 );
1293 }
1294
1295 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1296 async fn receipt_import_rejects_a_prune_config_that_bypasses_static_files() {
1297 let prune_modes = PruneModes {
1298 receipts_log_filter: ReceiptsLogPruneConfig(
1299 std::iter::once((Address::ZERO, PruneMode::Full)).collect(),
1300 ),
1301 ..Default::default()
1302 };
1303 let pf = create_test_provider_factory().with_prune_modes(prune_modes);
1304 init_genesis(&pf).unwrap();
1305 let folder = tempdir().unwrap();
1306 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1307
1308 let provider = pf.database_provider_rw().unwrap();
1309 provider.save_stage_checkpoint(StageId::Execution, StageCheckpoint::new(2)).unwrap();
1310 provider.commit().unwrap();
1311
1312 let stream = futures_util::stream::iter(vec![Ok(TestMeta { marked: Cell::new(false) })]);
1313 let result = import::<TestEraWithEmptyReceipts, _, _, _, Block, _, _>(
1314 stream,
1315 &pf,
1316 &mut hash_collector,
1317 None,
1318 true,
1319 &|_| true,
1320 );
1321
1322 assert!(result.is_err());
1323 }
1324
1325 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1326 async fn receipt_import_rejects_a_repair_starting_before_byzantium() {
1327 let pf = create_test_provider_factory();
1328 init_genesis(&pf).unwrap();
1329 let folder = tempdir().unwrap();
1330 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1331
1332 let provider = pf.database_provider_rw().unwrap();
1333 provider.save_stage_checkpoint(StageId::Execution, StageCheckpoint::new(2)).unwrap();
1334 provider.commit().unwrap();
1335
1336 let stream = futures_util::stream::iter(vec![Ok(TestMeta { marked: Cell::new(false) })]);
1338 let result = import::<TestEraWithEmptyReceipts, _, _, _, Block, _, _>(
1339 stream,
1340 &pf,
1341 &mut hash_collector,
1342 None,
1343 true,
1344 &|number| number >= 2,
1345 );
1346
1347 assert!(result.is_err());
1348 }
1349
1350 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1351 async fn receipt_import_errors_when_the_source_stops_short() {
1352 let pf = create_test_provider_factory();
1353 init_genesis(&pf).unwrap();
1354 let folder = tempdir().unwrap();
1355 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1356
1357 let provider = pf.database_provider_rw().unwrap();
1358 provider.save_stage_checkpoint(StageId::Execution, StageCheckpoint::new(10)).unwrap();
1359 provider.commit().unwrap();
1360
1361 let stream = futures_util::stream::iter(vec![Ok(TestMeta { marked: Cell::new(false) })]);
1363 let result = import::<TestEraWithEmptyReceipts, _, _, _, Block, _, _>(
1364 stream,
1365 &pf,
1366 &mut hash_collector,
1367 None,
1368 true,
1369 &|_| true,
1370 );
1371
1372 assert!(result.is_err());
1373 }
1374
1375 #[test]
1376 fn process_does_not_mark_partially_consumed_file_processed() {
1377 let pf = create_test_provider_factory();
1378 init_genesis(&pf).unwrap();
1379
1380 let static_file_provider = pf.static_file_provider();
1381 let mut writer = static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
1382 let provider = pf.database_provider_rw().unwrap();
1383 let folder = tempdir().unwrap();
1384 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1385 let meta = TestMeta { marked: Cell::new(false) };
1386
1387 let height = process::<TestEra, _, Block, _, _>(
1388 &meta,
1389 &mut writer,
1390 None,
1391 &provider,
1392 &mut hash_collector,
1393 0..=1,
1394 ImportPolicy { headers_tip: 0, is_receipt_verifiable: &|_| false },
1395 )
1396 .unwrap();
1397
1398 assert_eq!(height, 1);
1399 assert!(!meta.marked.get());
1400 }
1401
1402 #[test]
1403 fn process_iter_rejects_non_contiguous_blocks() {
1404 let pf = create_test_provider_factory();
1405 init_genesis(&pf).unwrap();
1406
1407 let static_file_provider = pf.static_file_provider();
1408 let mut writer = static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
1409 let provider = pf.database_provider_rw().unwrap();
1410 let folder = tempdir().unwrap();
1411 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1412
1413 let blocks = [5u64, 6].into_iter().map(|number| {
1416 Ok((Header { number, ..Default::default() }, BlockBody::default(), None))
1417 });
1418
1419 let result = process_iter::<_, Block, _, _>(
1420 blocks,
1421 &mut writer,
1422 None,
1423 &provider,
1424 &mut hash_collector,
1425 0..,
1426 ImportPolicy { headers_tip: 0, is_receipt_verifiable: &|_| false },
1427 );
1428
1429 assert!(result.is_err());
1430 }
1431
1432 #[test]
1433 fn process_writes_receipts_when_requested() {
1434 let pf = create_test_provider_factory();
1435 init_genesis(&pf).unwrap();
1436
1437 let static_file_provider = pf.static_file_provider();
1438 let mut writer = static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
1439 let mut receipts_writer =
1440 static_file_provider.latest_writer(StaticFileSegment::Receipts).unwrap();
1441 let provider = pf.database_provider_rw().unwrap();
1442 let folder = tempdir().unwrap();
1443 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1444 let meta = TestMeta { marked: Cell::new(false) };
1445
1446 let height = process::<TestEraWithEmptyReceipts, _, Block, _, _>(
1447 &meta,
1448 &mut writer,
1449 Some(&mut receipts_writer),
1450 &provider,
1451 &mut hash_collector,
1452 0..=1,
1453 ImportPolicy { headers_tip: 0, is_receipt_verifiable: &|_| false },
1454 )
1455 .unwrap();
1456 receipts_writer.commit().unwrap();
1457
1458 assert_eq!(height, 1);
1459 assert_eq!(
1460 static_file_provider.get_highest_static_file_block(StaticFileSegment::Receipts),
1461 Some(1)
1462 );
1463 }
1464
1465 #[test]
1466 fn process_iter_errors_when_receipts_missing() {
1467 let pf = create_test_provider_factory();
1468 init_genesis(&pf).unwrap();
1469
1470 let static_file_provider = pf.static_file_provider();
1471 let mut writer = static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
1472 let mut receipts_writer =
1473 static_file_provider.latest_writer(StaticFileSegment::Receipts).unwrap();
1474 let provider = pf.database_provider_rw().unwrap();
1475 let folder = tempdir().unwrap();
1476 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1477
1478 let blocks: Vec<
1481 eyre::Result<EraBlock<Header, BlockBody, reth_ethereum_primitives::Receipt>>,
1482 > = vec![Ok((Header { number: 1, ..Default::default() }, BlockBody::default(), None))];
1483
1484 let result = process_iter::<_, Block, _, _>(
1485 blocks.into_iter(),
1486 &mut writer,
1487 Some(&mut receipts_writer),
1488 &provider,
1489 &mut hash_collector,
1490 0..,
1491 ImportPolicy { headers_tip: 0, is_receipt_verifiable: &|_| false },
1492 );
1493
1494 assert!(result.is_err());
1495 }
1496
1497 #[test]
1498 fn process_errors_on_receipt_count_mismatch() {
1499 let pf = create_test_provider_factory();
1500 init_genesis(&pf).unwrap();
1501
1502 let static_file_provider = pf.static_file_provider();
1503 let mut writer = static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
1504 let mut receipts_writer =
1505 static_file_provider.latest_writer(StaticFileSegment::Receipts).unwrap();
1506 let provider = pf.database_provider_rw().unwrap();
1507 let folder = tempdir().unwrap();
1508 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1509 let meta = TestMeta { marked: Cell::new(false) };
1510
1511 let result = process::<TestEraWithMismatchedReceipts, _, Block, _, _>(
1514 &meta,
1515 &mut writer,
1516 Some(&mut receipts_writer),
1517 &provider,
1518 &mut hash_collector,
1519 0..,
1520 ImportPolicy { headers_tip: 0, is_receipt_verifiable: &|_| false },
1521 );
1522
1523 assert!(result.is_err());
1524 }
1525
1526 #[test]
1527 fn save_stage_checkpoints_leaves_execution_unset() {
1528 let pf = create_test_provider_factory();
1529 init_genesis(&pf).unwrap();
1530 let provider = pf.database_provider_rw().unwrap();
1531
1532 let execution_before = provider.get_stage_checkpoint(StageId::Execution).unwrap();
1533 save_stage_checkpoints(&provider, 0, 10, 10, 10).unwrap();
1534
1535 assert_eq!(
1536 provider.get_stage_checkpoint(StageId::Headers).unwrap().map(|c| c.block_number),
1537 Some(10)
1538 );
1539 assert_eq!(
1540 provider.get_stage_checkpoint(StageId::Bodies).unwrap().map(|c| c.block_number),
1541 Some(10)
1542 );
1543 assert_eq!(provider.get_stage_checkpoint(StageId::Execution).unwrap(), execution_before);
1545 }
1546
1547 #[test]
1548 fn backfills_receipts_onto_existing_headers() {
1549 let pf = create_test_provider_factory();
1550 init_genesis(&pf).unwrap();
1551 let static_file_provider = pf.static_file_provider();
1552 let folder = tempdir().unwrap();
1553 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1554
1555 {
1557 let provider = pf.database_provider_rw().unwrap();
1558 {
1561 let mut writer =
1562 static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
1563 let meta = TestMeta { marked: Cell::new(false) };
1564 process::<TestEra, _, Block, _, _>(
1565 &meta,
1566 &mut writer,
1567 None,
1568 &provider,
1569 &mut hash_collector,
1570 0..=2,
1571 ImportPolicy { headers_tip: 0, is_receipt_verifiable: &|_| false },
1572 )
1573 .unwrap();
1574 writer.commit().unwrap();
1575 }
1576 provider.commit().unwrap();
1577 }
1578 assert_eq!(
1579 static_file_provider.get_highest_static_file_block(StaticFileSegment::Headers),
1580 Some(2)
1581 );
1582 assert_eq!(
1584 static_file_provider.get_highest_static_file_block(StaticFileSegment::Receipts),
1585 Some(0)
1586 );
1587
1588 {
1590 let provider = pf.database_provider_rw().unwrap();
1591 {
1592 let mut writer =
1593 static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
1594 let mut receipts_writer =
1595 static_file_provider.latest_writer(StaticFileSegment::Receipts).unwrap();
1596 let meta = TestMeta { marked: Cell::new(false) };
1597 process::<TestEraWithEmptyReceipts, _, Block, _, _>(
1598 &meta,
1599 &mut writer,
1600 Some(&mut receipts_writer),
1601 &provider,
1602 &mut hash_collector,
1603 0..,
1604 ImportPolicy { headers_tip: 2, is_receipt_verifiable: &|_| false },
1605 )
1606 .unwrap();
1607 receipts_writer.commit().unwrap();
1608 writer.commit().unwrap();
1609 }
1610 provider.commit().unwrap();
1611 }
1612
1613 assert_eq!(
1615 static_file_provider.get_highest_static_file_block(StaticFileSegment::Receipts),
1616 Some(2)
1617 );
1618 assert_eq!(
1619 static_file_provider.get_highest_static_file_block(StaticFileSegment::Headers),
1620 Some(2)
1621 );
1622 }
1623
1624 #[test]
1625 fn process_iter_persists_verified_receipts() {
1626 let pf = create_test_provider_factory();
1627 init_genesis(&pf).unwrap();
1628 let static_file_provider = pf.static_file_provider();
1629 let folder = tempdir().unwrap();
1630 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1631
1632 let (header, body, receipt) = block_with_one_receipt(1);
1633 let provider = pf.database_provider_rw().unwrap();
1634 {
1635 let mut writer =
1636 static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
1637 let mut receipts_writer =
1638 static_file_provider.latest_writer(StaticFileSegment::Receipts).unwrap();
1639
1640 process_iter::<_, Block, _, _>(
1641 std::iter::once(Ok((header, body, Some(vec![receipt.clone()])))),
1642 &mut writer,
1643 Some(&mut receipts_writer),
1644 &provider,
1645 &mut hash_collector,
1646 0..,
1647 ImportPolicy { headers_tip: 0, is_receipt_verifiable: &|_| true },
1648 )
1649 .unwrap();
1650
1651 receipts_writer.commit().unwrap();
1652 writer.commit().unwrap();
1653 }
1654 provider.commit().unwrap();
1655
1656 let provider = pf.provider().unwrap();
1657 assert_eq!(provider.receipts_by_block(1.into()).unwrap(), Some(vec![receipt]));
1658 }
1659
1660 #[test]
1661 fn process_iter_rejects_receipts_the_header_does_not_commit_to() {
1662 let pf = create_test_provider_factory();
1663 init_genesis(&pf).unwrap();
1664 let static_file_provider = pf.static_file_provider();
1665 let folder = tempdir().unwrap();
1666 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1667
1668 let (header, body, receipt) = block_with_one_receipt(1);
1670 let tampered = Receipt { cumulative_gas_used: 42_000, ..receipt };
1671
1672 let provider = pf.database_provider_rw().unwrap();
1673 let mut writer = static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
1674 let mut receipts_writer =
1675 static_file_provider.latest_writer(StaticFileSegment::Receipts).unwrap();
1676
1677 let result = process_iter::<_, Block, _, _>(
1678 std::iter::once(Ok((header, body, Some(vec![tampered])))),
1679 &mut writer,
1680 Some(&mut receipts_writer),
1681 &provider,
1682 &mut hash_collector,
1683 0..,
1684 ImportPolicy { headers_tip: 0, is_receipt_verifiable: &|_| true },
1685 );
1686
1687 assert!(result.is_err());
1688 }
1689
1690 #[test]
1691 fn decodes_post_byzantium_era1_receipts() {
1692 let tuple = era1_block_tuple(4_370_000, vec![era1_receipt(Eip658Value::Eip658(true))]);
1693
1694 let (_, _, receipts) =
1695 decode_with_receipts::<Header, BlockBody, Receipt, E2sError>(Ok(tuple), true).unwrap();
1696
1697 let receipts = receipts.unwrap();
1698 assert_eq!(receipts.len(), 1);
1699 assert!(receipts[0].success);
1700 assert_eq!(receipts[0].cumulative_gas_used, 21_000);
1701 }
1702
1703 #[test]
1704 fn decodes_typed_ere_receipts_into_node_receipts() {
1705 let logs = vec![Log::new_unchecked(
1706 Address::repeat_byte(0x11),
1707 vec![B256::repeat_byte(0x22)],
1708 Bytes::from_static(b"typed"),
1709 )];
1710 let slim = vec![
1711 SlimReceipt {
1712 tx_type: TxType::Eip2930,
1713 status: Eip658Value::Eip658(true),
1714 cumulative_gas_used: 21_000,
1715 logs: logs.clone(),
1716 },
1717 SlimReceipt {
1718 tx_type: TxType::Eip1559,
1719 status: Eip658Value::Eip658(false),
1720 cumulative_gas_used: 42_000,
1721 logs: vec![],
1722 },
1723 ];
1724 let tuple = ere_block_tuple(12_965_000, &slim);
1725
1726 let (_, _, receipts) =
1727 Ere::decode_with_receipts::<Header, BlockBody, Receipt, E2sError>(Ok(tuple), true)
1728 .unwrap();
1729
1730 assert_eq!(
1731 receipts,
1732 Some(vec![
1733 Receipt {
1734 tx_type: TxType::Eip2930,
1735 success: true,
1736 cumulative_gas_used: 21_000,
1737 logs,
1738 },
1739 Receipt {
1740 tx_type: TxType::Eip1559,
1741 success: false,
1742 cumulative_gas_used: 42_000,
1743 logs: vec![],
1744 },
1745 ])
1746 );
1747 }
1748
1749 #[test]
1750 fn rejects_pre_byzantium_era1_receipts() {
1751 let tuple = era1_block_tuple(
1753 46_147,
1754 vec![era1_receipt(Eip658Value::PostState(B256::repeat_byte(1)))],
1755 );
1756
1757 let err = decode_with_receipts::<Header, BlockBody, Receipt, E2sError>(Ok(tuple), true)
1758 .unwrap_err()
1759 .to_string();
1760
1761 assert!(err.contains("pre-Byzantium"), "unexpected error: {err}");
1762 assert!(err.contains("46147"), "error should name the offending block: {err}");
1763 }
1764
1765 #[test]
1766 fn skips_pre_byzantium_receipts_when_not_requested() {
1767 let tuple = era1_block_tuple(
1769 46_147,
1770 vec![era1_receipt(Eip658Value::PostState(B256::repeat_byte(1)))],
1771 );
1772
1773 let (header, _, receipts) =
1774 decode_with_receipts::<Header, BlockBody, Receipt, E2sError>(Ok(tuple), false).unwrap();
1775
1776 assert_eq!(header.number, 46_147);
1777 assert!(receipts.is_none());
1778 }
1779
1780 #[test]
1781 fn backfill_rejects_a_header_that_differs_from_the_persisted_one() {
1782 let pf = create_test_provider_factory();
1783 init_genesis(&pf).unwrap();
1784 let static_file_provider = pf.static_file_provider();
1785 let folder = tempdir().unwrap();
1786 let mut hash_collector = Collector::new(4096, Some(folder.path().to_owned()));
1787
1788 let provider = pf.database_provider_rw().unwrap();
1789 {
1790 let mut writer =
1791 static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
1792 let meta = TestMeta { marked: Cell::new(false) };
1793 process::<TestEra, _, Block, _, _>(
1794 &meta,
1795 &mut writer,
1796 None,
1797 &provider,
1798 &mut hash_collector,
1799 0..=2,
1800 ImportPolicy { headers_tip: 0, is_receipt_verifiable: &|_| false },
1801 )
1802 .unwrap();
1803 writer.commit().unwrap();
1804 }
1805 provider.commit().unwrap();
1806
1807 let provider = pf.database_provider_rw().unwrap();
1809 let mut writer = static_file_provider.latest_writer(StaticFileSegment::Headers).unwrap();
1810 let mut receipts_writer =
1811 static_file_provider.latest_writer(StaticFileSegment::Receipts).unwrap();
1812
1813 let result = process_iter::<_, Block, _, _>(
1814 std::iter::once(Ok((
1815 Header { number: 1, gas_limit: 42, ..Default::default() },
1816 BlockBody::default(),
1817 Some(vec![]),
1818 ))),
1819 &mut writer,
1820 Some(&mut receipts_writer),
1821 &provider,
1822 &mut hash_collector,
1823 0..,
1824 ImportPolicy { headers_tip: 2, is_receipt_verifiable: &|_| false },
1825 );
1826
1827 assert!(result.is_err());
1828 }
1829
1830 #[test]
1831 fn verify_receipts_rejects_tampered_contents() {
1832 let receipts = vec![Receipt {
1833 tx_type: TxType::Legacy,
1834 success: true,
1835 cumulative_gas_used: 21_000,
1836 logs: vec![],
1837 }];
1838
1839 let with_bloom = receipts.iter().map(TxReceipt::with_bloom_ref).collect::<Vec<_>>();
1841 let header = Header {
1842 receipts_root: calculate_receipt_root(&with_bloom),
1843 logs_bloom: with_bloom.iter().fold(Bloom::ZERO, |bloom, r| bloom | r.bloom_ref()),
1844 ..Default::default()
1845 };
1846 verify_receipts(&header, &receipts, true).unwrap();
1847
1848 let tampered = vec![Receipt { cumulative_gas_used: 42_000, ..receipts[0].clone() }];
1850 assert!(verify_receipts(&header, &tampered, true).is_err());
1851 }
1852
1853 #[test]
1854 fn verify_receipts_checks_the_bloom_even_without_the_root() {
1855 let receipts = vec![Receipt {
1856 tx_type: TxType::Legacy,
1857 success: true,
1858 cumulative_gas_used: 21_000,
1859 logs: vec![Log::new_unchecked(
1860 Address::ZERO,
1861 vec![B256::repeat_byte(1)],
1862 Bytes::default(),
1863 )],
1864 }];
1865
1866 assert!(verify_receipts(&Header::default(), &receipts, false).is_err());
1868 }
1869}