1mod bodies;
3mod era;
4mod execution;
6mod finish;
8mod hashing_account;
10mod hashing_storage;
12mod headers;
14mod index_account_history;
16mod index_storage_history;
18mod merkle;
20mod prune;
21mod sender_recovery;
23mod tx_lookup;
25
26pub use bodies::*;
27pub use era::*;
28pub use execution::*;
29pub use finish::*;
30pub use hashing_account::*;
31pub use hashing_storage::*;
32pub use headers::*;
33pub use index_account_history::*;
34pub use index_storage_history::*;
35pub use merkle::*;
36pub use prune::*;
37pub use sender_recovery::*;
38pub use tx_lookup::*;
39
40mod utils;
41use utils::*;
42
43#[cfg(test)]
44mod tests {
45 use super::*;
46 use crate::test_utils::{StorageKind, TestStageDB};
47 use alloy_consensus::{SignableTransaction, TxLegacy};
48 use alloy_primitives::{
49 address, hex_literal::hex, keccak256, BlockNumber, Signature, B256, U256,
50 };
51 use alloy_rlp::Decodable;
52 use reth_chainspec::ChainSpecBuilder;
53 use reth_db::mdbx::{cursor::Cursor, RW};
54 use reth_db_api::{
55 cursor::{DbCursorRO, DbCursorRW},
56 models::StorageSettings,
57 table::Table,
58 tables,
59 transaction::{DbTx, DbTxMut},
60 AccountsHistory,
61 };
62 use reth_ethereum_consensus::EthBeaconConsensus;
63 use reth_ethereum_primitives::Block;
64 use reth_evm_ethereum::EthEvmConfig;
65 use reth_exex::ExExManagerHandle;
66 use reth_primitives_traits::{Account, Bytecode, SealedBlock};
67 use reth_provider::{
68 providers::{StaticFileProvider, StaticFileWriter},
69 test_utils::MockNodeTypesWithDB,
70 AccountExtReader, BlockBodyIndicesProvider, BlockWriter, DatabaseProviderFactory,
71 ProviderFactory, ProviderResult, PruneCheckpointWriter, ReceiptProvider,
72 StageCheckpointWriter, StaticFileProviderFactory, StorageReader,
73 };
74 use reth_prune_types::{PruneCheckpoint, PruneMode, PruneModes, PruneSegment};
75 use reth_stages_api::{
76 ExecInput, ExecutionStageThresholds, PipelineTarget, Stage, StageCheckpoint, StageId,
77 };
78 use reth_static_file_types::StaticFileSegment;
79 use reth_storage_api::StorageSettingsCache;
80 use reth_testing_utils::generators::{
81 self, random_block, random_block_range, random_receipt, BlockRangeParams,
82 };
83 use std::{io::Write, sync::Arc};
84
85 #[tokio::test]
86 #[ignore]
87 async fn test_prune() {
88 let test_db = TestStageDB::default();
89
90 let provider_rw = test_db.factory.provider_rw().unwrap();
91 let tip = 66;
92 let input = ExecInput { target: Some(tip), checkpoint: None };
93 let mut genesis_rlp = hex!("f901faf901f5a00000000000000000000000000000000000000000000000000000000000000000a01dcc4de8dec75d7aab85b567b6ccd41ad312451b948a7413f0a142fd40d49347942adc25665018aa1fe0e6bc666dac8fc2697ff9baa045571b40ae66ca7480791bbb2887286e4e4c4b1b298b191c889d6959023a32eda056e81f171bcc55a6ff8345e692c0f86e5b48e01b996cadc001622fb5e363b421a056e81f171bcc55a6ff8345e692c0f86e5b48e01b996cadc001622fb5e363b421b901000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000083020000808502540be400808000a00000000000000000000000000000000000000000000000000000000000000000880000000000000000c0c0").as_slice();
94 let genesis = SealedBlock::<Block>::decode(&mut genesis_rlp).unwrap();
95 let mut block_rlp = hex!("f90262f901f9a075c371ba45999d87f4542326910a11af515897aebce5265d3f6acd1f1161f82fa01dcc4de8dec75d7aab85b567b6ccd41ad312451b948a7413f0a142fd40d49347942adc25665018aa1fe0e6bc666dac8fc2697ff9baa098f2dcd87c8ae4083e7017a05456c14eea4b1db2032126e27b3b1563d57d7cc0a08151d548273f6683169524b66ca9fe338b9ce42bc3540046c828fd939ae23bcba03f4e5c2ec5b2170b711d97ee755c160457bb58d8daa338e835ec02ae6860bbabb901000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000083020000018502540be40082a8798203e800a00000000000000000000000000000000000000000000000000000000000000000880000000000000000f863f861800a8405f5e10094100000000000000000000000000000000000000080801ba07e09e26678ed4fac08a249ebe8ed680bf9051a5e14ad223e4b2b9d26e0208f37a05f6e3f188e3e6eab7d7d3b6568f5eac7d687b08d307d3154ccd8c87b4630509bc0").as_slice();
96 let block = SealedBlock::<Block>::decode(&mut block_rlp).unwrap();
97 let mut head = block.hash();
98 provider_rw.insert_block(&genesis.try_recover().unwrap()).unwrap();
99 provider_rw.insert_block(&block.try_recover().unwrap()).unwrap();
100
101 let mut rng = generators::rng();
103 for block_number in 2..=tip {
104 let nblock = random_block(
105 &mut rng,
106 block_number,
107 generators::BlockParams { parent: Some(head), ..Default::default() },
108 );
109 head = nblock.hash();
110 provider_rw.insert_block(&nblock.try_recover().unwrap()).unwrap();
111 }
112 provider_rw
113 .static_file_provider()
114 .latest_writer(StaticFileSegment::Headers)
115 .unwrap()
116 .commit()
117 .unwrap();
118 provider_rw.commit().unwrap();
119
120 let provider_rw = test_db.factory.provider_rw().unwrap();
122 let code = hex!("5a465a905090036002900360015500");
123 let code_hash = keccak256(hex!("5a465a905090036002900360015500"));
124 provider_rw
125 .tx_ref()
126 .put::<tables::PlainAccountState>(
127 address!("0x1000000000000000000000000000000000000000"),
128 Account { bytecode_hash: Some(code_hash), ..Default::default() },
129 )
130 .unwrap();
131 provider_rw
132 .tx_ref()
133 .put::<tables::PlainAccountState>(
134 address!("0xa94f5374fce5edbc8e2a8697c15331677e6ebf0b"),
135 Account { balance: U256::from(0x3635c9adc5dea00000u128), ..Default::default() },
136 )
137 .unwrap();
138 provider_rw
139 .tx_ref()
140 .put::<tables::Bytecodes>(code_hash, Bytecode::new_raw(code.to_vec().into()))
141 .unwrap();
142 provider_rw.commit().unwrap();
143
144 let check_pruning = |factory: ProviderFactory<MockNodeTypesWithDB>,
145 prune_modes: PruneModes,
146 expect_num_receipts: usize,
147 expect_num_acc_changesets: usize,
148 expect_num_storage_changesets: usize| async move {
149 let provider = factory.database_provider_rw().unwrap();
150
151 let mut execution_stage = ExecutionStage::new(
154 EthEvmConfig::ethereum(Arc::new(
155 ChainSpecBuilder::mainnet().berlin_activated().build(),
156 )),
157 Arc::new(EthBeaconConsensus::new(Arc::new(
158 ChainSpecBuilder::mainnet().berlin_activated().build(),
159 ))),
160 ExecutionStageThresholds {
161 max_blocks: Some(100),
162 max_changes: None,
163 max_cumulative_gas: None,
164 max_duration: None,
165 },
166 MERKLE_STAGE_DEFAULT_REBUILD_THRESHOLD,
167 ExExManagerHandle::empty(),
168 );
169
170 execution_stage.execute(&provider, input).unwrap();
171 assert_eq!(
172 provider.receipts_by_block(1.into()).unwrap().unwrap().len(),
173 expect_num_receipts
174 );
175
176 assert_eq!(
177 provider.changed_storages_and_blocks_with_range(0..=1000).unwrap().len(),
178 expect_num_storage_changesets
179 );
180
181 assert_eq!(
182 provider.changed_accounts_and_blocks_with_range(0..=1000).unwrap().len(),
183 expect_num_acc_changesets
184 );
185
186 let mut acc_indexing_stage = IndexAccountHistoryStage {
188 prune_mode: prune_modes.account_history,
189 ..Default::default()
190 };
191
192 if prune_modes.account_history == Some(PruneMode::Full) {
193 assert!(acc_indexing_stage.execute(&provider, input).is_err());
195 } else {
196 acc_indexing_stage.execute(&provider, input).unwrap();
197 let mut account_history: Cursor<RW, AccountsHistory> =
198 provider.tx_ref().cursor_read::<tables::AccountsHistory>().unwrap();
199 assert_eq!(account_history.walk(None).unwrap().count(), expect_num_acc_changesets);
200 }
201
202 let mut storage_indexing_stage = IndexStorageHistoryStage {
204 prune_mode: prune_modes.storage_history,
205 ..Default::default()
206 };
207
208 if prune_modes.storage_history == Some(PruneMode::Full) {
209 assert!(storage_indexing_stage.execute(&provider, input).is_err());
211 } else {
212 storage_indexing_stage.execute(&provider, input).unwrap();
213
214 let mut storage_history =
215 provider.tx_ref().cursor_read::<tables::StoragesHistory>().unwrap();
216 assert_eq!(
217 storage_history.walk(None).unwrap().count(),
218 expect_num_storage_changesets
219 );
220 }
221 };
222
223 let mut prune = PruneModes::default();
226 check_pruning(test_db.factory.clone(), prune.clone(), 1, 3, 1).await;
227
228 prune.receipts = Some(PruneMode::Full);
229 prune.account_history = Some(PruneMode::Full);
230 prune.storage_history = Some(PruneMode::Full);
231 check_pruning(test_db.factory.clone(), prune.clone(), 0, 0, 0).await;
233
234 prune.receipts = Some(PruneMode::Before(1));
235 prune.account_history = Some(PruneMode::Before(1));
236 prune.storage_history = Some(PruneMode::Before(1));
237 check_pruning(test_db.factory.clone(), prune.clone(), 1, 3, 1).await;
238
239 prune.receipts = Some(PruneMode::Before(2));
240 prune.account_history = Some(PruneMode::Before(2));
241 prune.storage_history = Some(PruneMode::Before(2));
242 check_pruning(test_db.factory.clone(), prune.clone(), 0, 1, 0).await;
244
245 prune.receipts = Some(PruneMode::Distance(66));
246 prune.account_history = Some(PruneMode::Distance(66));
247 prune.storage_history = Some(PruneMode::Distance(66));
248 check_pruning(test_db.factory.clone(), prune.clone(), 1, 3, 1).await;
249
250 prune.receipts = Some(PruneMode::Distance(64));
251 prune.account_history = Some(PruneMode::Distance(64));
252 prune.storage_history = Some(PruneMode::Distance(64));
253 check_pruning(test_db.factory.clone(), prune.clone(), 0, 1, 0).await;
255 }
256
257 fn seed_data(num_blocks: usize) -> ProviderResult<TestStageDB> {
260 let db = TestStageDB::default();
261 let mut rng = generators::rng();
262 let genesis_hash = B256::ZERO;
263 let tip = (num_blocks - 1) as u64;
264
265 let blocks = random_block_range(
266 &mut rng,
267 0..=tip,
268 BlockRangeParams { parent: Some(genesis_hash), tx_count: 2..3, ..Default::default() },
269 );
270 db.insert_blocks(blocks.iter(), StorageKind::Static)?;
271
272 let mut receipts = Vec::with_capacity(blocks.len());
273 let mut tx_num = 0u64;
274 for block in &blocks {
275 let mut block_receipts = Vec::with_capacity(block.transaction_count());
276 for transaction in &block.body().transactions {
277 block_receipts.push((tx_num, random_receipt(&mut rng, transaction, Some(0), None)));
278 tx_num += 1;
279 }
280 receipts.push((block.number, block_receipts));
281 }
282 db.insert_receipts_by_block(receipts, StorageKind::Static)?;
283
284 let provider_rw = db.factory.provider_rw()?;
286 for stage in StageId::ALL {
287 provider_rw.save_stage_checkpoint(stage, StageCheckpoint::new(tip))?;
288 }
289 provider_rw.commit()?;
290
291 Ok(db)
292 }
293
294 fn seed_v2_data() -> TestStageDB {
295 let db = seed_data(90).unwrap();
296 db.factory.set_storage_settings_cache(StorageSettings::v2());
297 db
298 }
299
300 fn save_prune_checkpoints(
301 db: &TestStageDB,
302 checkpoints: impl IntoIterator<Item = (PruneSegment, BlockNumber, PruneMode)>,
303 ) {
304 let provider_rw = db.factory.provider_rw().unwrap();
305 for (segment, block_number, prune_mode) in checkpoints {
306 provider_rw
307 .save_prune_checkpoint(
308 segment,
309 PruneCheckpoint {
310 block_number: Some(block_number),
311 tx_number: None,
312 prune_mode,
313 },
314 )
315 .unwrap();
316 }
317 provider_rw.commit().unwrap();
318 }
319
320 fn assert_consistency(db: &TestStageDB, expected: Option<PipelineTarget>) {
321 assert_eq!(
322 db.factory
323 .static_file_provider()
324 .check_consistency(&db.factory.database_provider_ro().unwrap())
325 .unwrap(),
326 expected
327 );
328 }
329
330 fn simulate_behind_checkpoint_corruption(
333 db: &TestStageDB,
334 prune_count: usize,
335 segment: StaticFileSegment,
336 expected: Option<PipelineTarget>,
337 ) {
338 let mut static_file_provider = db.factory.static_file_provider();
341 static_file_provider = StaticFileProvider::read_write(static_file_provider.path()).unwrap();
342
343 {
346 let mut headers_writer = static_file_provider.latest_writer(segment).unwrap();
347 let reader = headers_writer.inner().jar().open_data_reader().unwrap();
348 let columns = headers_writer.inner().jar().columns();
349 let data_file = headers_writer.inner().data_file();
350 let last_offset = reader.reverse_offset(prune_count * columns).unwrap();
351 data_file.get_mut().set_len(last_offset).unwrap();
352 data_file.flush().unwrap();
353 data_file.get_ref().sync_all().unwrap();
354 }
355
356 let mut static_file_provider = db.factory.static_file_provider();
359 static_file_provider = StaticFileProvider::read_write(static_file_provider.path()).unwrap();
360 assert!(matches!(
361 static_file_provider
362 .check_consistency(&db.factory.database_provider_ro().unwrap()),
363 Ok(e) if e == expected
364 ));
365 }
366
367 fn save_checkpoint_and_check(
370 db: &TestStageDB,
371 stage_id: StageId,
372 checkpoint_block_number: BlockNumber,
373 expected: Option<PipelineTarget>,
374 ) {
375 let provider_rw = db.factory.provider_rw().unwrap();
376 provider_rw
377 .save_stage_checkpoint(stage_id, StageCheckpoint::new(checkpoint_block_number))
378 .unwrap();
379 provider_rw.commit().unwrap();
380
381 assert!(matches!(
382 db.factory
383 .static_file_provider()
384 .check_consistency(&db.factory.database_provider_ro().unwrap(),),
385 Ok(e) if e == expected
386 ));
387 }
388
389 fn update_db_and_check<T: Table<Key = u64>>(
392 db: &TestStageDB,
393 key: u64,
394 expected: Option<PipelineTarget>,
395 ) where
396 <T as Table>::Value: Default,
397 {
398 update_db_with_and_check::<T>(db, key, expected, &Default::default());
399 }
400
401 fn update_db_with_and_check<T: Table<Key = u64>>(
404 db: &TestStageDB,
405 key: u64,
406 expected: Option<PipelineTarget>,
407 value: &T::Value,
408 ) {
409 let provider_rw = db.factory.provider_rw().unwrap();
410 let mut cursor = provider_rw.tx_ref().cursor_write::<T>().unwrap();
411 cursor.insert(key, value).unwrap();
412 provider_rw.commit().unwrap();
413
414 assert!(matches!(
415 db.factory
416 .static_file_provider()
417 .check_consistency(&db.factory.database_provider_ro().unwrap()),
418 Ok(e) if e == expected
419 ));
420 }
421
422 #[test]
423 fn test_consistency() {
424 let db = seed_data(90).unwrap();
425 let db_provider = db.factory.database_provider_ro().unwrap();
426
427 assert!(matches!(
428 db.factory.static_file_provider().check_consistency(&db_provider),
429 Ok(None)
430 ));
431 }
432
433 #[test]
434 fn test_consistency_no_commit_prune() {
435 let mut db_full = seed_data(90).unwrap();
437 db_full.factory = db_full.factory.with_prune_modes(PruneModes {
438 receipts: Some(PruneMode::Before(1)),
439 ..Default::default()
440 });
441
442 simulate_behind_checkpoint_corruption(&db_full, 1, StaticFileSegment::Receipts, None);
445
446 let db_archive = seed_data(90).unwrap();
448
449 simulate_behind_checkpoint_corruption(
452 &db_archive,
453 1,
454 StaticFileSegment::Receipts,
455 Some(PipelineTarget::Unwind(88)),
456 );
457
458 simulate_behind_checkpoint_corruption(
459 &db_archive,
460 3,
461 StaticFileSegment::Headers,
462 Some(PipelineTarget::Unwind(86)),
463 );
464 }
465
466 #[test]
467 fn test_consistency_checkpoints() {
468 let db = seed_data(90).unwrap();
469
470 let block = 87;
472 save_checkpoint_and_check(&db, StageId::Bodies, block, None);
473 assert_eq!(
474 db.factory
475 .static_file_provider()
476 .get_highest_static_file_block(StaticFileSegment::Transactions),
477 Some(block)
478 );
479 assert_eq!(
480 db.factory
481 .static_file_provider()
482 .get_highest_static_file_tx(StaticFileSegment::Transactions),
483 db.factory.block_body_indices(block).unwrap().map(|b| b.last_tx_num())
484 );
485
486 let block = 86;
487 save_checkpoint_and_check(&db, StageId::Execution, block, None);
488 assert_eq!(
489 db.factory
490 .static_file_provider()
491 .get_highest_static_file_block(StaticFileSegment::Receipts),
492 Some(block)
493 );
494 assert_eq!(
495 db.factory
496 .static_file_provider()
497 .get_highest_static_file_tx(StaticFileSegment::Receipts),
498 db.factory.block_body_indices(block).unwrap().map(|b| b.last_tx_num())
499 );
500
501 let block = 80;
502 save_checkpoint_and_check(&db, StageId::Headers, block, None);
503 assert_eq!(
504 db.factory
505 .static_file_provider()
506 .get_highest_static_file_block(StaticFileSegment::Headers),
507 Some(block)
508 );
509
510 save_checkpoint_and_check(&db, StageId::Headers, 91, Some(PipelineTarget::Unwind(block)));
512 }
513
514 #[test]
515 fn test_consistency_headers_gap() {
516 let db = seed_data(90).unwrap();
517 let current = db
518 .factory
519 .static_file_provider()
520 .get_highest_static_file_block(StaticFileSegment::Headers)
521 .unwrap();
522
523 update_db_and_check::<tables::Headers>(&db, current + 2, Some(PipelineTarget::Unwind(89)));
525
526 update_db_and_check::<tables::Headers>(&db, current + 1, None);
528 }
529
530 #[test]
531 fn test_consistency_tx_gap() {
532 let db = seed_data(90).unwrap();
533 let current = db
534 .factory
535 .static_file_provider()
536 .get_highest_static_file_tx(StaticFileSegment::Transactions)
537 .unwrap();
538
539 update_db_with_and_check::<tables::Transactions>(
541 &db,
542 current + 2,
543 Some(PipelineTarget::Unwind(89)),
544 &TxLegacy::default().into_signed(Signature::test_signature()).into(),
545 );
546
547 update_db_with_and_check::<tables::Transactions>(
549 &db,
550 current + 1,
551 None,
552 &TxLegacy::default().into_signed(Signature::test_signature()).into(),
553 );
554 }
555
556 #[test]
557 fn test_consistency_receipt_gap() {
558 let db = seed_data(90).unwrap();
559 let current = db
560 .factory
561 .static_file_provider()
562 .get_highest_static_file_tx(StaticFileSegment::Receipts)
563 .unwrap();
564
565 update_db_and_check::<tables::Receipts>(&db, current + 2, Some(PipelineTarget::Unwind(89)));
567
568 update_db_and_check::<tables::Receipts>(&db, current + 1, None);
570 }
571
572 #[test]
573 fn test_consistency_pruned_v2_segments() {
574 for prune_mode in [PruneMode::Distance(2_000_000), PruneMode::Before(90)] {
575 let db = seed_v2_data();
576
577 assert_consistency(&db, Some(PipelineTarget::Unwind(0)));
579
580 save_prune_checkpoints(
581 &db,
582 [
583 PruneSegment::SenderRecovery,
584 PruneSegment::AccountHistory,
585 PruneSegment::StorageHistory,
586 ]
587 .map(|segment| (segment, 89, prune_mode)),
588 );
589 assert_consistency(&db, None);
590 }
591 }
592
593 #[test]
594 fn test_consistency_unwind_bounded_by_prune_checkpoint() {
595 let db = seed_v2_data();
596
597 save_prune_checkpoints(
601 &db,
602 [
603 (PruneSegment::SenderRecovery, 50, PruneMode::Distance(39)),
604 (PruneSegment::AccountHistory, 89, PruneMode::Distance(2_000_000)),
605 (PruneSegment::StorageHistory, 89, PruneMode::Distance(2_000_000)),
606 ],
607 );
608 assert_consistency(&db, Some(PipelineTarget::Unwind(50)));
609 }
610
611 #[test]
612 fn test_consistency_receipts_distance_prune_checkpoint() {
613 let db = seed_data(90).unwrap();
614 let static_file_provider = db.factory.static_file_provider();
615
616 while let Some(block) =
619 static_file_provider.get_highest_static_file_block(StaticFileSegment::Receipts)
620 {
621 static_file_provider.delete_jar(StaticFileSegment::Receipts, block).unwrap();
622 }
623
624 assert_consistency(&db, Some(PipelineTarget::Unwind(0)));
625
626 save_prune_checkpoints(&db, [(PruneSegment::Receipts, 89, PruneMode::Distance(2_000_000))]);
627 assert_consistency(&db, None);
628 }
629}