1use super::SaveBlocksInput;
2use crate::{
3 changesets_utils::StorageRevertsIter,
4 providers::{
5 database::{chain::ChainStorage, metrics, DatabaseProviderMetrics},
6 rocksdb::{
7 OwnedRocksReadSnapshot, PendingRocksDBBatches, RocksDBProvider, RocksDBWriteCtx,
8 RocksReadSnapshot,
9 },
10 static_file::{StaticFileWriteCtx, StaticFileWriter},
11 NodeTypesForProvider, StaticFileProvider,
12 },
13 to_range,
14 traits::{
15 AccountExtReader, BlockSource, ChangeSetReader, ReceiptProvider, StageCheckpointWriter,
16 },
17 AccountReader, BlockBodyWriter, BlockExecutionWriter, BlockHashReader, BlockNumReader,
18 BlockReader, BlockWriter, BundleStateInit, ChainStateBlockReader, ChainStateBlockWriter,
19 DBProvider, DbTxProvider, EitherReader, EitherWriter, EitherWriterDestination, HashingWriter,
20 HeaderProvider, HeaderSyncGapProvider, HistoryWriter, LatestStateProviderRef,
21 OriginalValuesKnown, PersistenceFrontiers, ProviderError, PruneCheckpointReader,
22 PruneCheckpointWriter, RawRocksDBBatch, RevertsInit, RocksBatchArg, RocksDBProviderFactory,
23 StageCheckpointReader, StateWriter, StaticFileProviderFactory, StatsReader, StorageReader,
24 StorageTrieWriter, TransactionVariant, TransactionsProvider, TransactionsProviderExt,
25 TrieWriter,
26};
27use alloy_consensus::{
28 transaction::{SignerRecoverable, TransactionMeta, TxHashRef},
29 BlockHeader, TxReceipt,
30};
31use alloy_eips::BlockHashOrNumber;
32use alloy_primitives::{
33 keccak256,
34 map::{hash_map, AddressSet, B256Map, HashMap},
35 Address, BlockHash, BlockNumber, StorageKey, StorageValue, TxHash, TxNumber, B256,
36};
37use itertools::Itertools;
38use parking_lot::RwLock;
39use rayon::slice::ParallelSliceMut;
40use reth_chain_state::ExecutedBlock;
41use reth_chainspec::{ChainInfo, ChainSpecProvider, EthChainSpec};
42use reth_db_api::{
43 cursor::{DbCursorRO, DbCursorRW, DbDupCursorRO, DbDupCursorRW},
44 database::{Database, ReaderTxnTracker},
45 models::{
46 sharded_key, storage_sharded_key::StorageShardedKey, AccountBeforeTx, BlockNumberAddress,
47 BlockNumberAddressRange, ShardedKey, StorageBeforeTx, StorageSettings,
48 StoredBlockBodyIndices,
49 },
50 table::Table,
51 tables,
52 transaction::{DbTx, DbTxMut},
53 BlockNumberList,
54};
55use reth_execution_types::{
56 BlockExecutionOutput, BlockExecutionResult, Chain, ExecutionOutcome,
57 RecoveredBlockAndExecutionOutput,
58};
59use reth_node_types::{BlockTy, BodyTy, HeaderTy, NodeTypes, ReceiptTy, TxTy};
60use reth_primitives_traits::{
61 Account, Block as _, BlockBody as _, Bytecode, FastInstant as Instant, RecoveredBlock,
62 SealedHeader, StorageEntry,
63};
64use reth_prune_types::{
65 PruneCheckpoint, PruneMode, PruneModes, PruneSegment, MINIMUM_UNWIND_SAFE_DISTANCE,
66};
67use reth_stages_types::{FinishCheckpoint, StageCheckpoint, StageId};
68use reth_static_file_types::StaticFileSegment;
69use reth_storage_api::{
70 BlockBodyIndicesProvider, BlockBodyReader, HistoryInfo, HistoryReader, MetadataProvider,
71 MetadataWriter, NodePrimitivesProvider, StateProvider, StateWriteConfig,
72 StorageChangeSetReader, StoragePath, StorageSettingsCache, WriteStateInput,
73};
74use reth_storage_errors::provider::{ProviderResult, StaticFileWriterError};
75use reth_storage_overlay::OverlayManager;
76use reth_trie::{
77 updates::{StorageTrieUpdatesSorted, TrieUpdatesSorted},
78 HashedPostStateSorted,
79};
80use reth_trie_db::{DatabaseStorageTrieCursor, TrieTableAdapter};
81use revm::database::states::{
82 PlainStateReverts, PlainStorageChangeset, PlainStorageRevert, StateChangeset,
83};
84use smallvec::SmallVec;
85use std::{
86 cmp::Ordering,
87 collections::{BTreeMap, BTreeSet},
88 fmt::Debug,
89 ops::{Deref, DerefMut, Range, RangeBounds, RangeInclusive},
90 path::PathBuf,
91 sync::{Arc, OnceLock},
92};
93use tracing::{debug, instrument, trace};
94
95#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
97pub enum CommitOrder {
98 #[default]
100 Normal,
101 Unwind,
104}
105
106impl CommitOrder {
107 pub const fn is_unwind(&self) -> bool {
109 matches!(self, Self::Unwind)
110 }
111}
112
113pub type DatabaseProviderRO<DB, N> = DatabaseProvider<<DB as Database>::TX, N>;
115
116#[derive(Debug)]
121pub struct DatabaseProviderRW<DB: Database, N: NodeTypes>(
122 pub DatabaseProvider<<DB as Database>::TXMut, N>,
123);
124
125impl<DB: Database, N: NodeTypes> Deref for DatabaseProviderRW<DB, N> {
126 type Target = DatabaseProvider<<DB as Database>::TXMut, N>;
127
128 fn deref(&self) -> &Self::Target {
129 &self.0
130 }
131}
132
133impl<DB: Database, N: NodeTypes> DerefMut for DatabaseProviderRW<DB, N> {
134 fn deref_mut(&mut self) -> &mut Self::Target {
135 &mut self.0
136 }
137}
138
139impl<DB: Database, N: NodeTypes> AsRef<DatabaseProvider<<DB as Database>::TXMut, N>>
140 for DatabaseProviderRW<DB, N>
141{
142 fn as_ref(&self) -> &DatabaseProvider<<DB as Database>::TXMut, N> {
143 &self.0
144 }
145}
146
147impl<DB: Database, N: NodeTypes + 'static> DatabaseProviderRW<DB, N> {
148 pub fn commit(self) -> ProviderResult<()> {
150 self.0.commit()
151 }
152
153 pub fn into_tx(self) -> <DB as Database>::TXMut {
155 self.0.into_tx()
156 }
157
158 #[cfg(any(test, feature = "test-utils"))]
160 pub const fn with_minimum_pruning_distance(mut self, distance: u64) -> Self {
161 self.0.minimum_pruning_distance = distance;
162 self
163 }
164}
165
166impl<DB: Database, N: NodeTypes> From<DatabaseProviderRW<DB, N>>
167 for DatabaseProvider<<DB as Database>::TXMut, N>
168{
169 fn from(provider: DatabaseProviderRW<DB, N>) -> Self {
170 provider.0
171 }
172}
173
174#[derive(Debug, Clone, Copy, PartialEq, Eq)]
176enum SaveBlocksMode {
177 Full,
180 BlocksOnly,
184}
185
186impl SaveBlocksMode {
187 const fn with_state(self) -> bool {
189 matches!(self, Self::Full)
190 }
191}
192
193pub struct DatabaseProvider<TX, N: NodeTypes> {
196 tx: TX,
198 chain_spec: Arc<N::ChainSpec>,
200 static_file_provider: StaticFileProvider<N::Primitives>,
202 prune_modes: PruneModes,
204 storage: Arc<N::Storage>,
206 storage_settings: Arc<RwLock<StorageSettings>>,
208 rocksdb_provider: RocksDBProvider,
210 rocksdb_history_snapshot: OnceLock<Option<OwnedRocksReadSnapshot>>,
216 overlay_manager: OverlayManager<N::Primitives>,
218 runtime: reth_tasks::Runtime,
220 db_path: PathBuf,
222 pending_rocksdb_batches: PendingRocksDBBatches,
224 commit_order: CommitOrder,
226 minimum_pruning_distance: u64,
228 metrics: Arc<DatabaseProviderMetrics>,
230 reader_txn_tracker: Option<Arc<dyn ReaderTxnTracker>>,
232}
233
234impl<TX: Debug, N: NodeTypes> Debug for DatabaseProvider<TX, N> {
235 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
236 let mut s = f.debug_struct("DatabaseProvider");
237 s.field("tx", &self.tx)
238 .field("chain_spec", &self.chain_spec)
239 .field("static_file_provider", &self.static_file_provider)
240 .field("prune_modes", &self.prune_modes)
241 .field("storage", &self.storage)
242 .field("storage_settings", &self.storage_settings)
243 .field("rocksdb_provider", &self.rocksdb_provider)
244 .field("overlay_manager", &self.overlay_manager)
245 .field("runtime", &self.runtime)
246 .field("pending_rocksdb_batches", &"<pending batches>")
247 .field("commit_order", &self.commit_order)
248 .field("minimum_pruning_distance", &self.minimum_pruning_distance)
249 .field("reader_txn_tracker", &"<reader txn tracker>")
250 .finish()
251 }
252}
253
254impl<TX, N: NodeTypes> DatabaseProvider<TX, N> {
255 pub const fn prune_modes_ref(&self) -> &PruneModes {
257 &self.prune_modes
258 }
259
260 pub const fn with_minimum_pruning_distance(mut self, distance: u64) -> Self {
262 self.minimum_pruning_distance = distance;
263 self
264 }
265
266 pub(crate) fn with_reader_txn_tracker<T>(mut self, reader_txn_tracker: T) -> Self
268 where
269 T: ReaderTxnTracker + 'static,
270 {
271 self.reader_txn_tracker = Some(Arc::new(reader_txn_tracker));
272 self
273 }
274}
275
276impl<TX: DbTx + 'static, N: NodeTypes> DatabaseProvider<TX, N> {
277 fn commit_unwind(self) -> ProviderResult<()> {
288 let storage_v2 = self.cached_storage_settings().storage_v2;
289 let reader_txn_tracker = self.reader_txn_tracker.clone();
290 self.tx.commit()?;
291
292 if let Some(reader_txn_tracker) = reader_txn_tracker.as_ref() {
293 reader_txn_tracker.wait_for_pre_commit_readers();
294 }
295
296 if storage_v2 {
297 let batches = std::mem::take(&mut *self.pending_rocksdb_batches.lock());
298 for batch in batches {
299 self.rocksdb_provider.commit_batch(batch)?;
300 }
301 }
302
303 self.static_file_provider.commit()?;
304 Ok(())
305 }
306
307 pub fn latest<'a>(&'a self) -> Box<dyn StateProvider + 'a> {
309 trace!(target: "providers::db", "Returning latest state provider");
310 Box::new(LatestStateProviderRef::new(self))
311 }
312
313 #[cfg(feature = "test-utils")]
314 pub fn set_prune_modes(&mut self, prune_modes: PruneModes) {
316 self.prune_modes = prune_modes;
317 }
318}
319
320impl<TX, N: NodeTypes> NodePrimitivesProvider for DatabaseProvider<TX, N> {
321 type Primitives = N::Primitives;
322}
323
324impl<TX, N: NodeTypes> StaticFileProviderFactory for DatabaseProvider<TX, N> {
325 fn static_file_provider(&self) -> StaticFileProvider<Self::Primitives> {
327 self.static_file_provider.clone()
328 }
329
330 fn get_static_file_writer(
331 &self,
332 block: BlockNumber,
333 segment: StaticFileSegment,
334 ) -> ProviderResult<crate::providers::StaticFileProviderRWRefMut<'_, Self::Primitives>> {
335 self.static_file_provider.get_writer(block, segment)
336 }
337}
338
339impl<TX, N: NodeTypes> RocksDBProviderFactory for DatabaseProvider<TX, N> {
340 fn rocksdb_provider(&self) -> RocksDBProvider {
342 self.rocksdb_provider.clone()
343 }
344
345 fn set_pending_rocksdb_batch(&self, batch: rocksdb::WriteBatchWithTransaction<true>) {
346 self.pending_rocksdb_batches.lock().push(batch);
347 }
348
349 fn commit_pending_rocksdb_batches(&self) -> ProviderResult<()> {
350 let batches = std::mem::take(&mut *self.pending_rocksdb_batches.lock());
351 for batch in batches {
352 self.rocksdb_provider.commit_batch(batch)?;
353 }
354 Ok(())
355 }
356}
357
358impl<TX: Debug + Send, N: NodeTypes<ChainSpec: EthChainSpec + 'static>> ChainSpecProvider
359 for DatabaseProvider<TX, N>
360{
361 type ChainSpec = N::ChainSpec;
362
363 fn chain_spec(&self) -> Arc<Self::ChainSpec> {
364 self.chain_spec.clone()
365 }
366}
367
368impl<TX: DbTxMut, N: NodeTypes> DatabaseProvider<TX, N> {
369 #[expect(clippy::too_many_arguments)]
371 fn new_rw_inner(
372 tx: TX,
373 chain_spec: Arc<N::ChainSpec>,
374 static_file_provider: StaticFileProvider<N::Primitives>,
375 prune_modes: PruneModes,
376 storage: Arc<N::Storage>,
377 storage_settings: Arc<RwLock<StorageSettings>>,
378 rocksdb_provider: RocksDBProvider,
379 overlay_manager: OverlayManager<N::Primitives>,
380 runtime: reth_tasks::Runtime,
381 db_path: PathBuf,
382 commit_order: CommitOrder,
383 metrics: Arc<DatabaseProviderMetrics>,
384 ) -> Self {
385 Self {
386 tx,
387 chain_spec,
388 static_file_provider,
389 prune_modes,
390 storage,
391 storage_settings,
392 rocksdb_provider,
393 overlay_manager,
394 runtime,
395 db_path,
396 rocksdb_history_snapshot: OnceLock::new(),
397 pending_rocksdb_batches: Default::default(),
398 commit_order,
399 minimum_pruning_distance: MINIMUM_UNWIND_SAFE_DISTANCE,
400 metrics,
401 reader_txn_tracker: None,
402 }
403 }
404
405 #[expect(clippy::too_many_arguments)]
407 pub fn new_rw(
408 tx: TX,
409 chain_spec: Arc<N::ChainSpec>,
410 static_file_provider: StaticFileProvider<N::Primitives>,
411 prune_modes: PruneModes,
412 storage: Arc<N::Storage>,
413 storage_settings: Arc<RwLock<StorageSettings>>,
414 rocksdb_provider: RocksDBProvider,
415 overlay_manager: OverlayManager<N::Primitives>,
416 runtime: reth_tasks::Runtime,
417 db_path: PathBuf,
418 metrics: Arc<DatabaseProviderMetrics>,
419 ) -> Self {
420 Self::new_rw_inner(
421 tx,
422 chain_spec,
423 static_file_provider,
424 prune_modes,
425 storage,
426 storage_settings,
427 rocksdb_provider,
428 overlay_manager,
429 runtime,
430 db_path,
431 CommitOrder::Normal,
432 metrics,
433 )
434 }
435
436 #[expect(clippy::too_many_arguments)]
438 pub fn new_unwind_rw(
439 tx: TX,
440 chain_spec: Arc<N::ChainSpec>,
441 static_file_provider: StaticFileProvider<N::Primitives>,
442 prune_modes: PruneModes,
443 storage: Arc<N::Storage>,
444 storage_settings: Arc<RwLock<StorageSettings>>,
445 rocksdb_provider: RocksDBProvider,
446 overlay_manager: OverlayManager<N::Primitives>,
447 runtime: reth_tasks::Runtime,
448 db_path: PathBuf,
449 metrics: Arc<DatabaseProviderMetrics>,
450 ) -> Self {
451 Self::new_rw_inner(
452 tx,
453 chain_spec,
454 static_file_provider,
455 prune_modes,
456 storage,
457 storage_settings,
458 rocksdb_provider,
459 overlay_manager,
460 runtime,
461 db_path,
462 CommitOrder::Unwind,
463 metrics,
464 )
465 }
466}
467
468impl<TX, N: NodeTypes> AsRef<Self> for DatabaseProvider<TX, N> {
469 fn as_ref(&self) -> &Self {
470 self
471 }
472}
473
474impl<TX: DbTx + DbTxMut + 'static, N: NodeTypesForProvider> DatabaseProvider<TX, N> {
475 pub fn with_rocksdb_batch<F, R>(&self, f: F) -> ProviderResult<R>
479 where
480 F: FnOnce(RocksBatchArg<'_>) -> ProviderResult<(R, Option<RawRocksDBBatch>)>,
481 {
482 let rocksdb = self.rocksdb_provider();
483 let rocksdb_batch = rocksdb.batch();
484
485 let (result, raw_batch) = f(rocksdb_batch)?;
486
487 if let Some(batch) = raw_batch {
488 self.set_pending_rocksdb_batch(batch);
489 }
490 let _ = raw_batch; Ok(result)
493 }
494
495 fn static_file_write_ctx(
497 &self,
498 save_mode: SaveBlocksMode,
499 first_block: BlockNumber,
500 last_block: BlockNumber,
501 ) -> ProviderResult<StaticFileWriteCtx> {
502 let tip = self.last_block_number()?.max(last_block);
503 Ok(StaticFileWriteCtx {
504 write_senders: EitherWriterDestination::senders(self).is_static_file() &&
505 self.prune_modes.sender_recovery.is_none_or(|m| !m.is_full()),
506 write_receipts: save_mode.with_state() &&
507 EitherWriter::receipts_destination(self).is_static_file(),
508 write_account_changesets: save_mode.with_state() &&
509 EitherWriterDestination::account_changesets(self).is_static_file(),
510 write_storage_changesets: save_mode.with_state() &&
511 EitherWriterDestination::storage_changesets(self).is_static_file(),
512 tip,
513 receipts_prune_mode: self.prune_modes.receipts,
514 receipts_prunable: self
516 .static_file_provider
517 .get_highest_static_file_tx(StaticFileSegment::Receipts)
518 .is_none() &&
519 PruneMode::Distance(self.minimum_pruning_distance)
520 .should_prune(first_block, tip),
521 })
522 }
523
524 fn rocksdb_write_ctx(&self, first_block: BlockNumber) -> RocksDBWriteCtx {
526 RocksDBWriteCtx {
527 first_block_number: first_block,
528 prune_tx_lookup: self.prune_modes.transaction_lookup,
529 storage_settings: self.cached_storage_settings(),
530 pending_batches: self.pending_rocksdb_batches.clone(),
531 }
532 }
533
534 #[instrument(level = "debug", target = "providers::db", skip_all, fields(block_count = input.persist_rest_blocks().len()))]
539 pub fn save_blocks(&self, input: &SaveBlocksInput<N::Primitives>) -> ProviderResult<()> {
540 let (db_tip, partial_state_trie) = self
541 .get_stage_checkpoint(StageId::Finish)?
542 .map(|checkpoint| {
543 let partial_state_trie = checkpoint
544 .finish_stage_checkpoint()
545 .and_then(|finish| finish.partial_state_trie())
546 .unwrap_or(checkpoint.block_number);
547 (checkpoint.block_number, partial_state_trie)
548 })
549 .unwrap_or_default();
550
551 if db_tip != input.prev_db_tip() || partial_state_trie != input.prev_partial_state_trie() {
552 return Err(ProviderError::other(std::io::Error::other(format!(
553 "persistence frontiers do not match Finish checkpoint: expected database/state-trie tips #{}/{}, got #{}/{}",
554 input.prev_db_tip(),
555 input.prev_partial_state_trie(),
556 db_tip,
557 partial_state_trie,
558 ))))
559 }
560
561 self.save_blocks_inner(
562 input.persist_rest_blocks(),
563 input.state_trie_blocks(),
564 input.state_trie_masking_blocks(),
565 (input.new_partial_state_trie() < input.new_db_tip())
566 .then_some(input.new_partial_state_trie()),
567 SaveBlocksMode::Full,
568 )
569 }
570
571 fn save_blocks_inner(
572 &self,
573 blocks: &[ExecutedBlock<N::Primitives>],
574 state_trie_blocks: &[ExecutedBlock<N::Primitives>],
575 state_trie_masking_blocks: &[ExecutedBlock<N::Primitives>],
576 partial_state_trie: Option<BlockNumber>,
577 save_mode: SaveBlocksMode,
578 ) -> ProviderResult<()> {
579 let total_start = Instant::now();
580 let block_count = blocks.len() as u64;
581 let last_block_number = blocks
584 .last()
585 .or_else(|| state_trie_masking_blocks.last())
586 .or_else(|| state_trie_blocks.last())
587 .expect("at least one persistence range must be non-empty")
588 .recovered_block()
589 .number();
590 let first_number = blocks.first().map(|block| block.recovered_block().number());
591
592 debug!(target: "providers::db", block_count, "Writing blocks and execution data to storage");
593
594 let mut tx_nums: SmallVec<[TxNumber; 4]> = SmallVec::with_capacity(blocks.len());
596 if !blocks.is_empty() {
597 let first_tx_num = self
598 .tx
599 .cursor_read::<tables::TransactionBlocks>()?
600 .last()?
601 .map(|(n, _)| n + 1)
602 .unwrap_or_default();
603 let mut current = first_tx_num;
604 for block in blocks {
605 tx_nums.push(current);
606 current += block.recovered_block().body().transaction_count() as u64;
607 }
608 }
609
610 let mut timings =
611 metrics::SaveBlocksTimings { batch_size: block_count, ..Default::default() };
612
613 let sf_provider = &self.static_file_provider;
615 let rocksdb_provider = &self.rocksdb_provider;
616 let sf_ctx = first_number
617 .map(|first_number| {
618 self.static_file_write_ctx(save_mode, first_number, last_block_number)
619 })
620 .transpose()?;
621 let rocksdb_ctx = first_number.map(|first_number| self.rocksdb_write_ctx(first_number));
622 let rocksdb_enabled =
623 rocksdb_ctx.as_ref().is_some_and(|ctx| ctx.storage_settings.storage_v2);
624
625 let mut sf_result = None;
626 let mut rocksdb_result = None;
627
628 let runtime = &self.runtime;
630 let span = tracing::Span::current();
633 runtime.storage_pool().in_place_scope(|s| {
634 if sf_ctx.is_some() {
636 s.spawn(|_| {
637 let _guard = span.enter();
638 let start = Instant::now();
639 let sf_ctx =
640 sf_ctx.expect("static file context exists when blocks are persisted");
641 sf_result = Some(
642 sf_provider
643 .write_blocks_data(blocks, &tx_nums, sf_ctx, runtime)
644 .map(|()| start.elapsed()),
645 );
646 });
647 }
648
649 if rocksdb_enabled {
651 s.spawn(|_| {
652 let _guard = span.enter();
653 let start = Instant::now();
654 let rocksdb_ctx =
655 rocksdb_ctx.clone().expect("RocksDB context exists when enabled");
656 rocksdb_result = Some(
657 rocksdb_provider
658 .write_blocks_data(blocks, &tx_nums, rocksdb_ctx, runtime)
659 .map(|()| start.elapsed()),
660 );
661 });
662 }
663
664 let mdbx_start = Instant::now();
666
667 if !blocks.is_empty() &&
669 !self.cached_storage_settings().storage_v2 &&
670 self.prune_modes.transaction_lookup.is_none_or(|m| !m.is_full())
671 {
672 let start = Instant::now();
673 let total_tx_count: usize =
674 blocks.iter().map(|b| b.recovered_block().body().transaction_count()).sum();
675 let mut all_tx_hashes = Vec::with_capacity(total_tx_count);
676 for (i, block) in blocks.iter().enumerate() {
677 let recovered_block = block.recovered_block();
678 for (tx_num, transaction) in
679 (tx_nums[i]..).zip(recovered_block.body().transactions_iter())
680 {
681 all_tx_hashes.push((*transaction.tx_hash(), tx_num));
682 }
683 }
684
685 all_tx_hashes.sort_unstable_by_key(|(hash, _)| *hash);
687
688 self.with_rocksdb_batch(|batch| {
690 let mut tx_hash_writer =
691 EitherWriter::new_transaction_hash_numbers(self, batch)?;
692 tx_hash_writer.put_transaction_hash_numbers_batch(all_tx_hashes, false)?;
693 let raw_batch = tx_hash_writer.into_raw_rocksdb_batch();
694 Ok(((), raw_batch))
695 })?;
696 self.metrics.record_duration(
697 metrics::Action::InsertTransactionHashNumbers,
698 start.elapsed(),
699 );
700 }
701
702 for (i, block) in blocks.iter().enumerate() {
703 let recovered_block = block.recovered_block();
704
705 let start = Instant::now();
706 self.insert_block_mdbx_only(recovered_block, tx_nums[i])?;
707 timings.insert_block += start.elapsed();
708
709 if save_mode.with_state() {
710 let execution_output = block.execution_outcome();
711 let sf_ctx =
712 sf_ctx.expect("static file context exists when blocks are persisted");
713
714 let start = Instant::now();
718 self.write_state(
719 WriteStateInput::Single {
720 outcome: execution_output,
721 block: recovered_block.number(),
722 },
723 OriginalValuesKnown::No,
724 StateWriteConfig {
725 write_receipts: !sf_ctx.write_receipts,
726 write_account_changesets: !sf_ctx.write_account_changesets,
727 write_storage_changesets: !sf_ctx.write_storage_changesets,
728 },
729 )?;
730 timings.write_state += start.elapsed();
731 }
732 }
733
734 if save_mode.with_state() && !state_trie_blocks.is_empty() {
737 let start = Instant::now();
738 let batch = ExecutedBlock::hashed_state_refs(state_trie_blocks);
739 let mask = ExecutedBlock::hashed_state_refs(state_trie_masking_blocks);
740 let merged_hashed_state =
741 HashedPostStateSorted::disjointed_merge_batch(&batch, &mask);
742 if !merged_hashed_state.is_empty() {
743 self.write_hashed_state(&merged_hashed_state)?;
744 }
745 timings.write_hashed_state += start.elapsed();
746
747 let start = Instant::now();
748 let batch = ExecutedBlock::trie_updates_refs(state_trie_blocks);
749 let mask = ExecutedBlock::trie_updates_refs(state_trie_masking_blocks);
750 let merged_trie =
751 Arc::new(TrieUpdatesSorted::disjointed_merge_batch(&batch, &mask));
752 if !merged_trie.is_empty() {
753 self.write_trie_updates_sorted(&merged_trie)?;
754 }
755 timings.write_trie_updates += start.elapsed();
756 }
757
758 if save_mode.with_state() &&
760 let Some(first_number) = first_number
761 {
762 let start = Instant::now();
763 self.update_history_indices(first_number..=last_block_number)?;
764 timings.update_history_indices = start.elapsed();
765 }
766
767 let start = Instant::now();
769 if !blocks.is_empty() {
770 self.update_pipeline_stages(last_block_number, false)?;
771 }
772 if save_mode.with_state() {
773 let checkpoint = match partial_state_trie {
774 Some(partial_state_trie) => StageCheckpoint::new(last_block_number)
775 .with_finish_stage_checkpoint(FinishCheckpoint {
776 partial_state_trie: Some(partial_state_trie),
777 }),
778 None => StageCheckpoint::new(last_block_number),
779 };
780 self.save_stage_checkpoint(StageId::Finish, checkpoint)?;
781 }
782 timings.update_pipeline_stages = start.elapsed();
783
784 timings.mdbx = mdbx_start.elapsed();
785
786 Ok::<_, ProviderError>(())
787 })?;
788
789 if !blocks.is_empty() {
791 timings.sf = sf_result.ok_or(StaticFileWriterError::ThreadPanic("static file"))??;
792 }
793
794 if rocksdb_enabled {
795 timings.rocksdb = rocksdb_result.ok_or_else(|| {
796 ProviderError::Database(reth_db_api::DatabaseError::Other(
797 "RocksDB thread panicked".into(),
798 ))
799 })??;
800 }
801
802 timings.total = total_start.elapsed();
803
804 self.metrics.record_save_blocks(&timings);
805 if let Some(first_number) = first_number {
806 debug!(target: "providers::db", range = ?first_number..=last_block_number, "Appended block data");
807 }
808
809 Ok(())
810 }
811
812 #[instrument(level = "debug", target = "providers::db", skip_all)]
816 fn insert_block_mdbx_only(
817 &self,
818 block: &RecoveredBlock<BlockTy<N>>,
819 first_tx_num: TxNumber,
820 ) -> ProviderResult<StoredBlockBodyIndices> {
821 if self.prune_modes.sender_recovery.is_none_or(|m| !m.is_full()) &&
822 EitherWriterDestination::senders(self).is_database()
823 {
824 let start = Instant::now();
825 let tx_nums_iter = std::iter::successors(Some(first_tx_num), |n| Some(n + 1));
826 let mut cursor = self.tx.cursor_write::<tables::TransactionSenders>()?;
827 for (tx_num, sender) in tx_nums_iter.zip(block.senders_iter().copied()) {
828 cursor.append(tx_num, &sender)?;
829 }
830 self.metrics
831 .record_duration(metrics::Action::InsertTransactionSenders, start.elapsed());
832 }
833
834 let block_number = block.number();
835 let tx_count = block.body().transaction_count() as u64;
836
837 let start = Instant::now();
838 self.tx.put::<tables::HeaderNumbers>(block.hash(), block_number)?;
839 self.metrics.record_duration(metrics::Action::InsertHeaderNumbers, start.elapsed());
840
841 self.write_block_body_indices(block_number, block.body(), first_tx_num, tx_count)?;
842
843 Ok(StoredBlockBodyIndices { first_tx_num, tx_count })
844 }
845
846 fn write_block_body_indices(
849 &self,
850 block_number: BlockNumber,
851 body: &BodyTy<N>,
852 first_tx_num: TxNumber,
853 tx_count: u64,
854 ) -> ProviderResult<()> {
855 let start = Instant::now();
857 self.tx
858 .cursor_write::<tables::BlockBodyIndices>()?
859 .append(block_number, &StoredBlockBodyIndices { first_tx_num, tx_count })?;
860 self.metrics.record_duration(metrics::Action::InsertBlockBodyIndices, start.elapsed());
861
862 if tx_count > 0 {
864 let start = Instant::now();
865 self.tx
866 .cursor_write::<tables::TransactionBlocks>()?
867 .append(first_tx_num + tx_count - 1, &block_number)?;
868 self.metrics.record_duration(metrics::Action::InsertTransactionBlocks, start.elapsed());
869 }
870
871 self.storage.writer().write_block_bodies(self, vec![(block_number, Some(body))])?;
873
874 Ok(())
875 }
876
877 pub fn unwind_trie_state_from(&self, from: BlockNumber) -> ProviderResult<()> {
882 let changed_accounts = self.account_changesets_range(from..)?;
883
884 self.unwind_account_hashing(changed_accounts.iter())?;
886
887 self.unwind_account_history_indices(changed_accounts.iter())?;
889
890 let changed_storages = self.storage_changesets_range(from..)?;
891
892 self.unwind_storage_hashing(changed_storages.iter().copied())?;
894
895 self.unwind_storage_history_indices(changed_storages.iter().copied())?;
897
898 let db_tip_block = self
901 .get_stage_checkpoint(reth_stages_types::StageId::Finish)?
902 .as_ref()
903 .map(|chk| chk.block_number)
904 .ok_or_else(|| ProviderError::InsufficientChangesets {
905 requested: from,
906 available: 0..=0,
907 })?;
908
909 let trie_revert = self
910 .overlay_manager
911 .get_or_compute_cached_changesets_range(self, from..=db_tip_block)?;
912 self.write_trie_updates_sorted(&trie_revert)?;
913
914 Ok(())
915 }
916
917 fn remove_receipts_from(
919 &self,
920 from_tx: TxNumber,
921 last_block: BlockNumber,
922 ) -> ProviderResult<()> {
923 self.remove::<tables::Receipts<ReceiptTy<N>>>(from_tx..)?;
925
926 if EitherWriter::receipts_destination(self).is_static_file() {
927 let static_file_receipt_num =
928 self.static_file_provider.get_highest_static_file_tx(StaticFileSegment::Receipts);
929
930 let to_delete = static_file_receipt_num
931 .map(|static_num| (static_num + 1).saturating_sub(from_tx))
932 .unwrap_or_default();
933
934 self.static_file_provider
935 .latest_writer(StaticFileSegment::Receipts)?
936 .prune_receipts(to_delete, last_block)?;
937 }
938
939 Ok(())
940 }
941
942 fn write_bytecodes(
944 &self,
945 bytecodes: impl IntoIterator<Item = (B256, Bytecode)>,
946 ) -> ProviderResult<()> {
947 let mut bytecodes_cursor = self.tx_ref().cursor_write::<tables::Bytecodes>()?;
948 for (hash, bytecode) in bytecodes {
949 bytecodes_cursor.upsert(hash, &bytecode)?;
950 }
951 Ok(())
952 }
953}
954
955fn unwind_history_shards<S, T, C>(
970 cursor: &mut C,
971 start_key: T::Key,
972 block_number: BlockNumber,
973 mut shard_belongs_to_key: impl FnMut(&T::Key) -> bool,
974) -> ProviderResult<Vec<u64>>
975where
976 T: Table<Value = BlockNumberList>,
977 T::Key: AsRef<ShardedKey<S>>,
978 C: DbCursorRO<T> + DbCursorRW<T>,
979{
980 let mut item = cursor.seek_exact(start_key)?;
982 while let Some((sharded_key, list)) = item {
983 if !shard_belongs_to_key(&sharded_key) {
985 break
986 }
987
988 cursor.delete_current()?;
991
992 let first = list.iter().next().expect("List can't be empty");
995
996 if first >= block_number {
999 item = cursor.prev()?;
1000 continue
1001 }
1002 else if block_number <= sharded_key.as_ref().highest_block_number {
1005 return Ok(list.iter().take_while(|i| *i < block_number).collect::<Vec<_>>())
1008 }
1009 return Ok(list.iter().collect::<Vec<_>>())
1012 }
1013
1014 Ok(Vec::new())
1016}
1017
1018impl<TX: DbTx + 'static, N: NodeTypesForProvider> DatabaseProvider<TX, N> {
1019 #[expect(clippy::too_many_arguments)]
1021 pub fn new(
1022 tx: TX,
1023 chain_spec: Arc<N::ChainSpec>,
1024 static_file_provider: StaticFileProvider<N::Primitives>,
1025 prune_modes: PruneModes,
1026 storage: Arc<N::Storage>,
1027 storage_settings: Arc<RwLock<StorageSettings>>,
1028 rocksdb_provider: RocksDBProvider,
1029 overlay_manager: OverlayManager<N::Primitives>,
1030 runtime: reth_tasks::Runtime,
1031 db_path: PathBuf,
1032 metrics: Arc<DatabaseProviderMetrics>,
1033 ) -> Self {
1034 Self {
1035 tx,
1036 chain_spec,
1037 static_file_provider,
1038 prune_modes,
1039 storage,
1040 storage_settings,
1041 rocksdb_provider,
1042 overlay_manager,
1043 runtime,
1044 db_path,
1045 rocksdb_history_snapshot: OnceLock::new(),
1046 pending_rocksdb_batches: Default::default(),
1047 commit_order: CommitOrder::Normal,
1048 minimum_pruning_distance: MINIMUM_UNWIND_SAFE_DISTANCE,
1049 metrics,
1050 reader_txn_tracker: None,
1051 }
1052 }
1053
1054 pub fn into_tx(self) -> TX {
1056 self.tx
1057 }
1058
1059 pub const fn tx_mut(&mut self) -> &mut TX {
1061 &mut self.tx
1062 }
1063
1064 pub const fn tx_ref(&self) -> &TX {
1066 &self.tx
1067 }
1068
1069 pub fn chain_spec(&self) -> &N::ChainSpec {
1071 &self.chain_spec
1072 }
1073}
1074
1075impl<TX: DbTx + 'static, N: NodeTypesForProvider> DatabaseProvider<TX, N> {
1076 fn recovered_block<H, HF, B, BF>(
1077 &self,
1078 id: BlockHashOrNumber,
1079 _transaction_kind: TransactionVariant,
1080 header_by_number: HF,
1081 construct_block: BF,
1082 ) -> ProviderResult<Option<B>>
1083 where
1084 H: AsRef<HeaderTy<N>>,
1085 HF: FnOnce(BlockNumber) -> ProviderResult<Option<H>>,
1086 BF: FnOnce(H, BodyTy<N>, Vec<Address>) -> ProviderResult<Option<B>>,
1087 {
1088 let Some(block_number) = self.convert_hash_or_number(id)? else { return Ok(None) };
1089 let earliest_available = self.static_file_provider.earliest_history_height();
1090 if block_number < earliest_available {
1091 return Err(ProviderError::BlockExpired { requested: block_number, earliest_available })
1092 }
1093 let Some(header) = header_by_number(block_number)? else { return Ok(None) };
1094
1095 let Some(body) = self.block_body_indices(block_number)? else { return Ok(None) };
1102
1103 let tx_range = body.tx_num_range();
1104
1105 let transactions = if tx_range.is_empty() {
1106 vec![]
1107 } else {
1108 self.transactions_by_tx_range(tx_range.clone())?
1109 };
1110
1111 let body = self
1112 .storage
1113 .reader()
1114 .read_block_bodies(self, vec![(header.as_ref(), transactions)])?
1115 .pop()
1116 .ok_or(ProviderError::InvalidStorageOutput)?;
1117
1118 let senders = if tx_range.is_empty() {
1119 vec![]
1120 } else {
1121 let known_senders: HashMap<TxNumber, Address> =
1122 EitherReader::new_senders(self)?.senders_by_tx_range(tx_range.clone())?;
1123
1124 let mut senders = Vec::with_capacity(body.transactions().len());
1125 for (tx_num, tx) in tx_range.zip(body.transactions()) {
1126 match known_senders.get(&tx_num) {
1127 None => {
1128 let sender = tx.recover_signer_unchecked()?;
1129 senders.push(sender);
1130 }
1131 Some(sender) => senders.push(*sender),
1132 }
1133 }
1134 senders
1135 };
1136
1137 construct_block(header, body, senders)
1138 }
1139
1140 fn block_range<F, H, HF, R>(
1150 &self,
1151 range: RangeInclusive<BlockNumber>,
1152 headers_range: HF,
1153 mut assemble_block: F,
1154 ) -> ProviderResult<Vec<R>>
1155 where
1156 H: AsRef<HeaderTy<N>>,
1157 HF: FnOnce(RangeInclusive<BlockNumber>) -> ProviderResult<Vec<H>>,
1158 F: FnMut(H, BodyTy<N>, Range<TxNumber>) -> ProviderResult<R>,
1159 {
1160 if range.is_empty() {
1161 return Ok(Vec::new())
1162 }
1163
1164 let earliest_available = self.static_file_provider.earliest_history_height();
1167 if *range.start() < earliest_available {
1168 return Err(ProviderError::BlockExpired {
1169 requested: *range.start(),
1170 earliest_available,
1171 })
1172 }
1173
1174 let len = range.end().saturating_sub(*range.start()) as usize + 1;
1175 let mut blocks = Vec::with_capacity(len);
1176
1177 let headers = headers_range(range.clone())?;
1178
1179 let present_headers = self
1185 .block_body_indices_range(range)?
1186 .into_iter()
1187 .map(|b| b.tx_num_range())
1188 .zip(headers)
1189 .collect::<Vec<_>>();
1190
1191 let mut inputs = Vec::with_capacity(present_headers.len());
1192 for (tx_range, header) in &present_headers {
1193 let transactions = if tx_range.is_empty() {
1194 Vec::new()
1195 } else {
1196 self.transactions_by_tx_range(tx_range.clone())?
1197 };
1198
1199 inputs.push((header.as_ref(), transactions));
1200 }
1201
1202 let bodies = self.storage.reader().read_block_bodies(self, inputs)?;
1203
1204 for ((tx_range, header), body) in present_headers.into_iter().zip(bodies) {
1205 blocks.push(assemble_block(header, body, tx_range)?);
1206 }
1207
1208 Ok(blocks)
1209 }
1210
1211 fn block_with_senders_range<H, HF, B, BF>(
1222 &self,
1223 range: RangeInclusive<BlockNumber>,
1224 headers_range: HF,
1225 assemble_block: BF,
1226 ) -> ProviderResult<Vec<B>>
1227 where
1228 H: AsRef<HeaderTy<N>>,
1229 HF: Fn(RangeInclusive<BlockNumber>) -> ProviderResult<Vec<H>>,
1230 BF: Fn(H, BodyTy<N>, Vec<Address>) -> ProviderResult<B>,
1231 {
1232 self.block_range(range, headers_range, |header, body, tx_range| {
1233 let senders = if tx_range.is_empty() {
1234 Vec::new()
1235 } else {
1236 let known_senders: HashMap<TxNumber, Address> =
1237 EitherReader::new_senders(self)?.senders_by_tx_range(tx_range.clone())?;
1238
1239 let mut senders = Vec::with_capacity(body.transactions().len());
1240 for (tx_num, tx) in tx_range.zip(body.transactions()) {
1241 match known_senders.get(&tx_num) {
1242 None => {
1243 let sender = tx.recover_signer_unchecked()?;
1245 senders.push(sender);
1246 }
1247 Some(sender) => senders.push(*sender),
1248 }
1249 }
1250
1251 senders
1252 };
1253
1254 assemble_block(header, body, senders)
1255 })
1256 }
1257
1258 #[allow(clippy::clone_on_copy)]
1262 pub(crate) fn populate_bundle_state(
1263 &self,
1264 account_changeset: Vec<(u64, AccountBeforeTx)>,
1265 storage_changeset: Vec<(BlockNumberAddress, StorageEntry)>,
1266 mut get_account: impl FnMut(Address) -> ProviderResult<Option<Account>>,
1267 mut get_storage: impl FnMut(Address, StorageKey) -> ProviderResult<Option<StorageValue>>,
1268 ) -> ProviderResult<(BundleStateInit, RevertsInit)> {
1269 let mut state: BundleStateInit = HashMap::default();
1273
1274 let mut reverts: RevertsInit = HashMap::default();
1280
1281 for (block_number, account_before) in account_changeset.into_iter().rev() {
1283 let AccountBeforeTx { info: old_info, address } = account_before;
1284 match state.entry(address) {
1285 hash_map::Entry::Vacant(entry) => {
1286 let new_info = get_account(address)?;
1287 entry.insert((old_info.clone(), new_info, HashMap::default()));
1288 }
1289 hash_map::Entry::Occupied(mut entry) => {
1290 entry.get_mut().0 = old_info.clone();
1292 }
1293 }
1294 reverts.entry(block_number).or_default().entry(address).or_default().0 = Some(old_info);
1296 }
1297
1298 for (block_and_address, old_storage) in storage_changeset.into_iter().rev() {
1300 let BlockNumberAddress((block_number, address)) = block_and_address;
1301 let account_state = match state.entry(address) {
1303 hash_map::Entry::Vacant(entry) => {
1304 let present_info = get_account(address)?;
1305 entry.insert((present_info.clone(), present_info, HashMap::default()))
1306 }
1307 hash_map::Entry::Occupied(entry) => entry.into_mut(),
1308 };
1309
1310 match account_state.2.entry(old_storage.key) {
1312 hash_map::Entry::Vacant(entry) => {
1313 let new_storage = get_storage(address, old_storage.key)?.unwrap_or_default();
1314 entry.insert((old_storage.value, new_storage));
1315 }
1316 hash_map::Entry::Occupied(mut entry) => {
1317 entry.get_mut().0 = old_storage.value;
1318 }
1319 };
1320
1321 reverts
1322 .entry(block_number)
1323 .or_default()
1324 .entry(address)
1325 .or_default()
1326 .1
1327 .push(old_storage);
1328 }
1329
1330 Ok((state, reverts))
1331 }
1332
1333 fn populate_bundle_state_plain(
1336 &self,
1337 account_changeset: Vec<(u64, AccountBeforeTx)>,
1338 storage_changeset: Vec<(BlockNumberAddress, StorageEntry)>,
1339 plain_accounts_cursor: &mut impl DbCursorRO<tables::PlainAccountState>,
1340 plain_storage_cursor: &mut impl DbDupCursorRO<tables::PlainStorageState>,
1341 ) -> ProviderResult<(BundleStateInit, RevertsInit)> {
1342 self.populate_bundle_state(
1343 account_changeset,
1344 storage_changeset,
1345 |address| Ok(plain_accounts_cursor.seek_exact(address)?.map(|kv| kv.1)),
1346 |address, storage_key| {
1347 Ok(plain_storage_cursor
1348 .seek_by_key_subkey(address, storage_key)?
1349 .filter(|s| s.key == storage_key)
1350 .map(|s| s.value))
1351 },
1352 )
1353 }
1354
1355 fn populate_bundle_state_hashed(
1360 &self,
1361 account_changeset: Vec<(u64, AccountBeforeTx)>,
1362 storage_changeset: Vec<(BlockNumberAddress, StorageEntry)>,
1363 hashed_accounts_cursor: &mut impl DbCursorRO<tables::HashedAccounts>,
1364 hashed_storage_cursor: &mut impl DbDupCursorRO<tables::HashedStorages>,
1365 ) -> ProviderResult<(BundleStateInit, RevertsInit)> {
1366 self.populate_bundle_state(
1367 account_changeset,
1368 storage_changeset,
1369 |address| Ok(hashed_accounts_cursor.seek_exact(keccak256(address))?.map(|kv| kv.1)),
1370 |address, storage_key| {
1371 let hashed_storage_key = keccak256(storage_key);
1372 Ok(hashed_storage_cursor
1373 .seek_by_key_subkey(keccak256(address), hashed_storage_key)?
1374 .filter(|s| s.key == hashed_storage_key)
1375 .map(|s| s.value))
1376 },
1377 )
1378 }
1379}
1380
1381impl<TX: DbTxMut + DbTx + 'static, N: NodeTypes> DatabaseProvider<TX, N> {
1382 fn append_history_index<P, T>(
1390 &self,
1391 index_updates: impl IntoIterator<Item = (P, impl IntoIterator<Item = u64>)>,
1392 mut sharded_key_factory: impl FnMut(P, BlockNumber) -> T::Key,
1393 ) -> ProviderResult<()>
1394 where
1395 P: Copy,
1396 T: Table<Value = BlockNumberList>,
1397 {
1398 assert!(!T::DUPSORT, "append_history_index cannot be used with DUPSORT tables");
1401
1402 let mut cursor = self.tx.cursor_write::<T>()?;
1403
1404 for (partial_key, indices) in index_updates {
1405 let last_key = sharded_key_factory(partial_key, u64::MAX);
1406 let mut last_shard = cursor
1407 .seek_exact(last_key.clone())?
1408 .map(|(_, list)| list)
1409 .unwrap_or_else(BlockNumberList::empty);
1410
1411 last_shard.append(indices).map_err(ProviderError::other)?;
1412
1413 if last_shard.len() <= sharded_key::NUM_OF_INDICES_IN_SHARD as u64 {
1415 cursor.upsert(last_key, &last_shard)?;
1416 continue;
1417 }
1418
1419 let chunks = last_shard.iter().chunks(sharded_key::NUM_OF_INDICES_IN_SHARD);
1421 let mut chunks_peekable = chunks.into_iter().peekable();
1422
1423 while let Some(chunk) = chunks_peekable.next() {
1424 let shard = BlockNumberList::new_pre_sorted(chunk);
1425 let highest_block_number = if chunks_peekable.peek().is_some() {
1426 shard.iter().next_back().expect("`chunks` does not return empty list")
1427 } else {
1428 u64::MAX
1430 };
1431
1432 cursor.upsert(sharded_key_factory(partial_key, highest_block_number), &shard)?;
1433 }
1434 }
1435
1436 Ok(())
1437 }
1438}
1439
1440impl<TX: DbTx, N: NodeTypes> DatabaseProvider<TX, N> {
1441 fn ensure_finish_may_advance(&self, proposed: &StageCheckpoint) -> ProviderResult<()> {
1446 let current = self.get_stage_checkpoint(StageId::Finish)?.unwrap_or_default();
1447 let state_trie_frontier = |checkpoint: &StageCheckpoint| {
1448 checkpoint
1449 .finish_stage_checkpoint()
1450 .and_then(|finish| finish.partial_state_trie())
1451 .unwrap_or(checkpoint.block_number)
1452 };
1453 if proposed.block_number <= current.block_number &&
1454 state_trie_frontier(proposed) <= state_trie_frontier(¤t)
1455 {
1456 return Ok(())
1457 }
1458 match self.snap_attempt()? {
1459 Some(attempt) if !attempt.is_verified() => {
1460 Err(ProviderError::UnverifiedSnapState { attempt: attempt.id().into() })
1461 }
1462 _ => Ok(()),
1463 }
1464 }
1465
1466 pub fn ensure_snap_sync_layout(&self) -> ProviderResult<()> {
1468 if self.cached_storage_settings().use_hashed_state() {
1469 Ok(())
1470 } else {
1471 Err(ProviderError::SnapStorageLayoutUnsupported)
1472 }
1473 }
1474}
1475
1476impl<TX: DbTx, N: NodeTypes> AccountReader for DatabaseProvider<TX, N> {
1477 fn basic_account(&self, address: &Address) -> ProviderResult<Option<Account>> {
1478 if self.cached_storage_settings().use_hashed_state() {
1479 let hashed_address = keccak256(address);
1480 Ok(self.tx.get_by_encoded_key::<tables::HashedAccounts>(&hashed_address)?)
1481 } else {
1482 Ok(self.tx.get_by_encoded_key::<tables::PlainAccountState>(address)?)
1483 }
1484 }
1485}
1486
1487impl<TX: DbTx + 'static, N: NodeTypes> AccountExtReader for DatabaseProvider<TX, N> {
1488 fn changed_accounts_with_range(
1489 &self,
1490 range: RangeInclusive<BlockNumber>,
1491 ) -> ProviderResult<BTreeSet<Address>> {
1492 let mut reader = EitherReader::new_account_changesets(self)?;
1493
1494 reader.changed_accounts_with_range(range)
1495 }
1496
1497 fn basic_accounts(
1498 &self,
1499 iter: impl IntoIterator<Item = Address>,
1500 ) -> ProviderResult<Vec<(Address, Option<Account>)>> {
1501 if self.cached_storage_settings().use_hashed_state() {
1502 let mut hashed_accounts = self.tx.cursor_read::<tables::HashedAccounts>()?;
1503 Ok(iter
1504 .into_iter()
1505 .map(|address| {
1506 let hashed_address = keccak256(address);
1507 hashed_accounts.seek_exact(hashed_address).map(|a| (address, a.map(|(_, v)| v)))
1508 })
1509 .collect::<Result<Vec<_>, _>>()?)
1510 } else {
1511 let mut plain_accounts = self.tx.cursor_read::<tables::PlainAccountState>()?;
1512 Ok(iter
1513 .into_iter()
1514 .map(|address| {
1515 plain_accounts.seek_exact(address).map(|a| (address, a.map(|(_, v)| v)))
1516 })
1517 .collect::<Result<Vec<_>, _>>()?)
1518 }
1519 }
1520
1521 fn changed_accounts_and_blocks_with_range(
1522 &self,
1523 range: RangeInclusive<BlockNumber>,
1524 ) -> ProviderResult<BTreeMap<Address, Vec<u64>>> {
1525 let highest_static_block = self
1526 .static_file_provider
1527 .get_highest_static_file_block(StaticFileSegment::AccountChangeSets);
1528
1529 if let Some(highest) = highest_static_block &&
1530 self.cached_storage_settings().storage_v2
1531 {
1532 let start = *range.start();
1533 let static_end = (*range.end()).min(highest);
1534
1535 let mut changed_accounts_and_blocks: BTreeMap<_, Vec<u64>> = BTreeMap::default();
1536 if start <= static_end {
1537 for block in start..=static_end {
1538 let block_changesets = self.account_block_changeset(block)?;
1539 for changeset in block_changesets {
1540 changed_accounts_and_blocks
1541 .entry(changeset.address)
1542 .or_default()
1543 .push(block);
1544 }
1545 }
1546 }
1547
1548 Ok(changed_accounts_and_blocks)
1549 } else {
1550 let mut changeset_cursor = self.tx.cursor_read::<tables::AccountChangeSets>()?;
1551
1552 let account_transitions = changeset_cursor.walk_range(range)?.try_fold(
1553 BTreeMap::new(),
1554 |mut accounts: BTreeMap<Address, Vec<u64>>, entry| -> ProviderResult<_> {
1555 let (index, account) = entry?;
1556 accounts.entry(account.address).or_default().push(index);
1557 Ok(accounts)
1558 },
1559 )?;
1560
1561 Ok(account_transitions)
1562 }
1563 }
1564}
1565
1566impl<TX: DbTx, N: NodeTypes> StorageChangeSetReader for DatabaseProvider<TX, N> {
1567 fn storage_changeset(
1568 &self,
1569 block_number: BlockNumber,
1570 ) -> ProviderResult<Vec<(BlockNumberAddress, StorageEntry)>> {
1571 if self.cached_storage_settings().storage_v2 {
1572 self.static_file_provider.storage_changeset(block_number)
1573 } else {
1574 let range = block_number..=block_number;
1575 let storage_range = BlockNumberAddress::range(range);
1576 self.tx
1577 .cursor_dup_read::<tables::StorageChangeSets>()?
1578 .walk_range(storage_range)?
1579 .map(|r| {
1580 let (bna, entry) = r?;
1581 Ok((bna, entry))
1582 })
1583 .collect()
1584 }
1585 }
1586
1587 fn get_storage_before_block(
1588 &self,
1589 block_number: BlockNumber,
1590 address: Address,
1591 storage_key: B256,
1592 ) -> ProviderResult<Option<StorageEntry>> {
1593 if self.cached_storage_settings().storage_v2 {
1594 self.static_file_provider.get_storage_before_block(block_number, address, storage_key)
1595 } else {
1596 Ok(self
1597 .tx
1598 .cursor_dup_read::<tables::StorageChangeSets>()?
1599 .seek_by_key_subkey(BlockNumberAddress((block_number, address)), storage_key)?
1600 .filter(|entry| entry.key == storage_key))
1601 }
1602 }
1603
1604 fn storage_changesets_range(
1605 &self,
1606 range: impl RangeBounds<BlockNumber>,
1607 ) -> ProviderResult<Vec<(BlockNumberAddress, StorageEntry)>> {
1608 if self.cached_storage_settings().storage_v2 {
1609 self.static_file_provider.storage_changesets_range(range)
1610 } else {
1611 self.tx
1612 .cursor_dup_read::<tables::StorageChangeSets>()?
1613 .walk_range(BlockNumberAddressRange::from(range))?
1614 .map(|r| {
1615 let (bna, entry) = r?;
1616 Ok((bna, entry))
1617 })
1618 .collect()
1619 }
1620 }
1621}
1622
1623impl<TX: DbTx, N: NodeTypes> ChangeSetReader for DatabaseProvider<TX, N> {
1624 fn account_block_changeset(
1625 &self,
1626 block_number: BlockNumber,
1627 ) -> ProviderResult<Vec<AccountBeforeTx>> {
1628 if self.cached_storage_settings().storage_v2 {
1629 let static_changesets =
1630 self.static_file_provider.account_block_changeset(block_number)?;
1631 Ok(static_changesets)
1632 } else {
1633 let range = block_number..=block_number;
1634 self.tx
1635 .cursor_read::<tables::AccountChangeSets>()?
1636 .walk_range(range)?
1637 .map(|result| -> ProviderResult<_> {
1638 let (_, account_before) = result?;
1639 Ok(account_before)
1640 })
1641 .collect()
1642 }
1643 }
1644
1645 fn get_account_before_block(
1646 &self,
1647 block_number: BlockNumber,
1648 address: Address,
1649 ) -> ProviderResult<Option<AccountBeforeTx>> {
1650 if self.cached_storage_settings().storage_v2 {
1651 Ok(self.static_file_provider.get_account_before_block(block_number, address)?)
1652 } else {
1653 self.tx
1654 .cursor_dup_read::<tables::AccountChangeSets>()?
1655 .seek_by_key_subkey(block_number, address)?
1656 .filter(|acc| acc.address == address)
1657 .map(Ok)
1658 .transpose()
1659 }
1660 }
1661
1662 fn account_changesets_range(
1663 &self,
1664 range: impl core::ops::RangeBounds<BlockNumber>,
1665 ) -> ProviderResult<Vec<(BlockNumber, AccountBeforeTx)>> {
1666 if self.cached_storage_settings().storage_v2 {
1667 self.static_file_provider.account_changesets_range(range)
1668 } else {
1669 self.tx
1670 .cursor_read::<tables::AccountChangeSets>()?
1671 .walk_range(to_range(range))?
1672 .map(|r| r.map_err(Into::into))
1673 .collect()
1674 }
1675 }
1676}
1677
1678impl<TX: DbTx + 'static, N: NodeTypes> DatabaseProvider<TX, N> {
1679 fn history_rocksdb_snapshot(&self) -> Option<&RocksReadSnapshot<'_>> {
1692 self.rocksdb_history_snapshot
1693 .get_or_init(|| {
1694 self.cached_storage_settings()
1695 .storage_v2
1696 .then(|| self.rocksdb_provider.owned_snapshot())
1697 })
1698 .as_ref()
1699 .map(OwnedRocksReadSnapshot::as_snapshot)
1700 }
1701}
1702
1703impl<TX: DbTx + 'static, N: NodeTypes> HistoryReader for DatabaseProvider<TX, N> {
1704 fn account_history_info(
1705 &self,
1706 address: Address,
1707 block_number: BlockNumber,
1708 lowest_available_block_number: Option<BlockNumber>,
1709 ) -> ProviderResult<HistoryInfo> {
1710 let visible_tip = self.best_block_number()?;
1711 let mut reader = EitherReader::new_accounts_history(self, self.history_rocksdb_snapshot())?;
1712 reader
1713 .account_history_info(address, block_number, lowest_available_block_number, visible_tip)
1714 .map(Into::into)
1715 }
1716
1717 fn storage_history_info(
1718 &self,
1719 address: Address,
1720 storage_key: B256,
1721 block_number: BlockNumber,
1722 lowest_available_block_number: Option<BlockNumber>,
1723 ) -> ProviderResult<HistoryInfo> {
1724 let visible_tip = self.best_block_number()?;
1725 let mut reader = EitherReader::new_storages_history(self, self.history_rocksdb_snapshot())?;
1726 reader
1727 .storage_history_info(
1728 address,
1729 storage_key,
1730 block_number,
1731 lowest_available_block_number,
1732 visible_tip,
1733 )
1734 .map(Into::into)
1735 }
1736}
1737
1738impl<TX: DbTx + 'static, N: NodeTypesForProvider> HeaderSyncGapProvider
1739 for DatabaseProvider<TX, N>
1740{
1741 type Header = HeaderTy<N>;
1742
1743 fn local_tip_header(
1744 &self,
1745 highest_uninterrupted_block: BlockNumber,
1746 ) -> ProviderResult<SealedHeader<Self::Header>> {
1747 let static_file_provider = self.static_file_provider();
1748
1749 let next_static_file_block_num = static_file_provider
1752 .get_highest_static_file_block(StaticFileSegment::Headers)
1753 .map(|id| id + 1)
1754 .unwrap_or_default();
1755 let next_block = highest_uninterrupted_block + 1;
1756
1757 match next_static_file_block_num.cmp(&next_block) {
1758 Ordering::Greater => {
1761 let mut static_file_producer =
1762 static_file_provider.latest_writer(StaticFileSegment::Headers)?;
1763 static_file_producer.prune_headers(next_static_file_block_num - next_block)?;
1764 static_file_producer.commit()?
1767 }
1768 Ordering::Less => {
1769 return Err(ProviderError::HeaderNotFound(next_static_file_block_num.into()))
1771 }
1772 Ordering::Equal => {}
1773 }
1774
1775 let local_head = static_file_provider
1776 .sealed_header(highest_uninterrupted_block)?
1777 .ok_or_else(|| ProviderError::HeaderNotFound(highest_uninterrupted_block.into()))?;
1778
1779 Ok(local_head)
1780 }
1781}
1782
1783impl<TX: DbTx + 'static, N: NodeTypesForProvider> HeaderProvider for DatabaseProvider<TX, N> {
1784 type Header = HeaderTy<N>;
1785
1786 fn header(&self, block_hash: BlockHash) -> ProviderResult<Option<Self::Header>> {
1787 if let Some(num) = self.block_number(block_hash)? {
1788 Ok(self.header_by_number(num)?)
1789 } else {
1790 Ok(None)
1791 }
1792 }
1793
1794 fn header_by_number(&self, num: BlockNumber) -> ProviderResult<Option<Self::Header>> {
1795 self.static_file_provider.header_by_number(num)
1796 }
1797
1798 fn headers_range(
1799 &self,
1800 range: impl RangeBounds<BlockNumber>,
1801 ) -> ProviderResult<Vec<Self::Header>> {
1802 self.static_file_provider.headers_range(range)
1803 }
1804
1805 fn sealed_header(
1806 &self,
1807 number: BlockNumber,
1808 ) -> ProviderResult<Option<SealedHeader<Self::Header>>> {
1809 self.static_file_provider.sealed_header(number)
1810 }
1811
1812 fn sealed_headers_while(
1813 &self,
1814 range: impl RangeBounds<BlockNumber>,
1815 predicate: impl FnMut(&SealedHeader<Self::Header>) -> bool,
1816 ) -> ProviderResult<Vec<SealedHeader<Self::Header>>> {
1817 self.static_file_provider.sealed_headers_while(range, predicate)
1818 }
1819}
1820
1821impl<TX: DbTx + 'static, N: NodeTypes> BlockHashReader for DatabaseProvider<TX, N> {
1822 fn block_hash(&self, number: u64) -> ProviderResult<Option<B256>> {
1823 self.static_file_provider.block_hash(number)
1824 }
1825
1826 fn canonical_hashes_range(
1827 &self,
1828 start: BlockNumber,
1829 end: BlockNumber,
1830 ) -> ProviderResult<Vec<B256>> {
1831 self.static_file_provider.canonical_hashes_range(start, end)
1832 }
1833}
1834
1835impl<TX: DbTx + 'static, N: NodeTypes> BlockNumReader for DatabaseProvider<TX, N> {
1836 fn chain_info(&self) -> ProviderResult<ChainInfo> {
1837 let best_number = self.best_block_number()?;
1838 let best_hash = self.block_hash(best_number)?.unwrap_or_default();
1839 Ok(ChainInfo { best_hash, best_number })
1840 }
1841
1842 fn best_block_number(&self) -> ProviderResult<BlockNumber> {
1843 Ok(self
1846 .get_stage_checkpoint(StageId::Finish)?
1847 .map(|checkpoint| checkpoint.block_number)
1848 .unwrap_or_default())
1849 }
1850
1851 fn last_block_number(&self) -> ProviderResult<BlockNumber> {
1852 self.static_file_provider.last_block_number()
1853 }
1854
1855 fn block_number(&self, hash: B256) -> ProviderResult<Option<BlockNumber>> {
1856 Ok(self.tx.get::<tables::HeaderNumbers>(hash)?)
1857 }
1858}
1859
1860impl<TX: DbTx + 'static, N: NodeTypesForProvider> BlockReader for DatabaseProvider<TX, N> {
1861 type Block = BlockTy<N>;
1862
1863 fn find_block_by_hash(
1864 &self,
1865 hash: B256,
1866 source: BlockSource,
1867 ) -> ProviderResult<Option<Self::Block>> {
1868 if source.is_canonical() {
1869 self.block(hash.into())
1870 } else {
1871 Ok(None)
1872 }
1873 }
1874
1875 fn block(&self, id: BlockHashOrNumber) -> ProviderResult<Option<Self::Block>> {
1883 if let Some(number) = self.convert_hash_or_number(id)? {
1884 let earliest_available = self.static_file_provider.earliest_history_height();
1885 if number < earliest_available {
1886 return Err(ProviderError::BlockExpired { requested: number, earliest_available })
1887 }
1888
1889 let Some(header) = self.header_by_number(number)? else { return Ok(None) };
1890
1891 let Some(transactions) = self.transactions_by_block(number.into())? else {
1896 return Ok(None)
1897 };
1898
1899 let body = self
1900 .storage
1901 .reader()
1902 .read_block_bodies(self, vec![(&header, transactions)])?
1903 .pop()
1904 .ok_or(ProviderError::InvalidStorageOutput)?;
1905
1906 return Ok(Some(Self::Block::new(header, body)))
1907 }
1908
1909 Ok(None)
1910 }
1911
1912 fn pending_block(&self) -> ProviderResult<Option<Arc<RecoveredBlock<Self::Block>>>> {
1913 Ok(None)
1914 }
1915
1916 fn pending_block_and_receipts(
1917 &self,
1918 ) -> ProviderResult<Option<RecoveredBlockAndExecutionOutput<Self::Block, Self::Receipt>>> {
1919 Ok(None)
1920 }
1921
1922 fn recovered_block(
1931 &self,
1932 id: BlockHashOrNumber,
1933 transaction_kind: TransactionVariant,
1934 ) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
1935 self.recovered_block(
1936 id,
1937 transaction_kind,
1938 |block_number| self.header_by_number(block_number),
1939 |header, body, senders| {
1940 Self::Block::new(header, body)
1941 .try_into_recovered_unchecked(senders)
1945 .map(Some)
1946 .map_err(|_| ProviderError::SenderRecoveryError)
1947 },
1948 )
1949 }
1950
1951 fn sealed_block_with_senders(
1952 &self,
1953 id: BlockHashOrNumber,
1954 transaction_kind: TransactionVariant,
1955 ) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
1956 self.recovered_block(
1957 id,
1958 transaction_kind,
1959 |block_number| self.sealed_header(block_number),
1960 |header, body, senders| {
1961 Self::Block::new_sealed(header, body)
1962 .try_with_senders_unchecked(senders)
1966 .map(Some)
1967 .map_err(|_| ProviderError::SenderRecoveryError)
1968 },
1969 )
1970 }
1971
1972 fn block_range(&self, range: RangeInclusive<BlockNumber>) -> ProviderResult<Vec<Self::Block>> {
1973 self.block_range(
1974 range,
1975 |range| self.headers_range(range),
1976 |header, body, _| Ok(Self::Block::new(header, body)),
1977 )
1978 }
1979
1980 fn block_with_senders_range(
1981 &self,
1982 range: RangeInclusive<BlockNumber>,
1983 ) -> ProviderResult<Vec<RecoveredBlock<Self::Block>>> {
1984 self.block_with_senders_range(
1985 range,
1986 |range| self.headers_range(range),
1987 |header, body, senders| {
1988 Self::Block::new(header, body)
1989 .try_into_recovered_unchecked(senders)
1990 .map_err(|_| ProviderError::SenderRecoveryError)
1991 },
1992 )
1993 }
1994
1995 fn recovered_block_range(
1996 &self,
1997 range: RangeInclusive<BlockNumber>,
1998 ) -> ProviderResult<Vec<RecoveredBlock<Self::Block>>> {
1999 self.block_with_senders_range(
2000 range,
2001 |range| self.sealed_headers_range(range),
2002 |header, body, senders| {
2003 Self::Block::new_sealed(header, body)
2004 .try_with_senders(senders)
2005 .map_err(|_| ProviderError::SenderRecoveryError)
2006 },
2007 )
2008 }
2009
2010 fn block_by_transaction_id(&self, id: TxNumber) -> ProviderResult<Option<BlockNumber>> {
2011 Ok(self
2012 .tx
2013 .cursor_read::<tables::TransactionBlocks>()?
2014 .seek(id)
2015 .map(|b| b.map(|(_, bn)| bn))?)
2016 }
2017}
2018
2019impl<TX: DbTx + 'static, N: NodeTypesForProvider> TransactionsProviderExt
2020 for DatabaseProvider<TX, N>
2021{
2022 fn transaction_hashes_by_range(
2025 &self,
2026 tx_range: Range<TxNumber>,
2027 ) -> ProviderResult<Vec<(TxHash, TxNumber)>> {
2028 self.static_file_provider.transaction_hashes_by_range(tx_range)
2029 }
2030}
2031
2032impl<TX: DbTx + 'static, N: NodeTypesForProvider> TransactionsProvider for DatabaseProvider<TX, N> {
2034 type Transaction = TxTy<N>;
2035
2036 fn transaction_id(&self, tx_hash: TxHash) -> ProviderResult<Option<TxNumber>> {
2037 self.with_rocksdb_snapshot(|rocksdb_ref| {
2038 let mut reader = EitherReader::new_transaction_hash_numbers(self, rocksdb_ref)?;
2039 reader.get_transaction_hash_number(tx_hash)
2040 })
2041 }
2042
2043 fn transaction_by_id(&self, id: TxNumber) -> ProviderResult<Option<Self::Transaction>> {
2044 self.static_file_provider.transaction_by_id(id)
2045 }
2046
2047 fn transaction_by_id_unhashed(
2048 &self,
2049 id: TxNumber,
2050 ) -> ProviderResult<Option<Self::Transaction>> {
2051 self.static_file_provider.transaction_by_id_unhashed(id)
2052 }
2053
2054 fn transaction_by_hash(&self, hash: TxHash) -> ProviderResult<Option<Self::Transaction>> {
2055 if let Some(id) = self.transaction_id(hash)? {
2056 Ok(self.transaction_by_id_unhashed(id)?)
2057 } else {
2058 Ok(None)
2059 }
2060 }
2061
2062 fn transaction_by_hash_with_meta(
2063 &self,
2064 tx_hash: TxHash,
2065 ) -> ProviderResult<Option<(Self::Transaction, TransactionMeta)>> {
2066 if let Some(transaction_id) = self.transaction_id(tx_hash)? &&
2067 let Some(transaction) = self.transaction_by_id_unhashed(transaction_id)? &&
2068 let Some(block_number) = self.block_by_transaction_id(transaction_id)? &&
2069 let Some(sealed_header) = self.sealed_header(block_number)?
2070 {
2071 let (header, block_hash) = sealed_header.split();
2072 if let Some(block_body) = self.block_body_indices(block_number)? {
2073 let index = transaction_id - block_body.first_tx_num();
2078
2079 let meta = TransactionMeta {
2080 tx_hash,
2081 index,
2082 block_hash,
2083 block_number,
2084 base_fee: header.base_fee_per_gas(),
2085 excess_blob_gas: header.excess_blob_gas(),
2086 timestamp: header.timestamp(),
2087 };
2088
2089 return Ok(Some((transaction, meta)))
2090 }
2091 }
2092
2093 Ok(None)
2094 }
2095
2096 fn transactions_by_block(
2097 &self,
2098 id: BlockHashOrNumber,
2099 ) -> ProviderResult<Option<Vec<Self::Transaction>>> {
2100 if let Some(block_number) = self.convert_hash_or_number(id)? &&
2101 let Some(body) = self.block_body_indices(block_number)?
2102 {
2103 let tx_range = body.tx_num_range();
2104 return if tx_range.is_empty() {
2105 Ok(Some(Vec::new()))
2106 } else {
2107 self.transactions_by_tx_range(tx_range).map(Some)
2108 }
2109 }
2110 Ok(None)
2111 }
2112
2113 fn transactions_by_block_range(
2114 &self,
2115 range: impl RangeBounds<BlockNumber>,
2116 ) -> ProviderResult<Vec<Vec<Self::Transaction>>> {
2117 let range = to_range(range);
2118
2119 self.block_body_indices_range(range.start..=range.end.saturating_sub(1))?
2120 .into_iter()
2121 .map(|body| {
2122 let tx_num_range = body.tx_num_range();
2123 if tx_num_range.is_empty() {
2124 Ok(Vec::new())
2125 } else {
2126 self.transactions_by_tx_range(tx_num_range)
2127 }
2128 })
2129 .collect()
2130 }
2131
2132 fn transactions_by_tx_range(
2133 &self,
2134 range: impl RangeBounds<TxNumber>,
2135 ) -> ProviderResult<Vec<Self::Transaction>> {
2136 self.static_file_provider.transactions_by_tx_range(range)
2137 }
2138
2139 fn senders_by_tx_range(
2140 &self,
2141 range: impl RangeBounds<TxNumber>,
2142 ) -> ProviderResult<Vec<Address>> {
2143 if EitherWriterDestination::senders(self).is_static_file() {
2144 self.static_file_provider.senders_by_tx_range(range)
2145 } else {
2146 self.cursor_read_collect::<tables::TransactionSenders>(range)
2147 }
2148 }
2149
2150 fn transaction_sender(&self, id: TxNumber) -> ProviderResult<Option<Address>> {
2151 if EitherWriterDestination::senders(self).is_static_file() {
2152 self.static_file_provider.transaction_sender(id)
2153 } else {
2154 Ok(self.tx.get::<tables::TransactionSenders>(id)?)
2155 }
2156 }
2157}
2158
2159impl<TX: DbTx + 'static, N: NodeTypesForProvider> ReceiptProvider for DatabaseProvider<TX, N> {
2160 type Receipt = ReceiptTy<N>;
2161
2162 fn receipt(&self, id: TxNumber) -> ProviderResult<Option<Self::Receipt>> {
2163 self.static_file_provider.get_with_static_file_or_database(
2164 StaticFileSegment::Receipts,
2165 id,
2166 |static_file| static_file.receipt(id),
2167 || Ok(self.tx.get::<tables::Receipts<Self::Receipt>>(id)?),
2168 )
2169 }
2170
2171 fn receipt_by_hash(&self, hash: TxHash) -> ProviderResult<Option<Self::Receipt>> {
2172 if let Some(id) = self.transaction_id(hash)? {
2173 self.receipt(id)
2174 } else {
2175 Ok(None)
2176 }
2177 }
2178
2179 fn receipts_by_block(
2180 &self,
2181 block: BlockHashOrNumber,
2182 ) -> ProviderResult<Option<Vec<Self::Receipt>>> {
2183 if let Some(number) = self.convert_hash_or_number(block)? &&
2184 let Some(body) = self.block_body_indices(number)?
2185 {
2186 let tx_range = body.tx_num_range();
2187 return if tx_range.is_empty() {
2188 Ok(Some(Vec::new()))
2189 } else {
2190 let receipts = self.receipts_by_tx_range(tx_range)?;
2191
2192 if receipts.len() != body.tx_count as usize {
2193 return Ok(None)
2194 }
2195
2196 Ok(Some(receipts))
2197 }
2198 }
2199 Ok(None)
2200 }
2201
2202 fn receipts_by_tx_range(
2203 &self,
2204 range: impl RangeBounds<TxNumber>,
2205 ) -> ProviderResult<Vec<Self::Receipt>> {
2206 self.static_file_provider.get_range_with_static_file_or_database(
2207 StaticFileSegment::Receipts,
2208 to_range(range),
2209 |static_file, range, _| static_file.receipts_by_tx_range(range),
2210 |range, _| self.cursor_read_collect::<tables::Receipts<Self::Receipt>>(range),
2211 |_| true,
2212 )
2213 }
2214
2215 fn receipts_by_block_range(
2216 &self,
2217 block_range: RangeInclusive<BlockNumber>,
2218 ) -> ProviderResult<Vec<Vec<Self::Receipt>>> {
2219 if block_range.is_empty() {
2220 return Ok(Vec::new());
2221 }
2222
2223 let range_len = block_range.end().saturating_sub(*block_range.start()) as usize + 1;
2225 let mut block_body_indices = Vec::with_capacity(range_len);
2226 for block_num in block_range {
2227 if let Some(indices) = self.block_body_indices(block_num)? {
2228 block_body_indices.push(indices);
2229 } else {
2230 block_body_indices.push(StoredBlockBodyIndices::default());
2232 }
2233 }
2234
2235 if block_body_indices.is_empty() {
2236 return Ok(Vec::new());
2237 }
2238
2239 let non_empty_blocks: Vec<_> =
2241 block_body_indices.iter().filter(|indices| indices.tx_count > 0).collect();
2242
2243 if non_empty_blocks.is_empty() {
2244 return Ok(vec![Vec::new(); block_body_indices.len()]);
2246 }
2247
2248 let first_tx = non_empty_blocks[0].first_tx_num();
2250 let last_tx = non_empty_blocks[non_empty_blocks.len() - 1].last_tx_num();
2251
2252 let all_receipts = self.receipts_by_tx_range(first_tx..=last_tx)?;
2254 let mut receipts_iter = all_receipts.into_iter();
2255
2256 let mut result = Vec::with_capacity(block_body_indices.len());
2258 for indices in &block_body_indices {
2259 if indices.tx_count == 0 {
2260 result.push(Vec::new());
2261 } else {
2262 let block_receipts =
2263 receipts_iter.by_ref().take(indices.tx_count as usize).collect();
2264 result.push(block_receipts);
2265 }
2266 }
2267
2268 Ok(result)
2269 }
2270}
2271
2272impl<TX: DbTx + 'static, N: NodeTypesForProvider> BlockBodyIndicesProvider
2273 for DatabaseProvider<TX, N>
2274{
2275 fn block_body_indices(&self, num: u64) -> ProviderResult<Option<StoredBlockBodyIndices>> {
2276 Ok(self.tx.get::<tables::BlockBodyIndices>(num)?)
2277 }
2278
2279 fn block_body_indices_range(
2280 &self,
2281 range: RangeInclusive<BlockNumber>,
2282 ) -> ProviderResult<Vec<StoredBlockBodyIndices>> {
2283 self.cursor_read_collect::<tables::BlockBodyIndices>(range)
2284 }
2285}
2286
2287impl<TX: DbTx, N: NodeTypes> StageCheckpointReader for DatabaseProvider<TX, N> {
2288 fn get_stage_checkpoint(&self, id: StageId) -> ProviderResult<Option<StageCheckpoint>> {
2289 Ok(if let Some(encoded) = id.get_pre_encoded() {
2290 self.tx.get_by_encoded_key::<tables::StageCheckpoints>(encoded)?
2291 } else {
2292 self.tx.get::<tables::StageCheckpoints>(id.to_string())?
2293 })
2294 }
2295
2296 fn get_stage_checkpoint_progress(&self, id: StageId) -> ProviderResult<Option<Vec<u8>>> {
2298 Ok(self.tx.get::<tables::StageCheckpointProgresses>(id.to_string())?)
2299 }
2300
2301 fn get_all_checkpoints(&self) -> ProviderResult<Vec<(String, StageCheckpoint)>> {
2302 self.tx
2303 .cursor_read::<tables::StageCheckpoints>()?
2304 .walk(None)?
2305 .collect::<Result<Vec<(String, StageCheckpoint)>, _>>()
2306 .map_err(ProviderError::Database)
2307 }
2308}
2309
2310impl<TX: DbTxMut + DbTx, N: NodeTypes> StageCheckpointWriter for DatabaseProvider<TX, N> {
2311 fn save_stage_checkpoint(
2313 &self,
2314 id: StageId,
2315 checkpoint: StageCheckpoint,
2316 ) -> ProviderResult<()> {
2317 if id == StageId::Finish {
2318 self.ensure_finish_may_advance(&checkpoint)?;
2319 }
2320 Ok(self.tx.put::<tables::StageCheckpoints>(id.to_string(), checkpoint)?)
2321 }
2322
2323 fn save_stage_checkpoint_progress(
2325 &self,
2326 id: StageId,
2327 checkpoint: Vec<u8>,
2328 ) -> ProviderResult<()> {
2329 Ok(self.tx.put::<tables::StageCheckpointProgresses>(id.to_string(), checkpoint)?)
2330 }
2331
2332 #[instrument(level = "debug", target = "providers::db", skip_all)]
2333 fn update_pipeline_stages(
2334 &self,
2335 block_number: BlockNumber,
2336 drop_stage_checkpoint: bool,
2337 ) -> ProviderResult<()> {
2338 let current = self.get_stage_checkpoint(StageId::Finish)?.unwrap_or_default();
2339 let proposed = StageCheckpoint {
2340 block_number,
2341 ..if drop_stage_checkpoint { Default::default() } else { current }
2342 };
2343 self.ensure_finish_may_advance(&proposed)?;
2344
2345 let mut cursor = self.tx.cursor_write::<tables::StageCheckpoints>()?;
2347 for stage_id in StageId::ALL {
2348 let (_, checkpoint) = cursor.seek_exact(stage_id.to_string())?.unwrap_or_default();
2349 cursor.upsert(
2350 stage_id.to_string(),
2351 &StageCheckpoint {
2352 block_number,
2353 ..if drop_stage_checkpoint { Default::default() } else { checkpoint }
2354 },
2355 )?;
2356 }
2357
2358 Ok(())
2359 }
2360}
2361
2362impl<TX: DbTxMut + DbTx, N: NodeTypes> DatabaseProvider<TX, N> {
2363 fn update_pipeline_stages_after_unwind(
2366 &self,
2367 block_number: BlockNumber,
2368 ) -> ProviderResult<PersistenceFrontiers> {
2369 let partial_state_trie = self
2370 .get_stage_checkpoint(StageId::Finish)?
2371 .map(|checkpoint| {
2372 checkpoint
2373 .finish_stage_checkpoint()
2374 .and_then(|finish| finish.partial_state_trie())
2375 .unwrap_or(checkpoint.block_number)
2376 })
2377 .unwrap_or(block_number)
2378 .min(block_number);
2379
2380 self.update_pipeline_stages(block_number, true)?;
2381 if partial_state_trie < block_number {
2382 self.save_stage_checkpoint(
2383 StageId::Finish,
2384 StageCheckpoint::new(block_number).with_finish_stage_checkpoint(FinishCheckpoint {
2385 partial_state_trie: Some(partial_state_trie),
2386 }),
2387 )?;
2388 }
2389
2390 Ok(PersistenceFrontiers { db_tip: block_number, partial_state_trie })
2391 }
2392}
2393
2394impl<TX: DbTx + 'static, N: NodeTypes> StorageReader for DatabaseProvider<TX, N> {
2395 fn plain_state_storages(
2396 &self,
2397 addresses_with_keys: impl IntoIterator<Item = (Address, impl IntoIterator<Item = B256>)>,
2398 ) -> ProviderResult<Vec<(Address, Vec<StorageEntry>)>> {
2399 if self.cached_storage_settings().use_hashed_state() {
2400 let mut hashed_storage = self.tx.cursor_dup_read::<tables::HashedStorages>()?;
2401
2402 addresses_with_keys
2403 .into_iter()
2404 .map(|(address, storage)| {
2405 let hashed_address = keccak256(address);
2406 storage
2407 .into_iter()
2408 .map(|key| -> ProviderResult<_> {
2409 let hashed_key = keccak256(key);
2410 let value = hashed_storage
2411 .seek_by_key_subkey(hashed_address, hashed_key)?
2412 .filter(|v| v.key == hashed_key)
2413 .map(|v| v.value)
2414 .unwrap_or_default();
2415 Ok(StorageEntry { key, value })
2416 })
2417 .collect::<ProviderResult<Vec<_>>>()
2418 .map(|storage| (address, storage))
2419 })
2420 .collect::<ProviderResult<Vec<(_, _)>>>()
2421 } else {
2422 let mut plain_storage = self.tx.cursor_dup_read::<tables::PlainStorageState>()?;
2423
2424 addresses_with_keys
2425 .into_iter()
2426 .map(|(address, storage)| {
2427 storage
2428 .into_iter()
2429 .map(|key| -> ProviderResult<_> {
2430 Ok(plain_storage
2431 .seek_by_key_subkey(address, key)?
2432 .filter(|v| v.key == key)
2433 .unwrap_or_else(|| StorageEntry { key, value: Default::default() }))
2434 })
2435 .collect::<ProviderResult<Vec<_>>>()
2436 .map(|storage| (address, storage))
2437 })
2438 .collect::<ProviderResult<Vec<(_, _)>>>()
2439 }
2440 }
2441
2442 fn changed_storages_with_range(
2443 &self,
2444 range: RangeInclusive<BlockNumber>,
2445 ) -> ProviderResult<BTreeMap<Address, BTreeSet<B256>>> {
2446 if self.cached_storage_settings().storage_v2 {
2447 self.storage_changesets_range(range)?.into_iter().try_fold(
2448 BTreeMap::new(),
2449 |mut accounts: BTreeMap<Address, BTreeSet<B256>>, entry| {
2450 let (BlockNumberAddress((_, address)), storage_entry) = entry;
2451 accounts.entry(address).or_default().insert(storage_entry.key);
2452 Ok(accounts)
2453 },
2454 )
2455 } else {
2456 self.tx
2457 .cursor_read::<tables::StorageChangeSets>()?
2458 .walk_range(BlockNumberAddress::range(range))?
2459 .try_fold(
2462 BTreeMap::new(),
2463 |mut accounts: BTreeMap<Address, BTreeSet<B256>>, entry| {
2464 let (BlockNumberAddress((_, address)), storage_entry) = entry?;
2465 accounts.entry(address).or_default().insert(storage_entry.key);
2466 Ok(accounts)
2467 },
2468 )
2469 }
2470 }
2471
2472 fn changed_storages_and_blocks_with_range(
2473 &self,
2474 range: RangeInclusive<BlockNumber>,
2475 ) -> ProviderResult<BTreeMap<(Address, B256), Vec<u64>>> {
2476 if self.cached_storage_settings().storage_v2 {
2477 self.storage_changesets_range(range)?.into_iter().try_fold(
2478 BTreeMap::new(),
2479 |mut storages: BTreeMap<(Address, B256), Vec<u64>>, (index, storage)| {
2480 storages
2481 .entry((index.address(), storage.key))
2482 .or_default()
2483 .push(index.block_number());
2484 Ok(storages)
2485 },
2486 )
2487 } else {
2488 let mut changeset_cursor = self.tx.cursor_read::<tables::StorageChangeSets>()?;
2489
2490 let storage_changeset_lists =
2491 changeset_cursor.walk_range(BlockNumberAddress::range(range))?.try_fold(
2492 BTreeMap::new(),
2493 |mut storages: BTreeMap<(Address, B256), Vec<u64>>,
2494 entry|
2495 -> ProviderResult<_> {
2496 let (index, storage) = entry?;
2497 storages
2498 .entry((index.address(), storage.key))
2499 .or_default()
2500 .push(index.block_number());
2501 Ok(storages)
2502 },
2503 )?;
2504
2505 Ok(storage_changeset_lists)
2506 }
2507 }
2508}
2509
2510impl<TX: DbTxMut + DbTx + 'static, N: NodeTypesForProvider> StateWriter
2511 for DatabaseProvider<TX, N>
2512{
2513 type Receipt = ReceiptTy<N>;
2514
2515 #[instrument(level = "debug", target = "providers::db", skip_all)]
2516 fn write_state<'a>(
2517 &self,
2518 execution_outcome: impl Into<WriteStateInput<'a, Self::Receipt>>,
2519 is_value_known: OriginalValuesKnown,
2520 config: StateWriteConfig,
2521 ) -> ProviderResult<()> {
2522 let execution_outcome = execution_outcome.into();
2523
2524 if self.cached_storage_settings().use_hashed_state() &&
2525 !config.write_receipts &&
2526 !config.write_account_changesets &&
2527 !config.write_storage_changesets
2528 {
2529 self.write_bytecodes(
2533 execution_outcome.state().contracts.iter().map(|(h, b)| (*h, Bytecode(b.clone()))),
2534 )?;
2535 return Ok(());
2536 }
2537
2538 let first_block = execution_outcome.first_block();
2539 let (plain_state, reverts) =
2540 execution_outcome.state().to_plain_state_and_reverts(is_value_known);
2541
2542 self.write_state_reverts(reverts, first_block, config)?;
2543 self.write_state_changes(plain_state)?;
2544
2545 if !config.write_receipts {
2546 return Ok(());
2547 }
2548
2549 let block_count = execution_outcome.len() as u64;
2550 let last_block = execution_outcome.last_block();
2551 let block_range = first_block..=last_block;
2552
2553 let tip = self.last_block_number()?.max(last_block);
2554
2555 let block_indices: Vec<_> = self
2557 .block_body_indices_range(block_range)?
2558 .into_iter()
2559 .map(|b| b.first_tx_num)
2560 .collect();
2561
2562 if block_indices.len() < block_count as usize {
2564 let missing_blocks = block_count - block_indices.len() as u64;
2565 return Err(ProviderError::BlockBodyIndicesNotFound(
2566 last_block.saturating_sub(missing_blocks - 1),
2567 ));
2568 }
2569
2570 let mut receipts_writer = EitherWriter::new_receipts(self, first_block)?;
2571
2572 let has_contract_log_filter = !self.prune_modes.receipts_log_filter.is_empty();
2573 let contract_log_pruner = self.prune_modes.receipts_log_filter.group_by_block(tip, None)?;
2574
2575 let prunable_receipts = (EitherWriter::receipts_destination(self).is_database() ||
2583 self.static_file_provider()
2584 .get_highest_static_file_tx(StaticFileSegment::Receipts)
2585 .is_none()) &&
2586 PruneMode::Distance(self.minimum_pruning_distance).should_prune(first_block, tip);
2587
2588 let mut allowed_addresses: AddressSet = AddressSet::default();
2590 for (_, addresses) in contract_log_pruner.range(..first_block) {
2591 allowed_addresses.extend(addresses.iter().copied());
2592 }
2593
2594 for (idx, (receipts, first_tx_index)) in
2595 execution_outcome.receipts().zip(block_indices).enumerate()
2596 {
2597 let block_number = first_block + idx as u64;
2598
2599 receipts_writer.increment_block(block_number)?;
2601
2602 if prunable_receipts &&
2604 self.prune_modes
2605 .receipts
2606 .is_some_and(|mode| mode.should_prune(block_number, tip))
2607 {
2608 continue
2609 }
2610
2611 if let Some(new_addresses) = contract_log_pruner.get(&block_number) {
2613 allowed_addresses.extend(new_addresses.iter().copied());
2614 }
2615
2616 for (idx, receipt) in receipts.iter().enumerate() {
2617 let receipt_idx = first_tx_index + idx as u64;
2618 if prunable_receipts &&
2621 has_contract_log_filter &&
2622 !receipt.logs().iter().any(|log| allowed_addresses.contains(&log.address))
2623 {
2624 continue
2625 }
2626
2627 receipts_writer.append_receipt(receipt_idx, receipt)?;
2628 }
2629 }
2630
2631 Ok(())
2632 }
2633
2634 fn write_state_reverts(
2635 &self,
2636 reverts: PlainStateReverts,
2637 first_block: BlockNumber,
2638 config: StateWriteConfig,
2639 ) -> ProviderResult<()> {
2640 if config.write_storage_changesets {
2642 tracing::trace!("Writing storage changes");
2643 let mut storages_cursor =
2644 self.tx_ref().cursor_dup_write::<tables::PlainStorageState>()?;
2645 for (block_index, mut storage_changes) in reverts.storage.into_iter().enumerate() {
2646 let block_number = first_block + block_index as BlockNumber;
2647
2648 tracing::trace!(block_number, "Writing block change");
2649 storage_changes.par_sort_unstable_by_key(|a| a.address);
2651 let total_changes =
2652 storage_changes.iter().map(|change| change.storage_revert.len()).sum();
2653 let mut changeset = Vec::with_capacity(total_changes);
2654 for PlainStorageRevert { address, wiped, storage_revert } in storage_changes {
2655 let mut storage = storage_revert
2656 .into_iter()
2657 .map(|(k, v)| (B256::from(k.to_be_bytes()), v))
2658 .collect::<Vec<_>>();
2659 storage.par_sort_unstable_by_key(|a| a.0);
2661
2662 let mut wiped_storage = Vec::new();
2670 if wiped {
2671 tracing::trace!(?address, "Wiping storage");
2672 if let Some((_, entry)) = storages_cursor.seek_exact(address)? {
2673 wiped_storage.push((entry.key, entry.value));
2674 while let Some(entry) = storages_cursor.next_dup_val()? {
2675 wiped_storage.push((entry.key, entry.value))
2676 }
2677 }
2678 }
2679
2680 tracing::trace!(?address, ?storage, "Writing storage reverts");
2681 for (key, value) in StorageRevertsIter::new(storage, wiped_storage) {
2682 changeset.push(StorageBeforeTx { address, key, value });
2683 }
2684 }
2685
2686 let mut storage_changesets_writer =
2687 EitherWriter::new_storage_changesets(self, block_number)?;
2688 storage_changesets_writer.append_storage_changeset(block_number, changeset)?;
2689 }
2690 }
2691
2692 if !config.write_account_changesets {
2693 return Ok(());
2694 }
2695
2696 tracing::trace!(?first_block, "Writing account changes");
2698 for (block_index, account_block_reverts) in reverts.accounts.into_iter().enumerate() {
2699 let block_number = first_block + block_index as BlockNumber;
2700 let changeset = account_block_reverts
2701 .into_iter()
2702 .map(|(address, info)| AccountBeforeTx { address, info: info.map(Into::into) })
2703 .collect::<Vec<_>>();
2704 let mut account_changesets_writer =
2705 EitherWriter::new_account_changesets(self, block_number)?;
2706
2707 account_changesets_writer.append_account_changeset(block_number, changeset)?;
2708 }
2709
2710 Ok(())
2711 }
2712
2713 fn write_state_changes(&self, mut changes: StateChangeset) -> ProviderResult<()> {
2714 changes.accounts.par_sort_by_key(|a| a.0);
2717 changes.storage.par_sort_by_key(|a| a.address);
2718 changes.contracts.par_sort_by_key(|a| a.0);
2719
2720 if !self.cached_storage_settings().use_hashed_state() {
2721 tracing::trace!(len = changes.accounts.len(), "Writing new account state");
2723 let mut accounts_cursor = self.tx_ref().cursor_write::<tables::PlainAccountState>()?;
2724 for (address, account) in changes.accounts {
2726 if let Some(account) = account {
2727 tracing::trace!(?address, "Updating plain state account");
2728 accounts_cursor.upsert(address, &account.into())?;
2729 } else if accounts_cursor.seek_exact(address)?.is_some() {
2730 tracing::trace!(?address, "Deleting plain state account");
2731 accounts_cursor.delete_current()?;
2732 }
2733 }
2734
2735 tracing::trace!(len = changes.storage.len(), "Writing new storage state");
2737 let mut storages_cursor =
2738 self.tx_ref().cursor_dup_write::<tables::PlainStorageState>()?;
2739 for PlainStorageChangeset { address, wipe_storage, storage } in changes.storage {
2740 if wipe_storage && storages_cursor.seek_exact(address)?.is_some() {
2742 storages_cursor.delete_current_duplicates()?;
2743 }
2744 let mut storage = storage
2746 .into_iter()
2747 .map(|(k, value)| StorageEntry { key: k.into(), value })
2748 .collect::<Vec<_>>();
2749 storage.par_sort_unstable_by_key(|a| a.key);
2751
2752 for entry in storage {
2753 tracing::trace!(?address, ?entry.key, "Updating plain state storage");
2754 if let Some(db_entry) =
2755 storages_cursor.seek_by_key_subkey(address, entry.key)? &&
2756 db_entry.key == entry.key
2757 {
2758 storages_cursor.delete_current()?;
2759 }
2760
2761 if !entry.value.is_zero() {
2762 storages_cursor.upsert(address, &entry)?;
2763 }
2764 }
2765 }
2766 }
2767
2768 tracing::trace!(len = changes.contracts.len(), "Writing bytecodes");
2770 self.write_bytecodes(
2771 changes.contracts.into_iter().map(|(hash, bytecode)| (hash, Bytecode(bytecode))),
2772 )?;
2773
2774 Ok(())
2775 }
2776
2777 #[instrument(level = "debug", target = "providers::db", skip_all)]
2778 fn write_hashed_state(&self, hashed_state: &HashedPostStateSorted) -> ProviderResult<()> {
2779 let mut hashed_accounts_cursor = self.tx_ref().cursor_write::<tables::HashedAccounts>()?;
2781 for (hashed_address, account) in hashed_state.accounts() {
2782 if let Some(account) = account {
2783 hashed_accounts_cursor.upsert(*hashed_address, account)?;
2784 } else if hashed_accounts_cursor.seek_exact(*hashed_address)?.is_some() {
2785 hashed_accounts_cursor.delete_current()?;
2786 }
2787 }
2788
2789 let sorted_storages = hashed_state.account_storages().iter().sorted_by_key(|(key, _)| *key);
2791 let mut hashed_storage_cursor =
2792 self.tx_ref().cursor_dup_write::<tables::HashedStorages>()?;
2793 for (hashed_address, storage) in sorted_storages {
2794 for (hashed_slot, value) in storage.storage_slots_ref() {
2795 let entry = StorageEntry { key: *hashed_slot, value: *value };
2796
2797 if let Some(db_entry) =
2798 hashed_storage_cursor.seek_by_key_subkey(*hashed_address, entry.key)? &&
2799 db_entry.key == entry.key
2800 {
2801 hashed_storage_cursor.delete_current()?;
2802 }
2803
2804 if !entry.value.is_zero() {
2805 hashed_storage_cursor.upsert(*hashed_address, &entry)?;
2806 }
2807 }
2808 }
2809
2810 Ok(())
2811 }
2812
2813 fn remove_state_above(&self, block: BlockNumber) -> ProviderResult<()> {
2835 let range = block + 1..=self.last_block_number()?;
2836
2837 if range.is_empty() {
2838 return Ok(());
2839 }
2840
2841 let from_transaction_num = self
2843 .block_body_indices(block)?
2844 .map(|b| b.next_tx_num())
2845 .ok_or(ProviderError::BlockBodyIndicesNotFound(block))?;
2846
2847 let storage_range = BlockNumberAddress::range(range.clone());
2848 let storage_changeset = if self.cached_storage_settings().storage_v2 {
2849 let changesets = self.storage_changesets_range(range.clone())?;
2850 let mut changeset_writer =
2851 self.static_file_provider.latest_writer(StaticFileSegment::StorageChangeSets)?;
2852 changeset_writer.prune_storage_changesets(block)?;
2853 changesets
2854 } else {
2855 self.take::<tables::StorageChangeSets>(storage_range)?.into_iter().collect()
2856 };
2857 let account_changeset = if self.cached_storage_settings().storage_v2 {
2858 let changesets = self.account_changesets_range(range)?;
2859 let mut changeset_writer =
2860 self.static_file_provider.latest_writer(StaticFileSegment::AccountChangeSets)?;
2861 changeset_writer.prune_account_changesets(block)?;
2862 changesets
2863 } else {
2864 self.take::<tables::AccountChangeSets>(range)?
2865 };
2866
2867 if self.cached_storage_settings().use_hashed_state() {
2868 let mut hashed_accounts_cursor = self.tx.cursor_write::<tables::HashedAccounts>()?;
2869 let mut hashed_storage_cursor = self.tx.cursor_dup_write::<tables::HashedStorages>()?;
2870
2871 let (state, _) = self.populate_bundle_state_hashed(
2872 account_changeset,
2873 storage_changeset,
2874 &mut hashed_accounts_cursor,
2875 &mut hashed_storage_cursor,
2876 )?;
2877
2878 for (address, (old_account, new_account, storage)) in &state {
2879 if old_account != new_account {
2880 let hashed_address = keccak256(address);
2881 let existing_entry = hashed_accounts_cursor.seek_exact(hashed_address)?;
2882 if let Some(account) = old_account {
2883 hashed_accounts_cursor.upsert(hashed_address, account)?;
2884 } else if existing_entry.is_some() {
2885 hashed_accounts_cursor.delete_current()?;
2886 }
2887 }
2888
2889 for (storage_key, (old_storage_value, _new_storage_value)) in storage {
2890 let hashed_address = keccak256(address);
2891 let hashed_storage_key = keccak256(storage_key);
2892 let storage_entry =
2893 StorageEntry { key: hashed_storage_key, value: *old_storage_value };
2894 if hashed_storage_cursor
2895 .seek_by_key_subkey(hashed_address, hashed_storage_key)?
2896 .is_some_and(|s| s.key == hashed_storage_key)
2897 {
2898 hashed_storage_cursor.delete_current()?
2899 }
2900
2901 if !old_storage_value.is_zero() {
2902 hashed_storage_cursor.upsert(hashed_address, &storage_entry)?;
2903 }
2904 }
2905 }
2906 } else {
2907 let mut plain_accounts_cursor = self.tx.cursor_write::<tables::PlainAccountState>()?;
2912 let mut plain_storage_cursor =
2913 self.tx.cursor_dup_write::<tables::PlainStorageState>()?;
2914
2915 let (state, _) = self.populate_bundle_state_plain(
2916 account_changeset,
2917 storage_changeset,
2918 &mut plain_accounts_cursor,
2919 &mut plain_storage_cursor,
2920 )?;
2921
2922 for (address, (old_account, new_account, storage)) in &state {
2923 if old_account != new_account {
2924 let existing_entry = plain_accounts_cursor.seek_exact(*address)?;
2925 if let Some(account) = old_account {
2926 plain_accounts_cursor.upsert(*address, account)?;
2927 } else if existing_entry.is_some() {
2928 plain_accounts_cursor.delete_current()?;
2929 }
2930 }
2931
2932 for (storage_key, (old_storage_value, _new_storage_value)) in storage {
2933 let storage_entry =
2934 StorageEntry { key: *storage_key, value: *old_storage_value };
2935 if plain_storage_cursor
2936 .seek_by_key_subkey(*address, *storage_key)?
2937 .is_some_and(|s| s.key == *storage_key)
2938 {
2939 plain_storage_cursor.delete_current()?
2940 }
2941
2942 if !old_storage_value.is_zero() {
2943 plain_storage_cursor.upsert(*address, &storage_entry)?;
2944 }
2945 }
2946 }
2947 }
2948
2949 self.remove_receipts_from(from_transaction_num, block)?;
2950
2951 Ok(())
2952 }
2953
2954 fn take_state_above(
2976 &self,
2977 block: BlockNumber,
2978 ) -> ProviderResult<ExecutionOutcome<Self::Receipt>> {
2979 let range = block + 1..=self.last_block_number()?;
2980
2981 if range.is_empty() {
2982 return Ok(ExecutionOutcome::default())
2983 }
2984 let start_block_number = *range.start();
2985
2986 let block_bodies = self.block_body_indices_range(range.clone())?;
2988
2989 let from_transaction_num =
2991 block_bodies.first().expect("already checked if there are blocks").first_tx_num();
2992 let to_transaction_num =
2993 block_bodies.last().expect("already checked if there are blocks").last_tx_num();
2994
2995 let storage_range = BlockNumberAddress::range(range.clone());
2996 let storage_changeset = if let Some(highest_block) = self
2997 .static_file_provider
2998 .get_highest_static_file_block(StaticFileSegment::StorageChangeSets) &&
2999 self.cached_storage_settings().storage_v2
3000 {
3001 let changesets = self.storage_changesets_range(block + 1..=highest_block)?;
3002 let mut changeset_writer =
3003 self.static_file_provider.latest_writer(StaticFileSegment::StorageChangeSets)?;
3004 changeset_writer.prune_storage_changesets(block)?;
3005 changesets
3006 } else {
3007 self.take::<tables::StorageChangeSets>(storage_range)?.into_iter().collect()
3008 };
3009
3010 let highest_changeset_block = self
3012 .static_file_provider
3013 .get_highest_static_file_block(StaticFileSegment::AccountChangeSets);
3014 let account_changeset = if let Some(highest_block) = highest_changeset_block &&
3015 self.cached_storage_settings().storage_v2
3016 {
3017 let changesets = self.account_changesets_range(block + 1..highest_block + 1)?;
3019 let mut changeset_writer =
3020 self.static_file_provider.latest_writer(StaticFileSegment::AccountChangeSets)?;
3021 changeset_writer.prune_account_changesets(block)?;
3022
3023 changesets
3024 } else {
3025 self.take::<tables::AccountChangeSets>(range)?
3028 };
3029
3030 let (state, reverts) = if self.cached_storage_settings().use_hashed_state() {
3031 let mut hashed_accounts_cursor = self.tx.cursor_write::<tables::HashedAccounts>()?;
3032 let mut hashed_storage_cursor = self.tx.cursor_dup_write::<tables::HashedStorages>()?;
3033
3034 let (state, reverts) = self.populate_bundle_state_hashed(
3035 account_changeset,
3036 storage_changeset,
3037 &mut hashed_accounts_cursor,
3038 &mut hashed_storage_cursor,
3039 )?;
3040
3041 for (address, (old_account, new_account, storage)) in &state {
3042 if old_account != new_account {
3043 let hashed_address = keccak256(address);
3044 let existing_entry = hashed_accounts_cursor.seek_exact(hashed_address)?;
3045 if let Some(account) = old_account {
3046 hashed_accounts_cursor.upsert(hashed_address, account)?;
3047 } else if existing_entry.is_some() {
3048 hashed_accounts_cursor.delete_current()?;
3049 }
3050 }
3051
3052 for (storage_key, (old_storage_value, _new_storage_value)) in storage {
3053 let hashed_address = keccak256(address);
3054 let hashed_storage_key = keccak256(storage_key);
3055 let storage_entry =
3056 StorageEntry { key: hashed_storage_key, value: *old_storage_value };
3057 if hashed_storage_cursor
3058 .seek_by_key_subkey(hashed_address, hashed_storage_key)?
3059 .is_some_and(|s| s.key == hashed_storage_key)
3060 {
3061 hashed_storage_cursor.delete_current()?
3062 }
3063
3064 if !old_storage_value.is_zero() {
3065 hashed_storage_cursor.upsert(hashed_address, &storage_entry)?;
3066 }
3067 }
3068 }
3069
3070 (state, reverts)
3071 } else {
3072 let mut plain_accounts_cursor = self.tx.cursor_write::<tables::PlainAccountState>()?;
3077 let mut plain_storage_cursor =
3078 self.tx.cursor_dup_write::<tables::PlainStorageState>()?;
3079
3080 let (state, reverts) = self.populate_bundle_state_plain(
3081 account_changeset,
3082 storage_changeset,
3083 &mut plain_accounts_cursor,
3084 &mut plain_storage_cursor,
3085 )?;
3086
3087 for (address, (old_account, new_account, storage)) in &state {
3088 if old_account != new_account {
3089 let existing_entry = plain_accounts_cursor.seek_exact(*address)?;
3090 if let Some(account) = old_account {
3091 plain_accounts_cursor.upsert(*address, account)?;
3092 } else if existing_entry.is_some() {
3093 plain_accounts_cursor.delete_current()?;
3094 }
3095 }
3096
3097 for (storage_key, (old_storage_value, _new_storage_value)) in storage {
3098 let storage_entry =
3099 StorageEntry { key: *storage_key, value: *old_storage_value };
3100 if plain_storage_cursor
3101 .seek_by_key_subkey(*address, *storage_key)?
3102 .is_some_and(|s| s.key == *storage_key)
3103 {
3104 plain_storage_cursor.delete_current()?
3105 }
3106
3107 if !old_storage_value.is_zero() {
3108 plain_storage_cursor.upsert(*address, &storage_entry)?;
3109 }
3110 }
3111 }
3112
3113 (state, reverts)
3114 };
3115
3116 let mut receipts_iter = self
3118 .static_file_provider
3119 .get_range_with_static_file_or_database(
3120 StaticFileSegment::Receipts,
3121 from_transaction_num..to_transaction_num + 1,
3122 |static_file, range, _| {
3123 static_file
3124 .receipts_by_tx_range(range.clone())
3125 .map(|r| range.into_iter().zip(r).collect())
3126 },
3127 |range, _| {
3128 self.tx
3129 .cursor_read::<tables::Receipts<Self::Receipt>>()?
3130 .walk_range(range)?
3131 .map(|r| r.map_err(Into::into))
3132 .collect()
3133 },
3134 |_| true,
3135 )?
3136 .into_iter()
3137 .peekable();
3138
3139 let mut receipts = Vec::with_capacity(block_bodies.len());
3140 for block_body in block_bodies {
3142 let mut block_receipts = Vec::with_capacity(block_body.tx_count as usize);
3143 for num in block_body.tx_num_range() {
3144 if receipts_iter.peek().is_some_and(|(n, _)| *n == num) {
3145 block_receipts.push(receipts_iter.next().unwrap().1);
3146 }
3147 }
3148 receipts.push(block_receipts);
3149 }
3150
3151 self.remove_receipts_from(from_transaction_num, block)?;
3152
3153 Ok(ExecutionOutcome::new_init(
3154 state,
3155 reverts,
3156 Vec::new(),
3157 receipts,
3158 start_block_number,
3159 Vec::new(),
3160 ))
3161 }
3162}
3163
3164impl<TX: DbTxMut + DbTx + 'static, N: NodeTypes> DatabaseProvider<TX, N> {
3165 fn write_account_trie_updates<A: TrieTableAdapter>(
3166 tx: &TX,
3167 trie_updates: &TrieUpdatesSorted,
3168 num_entries: &mut usize,
3169 ) -> ProviderResult<()>
3170 where
3171 TX: DbTxMut,
3172 {
3173 let mut account_trie_cursor = tx.cursor_write::<A::AccountTrieTable>()?;
3174 for (key, updated_node) in trie_updates.account_nodes_ref() {
3176 let nibbles = A::AccountKey::from(*key);
3177 match updated_node {
3178 Some(node) => {
3179 if !key.is_empty() {
3180 *num_entries += 1;
3181 account_trie_cursor.upsert(nibbles, node)?;
3182 }
3183 }
3184 None => {
3185 *num_entries += 1;
3186 if account_trie_cursor.seek_exact(nibbles)?.is_some() {
3187 account_trie_cursor.delete_current()?;
3188 }
3189 }
3190 }
3191 }
3192 Ok(())
3193 }
3194
3195 fn write_storage_tries<A: TrieTableAdapter>(
3196 tx: &TX,
3197 storage_tries: Vec<(&B256, &StorageTrieUpdatesSorted)>,
3198 num_entries: &mut usize,
3199 ) -> ProviderResult<()>
3200 where
3201 TX: DbTxMut,
3202 {
3203 let mut cursor = tx.cursor_dup_write::<A::StorageTrieTable>()?;
3204 for (hashed_address, storage_trie_updates) in storage_tries {
3205 let mut db_storage_trie_cursor: DatabaseStorageTrieCursor<_, A> =
3206 DatabaseStorageTrieCursor::new(cursor, *hashed_address);
3207 *num_entries +=
3208 db_storage_trie_cursor.write_storage_trie_updates_sorted(storage_trie_updates)?;
3209 cursor = db_storage_trie_cursor.cursor;
3210 }
3211 Ok(())
3212 }
3213}
3214
3215impl<TX: DbTxMut + DbTx + 'static, N: NodeTypes> TrieWriter for DatabaseProvider<TX, N> {
3216 #[instrument(level = "debug", target = "providers::db", skip_all)]
3220 fn write_trie_updates_sorted(&self, trie_updates: &TrieUpdatesSorted) -> ProviderResult<usize> {
3221 if trie_updates.is_empty() {
3222 return Ok(0)
3223 }
3224
3225 let mut num_entries = 0;
3227
3228 reth_trie_db::with_adapter!(self, |A| {
3229 Self::write_account_trie_updates::<A>(self.tx_ref(), trie_updates, &mut num_entries)?;
3230 });
3231
3232 num_entries +=
3233 self.write_storage_trie_updates_sorted(trie_updates.storage_tries_ref().iter())?;
3234
3235 Ok(num_entries)
3236 }
3237}
3238
3239impl<TX: DbTxMut + DbTx + 'static, N: NodeTypes> StorageTrieWriter for DatabaseProvider<TX, N> {
3240 fn write_storage_trie_updates_sorted<'a>(
3246 &self,
3247 storage_tries: impl Iterator<Item = (&'a B256, &'a StorageTrieUpdatesSorted)>,
3248 ) -> ProviderResult<usize> {
3249 let mut num_entries = 0;
3250 let mut storage_tries = storage_tries.collect::<Vec<_>>();
3251 storage_tries.sort_unstable_by(|a, b| a.0.cmp(b.0));
3252 reth_trie_db::with_adapter!(self, |A| {
3253 Self::write_storage_tries::<A>(self.tx_ref(), storage_tries, &mut num_entries)?;
3254 });
3255 Ok(num_entries)
3256 }
3257}
3258
3259impl<TX: DbTxMut + DbTx + 'static, N: NodeTypes> HashingWriter for DatabaseProvider<TX, N> {
3260 #[allow(clippy::clone_on_copy)]
3261 fn unwind_account_hashing<'a>(
3262 &self,
3263 changesets: impl Iterator<Item = &'a (BlockNumber, AccountBeforeTx)>,
3264 ) -> ProviderResult<BTreeMap<B256, Option<Account>>> {
3265 let hashed_accounts = changesets
3269 .into_iter()
3270 .map(|(_, e)| (keccak256(e.address), e.info.clone()))
3271 .collect::<Vec<_>>()
3272 .into_iter()
3273 .rev()
3274 .collect::<BTreeMap<_, _>>();
3275
3276 let mut hashed_accounts_cursor = self.tx.cursor_write::<tables::HashedAccounts>()?;
3278 for (hashed_address, account) in &hashed_accounts {
3279 if let Some(account) = account {
3280 hashed_accounts_cursor.upsert(*hashed_address, account)?;
3281 } else if hashed_accounts_cursor.seek_exact(*hashed_address)?.is_some() {
3282 hashed_accounts_cursor.delete_current()?;
3283 }
3284 }
3285
3286 Ok(hashed_accounts)
3287 }
3288
3289 fn unwind_account_hashing_range(
3290 &self,
3291 range: impl RangeBounds<BlockNumber>,
3292 ) -> ProviderResult<BTreeMap<B256, Option<Account>>> {
3293 let changesets = self.account_changesets_range(range)?;
3294 self.unwind_account_hashing(changesets.iter())
3295 }
3296
3297 fn insert_account_for_hashing(
3298 &self,
3299 changesets: impl IntoIterator<Item = (Address, Option<Account>)>,
3300 ) -> ProviderResult<BTreeMap<B256, Option<Account>>> {
3301 let mut hashed_accounts_cursor = self.tx.cursor_write::<tables::HashedAccounts>()?;
3302 let hashed_accounts =
3303 changesets.into_iter().map(|(ad, ac)| (keccak256(ad), ac)).collect::<BTreeMap<_, _>>();
3304 for (hashed_address, account) in &hashed_accounts {
3305 if let Some(account) = account {
3306 hashed_accounts_cursor.upsert(*hashed_address, account)?;
3307 } else if hashed_accounts_cursor.seek_exact(*hashed_address)?.is_some() {
3308 hashed_accounts_cursor.delete_current()?;
3309 }
3310 }
3311 Ok(hashed_accounts)
3312 }
3313
3314 fn unwind_storage_hashing(
3315 &self,
3316 changesets: impl Iterator<Item = (BlockNumberAddress, StorageEntry)>,
3317 ) -> ProviderResult<B256Map<BTreeSet<B256>>> {
3318 let mut hashed_storages = changesets
3320 .into_iter()
3321 .map(|(BlockNumberAddress((_, address)), storage_entry)| {
3322 let hashed_key = keccak256(storage_entry.key);
3323 (keccak256(address), hashed_key, storage_entry.value)
3324 })
3325 .collect::<Vec<_>>();
3326 hashed_storages.sort_by_key(|(ha, hk, _)| (*ha, *hk));
3327
3328 let mut hashed_storage_keys: B256Map<BTreeSet<B256>> =
3330 B256Map::with_capacity_and_hasher(hashed_storages.len(), Default::default());
3331 let mut hashed_storage = self.tx.cursor_dup_write::<tables::HashedStorages>()?;
3332 for (hashed_address, key, value) in hashed_storages.into_iter().rev() {
3333 hashed_storage_keys.entry(hashed_address).or_default().insert(key);
3334
3335 if hashed_storage
3336 .seek_by_key_subkey(hashed_address, key)?
3337 .is_some_and(|entry| entry.key == key)
3338 {
3339 hashed_storage.delete_current()?;
3340 }
3341
3342 if !value.is_zero() {
3343 hashed_storage.upsert(hashed_address, &StorageEntry { key, value })?;
3344 }
3345 }
3346 Ok(hashed_storage_keys)
3347 }
3348
3349 fn unwind_storage_hashing_range(
3350 &self,
3351 range: impl RangeBounds<BlockNumber>,
3352 ) -> ProviderResult<B256Map<BTreeSet<B256>>> {
3353 let changesets = self.storage_changesets_range(range)?;
3354 self.unwind_storage_hashing(changesets.into_iter())
3355 }
3356
3357 fn insert_storage_for_hashing(
3358 &self,
3359 storages: impl IntoIterator<Item = (Address, impl IntoIterator<Item = StorageEntry>)>,
3360 ) -> ProviderResult<B256Map<BTreeSet<B256>>> {
3361 let hashed_storages =
3363 storages.into_iter().fold(BTreeMap::new(), |mut map, (address, storage)| {
3364 let storage = storage.into_iter().fold(BTreeMap::new(), |mut map, entry| {
3365 map.insert(keccak256(entry.key), entry.value);
3366 map
3367 });
3368 map.insert(keccak256(address), storage);
3369 map
3370 });
3371
3372 let hashed_storage_keys = hashed_storages
3373 .iter()
3374 .map(|(hashed_address, entries)| (*hashed_address, entries.keys().copied().collect()))
3375 .collect();
3376
3377 let mut hashed_storage_cursor = self.tx.cursor_dup_write::<tables::HashedStorages>()?;
3378 hashed_storages.into_iter().try_for_each(|(hashed_address, storage)| {
3381 storage.into_iter().try_for_each(|(key, value)| -> ProviderResult<()> {
3382 if hashed_storage_cursor
3383 .seek_by_key_subkey(hashed_address, key)?
3384 .is_some_and(|entry| entry.key == key)
3385 {
3386 hashed_storage_cursor.delete_current()?;
3387 }
3388
3389 if !value.is_zero() {
3390 hashed_storage_cursor.upsert(hashed_address, &StorageEntry { key, value })?;
3391 }
3392 Ok(())
3393 })
3394 })?;
3395
3396 Ok(hashed_storage_keys)
3397 }
3398}
3399
3400impl<TX: DbTxMut + DbTx + 'static, N: NodeTypes> HistoryWriter for DatabaseProvider<TX, N> {
3401 fn unwind_account_history_indices<'a>(
3402 &self,
3403 changesets: impl Iterator<Item = &'a (BlockNumber, AccountBeforeTx)>,
3404 ) -> ProviderResult<usize> {
3405 let mut last_indices = changesets
3406 .into_iter()
3407 .map(|(index, account)| (account.address, *index))
3408 .collect::<Vec<_>>();
3409 last_indices.sort_unstable_by_key(|(a, _)| *a);
3410
3411 if self.cached_storage_settings().storage_v2 {
3412 let batch = self.rocksdb_provider.unwind_account_history_indices(&last_indices)?;
3413 self.pending_rocksdb_batches.lock().push(batch);
3414 } else {
3415 let mut cursor = self.tx.cursor_write::<tables::AccountsHistory>()?;
3417 for &(address, rem_index) in &last_indices {
3418 let partial_shard = unwind_history_shards::<_, tables::AccountsHistory, _>(
3419 &mut cursor,
3420 ShardedKey::last(address),
3421 rem_index,
3422 |sharded_key| sharded_key.key == address,
3423 )?;
3424
3425 if !partial_shard.is_empty() {
3428 cursor.insert(
3429 ShardedKey::last(address),
3430 &BlockNumberList::new_pre_sorted(partial_shard),
3431 )?;
3432 }
3433 }
3434 }
3435
3436 let changesets = last_indices.len();
3437 Ok(changesets)
3438 }
3439
3440 fn unwind_account_history_indices_range(
3441 &self,
3442 range: impl RangeBounds<BlockNumber>,
3443 ) -> ProviderResult<usize> {
3444 let changesets = self.account_changesets_range(range)?;
3445 self.unwind_account_history_indices(changesets.iter())
3446 }
3447
3448 fn insert_account_history_index(
3449 &self,
3450 account_transitions: impl IntoIterator<Item = (Address, impl IntoIterator<Item = u64>)>,
3451 ) -> ProviderResult<()> {
3452 self.append_history_index::<_, tables::AccountsHistory>(
3453 account_transitions,
3454 ShardedKey::new,
3455 )
3456 }
3457
3458 fn unwind_storage_history_indices(
3459 &self,
3460 changesets: impl Iterator<Item = (BlockNumberAddress, StorageEntry)>,
3461 ) -> ProviderResult<usize> {
3462 let mut storage_changesets = changesets
3463 .into_iter()
3464 .map(|(BlockNumberAddress((bn, address)), storage)| (address, storage.key, bn))
3465 .collect::<Vec<_>>();
3466 storage_changesets.sort_unstable_by_key(|(address, key, _)| (*address, *key));
3467
3468 if self.cached_storage_settings().storage_v2 {
3469 let batch =
3470 self.rocksdb_provider.unwind_storage_history_indices(&storage_changesets)?;
3471 self.pending_rocksdb_batches.lock().push(batch);
3472 } else {
3473 let mut cursor = self.tx.cursor_write::<tables::StoragesHistory>()?;
3475 for &(address, storage_key, rem_index) in &storage_changesets {
3476 let partial_shard = unwind_history_shards::<_, tables::StoragesHistory, _>(
3477 &mut cursor,
3478 StorageShardedKey::last(address, storage_key),
3479 rem_index,
3480 |storage_sharded_key| {
3481 storage_sharded_key.address == address &&
3482 storage_sharded_key.sharded_key.key == storage_key
3483 },
3484 )?;
3485
3486 if !partial_shard.is_empty() {
3489 cursor.insert(
3490 StorageShardedKey::last(address, storage_key),
3491 &BlockNumberList::new_pre_sorted(partial_shard),
3492 )?;
3493 }
3494 }
3495 }
3496
3497 let changesets = storage_changesets.len();
3498 Ok(changesets)
3499 }
3500
3501 fn unwind_storage_history_indices_range(
3502 &self,
3503 range: impl RangeBounds<BlockNumber>,
3504 ) -> ProviderResult<usize> {
3505 let changesets = self.storage_changesets_range(range)?;
3506 self.unwind_storage_history_indices(changesets.into_iter())
3507 }
3508
3509 fn insert_storage_history_index(
3510 &self,
3511 storage_transitions: impl IntoIterator<Item = ((Address, B256), impl IntoIterator<Item = u64>)>,
3512 ) -> ProviderResult<()> {
3513 self.append_history_index::<_, tables::StoragesHistory>(
3514 storage_transitions,
3515 |(address, storage_key), highest_block_number| {
3516 StorageShardedKey::new(address, storage_key, highest_block_number)
3517 },
3518 )
3519 }
3520
3521 #[instrument(level = "debug", target = "providers::db", skip_all)]
3522 fn update_history_indices(&self, range: RangeInclusive<BlockNumber>) -> ProviderResult<()> {
3523 let storage_settings = self.cached_storage_settings();
3524 if !storage_settings.storage_v2 {
3525 let indices = self.changed_accounts_and_blocks_with_range(range.clone())?;
3526 self.insert_account_history_index(indices)?;
3527 }
3528
3529 if !storage_settings.storage_v2 {
3530 let indices = self.changed_storages_and_blocks_with_range(range)?;
3531 self.insert_storage_history_index(indices)?;
3532 }
3533
3534 Ok(())
3535 }
3536}
3537
3538impl<TX: DbTxMut + DbTx + 'static, N: NodeTypesForProvider> BlockExecutionWriter
3539 for DatabaseProvider<TX, N>
3540{
3541 fn take_block_and_execution_above(
3542 &self,
3543 block: BlockNumber,
3544 ) -> ProviderResult<Chain<Self::Primitives>> {
3545 let range = block + 1..=self.last_block_number()?;
3546
3547 self.unwind_trie_state_from(block + 1)?;
3548
3549 let execution_state = self.take_state_above(block)?;
3551
3552 let blocks = self.recovered_block_range(range)?;
3553
3554 self.remove_blocks_above(block)?;
3557
3558 self.update_pipeline_stages_after_unwind(block)?;
3560
3561 Ok(Chain::new(blocks, execution_state, BTreeMap::new()))
3562 }
3563
3564 fn remove_block_and_execution_above(
3565 &self,
3566 block: BlockNumber,
3567 ) -> ProviderResult<PersistenceFrontiers> {
3568 self.unwind_trie_state_from(block + 1)?;
3569
3570 self.remove_state_above(block)?;
3572
3573 self.remove_blocks_above(block)?;
3576
3577 self.update_pipeline_stages_after_unwind(block)
3579 }
3580}
3581
3582impl<TX: DbTxMut + DbTx + 'static, N: NodeTypesForProvider> BlockWriter
3583 for DatabaseProvider<TX, N>
3584{
3585 type Block = BlockTy<N>;
3586 type Receipt = ReceiptTy<N>;
3587
3588 fn insert_block(
3593 &self,
3594 block: &RecoveredBlock<Self::Block>,
3595 ) -> ProviderResult<StoredBlockBodyIndices> {
3596 let block_number = block.number();
3597
3598 let executed_block = ExecutedBlock::new(
3600 Arc::new(block.clone()),
3601 Arc::new(BlockExecutionOutput {
3602 result: BlockExecutionResult {
3603 receipts: Default::default(),
3604 requests: Default::default(),
3605 gas_used: 0,
3606 blob_gas_used: 0,
3607 },
3608 state: Default::default(),
3609 }),
3610 Default::default(),
3611 Default::default(),
3612 );
3613
3614 self.save_blocks_inner(
3615 std::slice::from_ref(&executed_block),
3616 &[],
3617 &[],
3618 None,
3619 SaveBlocksMode::BlocksOnly,
3620 )?;
3621
3622 self.block_body_indices(block_number)?
3624 .ok_or(ProviderError::BlockBodyIndicesNotFound(block_number))
3625 }
3626
3627 fn append_block_bodies(
3628 &self,
3629 bodies: Vec<(BlockNumber, Option<&BodyTy<N>>)>,
3630 ) -> ProviderResult<()> {
3631 let Some(from_block) = bodies.first().map(|(block, _)| *block) else { return Ok(()) };
3632
3633 let mut tx_writer =
3635 self.static_file_provider.get_writer(from_block, StaticFileSegment::Transactions)?;
3636
3637 let mut block_indices_cursor = self.tx.cursor_write::<tables::BlockBodyIndices>()?;
3638 let mut tx_block_cursor = self.tx.cursor_write::<tables::TransactionBlocks>()?;
3639
3640 let mut next_tx_num = tx_block_cursor.last()?.map(|(id, _)| id + 1).unwrap_or_default();
3642
3643 for (block_number, body) in &bodies {
3644 tx_writer.increment_block(*block_number)?;
3646
3647 let tx_count = body.as_ref().map(|b| b.transactions().len() as u64).unwrap_or_default();
3648 let block_indices = StoredBlockBodyIndices { first_tx_num: next_tx_num, tx_count };
3649
3650 let mut durations_recorder = metrics::DurationsRecorder::new(&self.metrics);
3651
3652 block_indices_cursor.append(*block_number, &block_indices)?;
3654
3655 durations_recorder.record_relative(metrics::Action::InsertBlockBodyIndices);
3656
3657 let Some(body) = body else { continue };
3658
3659 if !body.transactions().is_empty() {
3661 tx_block_cursor.append(block_indices.last_tx_num(), block_number)?;
3662 durations_recorder.record_relative(metrics::Action::InsertTransactionBlocks);
3663 }
3664
3665 for transaction in body.transactions() {
3667 tx_writer.append_transaction(next_tx_num, transaction)?;
3668
3669 next_tx_num += 1;
3671 }
3672 }
3673
3674 self.storage.writer().write_block_bodies(self, bodies)?;
3675
3676 Ok(())
3677 }
3678
3679 fn remove_blocks_above(&self, block: BlockNumber) -> ProviderResult<()> {
3680 let last_block_number = self.last_block_number()?;
3681 for hash in self.canonical_hashes_range(block + 1, last_block_number + 1)? {
3683 self.tx.delete::<tables::HeaderNumbers>(hash, None)?;
3684 }
3685
3686 let highest_static_file_block = self
3688 .static_file_provider()
3689 .get_highest_static_file_block(StaticFileSegment::Headers)
3690 .expect("todo: error handling, headers should exist");
3691
3692 debug!(target: "providers::db", ?block, "Removing static file blocks above block_number");
3698 self.static_file_provider()
3699 .get_writer(block, StaticFileSegment::Headers)?
3700 .prune_headers(highest_static_file_block.saturating_sub(block))?;
3701
3702 let unwind_tx_from = self
3704 .block_body_indices(block)?
3705 .map(|b| b.next_tx_num())
3706 .ok_or(ProviderError::BlockBodyIndicesNotFound(block))?;
3707
3708 let unwind_tx_to = self
3710 .tx
3711 .cursor_read::<tables::BlockBodyIndices>()?
3712 .last()?
3713 .ok_or(ProviderError::BlockBodyIndicesNotFound(block))?
3715 .1
3716 .last_tx_num();
3717
3718 if unwind_tx_from <= unwind_tx_to {
3719 let hashes = self.transaction_hashes_by_range(unwind_tx_from..(unwind_tx_to + 1))?;
3720 self.with_rocksdb_batch(|batch| {
3721 let mut writer = EitherWriter::new_transaction_hash_numbers(self, batch)?;
3722 for (hash, _) in hashes {
3723 writer.delete_transaction_hash_number(hash)?;
3724 }
3725 Ok(((), writer.into_raw_rocksdb_batch()))
3726 })?;
3727 }
3728
3729 if self.prune_modes.sender_recovery.is_none_or(|m| !m.is_full()) {
3732 EitherWriter::new_senders(self, last_block_number)?
3733 .prune_senders(unwind_tx_from, block)?;
3734 }
3735
3736 self.remove_bodies_above(block)?;
3737
3738 Ok(())
3739 }
3740
3741 fn remove_bodies_above(&self, block: BlockNumber) -> ProviderResult<()> {
3742 self.storage.writer().remove_block_bodies_above(self, block)?;
3743
3744 let unwind_tx_from = self
3746 .block_body_indices(block)?
3747 .map(|b| b.next_tx_num())
3748 .ok_or(ProviderError::BlockBodyIndicesNotFound(block))?;
3749
3750 self.remove::<tables::BlockBodyIndices>(block + 1..)?;
3751 self.remove::<tables::TransactionBlocks>(unwind_tx_from..)?;
3752
3753 let static_file_tx_num =
3754 self.static_file_provider.get_highest_static_file_tx(StaticFileSegment::Transactions);
3755
3756 let to_delete = static_file_tx_num
3757 .map(|static_tx| (static_tx + 1).saturating_sub(unwind_tx_from))
3758 .unwrap_or_default();
3759
3760 self.static_file_provider
3761 .latest_writer(StaticFileSegment::Transactions)?
3762 .prune_transactions(to_delete, block)?;
3763
3764 Ok(())
3765 }
3766
3767 fn append_blocks_with_state(
3776 &self,
3777 blocks: Vec<RecoveredBlock<Self::Block>>,
3778 execution_outcome: &ExecutionOutcome<Self::Receipt>,
3779 hashed_state: HashedPostStateSorted,
3780 ) -> ProviderResult<()> {
3781 if blocks.is_empty() {
3782 debug!(target: "providers::db", "Attempted to append empty block range");
3783 return Ok(())
3784 }
3785
3786 let first_number = blocks[0].number();
3789
3790 let last_block_number = blocks[blocks.len() - 1].number();
3793
3794 let mut durations_recorder = metrics::DurationsRecorder::new(&self.metrics);
3795
3796 let (account_transitions, storage_transitions) = {
3801 let mut account_transitions: BTreeMap<Address, Vec<u64>> = BTreeMap::new();
3802 let mut storage_transitions: BTreeMap<(Address, B256), Vec<u64>> = BTreeMap::new();
3803 for (block_idx, block_reverts) in execution_outcome.bundle.reverts.iter().enumerate() {
3804 let block_number = first_number + block_idx as u64;
3805 for (address, account_revert) in block_reverts {
3806 account_transitions.entry(*address).or_default().push(block_number);
3807 for storage_key in account_revert.storage.keys() {
3808 let key = B256::from(storage_key.to_be_bytes());
3809 storage_transitions.entry((*address, key)).or_default().push(block_number);
3810 }
3811 }
3812 }
3813 (account_transitions, storage_transitions)
3814 };
3815
3816 for block in blocks {
3818 self.insert_block(&block)?;
3819 durations_recorder.record_relative(metrics::Action::InsertBlock);
3820 }
3821
3822 self.write_state(execution_outcome, OriginalValuesKnown::No, StateWriteConfig::default())?;
3823 durations_recorder.record_relative(metrics::Action::InsertState);
3824
3825 self.write_hashed_state(&hashed_state)?;
3827 durations_recorder.record_relative(metrics::Action::InsertHashes);
3828
3829 let storage_settings = self.cached_storage_settings();
3834 if storage_settings.storage_v2 {
3835 self.with_rocksdb_batch(|mut batch| {
3836 for (address, blocks) in account_transitions {
3837 batch.append_account_history_shard(address, blocks)?;
3838 }
3839 Ok(((), Some(batch.into_inner())))
3840 })?;
3841 } else {
3842 self.insert_account_history_index(account_transitions)?;
3843 }
3844 if storage_settings.storage_v2 {
3845 self.with_rocksdb_batch(|mut batch| {
3846 for ((address, key), blocks) in storage_transitions {
3847 batch.append_storage_history_shard(address, key, blocks)?;
3848 }
3849 Ok(((), Some(batch.into_inner())))
3850 })?;
3851 } else {
3852 self.insert_storage_history_index(storage_transitions)?;
3853 }
3854 durations_recorder.record_relative(metrics::Action::InsertHistoryIndices);
3855
3856 self.update_pipeline_stages(last_block_number, false)?;
3858 durations_recorder.record_relative(metrics::Action::UpdatePipelineStages);
3859
3860 debug!(target: "providers::db", range = ?first_number..=last_block_number, actions = ?durations_recorder.actions, "Appended blocks");
3861
3862 Ok(())
3863 }
3864}
3865
3866impl<TX: DbTx + 'static, N: NodeTypes> PruneCheckpointReader for DatabaseProvider<TX, N> {
3867 fn get_prune_checkpoint(
3868 &self,
3869 segment: PruneSegment,
3870 ) -> ProviderResult<Option<PruneCheckpoint>> {
3871 Ok(self.tx.get::<tables::PruneCheckpoints>(segment)?)
3872 }
3873
3874 fn get_prune_checkpoints(&self) -> ProviderResult<Vec<(PruneSegment, PruneCheckpoint)>> {
3875 Ok(PruneSegment::variants()
3876 .filter_map(|segment| {
3877 self.tx
3878 .get::<tables::PruneCheckpoints>(segment)
3879 .transpose()
3880 .map(|chk| chk.map(|chk| (segment, chk)))
3881 })
3882 .collect::<Result<_, _>>()?)
3883 }
3884}
3885
3886impl<TX: DbTxMut, N: NodeTypes> PruneCheckpointWriter for DatabaseProvider<TX, N> {
3887 fn save_prune_checkpoint(
3888 &self,
3889 segment: PruneSegment,
3890 checkpoint: PruneCheckpoint,
3891 ) -> ProviderResult<()> {
3892 Ok(self.tx.put::<tables::PruneCheckpoints>(segment, checkpoint)?)
3893 }
3894}
3895
3896impl<TX: DbTx + 'static, N: NodeTypesForProvider> StatsReader for DatabaseProvider<TX, N> {
3897 fn count_entries<T: Table>(&self) -> ProviderResult<usize> {
3898 let db_entries = self.tx.entries::<T>()?;
3899 let static_file_entries = match self.static_file_provider.count_entries::<T>() {
3900 Ok(entries) => entries,
3901 Err(ProviderError::UnsupportedProvider) => 0,
3902 Err(err) => return Err(err),
3903 };
3904
3905 Ok(db_entries + static_file_entries)
3906 }
3907}
3908
3909impl<TX: DbTx + 'static, N: NodeTypes> ChainStateBlockReader for DatabaseProvider<TX, N> {
3910 fn last_finalized_block_number(&self) -> ProviderResult<Option<BlockNumber>> {
3911 let mut finalized_blocks = self
3912 .tx
3913 .cursor_read::<tables::ChainState>()?
3914 .walk(Some(tables::ChainStateKey::LastFinalizedBlock))?
3915 .take(1)
3916 .collect::<Result<BTreeMap<tables::ChainStateKey, BlockNumber>, _>>()?;
3917
3918 let last_finalized_block_number = finalized_blocks.pop_first().map(|pair| pair.1);
3919 Ok(last_finalized_block_number)
3920 }
3921
3922 fn last_safe_block_number(&self) -> ProviderResult<Option<BlockNumber>> {
3923 let mut finalized_blocks = self
3924 .tx
3925 .cursor_read::<tables::ChainState>()?
3926 .walk(Some(tables::ChainStateKey::LastSafeBlock))?
3927 .take(1)
3928 .collect::<Result<BTreeMap<tables::ChainStateKey, BlockNumber>, _>>()?;
3929
3930 let last_finalized_block_number = finalized_blocks.pop_first().map(|pair| pair.1);
3931 Ok(last_finalized_block_number)
3932 }
3933}
3934
3935impl<TX: DbTxMut, N: NodeTypes> ChainStateBlockWriter for DatabaseProvider<TX, N> {
3936 fn save_finalized_block_number(&self, block_number: BlockNumber) -> ProviderResult<()> {
3937 Ok(self
3938 .tx
3939 .put::<tables::ChainState>(tables::ChainStateKey::LastFinalizedBlock, block_number)?)
3940 }
3941
3942 fn save_safe_block_number(&self, block_number: BlockNumber) -> ProviderResult<()> {
3943 Ok(self.tx.put::<tables::ChainState>(tables::ChainStateKey::LastSafeBlock, block_number)?)
3944 }
3945}
3946
3947impl<TX: DbTx + 'static, N: NodeTypes + 'static> DbTxProvider for DatabaseProvider<TX, N> {
3948 type Tx = TX;
3949
3950 fn tx(&self) -> &Self::Tx {
3951 &self.tx
3952 }
3953}
3954
3955impl<TX: DbTx + 'static, N: NodeTypes + 'static> DBProvider for DatabaseProvider<TX, N> {
3956 fn tx_mut(&mut self) -> &mut Self::Tx {
3957 &mut self.tx
3958 }
3959
3960 fn into_tx(self) -> Self::Tx {
3961 self.tx
3962 }
3963
3964 fn prune_modes_ref(&self) -> &PruneModes {
3965 self.prune_modes_ref()
3966 }
3967
3968 #[instrument(
3970 name = "DatabaseProvider::commit",
3971 level = "debug",
3972 target = "providers::db",
3973 skip_all
3974 )]
3975 fn commit(self) -> ProviderResult<()> {
3976 if self.static_file_provider.has_unwind_queued() || self.commit_order.is_unwind() {
3977 self.commit_unwind()?;
3978 } else {
3979 let mut timings = metrics::CommitTimings::default();
3981
3982 let start = Instant::now();
3983 self.static_file_provider.finalize()?;
3984 timings.sf = start.elapsed();
3985
3986 let start = Instant::now();
3987 let batches = std::mem::take(&mut *self.pending_rocksdb_batches.lock());
3988 for batch in batches {
3989 self.rocksdb_provider.commit_batch(batch)?;
3990 }
3991 timings.rocksdb = start.elapsed();
3992
3993 let start = Instant::now();
3994 self.tx.commit()?;
3995 timings.mdbx = start.elapsed();
3996
3997 self.metrics.record_commit(&timings);
3998 }
3999
4000 Ok(())
4001 }
4002}
4003
4004impl<TX: DbTx, N: NodeTypes> MetadataProvider for DatabaseProvider<TX, N> {
4005 fn get_metadata(&self, key: &str) -> ProviderResult<Option<Vec<u8>>> {
4006 self.tx.get::<tables::Metadata>(key.to_string()).map_err(Into::into)
4007 }
4008}
4009
4010impl<TX: DbTxMut, N: NodeTypes> MetadataWriter for DatabaseProvider<TX, N> {
4011 fn write_metadata(&self, key: &str, value: Vec<u8>) -> ProviderResult<()> {
4012 self.tx.put::<tables::Metadata>(key.to_string(), value).map_err(Into::into)
4013 }
4014
4015 fn delete_metadata(&self, key: &str) -> ProviderResult<()> {
4016 self.tx.delete::<tables::Metadata>(key.to_string(), None)?;
4017 Ok(())
4018 }
4019}
4020
4021impl<TX: Send, N: NodeTypes> StorageSettingsCache for DatabaseProvider<TX, N> {
4022 fn cached_storage_settings(&self) -> StorageSettings {
4023 *self.storage_settings.read()
4024 }
4025
4026 fn set_storage_settings_cache(&self, settings: StorageSettings) {
4027 *self.storage_settings.write() = settings;
4028 }
4029}
4030
4031impl<TX: Send, N: NodeTypes> StoragePath for DatabaseProvider<TX, N> {
4032 fn storage_path(&self) -> PathBuf {
4033 self.db_path.clone()
4034 }
4035}
4036
4037#[cfg(test)]
4038mod tests {
4039 use super::*;
4040 use crate::{
4041 test_utils::{blocks::BlockchainTestData, create_test_provider_factory},
4042 BlockWriter,
4043 };
4044 use alloy_consensus::Header;
4045 use alloy_primitives::{
4046 map::{AddressMap, B256Map},
4047 U256,
4048 };
4049 use reth_chain_state::{test_utils::TestBlockBuilder, ExecutedBlock};
4050 use reth_db_api::models::StorageSettings;
4051 use reth_ethereum_primitives::Receipt;
4052 use reth_execution_types::{AccountRevertInit, BlockExecutionOutput, BlockExecutionResult};
4053 use reth_primitives_traits::SealedBlock;
4054 use reth_storage_api::{DatabaseProviderFactory, MetadataProvider, MetadataWriter};
4055 use reth_testing_utils::generators::{self, random_block, BlockParams};
4056 use reth_trie::{
4057 HashedPostState, KeccakKeyHasher, Nibbles, StoredNibbles, StoredNibblesSubKey,
4058 };
4059 use revm::{database::BundleState, state::AccountInfo};
4060 use std::{sync::mpsc, time::Duration};
4061
4062 fn save_genesis<TX, N>(
4065 provider: &DatabaseProvider<TX, N>,
4066 genesis: &ExecutedBlock<N::Primitives>,
4067 ) -> ProviderResult<()>
4068 where
4069 TX: DbTx + DbTxMut + 'static,
4070 N: NodeTypesForProvider,
4071 {
4072 assert_eq!(genesis.recovered_block().number(), 0);
4073 provider.save_blocks_inner(
4074 std::slice::from_ref(genesis),
4075 std::slice::from_ref(genesis),
4076 &[],
4077 None,
4078 SaveBlocksMode::Full,
4079 )
4080 }
4081
4082 #[test]
4083 fn snap_attempt_guards_finish_checkpoint_writers() {
4084 use alloy_eips::BlockNumHash;
4085 use reth_db_api::models::SnapAttempt;
4086
4087 let factory = create_test_provider_factory();
4088 let provider = factory.database_provider_rw().unwrap();
4089 provider.update_pipeline_stages(5, false).unwrap();
4090 let mut attempt = SnapAttempt::start(
4091 None,
4092 BlockNumHash::new(10, B256::repeat_byte(1)),
4093 B256::repeat_byte(2),
4094 );
4095 provider.write_snap_attempt(&attempt).unwrap();
4096 provider.commit().unwrap();
4097
4098 let provider = factory.database_provider_rw().unwrap();
4099 let refused = provider.save_stage_checkpoint(StageId::Finish, StageCheckpoint::new(10));
4100 assert!(matches!(refused, Err(ProviderError::UnverifiedSnapState { attempt: 0 })));
4101 let refused = provider.update_pipeline_stages(10, false);
4102 assert!(matches!(refused, Err(ProviderError::UnverifiedSnapState { attempt: 0 })));
4103 for stage in [StageId::Finish, StageId::Headers] {
4104 assert_eq!(provider.get_stage_checkpoint(stage).unwrap().unwrap().block_number, 5);
4105 }
4106
4107 provider.update_pipeline_stages(4, true).unwrap();
4109 provider.save_stage_checkpoint(StageId::Finish, StageCheckpoint::new(3)).unwrap();
4110 assert_eq!(
4111 provider.get_stage_checkpoint(StageId::Finish).unwrap().unwrap().block_number,
4112 3
4113 );
4114
4115 attempt.abandon();
4117 provider.write_snap_attempt(&attempt).unwrap();
4118 let refused = provider.save_stage_checkpoint(StageId::Finish, StageCheckpoint::new(10));
4119 assert!(matches!(refused, Err(ProviderError::UnverifiedSnapState { attempt: 0 })));
4120
4121 provider.save_stage_checkpoint(StageId::Headers, StageCheckpoint::new(10)).unwrap();
4123 attempt.verify();
4124 provider.write_snap_attempt(&attempt).unwrap();
4125 provider.save_stage_checkpoint(StageId::Finish, StageCheckpoint::new(10)).unwrap();
4126 provider.update_pipeline_stages(11, false).unwrap();
4127 provider.commit().unwrap();
4128
4129 let provider = factory.database_provider_ro().unwrap();
4130 assert_eq!(
4131 provider.get_stage_checkpoint(StageId::Finish).unwrap().unwrap().block_number,
4132 11
4133 );
4134 }
4135
4136 #[test]
4137 fn snap_attempt_guards_partial_state_trie_advances() {
4138 use alloy_eips::BlockNumHash;
4139 use reth_db_api::models::SnapAttempt;
4140
4141 for abandoned in [false, true] {
4142 let factory = create_test_provider_factory();
4143 let provider = factory.database_provider_rw().unwrap();
4144 provider.update_pipeline_stages(100, false).unwrap();
4145 let checkpoint = StageCheckpoint::new(100)
4146 .with_finish_stage_checkpoint(FinishCheckpoint { partial_state_trie: Some(50) });
4147 provider.save_stage_checkpoint(StageId::Finish, checkpoint).unwrap();
4148 let mut attempt = SnapAttempt::start(None, BlockNumHash::default(), B256::ZERO);
4149 if abandoned {
4150 attempt.abandon();
4151 }
4152 provider.write_snap_attempt(&attempt).unwrap();
4153
4154 for block in [100, 90] {
4155 let result =
4156 provider.save_stage_checkpoint(StageId::Finish, StageCheckpoint::new(block));
4157 assert!(matches!(result, Err(ProviderError::UnverifiedSnapState { .. })));
4158 let result = provider.update_pipeline_stages(block, true);
4159 assert!(matches!(result, Err(ProviderError::UnverifiedSnapState { .. })));
4160 for stage in [StageId::Finish, StageId::Headers] {
4161 assert_eq!(
4162 provider.get_stage_checkpoint(stage).unwrap().unwrap().block_number,
4163 100
4164 );
4165 }
4166 assert_eq!(
4167 provider.get_stage_checkpoint(StageId::Finish).unwrap(),
4168 Some(checkpoint)
4169 );
4170 }
4171
4172 let advanced = StageCheckpoint::new(100)
4173 .with_finish_stage_checkpoint(FinishCheckpoint { partial_state_trie: Some(60) });
4174 assert!(matches!(
4175 provider.save_stage_checkpoint(StageId::Finish, advanced),
4176 Err(ProviderError::UnverifiedSnapState { .. })
4177 ));
4178
4179 provider.update_pipeline_stages(90, false).unwrap();
4181 let rewind = StageCheckpoint::new(80)
4182 .with_finish_stage_checkpoint(FinishCheckpoint { partial_state_trie: Some(40) });
4183 provider.save_stage_checkpoint(StageId::Finish, rewind).unwrap();
4184 provider.update_pipeline_stages(30, true).unwrap();
4185
4186 attempt.verify();
4187 provider.write_snap_attempt(&attempt).unwrap();
4188 provider.save_stage_checkpoint(StageId::Finish, advanced).unwrap();
4189 provider.update_pipeline_stages(100, true).unwrap();
4190 }
4191 }
4192
4193 #[test]
4194 fn test_receipts_by_block_range_empty_range() {
4195 let factory = create_test_provider_factory();
4196 let provider = factory.provider().unwrap();
4197
4198 let start = 10u64;
4200 let end = 9u64;
4201 let result = provider.receipts_by_block_range(start..=end).unwrap();
4202 assert_eq!(result, Vec::<Vec<reth_ethereum_primitives::Receipt>>::new());
4203 }
4204
4205 #[test]
4206 fn metadata_can_be_deleted() {
4207 let factory = create_test_provider_factory();
4208 let key = "metadata-delete-test";
4209
4210 let provider_rw = factory.provider_rw().unwrap();
4211 provider_rw.write_metadata(key, vec![1]).unwrap();
4212 provider_rw.commit().unwrap();
4213 assert_eq!(factory.provider().unwrap().get_metadata(key).unwrap(), Some(vec![1]));
4214
4215 let provider_rw = factory.provider_rw().unwrap();
4216 provider_rw.delete_metadata(key).unwrap();
4217 provider_rw.commit().unwrap();
4218 assert_eq!(factory.provider().unwrap().get_metadata(key).unwrap(), None);
4219 }
4220
4221 #[test]
4222 fn unwind_commit_waits_for_pre_commit_readers() {
4223 let factory = create_test_provider_factory();
4224
4225 let reader = factory.provider().unwrap();
4226 let provider_rw = factory.unwind_provider_rw().unwrap();
4227 provider_rw.write_metadata("unwind-wait-test", vec![1]).unwrap();
4228 let (done_tx, done_rx) = mpsc::channel();
4229
4230 let handle = std::thread::spawn(move || {
4231 let result = provider_rw.commit();
4232 done_tx.send(result).unwrap();
4233 });
4234
4235 assert!(
4236 done_rx.recv_timeout(Duration::from_millis(50)).is_err(),
4237 "unwind commit should wait while an older read transaction is still open"
4238 );
4239
4240 drop(reader);
4241
4242 done_rx.recv_timeout(Duration::from_secs(1)).unwrap().unwrap();
4243 handle.join().unwrap();
4244 }
4245
4246 #[test]
4247 fn test_receipts_by_block_range_nonexistent_blocks() {
4248 let factory = create_test_provider_factory();
4249 let provider = factory.provider().unwrap();
4250
4251 let result = provider.receipts_by_block_range(10..=12).unwrap();
4253 assert_eq!(result, vec![vec![], vec![], vec![]]);
4254 }
4255
4256 #[test]
4257 fn test_receipts_by_block_range_single_block() {
4258 let factory = create_test_provider_factory();
4259 let data = BlockchainTestData::default();
4260
4261 let provider_rw = factory.provider_rw().unwrap();
4262 provider_rw.insert_block(&data.genesis.try_recover().unwrap()).unwrap();
4263 provider_rw
4264 .write_state(
4265 &ExecutionOutcome { first_block: 0, receipts: vec![vec![]], ..Default::default() },
4266 crate::OriginalValuesKnown::No,
4267 StateWriteConfig::default(),
4268 )
4269 .unwrap();
4270 provider_rw.insert_block(&data.blocks[0].0).unwrap();
4271 provider_rw
4272 .write_state(
4273 &data.blocks[0].1,
4274 crate::OriginalValuesKnown::No,
4275 StateWriteConfig::default(),
4276 )
4277 .unwrap();
4278 provider_rw.commit().unwrap();
4279
4280 let provider = factory.provider().unwrap();
4281 let result = provider.receipts_by_block_range(1..=1).unwrap();
4282
4283 assert_eq!(result.len(), 1);
4285 assert_eq!(result[0].len(), 1);
4286 assert_eq!(result[0][0], data.blocks[0].1.receipts()[0][0]);
4287 }
4288
4289 #[test]
4290 fn test_receipts_by_block_range_multiple_blocks() {
4291 let factory = create_test_provider_factory();
4292 let data = BlockchainTestData::default();
4293
4294 let provider_rw = factory.provider_rw().unwrap();
4295 provider_rw.insert_block(&data.genesis.try_recover().unwrap()).unwrap();
4296 provider_rw
4297 .write_state(
4298 &ExecutionOutcome { first_block: 0, receipts: vec![vec![]], ..Default::default() },
4299 crate::OriginalValuesKnown::No,
4300 StateWriteConfig::default(),
4301 )
4302 .unwrap();
4303 for (block, outcome) in data.blocks.iter().take(3) {
4304 provider_rw.insert_block(block).unwrap();
4305 provider_rw
4306 .write_state(outcome, crate::OriginalValuesKnown::No, StateWriteConfig::default())
4307 .unwrap();
4308 }
4309 provider_rw.commit().unwrap();
4310
4311 let provider = factory.provider().unwrap();
4312 let result = provider.receipts_by_block_range(1..=3).unwrap();
4313
4314 assert_eq!(result.len(), 3);
4316 for (i, block_receipts) in result.iter().enumerate() {
4317 assert_eq!(block_receipts.len(), 1);
4318 assert_eq!(block_receipts[0], data.blocks[i].1.receipts()[0][0]);
4319 }
4320 }
4321
4322 #[test]
4323 fn test_receipts_by_block_range_blocks_with_varying_tx_counts() {
4324 let factory = create_test_provider_factory();
4325 let data = BlockchainTestData::default();
4326
4327 let provider_rw = factory.provider_rw().unwrap();
4328 provider_rw.insert_block(&data.genesis.try_recover().unwrap()).unwrap();
4329 provider_rw
4330 .write_state(
4331 &ExecutionOutcome { first_block: 0, receipts: vec![vec![]], ..Default::default() },
4332 crate::OriginalValuesKnown::No,
4333 StateWriteConfig::default(),
4334 )
4335 .unwrap();
4336
4337 for (block, outcome) in data.blocks.iter().take(3) {
4339 provider_rw.insert_block(block).unwrap();
4340 provider_rw
4341 .write_state(outcome, crate::OriginalValuesKnown::No, StateWriteConfig::default())
4342 .unwrap();
4343 }
4344 provider_rw.commit().unwrap();
4345
4346 let provider = factory.provider().unwrap();
4347 let result = provider.receipts_by_block_range(1..=3).unwrap();
4348
4349 assert_eq!(result.len(), 3);
4351 for block_receipts in &result {
4352 assert_eq!(block_receipts.len(), 1);
4353 }
4354 }
4355
4356 #[test]
4357 fn test_receipts_by_block_range_partial_range() {
4358 let factory = create_test_provider_factory();
4359 let data = BlockchainTestData::default();
4360
4361 let provider_rw = factory.provider_rw().unwrap();
4362 provider_rw.insert_block(&data.genesis.try_recover().unwrap()).unwrap();
4363 provider_rw
4364 .write_state(
4365 &ExecutionOutcome { first_block: 0, receipts: vec![vec![]], ..Default::default() },
4366 crate::OriginalValuesKnown::No,
4367 StateWriteConfig::default(),
4368 )
4369 .unwrap();
4370 for (block, outcome) in data.blocks.iter().take(3) {
4371 provider_rw.insert_block(block).unwrap();
4372 provider_rw
4373 .write_state(outcome, crate::OriginalValuesKnown::No, StateWriteConfig::default())
4374 .unwrap();
4375 }
4376 provider_rw.commit().unwrap();
4377
4378 let provider = factory.provider().unwrap();
4379
4380 let result = provider.receipts_by_block_range(2..=5).unwrap();
4382 assert_eq!(result.len(), 4);
4383
4384 assert_eq!(result[0].len(), 1); assert_eq!(result[1].len(), 1); assert_eq!(result[2].len(), 0); assert_eq!(result[3].len(), 0); assert_eq!(result[0][0], data.blocks[1].1.receipts()[0][0]);
4391 assert_eq!(result[1][0], data.blocks[2].1.receipts()[0][0]);
4392 }
4393
4394 #[test]
4395 fn test_receipts_by_block_range_all_empty_blocks() {
4396 let factory = create_test_provider_factory();
4397 let mut rng = generators::rng();
4398
4399 let mut blocks = Vec::new();
4401 for i in 0..3 {
4402 let block =
4403 random_block(&mut rng, i, BlockParams { tx_count: Some(0), ..Default::default() });
4404 blocks.push(block);
4405 }
4406
4407 let provider_rw = factory.provider_rw().unwrap();
4408 for block in blocks {
4409 provider_rw.insert_block(&block.try_recover().unwrap()).unwrap();
4410 }
4411 provider_rw.commit().unwrap();
4412
4413 let provider = factory.provider().unwrap();
4414 let result = provider.receipts_by_block_range(1..=3).unwrap();
4415
4416 assert_eq!(result.len(), 3);
4417 for block_receipts in result {
4418 assert_eq!(block_receipts.len(), 0);
4419 }
4420 }
4421
4422 #[test]
4423 fn test_receipts_by_block_range_consistency_with_individual_calls() {
4424 let factory = create_test_provider_factory();
4425 let data = BlockchainTestData::default();
4426
4427 let provider_rw = factory.provider_rw().unwrap();
4428 provider_rw.insert_block(&data.genesis.try_recover().unwrap()).unwrap();
4429 provider_rw
4430 .write_state(
4431 &ExecutionOutcome { first_block: 0, receipts: vec![vec![]], ..Default::default() },
4432 crate::OriginalValuesKnown::No,
4433 StateWriteConfig::default(),
4434 )
4435 .unwrap();
4436 for (block, outcome) in data.blocks.iter().take(3) {
4437 provider_rw.insert_block(block).unwrap();
4438 provider_rw
4439 .write_state(outcome, crate::OriginalValuesKnown::No, StateWriteConfig::default())
4440 .unwrap();
4441 }
4442 provider_rw.commit().unwrap();
4443
4444 let provider = factory.provider().unwrap();
4445
4446 let range_result = provider.receipts_by_block_range(1..=3).unwrap();
4448
4449 let mut individual_results = Vec::new();
4451 for block_num in 1..=3 {
4452 let receipts =
4453 provider.receipts_by_block(block_num.into()).unwrap().unwrap_or_default();
4454 individual_results.push(receipts);
4455 }
4456
4457 assert_eq!(range_result, individual_results);
4458 }
4459
4460 #[test]
4461 fn test_receipts_by_block_returns_none_for_missing_unpruned_receipts() {
4462 let factory = create_test_provider_factory();
4463 let data = BlockchainTestData::default();
4464
4465 let provider_rw = factory.provider_rw().unwrap();
4466 provider_rw.insert_block(&data.genesis.try_recover().unwrap()).unwrap();
4467 provider_rw.insert_block(&data.blocks[0].0).unwrap();
4468 provider_rw.commit().unwrap();
4469
4470 let provider = factory.provider().unwrap();
4471 assert!(provider.receipts_by_block(1.into()).unwrap().is_none());
4472 }
4473
4474 #[test]
4475 fn test_write_trie_updates_sorted() {
4476 use reth_trie::{
4477 updates::{StorageTrieUpdatesSorted, TrieUpdatesSorted},
4478 BranchNodeCompact, StorageTrieEntry,
4479 };
4480
4481 let factory = create_test_provider_factory();
4482 let provider_rw = factory.provider_rw().unwrap();
4483
4484 {
4486 let tx = provider_rw.tx_ref();
4487 let mut cursor = tx.cursor_write::<tables::AccountsTrie>().unwrap();
4488
4489 let to_delete = StoredNibbles(Nibbles::from_nibbles([0x3, 0x4]));
4491 cursor
4492 .upsert(
4493 to_delete,
4494 &BranchNodeCompact::new(
4495 0b1010_1010_1010_1010, 0b0000_0000_0000_0000, 0b0000_0000_0000_0000, vec![],
4499 None,
4500 ),
4501 )
4502 .unwrap();
4503
4504 let to_update = StoredNibbles(Nibbles::from_nibbles([0x1, 0x2]));
4506 cursor
4507 .upsert(
4508 to_update,
4509 &BranchNodeCompact::new(
4510 0b0101_0101_0101_0101, 0b0000_0000_0000_0000, 0b0000_0000_0000_0000, vec![],
4514 None,
4515 ),
4516 )
4517 .unwrap();
4518 }
4519
4520 let storage_address1 = B256::from([1u8; 32]);
4522 let storage_address2 = B256::from([2u8; 32]);
4523 {
4524 let tx = provider_rw.tx_ref();
4525 let mut storage_cursor = tx.cursor_dup_write::<tables::StoragesTrie>().unwrap();
4526
4527 storage_cursor
4529 .upsert(
4530 storage_address1,
4531 &StorageTrieEntry {
4532 nibbles: StoredNibblesSubKey(Nibbles::from_nibbles([0x2, 0x0])),
4533 node: BranchNodeCompact::new(
4534 0b0011_0011_0011_0011, 0b0000_0000_0000_0000,
4536 0b0000_0000_0000_0000,
4537 vec![],
4538 None,
4539 ),
4540 },
4541 )
4542 .unwrap();
4543
4544 storage_cursor
4546 .upsert(
4547 storage_address2,
4548 &StorageTrieEntry {
4549 nibbles: StoredNibblesSubKey(Nibbles::from_nibbles([0xa, 0xb])),
4550 node: BranchNodeCompact::new(
4551 0b1100_1100_1100_1100,
4552 0b0000_0000_0000_0000,
4553 0b0000_0000_0000_0000,
4554 vec![],
4555 None,
4556 ),
4557 },
4558 )
4559 .unwrap();
4560 storage_cursor
4561 .upsert(
4562 storage_address2,
4563 &StorageTrieEntry {
4564 nibbles: StoredNibblesSubKey(Nibbles::from_nibbles([0xc, 0xd])),
4565 node: BranchNodeCompact::new(
4566 0b0011_1100_0011_1100,
4567 0b0000_0000_0000_0000,
4568 0b0000_0000_0000_0000,
4569 vec![],
4570 None,
4571 ),
4572 },
4573 )
4574 .unwrap();
4575 }
4576
4577 let account_nodes = vec![
4579 (
4580 Nibbles::from_nibbles([0x1, 0x2]),
4581 Some(BranchNodeCompact::new(
4582 0b1111_1111_1111_1111, 0b0000_0000_0000_0000, 0b0000_0000_0000_0000, vec![],
4586 None,
4587 )),
4588 ),
4589 (Nibbles::from_nibbles([0x3, 0x4]), None), (
4591 Nibbles::from_nibbles([0x5, 0x6]),
4592 Some(BranchNodeCompact::new(
4593 0b1111_1111_1111_1111, 0b0000_0000_0000_0000, 0b0000_0000_0000_0000, vec![],
4597 None,
4598 )),
4599 ),
4600 ];
4601
4602 let storage_trie1 = StorageTrieUpdatesSorted {
4604 storage_nodes: vec![
4605 (
4606 Nibbles::from_nibbles([0x1, 0x0]),
4607 Some(BranchNodeCompact::new(
4608 0b1111_0000_0000_0000, 0b0000_0000_0000_0000, 0b0000_0000_0000_0000, vec![],
4612 None,
4613 )),
4614 ),
4615 (Nibbles::from_nibbles([0x2, 0x0]), None), ],
4617 };
4618
4619 let storage_trie2 = StorageTrieUpdatesSorted {
4620 storage_nodes: vec![
4621 (Nibbles::from_nibbles([0xa, 0xb]), None),
4622 (Nibbles::from_nibbles([0xc, 0xd]), None),
4623 ],
4624 };
4625
4626 let mut storage_tries = B256Map::default();
4627 storage_tries.insert(storage_address1, storage_trie1);
4628 storage_tries.insert(storage_address2, storage_trie2);
4629
4630 let trie_updates = TrieUpdatesSorted::new(account_nodes, storage_tries);
4631
4632 let num_entries = provider_rw.write_trie_updates_sorted(&trie_updates).unwrap();
4634
4635 assert_eq!(num_entries, 7);
4638
4639 let tx = provider_rw.tx_ref();
4641 let mut cursor = tx.cursor_read::<tables::AccountsTrie>().unwrap();
4642
4643 let nibbles1 = StoredNibbles(Nibbles::from_nibbles([0x1, 0x2]));
4645 let entry1 = cursor.seek_exact(nibbles1).unwrap();
4646 assert!(entry1.is_some(), "Updated account node should exist");
4647 let expected_mask = reth_trie::TrieMask::new(0b1111_1111_1111_1111);
4648 assert_eq!(
4649 entry1.unwrap().1.state_mask,
4650 expected_mask,
4651 "Account node should have updated state_mask"
4652 );
4653
4654 let nibbles2 = StoredNibbles(Nibbles::from_nibbles([0x3, 0x4]));
4656 let entry2 = cursor.seek_exact(nibbles2).unwrap();
4657 assert!(entry2.is_none(), "Deleted account node should not exist");
4658
4659 let nibbles3 = StoredNibbles(Nibbles::from_nibbles([0x5, 0x6]));
4661 let entry3 = cursor.seek_exact(nibbles3).unwrap();
4662 assert!(entry3.is_some(), "New account node should exist");
4663
4664 let mut storage_cursor = tx.cursor_dup_read::<tables::StoragesTrie>().unwrap();
4666
4667 let storage_entries1: Vec<_> = storage_cursor
4669 .walk_dup(Some(storage_address1), None)
4670 .unwrap()
4671 .collect::<Result<Vec<_>, _>>()
4672 .unwrap();
4673 assert_eq!(
4674 storage_entries1.len(),
4675 1,
4676 "Storage address1 should have 1 entry after deletion"
4677 );
4678 assert_eq!(
4679 storage_entries1[0].1.nibbles.0,
4680 Nibbles::from_nibbles([0x1, 0x0]),
4681 "Remaining entry should be [0x1, 0x0]"
4682 );
4683
4684 let storage_entries2: Vec<_> = storage_cursor
4686 .walk_dup(Some(storage_address2), None)
4687 .unwrap()
4688 .collect::<Result<Vec<_>, _>>()
4689 .unwrap();
4690 assert_eq!(storage_entries2.len(), 0, "Storage address2 should be empty after removal");
4691
4692 provider_rw.commit().unwrap();
4693 }
4694
4695 #[test]
4696 fn test_save_blocks_only_masks_trie_with_deferred_blocks() {
4697 use reth_trie::{
4698 updates::{StorageTrieUpdatesSorted, TrieUpdatesSorted},
4699 BranchNodeCompact, HashedPostStateSorted, HashedStorageSorted,
4700 };
4701
4702 fn branch(mask: u16) -> BranchNodeCompact {
4703 BranchNodeCompact::new(mask, 0, 0, vec![], None)
4704 }
4705
4706 let factory = create_test_provider_factory();
4707 factory.set_storage_settings_cache(StorageSettings::v1());
4708
4709 let mut test_block_builder = TestBlockBuilder::eth().with_state();
4710 let genesis = test_block_builder.get_executed_blocks(0..1).next().unwrap();
4711 let blocks: Vec<_> = test_block_builder.get_executed_blocks(1..3).collect();
4712
4713 let provider_rw = factory.provider_rw().unwrap();
4714 save_genesis(&provider_rw, &genesis).unwrap();
4715 provider_rw.commit().unwrap();
4716
4717 let kept_account = B256::with_last_byte(0x11);
4718 let masked_account = B256::with_last_byte(0x12);
4719 let kept_storage = B256::with_last_byte(0x21);
4720 let masked_storage = B256::with_last_byte(0x22);
4721 let kept_slot = B256::with_last_byte(0x31);
4722 let masked_slot = B256::with_last_byte(0x32);
4723 let kept_account_node = Nibbles::from_nibbles([0x1, 0x2]);
4724 let masked_account_node = Nibbles::from_nibbles([0x1, 0x3]);
4725 let kept_storage_node = Nibbles::from_nibbles([0x2, 0x1]);
4726 let masked_storage_node = Nibbles::from_nibbles([0x2, 0x2]);
4727 let full_persist_base = &blocks[0];
4728 let deferred_trie_base = &blocks[1];
4729
4730 let full_persist_hashed_state = HashedPostStateSorted::new(
4731 vec![
4732 (kept_account, Some(Account::default())),
4733 (masked_account, Some(Account { nonce: 1, ..Default::default() })),
4734 ],
4735 B256Map::from_iter([
4736 (
4737 kept_storage,
4738 HashedStorageSorted { storage_slots: vec![(kept_slot, U256::from(1))] },
4739 ),
4740 (
4741 masked_storage,
4742 HashedStorageSorted { storage_slots: vec![(masked_slot, U256::from(2))] },
4743 ),
4744 ]),
4745 );
4746 let full_persist_trie_updates = TrieUpdatesSorted::new(
4747 vec![
4748 (kept_account_node, Some(branch(0b0000_1111_0000_1111))),
4749 (masked_account_node, Some(branch(0b1111_0000_1111_0000))),
4750 ],
4751 B256Map::from_iter([
4752 (
4753 kept_storage,
4754 StorageTrieUpdatesSorted {
4755 storage_nodes: vec![(kept_storage_node, Some(branch(0b1010)))],
4756 },
4757 ),
4758 (
4759 masked_storage,
4760 StorageTrieUpdatesSorted {
4761 storage_nodes: vec![(masked_storage_node, Some(branch(0b0101)))],
4762 },
4763 ),
4764 ]),
4765 );
4766
4767 let full_persist_block = ExecutedBlock::new(
4768 Arc::clone(&full_persist_base.recovered_block),
4769 Arc::clone(&full_persist_base.execution_output),
4770 Arc::new(full_persist_hashed_state),
4771 Arc::new(full_persist_trie_updates),
4772 );
4773
4774 let deferred_trie_hashed_state = HashedPostStateSorted::new(
4775 vec![(masked_account, Some(Account { nonce: 3, ..Default::default() }))],
4776 B256Map::from_iter([(
4777 masked_storage,
4778 HashedStorageSorted { storage_slots: vec![(masked_slot, U256::from(4))] },
4779 )]),
4780 );
4781 let deferred_trie_updates = TrieUpdatesSorted::new(
4782 vec![(masked_account_node, Some(branch(0b0011_0011)))],
4783 B256Map::from_iter([(
4784 masked_storage,
4785 StorageTrieUpdatesSorted {
4786 storage_nodes: vec![(masked_storage_node, Some(branch(0b1100)))],
4787 },
4788 )]),
4789 );
4790 let deferred_trie_block = ExecutedBlock::new(
4791 Arc::clone(&deferred_trie_base.recovered_block),
4792 Arc::clone(&deferred_trie_base.execution_output),
4793 Arc::new(deferred_trie_hashed_state),
4794 Arc::new(deferred_trie_updates),
4795 );
4796
4797 let provider_rw = factory.provider_rw().unwrap();
4798 let input = SaveBlocksInput::new(
4799 vec![full_persist_block.clone(), deferred_trie_block.clone()],
4800 0,
4801 0,
4802 2,
4803 0,
4804 );
4805 provider_rw.save_blocks(&input).unwrap();
4806 provider_rw.commit().unwrap();
4807
4808 let provider_rw = factory.provider_rw().unwrap();
4809 let input =
4810 SaveBlocksInput::new(vec![full_persist_block, deferred_trie_block.clone()], 2, 0, 2, 1);
4811 assert!(input.persist_rest_blocks().is_empty());
4812 provider_rw.save_blocks(&input).unwrap();
4813 provider_rw.commit().unwrap();
4814
4815 let provider = factory.provider().unwrap();
4816 let tx = provider.tx_ref();
4817 let finish_checkpoint = provider.get_stage_checkpoint(StageId::Finish).unwrap().unwrap();
4818 assert_eq!(finish_checkpoint.block_number, 2);
4819 assert_eq!(
4820 finish_checkpoint.finish_stage_checkpoint().unwrap().partial_state_trie,
4821 Some(1)
4822 );
4823 assert!(provider.block_hash(2).unwrap().is_some());
4824
4825 let mut hashed_accounts = tx.cursor_read::<tables::HashedAccounts>().unwrap();
4826 assert!(hashed_accounts.seek_exact(kept_account).unwrap().is_some());
4827 assert!(hashed_accounts.seek_exact(masked_account).unwrap().is_none());
4828
4829 let mut hashed_storages = tx.cursor_dup_read::<tables::HashedStorages>().unwrap();
4830 assert!(hashed_storages.seek_by_key_subkey(kept_storage, kept_slot).unwrap().is_some());
4831 assert!(hashed_storages
4832 .walk_dup(Some(masked_storage), None)
4833 .unwrap()
4834 .next()
4835 .transpose()
4836 .unwrap()
4837 .is_none());
4838
4839 let mut account_trie = tx.cursor_read::<tables::AccountsTrie>().unwrap();
4840 assert!(account_trie.seek_exact(StoredNibbles(kept_account_node)).unwrap().is_some());
4841 assert!(account_trie.seek_exact(StoredNibbles(masked_account_node)).unwrap().is_none());
4842
4843 let mut storage_trie = tx.cursor_dup_read::<tables::StoragesTrie>().unwrap();
4844 let kept_entries: Vec<_> = storage_trie
4845 .walk_dup(Some(kept_storage), None)
4846 .unwrap()
4847 .collect::<Result<Vec<_>, _>>()
4848 .unwrap();
4849 assert_eq!(kept_entries.len(), 1);
4850 assert_eq!(kept_entries[0].1.nibbles.0, kept_storage_node);
4851
4852 let masked_entries: Vec<_> = storage_trie
4853 .walk_dup(Some(masked_storage), None)
4854 .unwrap()
4855 .collect::<Result<Vec<_>, _>>()
4856 .unwrap();
4857 assert!(masked_entries.is_empty());
4858
4859 drop(storage_trie);
4860 drop(account_trie);
4861 drop(hashed_storages);
4862 drop(hashed_accounts);
4863 drop(provider);
4864
4865 let provider_rw = factory.provider_rw().unwrap();
4866 let input = SaveBlocksInput::new(vec![deferred_trie_block], 2, 1, 2, 2);
4867 assert!(input.persist_rest_blocks().is_empty());
4868 provider_rw.save_blocks(&input).unwrap();
4869 provider_rw.commit().unwrap();
4870
4871 let provider = factory.provider().unwrap();
4872 let finish_checkpoint = provider.get_stage_checkpoint(StageId::Finish).unwrap().unwrap();
4873 assert_eq!(finish_checkpoint.block_number, 2);
4874 assert!(finish_checkpoint.finish_stage_checkpoint().is_none());
4875
4876 let mut hashed_accounts =
4877 provider.tx_ref().cursor_read::<tables::HashedAccounts>().unwrap();
4878 let (_, account) = hashed_accounts.seek_exact(masked_account).unwrap().unwrap();
4879 assert_eq!(account.nonce, 3);
4880
4881 let mut hashed_storages =
4882 provider.tx_ref().cursor_dup_read::<tables::HashedStorages>().unwrap();
4883 let storage =
4884 hashed_storages.seek_by_key_subkey(masked_storage, masked_slot).unwrap().unwrap();
4885 assert_eq!(storage.value, U256::from(4));
4886
4887 let mut account_trie = provider.tx_ref().cursor_read::<tables::AccountsTrie>().unwrap();
4888 assert!(account_trie.seek_exact(StoredNibbles(masked_account_node)).unwrap().is_some());
4889
4890 let mut storage_trie = provider.tx_ref().cursor_dup_read::<tables::StoragesTrie>().unwrap();
4891 let masked_entries: Vec<_> = storage_trie
4892 .walk_dup(Some(masked_storage), None)
4893 .unwrap()
4894 .collect::<Result<Vec<_>, _>>()
4895 .unwrap();
4896 assert_eq!(masked_entries.len(), 1);
4897 assert_eq!(masked_entries[0].1.nibbles.0, masked_storage_node);
4898 }
4899
4900 #[test]
4901 fn test_save_blocks_partial_cycles_do_not_duplicate_static_file_writes() {
4902 let factory = create_test_provider_factory();
4903 let mut test_block_builder = TestBlockBuilder::eth().with_state();
4904
4905 let genesis = test_block_builder.get_executed_blocks(0..1).next().unwrap();
4906 let blocks: Vec<_> = test_block_builder.get_executed_blocks(1..5).collect();
4907
4908 let provider_rw = factory.provider_rw().unwrap();
4909 save_genesis(&provider_rw, &genesis).unwrap();
4910 provider_rw.commit().unwrap();
4911
4912 let provider_rw = factory.provider_rw().unwrap();
4913 let input = SaveBlocksInput::new(blocks[..2].to_vec(), 0, 0, 2, 2);
4914 provider_rw.save_blocks(&input).unwrap();
4915 provider_rw.commit().unwrap();
4916
4917 let provider_rw = factory.provider_rw().unwrap();
4918 let input = SaveBlocksInput::new(blocks[2..].to_vec(), 2, 2, 4, 2);
4919 provider_rw.save_blocks(&input).unwrap();
4920 provider_rw.commit().unwrap();
4921
4922 let provider_rw = factory.provider_rw().unwrap();
4923 let stale_input = SaveBlocksInput::new(vec![blocks[0].clone()], 0, 0, 1, 1);
4924 let err = provider_rw.save_blocks(&stale_input).unwrap_err();
4925 assert!(err.to_string().contains("persistence frontiers do not match Finish checkpoint"));
4926 drop(provider_rw);
4927
4928 let provider = factory.provider().unwrap();
4929 let finish_checkpoint = provider.get_stage_checkpoint(StageId::Finish).unwrap().unwrap();
4930 assert_eq!(finish_checkpoint.block_number, 4);
4931 assert_eq!(
4932 finish_checkpoint.finish_stage_checkpoint().unwrap().partial_state_trie,
4933 Some(2)
4934 );
4935
4936 let static_files = factory.static_file_provider();
4937 assert_eq!(static_files.get_highest_static_file_block(StaticFileSegment::Headers), Some(4));
4938 assert_eq!(
4939 static_files.get_highest_static_file_block(StaticFileSegment::Transactions),
4940 Some(4)
4941 );
4942 assert_eq!(
4943 static_files.get_highest_static_file_block(StaticFileSegment::Receipts),
4944 Some(4)
4945 );
4946 }
4947
4948 #[test]
4949 fn remove_block_and_execution_above_returns_persistence_frontiers() {
4950 let factory = create_test_provider_factory();
4951 let mut test_block_builder = TestBlockBuilder::eth().with_state();
4952
4953 let genesis = test_block_builder.get_executed_blocks(0..1).next().unwrap();
4954 let blocks: Vec<_> = test_block_builder.get_executed_blocks(1..5).collect();
4955
4956 let provider_rw = factory.provider_rw().unwrap();
4957 save_genesis(&provider_rw, &genesis).unwrap();
4958 provider_rw.commit().unwrap();
4959
4960 for block in &blocks[2..] {
4961 factory.overlay_manager().insert_block(block.clone());
4962 }
4963
4964 let provider_rw = factory.provider_rw().unwrap();
4965 let input = SaveBlocksInput::new(blocks, 0, 0, 4, 2);
4966 provider_rw.save_blocks(&input).unwrap();
4967 provider_rw.commit().unwrap();
4968
4969 let provider_rw = factory.provider_rw().unwrap();
4970 let frontiers = provider_rw.remove_block_and_execution_above(3).unwrap();
4971 assert_eq!(frontiers, PersistenceFrontiers { db_tip: 3, partial_state_trie: 2 });
4972 provider_rw.commit().unwrap();
4973
4974 let provider = factory.provider().unwrap();
4975 let checkpoint = provider.get_stage_checkpoint(StageId::Finish).unwrap().unwrap();
4976 assert_eq!(
4977 checkpoint.finish_stage_checkpoint().and_then(|finish| finish.partial_state_trie()),
4978 Some(2)
4979 );
4980 }
4981
4982 #[test]
4983 fn test_prunable_receipts_logic() {
4984 let insert_blocks =
4985 |provider_rw: &DatabaseProviderRW<_, _>, tip_block: u64, tx_count: u8| {
4986 let mut rng = generators::rng();
4987 for block_num in 0..=tip_block {
4988 let block = random_block(
4989 &mut rng,
4990 block_num,
4991 BlockParams { tx_count: Some(tx_count), ..Default::default() },
4992 );
4993 provider_rw.insert_block(&block.try_recover().unwrap()).unwrap();
4994 }
4995 };
4996
4997 let write_receipts = |provider_rw: DatabaseProviderRW<_, _>, block: u64| {
4998 let outcome = ExecutionOutcome {
4999 first_block: block,
5000 receipts: vec![vec![Receipt {
5001 tx_type: Default::default(),
5002 success: true,
5003 cumulative_gas_used: block, logs: vec![],
5005 }]],
5006 ..Default::default()
5007 };
5008 provider_rw
5009 .write_state(&outcome, crate::OriginalValuesKnown::No, StateWriteConfig::default())
5010 .unwrap();
5011 provider_rw.commit().unwrap();
5012 };
5013
5014 {
5016 let factory = create_test_provider_factory();
5017 let storage_settings = StorageSettings::v1();
5018 factory.set_storage_settings_cache(storage_settings);
5019 let factory = factory.with_prune_modes(PruneModes {
5020 receipts: Some(PruneMode::Before(100)),
5021 ..Default::default()
5022 });
5023
5024 let tip_block = 200u64;
5025 let first_block = 1u64;
5026
5027 let provider_rw = factory.provider_rw().unwrap();
5029 insert_blocks(&provider_rw, tip_block, 1);
5030 provider_rw.commit().unwrap();
5031
5032 write_receipts(
5033 factory.provider_rw().unwrap().with_minimum_pruning_distance(100),
5034 first_block,
5035 );
5036 write_receipts(
5037 factory.provider_rw().unwrap().with_minimum_pruning_distance(100),
5038 tip_block - 1,
5039 );
5040
5041 let provider = factory.provider().unwrap();
5042
5043 assert!(provider.receipts_by_block(0.into()).unwrap().is_none());
5044 assert!(provider
5045 .receipts_by_block((tip_block - 1).into())
5046 .unwrap()
5047 .is_some_and(|r| r.len() == 1));
5048 }
5049
5050 {
5052 let factory = create_test_provider_factory();
5053 let storage_settings = StorageSettings::v2();
5054 factory.set_storage_settings_cache(storage_settings);
5055 let factory = factory.with_prune_modes(PruneModes {
5056 receipts: Some(PruneMode::Before(2)),
5057 ..Default::default()
5058 });
5059
5060 let tip_block = 200u64;
5061
5062 let provider_rw = factory.provider_rw().unwrap();
5064 insert_blocks(&provider_rw, tip_block, 1);
5065 provider_rw.commit().unwrap();
5066
5067 write_receipts(factory.provider_rw().unwrap().with_minimum_pruning_distance(100), 0);
5069 write_receipts(factory.provider_rw().unwrap().with_minimum_pruning_distance(100), 1);
5070
5071 assert!(factory
5072 .static_file_provider()
5073 .get_highest_static_file_tx(StaticFileSegment::Receipts)
5074 .is_none(),);
5075 assert!(factory
5076 .static_file_provider()
5077 .get_highest_static_file_block(StaticFileSegment::Receipts)
5078 .is_some_and(|b| b == 1),);
5079
5080 write_receipts(factory.provider_rw().unwrap().with_minimum_pruning_distance(100), 2);
5083 assert!(factory
5084 .static_file_provider()
5085 .get_highest_static_file_tx(StaticFileSegment::Receipts)
5086 .is_some_and(|num| num == 2),);
5087
5088 let factory = factory.with_prune_modes(PruneModes {
5092 receipts: Some(PruneMode::Before(100)),
5093 ..Default::default()
5094 });
5095 let provider_rw = factory.provider_rw().unwrap().with_minimum_pruning_distance(1);
5096 assert!(PruneMode::Distance(1).should_prune(3, tip_block));
5097 write_receipts(provider_rw, 3);
5098
5099 let provider = factory.provider().unwrap();
5104 assert!(EitherWriter::receipts_destination(&provider).is_static_file());
5105 for (num, has_receipt) in [(0, false), (1, false), (2, true), (3, true)] {
5106 let receipts = provider.receipts_by_block(num.into()).unwrap();
5107 if has_receipt {
5108 assert!(receipts.is_some_and(|r| r.len() == 1));
5109 } else {
5110 assert!(receipts.is_none());
5111 }
5112
5113 let receipt = provider.receipt(num).unwrap();
5114 if has_receipt {
5115 assert!(receipt.is_some_and(|r| r.cumulative_gas_used == num));
5116 } else {
5117 assert!(receipt.is_none());
5118 }
5119 }
5120 }
5121 }
5122
5123 #[test]
5124 fn test_unwind_storage_hashing_with_hashed_state() {
5125 let factory = create_test_provider_factory();
5126 let storage_settings = StorageSettings::v2();
5127 factory.set_storage_settings_cache(storage_settings);
5128
5129 let address = Address::random();
5130 let hashed_address = keccak256(address);
5131
5132 let plain_slot = B256::random();
5133 let hashed_slot = keccak256(plain_slot);
5134
5135 let current_value = U256::from(100);
5136 let old_value = U256::from(42);
5137
5138 let provider_rw = factory.provider_rw().unwrap();
5139 provider_rw
5140 .tx
5141 .cursor_dup_write::<tables::HashedStorages>()
5142 .unwrap()
5143 .upsert(hashed_address, &StorageEntry { key: hashed_slot, value: current_value })
5144 .unwrap();
5145
5146 let changesets = vec![(
5147 BlockNumberAddress((1, address)),
5148 StorageEntry { key: plain_slot, value: old_value },
5149 )];
5150
5151 let result = provider_rw.unwind_storage_hashing(changesets.into_iter()).unwrap();
5152
5153 assert_eq!(result.len(), 1);
5154 assert!(result.contains_key(&hashed_address));
5155 assert!(result[&hashed_address].contains(&hashed_slot));
5156
5157 let mut cursor = provider_rw.tx.cursor_dup_read::<tables::HashedStorages>().unwrap();
5158 let entry = cursor
5159 .seek_by_key_subkey(hashed_address, hashed_slot)
5160 .unwrap()
5161 .expect("entry should exist");
5162 assert_eq!(entry.key, hashed_slot);
5163 assert_eq!(entry.value, old_value);
5164 }
5165
5166 #[test]
5167 fn test_write_and_remove_state_roundtrip_legacy() {
5168 let factory = create_test_provider_factory();
5169 let storage_settings = StorageSettings::v1();
5170 assert!(!storage_settings.use_hashed_state());
5171 factory.set_storage_settings_cache(storage_settings);
5172
5173 let address = Address::with_last_byte(1);
5174 let hashed_address = keccak256(address);
5175 let slot = U256::from(5);
5176 let slot_key = B256::from(slot);
5177 let hashed_slot = keccak256(slot_key);
5178
5179 let mut rng = generators::rng();
5180 let block0 =
5181 random_block(&mut rng, 0, BlockParams { tx_count: Some(0), ..Default::default() });
5182 let block1 =
5183 random_block(&mut rng, 1, BlockParams { tx_count: Some(0), ..Default::default() });
5184
5185 {
5186 let provider_rw = factory.provider_rw().unwrap();
5187 provider_rw.insert_block(&block0.try_recover().unwrap()).unwrap();
5188 provider_rw.insert_block(&block1.try_recover().unwrap()).unwrap();
5189 provider_rw
5190 .tx
5191 .cursor_write::<tables::PlainAccountState>()
5192 .unwrap()
5193 .upsert(address, &Account::default())
5194 .unwrap();
5195 provider_rw.commit().unwrap();
5196 }
5197
5198 let provider_rw = factory.provider_rw().unwrap();
5199
5200 let mut state_init: BundleStateInit = AddressMap::default();
5201 let mut storage_map: B256Map<(U256, U256)> = B256Map::default();
5202 storage_map.insert(slot_key, (U256::ZERO, U256::from(10)));
5203 state_init.insert(
5204 address,
5205 (
5206 Some(Account::default()),
5207 Some(Account { nonce: 1, ..Default::default() }),
5208 storage_map,
5209 ),
5210 );
5211
5212 let mut reverts_init: RevertsInit = HashMap::default();
5213 let mut block_reverts: AddressMap<AccountRevertInit> = AddressMap::default();
5214 block_reverts.insert(
5215 address,
5216 (
5217 Some(Some(Account::default())),
5218 vec![StorageEntry { key: slot_key, value: U256::ZERO }],
5219 ),
5220 );
5221 reverts_init.insert(1, block_reverts);
5222
5223 let execution_outcome =
5224 ExecutionOutcome::new_init(state_init, reverts_init, [], vec![vec![]], 1, vec![]);
5225
5226 provider_rw
5227 .write_state(
5228 &execution_outcome,
5229 OriginalValuesKnown::Yes,
5230 StateWriteConfig {
5231 write_receipts: false,
5232 write_account_changesets: true,
5233 write_storage_changesets: true,
5234 },
5235 )
5236 .unwrap();
5237
5238 let hashed_state =
5239 execution_outcome.hash_state_slow::<reth_trie::KeccakKeyHasher>().into_sorted();
5240 provider_rw.write_hashed_state(&hashed_state).unwrap();
5241
5242 let account = provider_rw
5243 .tx
5244 .cursor_read::<tables::PlainAccountState>()
5245 .unwrap()
5246 .seek_exact(address)
5247 .unwrap()
5248 .unwrap()
5249 .1;
5250 assert_eq!(account.nonce, 1);
5251
5252 let storage_entry = provider_rw
5253 .tx
5254 .cursor_dup_read::<tables::PlainStorageState>()
5255 .unwrap()
5256 .seek_by_key_subkey(address, slot_key)
5257 .unwrap()
5258 .unwrap();
5259 assert_eq!(storage_entry.key, slot_key);
5260 assert_eq!(storage_entry.value, U256::from(10));
5261
5262 let hashed_entry = provider_rw
5263 .tx
5264 .cursor_dup_read::<tables::HashedStorages>()
5265 .unwrap()
5266 .seek_by_key_subkey(hashed_address, hashed_slot)
5267 .unwrap()
5268 .unwrap();
5269 assert_eq!(hashed_entry.key, hashed_slot);
5270 assert_eq!(hashed_entry.value, U256::from(10));
5271
5272 let account_cs_entries = provider_rw
5273 .tx
5274 .cursor_dup_read::<tables::AccountChangeSets>()
5275 .unwrap()
5276 .walk(Some(1))
5277 .unwrap()
5278 .collect::<Result<Vec<_>, _>>()
5279 .unwrap();
5280 assert!(!account_cs_entries.is_empty());
5281
5282 let storage_cs_entries = provider_rw
5283 .tx
5284 .cursor_read::<tables::StorageChangeSets>()
5285 .unwrap()
5286 .walk(Some(BlockNumberAddress((1, address))))
5287 .unwrap()
5288 .collect::<Result<Vec<_>, _>>()
5289 .unwrap();
5290 assert!(!storage_cs_entries.is_empty());
5291 assert_eq!(storage_cs_entries[0].1.key, slot_key);
5292
5293 provider_rw.remove_state_above(0).unwrap();
5294
5295 let restored_account = provider_rw
5296 .tx
5297 .cursor_read::<tables::PlainAccountState>()
5298 .unwrap()
5299 .seek_exact(address)
5300 .unwrap()
5301 .unwrap()
5302 .1;
5303 assert_eq!(restored_account.nonce, 0);
5304
5305 let storage_gone = provider_rw
5306 .tx
5307 .cursor_dup_read::<tables::PlainStorageState>()
5308 .unwrap()
5309 .seek_by_key_subkey(address, slot_key)
5310 .unwrap();
5311 assert!(storage_gone.is_none() || storage_gone.unwrap().key != slot_key);
5312
5313 let account_cs_after = provider_rw
5314 .tx
5315 .cursor_dup_read::<tables::AccountChangeSets>()
5316 .unwrap()
5317 .walk(Some(1))
5318 .unwrap()
5319 .collect::<Result<Vec<_>, _>>()
5320 .unwrap();
5321 assert!(account_cs_after.is_empty());
5322
5323 let storage_cs_after = provider_rw
5324 .tx
5325 .cursor_read::<tables::StorageChangeSets>()
5326 .unwrap()
5327 .walk(Some(BlockNumberAddress((1, address))))
5328 .unwrap()
5329 .collect::<Result<Vec<_>, _>>()
5330 .unwrap();
5331 assert!(storage_cs_after.is_empty());
5332 }
5333
5334 #[test]
5335 fn test_unwind_storage_hashing_legacy() {
5336 let factory = create_test_provider_factory();
5337 let storage_settings = StorageSettings::v1();
5338 assert!(!storage_settings.use_hashed_state());
5339 factory.set_storage_settings_cache(storage_settings);
5340
5341 let address = Address::random();
5342 let hashed_address = keccak256(address);
5343
5344 let plain_slot = B256::random();
5345 let hashed_slot = keccak256(plain_slot);
5346
5347 let current_value = U256::from(100);
5348 let old_value = U256::from(42);
5349
5350 let provider_rw = factory.provider_rw().unwrap();
5351 provider_rw
5352 .tx
5353 .cursor_dup_write::<tables::HashedStorages>()
5354 .unwrap()
5355 .upsert(hashed_address, &StorageEntry { key: hashed_slot, value: current_value })
5356 .unwrap();
5357
5358 let changesets = vec![(
5359 BlockNumberAddress((1, address)),
5360 StorageEntry { key: plain_slot, value: old_value },
5361 )];
5362
5363 let result = provider_rw.unwind_storage_hashing(changesets.into_iter()).unwrap();
5364
5365 assert_eq!(result.len(), 1);
5366 assert!(result.contains_key(&hashed_address));
5367 assert!(result[&hashed_address].contains(&hashed_slot));
5368
5369 let mut cursor = provider_rw.tx.cursor_dup_read::<tables::HashedStorages>().unwrap();
5370 let entry = cursor
5371 .seek_by_key_subkey(hashed_address, hashed_slot)
5372 .unwrap()
5373 .expect("entry should exist");
5374 assert_eq!(entry.key, hashed_slot);
5375 assert_eq!(entry.value, old_value);
5376 }
5377
5378 #[test]
5379 fn test_write_state_hashed() {
5380 use reth_trie::{HashedPostState, KeccakKeyHasher};
5381 use revm::{database::BundleState, state::AccountInfo};
5382
5383 let factory = create_test_provider_factory();
5384 factory.set_storage_settings_cache(StorageSettings::v2());
5385
5386 let address = Address::with_last_byte(1);
5387 let slot = U256::from(5);
5388 let slot_key = B256::from(slot);
5389 let hashed_address = keccak256(address);
5390 let hashed_slot = keccak256(slot_key);
5391
5392 {
5393 let sf = factory.static_file_provider();
5394 let mut hw = sf.latest_writer(StaticFileSegment::Headers).unwrap();
5395 let h0 = alloy_consensus::Header { number: 0, ..Default::default() };
5396 hw.append_header(&h0, &B256::ZERO).unwrap();
5397 let h1 = alloy_consensus::Header { number: 1, ..Default::default() };
5398 hw.append_header(&h1, &B256::ZERO).unwrap();
5399 hw.commit().unwrap();
5400
5401 let mut aw = sf.latest_writer(StaticFileSegment::AccountChangeSets).unwrap();
5402 aw.append_account_changeset(vec![], 0).unwrap();
5403 aw.commit().unwrap();
5404
5405 let mut sw = sf.latest_writer(StaticFileSegment::StorageChangeSets).unwrap();
5406 sw.append_storage_changeset(vec![], 0).unwrap();
5407 sw.commit().unwrap();
5408 }
5409
5410 let provider_rw = factory.provider_rw().unwrap();
5411
5412 let bundle = BundleState::builder(1..=1)
5413 .state_present_account_info(
5414 address,
5415 AccountInfo { nonce: 1, balance: U256::from(10), ..Default::default() },
5416 )
5417 .state_storage(address, HashMap::from_iter([(slot, (U256::ZERO, U256::from(10)))]))
5418 .revert_account_info(1, address, Some(None))
5419 .revert_storage(1, address, vec![(slot, U256::ZERO)])
5420 .build();
5421
5422 let execution_outcome = ExecutionOutcome::new(bundle.clone(), vec![vec![]], 1, Vec::new());
5423
5424 provider_rw
5425 .tx
5426 .put::<tables::BlockBodyIndices>(
5427 1,
5428 StoredBlockBodyIndices { first_tx_num: 0, tx_count: 0 },
5429 )
5430 .unwrap();
5431
5432 provider_rw
5433 .write_state(
5434 &execution_outcome,
5435 OriginalValuesKnown::Yes,
5436 StateWriteConfig {
5437 write_receipts: false,
5438 write_account_changesets: true,
5439 write_storage_changesets: true,
5440 },
5441 )
5442 .unwrap();
5443
5444 let hashed_state =
5445 HashedPostState::from_bundle_state::<KeccakKeyHasher>(bundle.state()).into_sorted();
5446 provider_rw.write_hashed_state(&hashed_state).unwrap();
5447
5448 let plain_storage_entries = provider_rw
5449 .tx
5450 .cursor_dup_read::<tables::PlainStorageState>()
5451 .unwrap()
5452 .walk(None)
5453 .unwrap()
5454 .collect::<Result<Vec<_>, _>>()
5455 .unwrap();
5456 assert!(plain_storage_entries.is_empty());
5457
5458 let hashed_entry = provider_rw
5459 .tx
5460 .cursor_dup_read::<tables::HashedStorages>()
5461 .unwrap()
5462 .seek_by_key_subkey(hashed_address, hashed_slot)
5463 .unwrap()
5464 .unwrap();
5465 assert_eq!(hashed_entry.key, hashed_slot);
5466 assert_eq!(hashed_entry.value, U256::from(10));
5467
5468 provider_rw.static_file_provider().commit().unwrap();
5469
5470 let sf = factory.static_file_provider();
5471 let storage_cs = sf.storage_changeset(1).unwrap();
5472 assert!(!storage_cs.is_empty());
5473 assert_eq!(storage_cs[0].1.key, slot_key);
5474
5475 let account_cs = sf.account_block_changeset(1).unwrap();
5476 assert!(!account_cs.is_empty());
5477 assert_eq!(account_cs[0].address, address);
5478 }
5479
5480 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
5481 enum StorageMode {
5482 V1,
5483 V2,
5484 }
5485
5486 fn run_save_blocks_and_verify(mode: StorageMode) {
5487 use alloy_primitives::map::{FbBuildHasher, HashMap};
5488
5489 let factory = create_test_provider_factory();
5490
5491 match mode {
5492 StorageMode::V1 => factory.set_storage_settings_cache(StorageSettings::v1()),
5493 StorageMode::V2 => factory.set_storage_settings_cache(StorageSettings::v2()),
5494 }
5495
5496 let num_blocks = 3u64;
5497 let accounts_per_block = 5usize;
5498 let slots_per_account = 3usize;
5499
5500 let genesis = SealedBlock::<reth_ethereum_primitives::Block>::from_sealed_parts(
5501 SealedHeader::new(
5502 Header { number: 0, difficulty: U256::from(1), ..Default::default() },
5503 B256::ZERO,
5504 ),
5505 Default::default(),
5506 );
5507
5508 let genesis_executed: ExecutedBlock = ExecutedBlock::new(
5509 Arc::new(genesis.try_recover().unwrap()),
5510 Arc::new(BlockExecutionOutput {
5511 result: BlockExecutionResult {
5512 receipts: vec![],
5513 requests: Default::default(),
5514 gas_used: 0,
5515 blob_gas_used: 0,
5516 },
5517 state: Default::default(),
5518 }),
5519 Default::default(),
5520 Default::default(),
5521 );
5522 let provider_rw = factory.provider_rw().unwrap();
5523 save_genesis(&provider_rw, &genesis_executed).unwrap();
5524 provider_rw.commit().unwrap();
5525
5526 let mut blocks: Vec<ExecutedBlock> = Vec::new();
5527 let mut parent_hash = B256::ZERO;
5528
5529 for block_num in 1..=num_blocks {
5530 let mut builder = BundleState::builder(block_num..=block_num);
5531
5532 for acct_idx in 0..accounts_per_block {
5533 let address = Address::with_last_byte((block_num * 10 + acct_idx as u64) as u8);
5534 let info = AccountInfo {
5535 nonce: block_num,
5536 balance: U256::from(block_num * 100 + acct_idx as u64),
5537 ..Default::default()
5538 };
5539
5540 let storage: HashMap<U256, (U256, U256), FbBuildHasher<32>> = (1..=
5541 slots_per_account as u64)
5542 .map(|s| {
5543 (
5544 U256::from(s + acct_idx as u64 * 100),
5545 (U256::ZERO, U256::from(block_num * 1000 + s)),
5546 )
5547 })
5548 .collect();
5549
5550 let revert_storage: Vec<(U256, U256)> = (1..=slots_per_account as u64)
5551 .map(|s| (U256::from(s + acct_idx as u64 * 100), U256::ZERO))
5552 .collect();
5553
5554 builder = builder
5555 .state_present_account_info(address, info)
5556 .revert_account_info(block_num, address, Some(None))
5557 .state_storage(address, storage)
5558 .revert_storage(block_num, address, revert_storage);
5559 }
5560
5561 let bundle = builder.build();
5562
5563 let hashed_state =
5564 HashedPostState::from_bundle_state::<KeccakKeyHasher>(bundle.state()).into_sorted();
5565
5566 let header = Header {
5567 number: block_num,
5568 parent_hash,
5569 difficulty: U256::from(1),
5570 ..Default::default()
5571 };
5572 let block = SealedBlock::<reth_ethereum_primitives::Block>::seal_parts(
5573 header,
5574 Default::default(),
5575 );
5576 parent_hash = block.hash();
5577
5578 let executed = ExecutedBlock::new(
5579 Arc::new(block.try_recover().unwrap()),
5580 Arc::new(BlockExecutionOutput {
5581 result: BlockExecutionResult {
5582 receipts: vec![],
5583 requests: Default::default(),
5584 gas_used: 0,
5585 blob_gas_used: 0,
5586 },
5587 state: bundle,
5588 }),
5589 Arc::new(hashed_state),
5590 Default::default(),
5591 );
5592 blocks.push(executed);
5593 }
5594
5595 let provider_rw = factory.provider_rw().unwrap();
5596 let input = SaveBlocksInput::new(blocks, 0, 0, num_blocks, num_blocks);
5597 provider_rw.save_blocks(&input).unwrap();
5598 provider_rw.commit().unwrap();
5599
5600 let provider = factory.provider().unwrap();
5601
5602 for block_num in 1..=num_blocks {
5603 for acct_idx in 0..accounts_per_block {
5604 let address = Address::with_last_byte((block_num * 10 + acct_idx as u64) as u8);
5605 let hashed_address = keccak256(address);
5606
5607 let ha_entry = provider
5608 .tx_ref()
5609 .cursor_read::<tables::HashedAccounts>()
5610 .unwrap()
5611 .seek_exact(hashed_address)
5612 .unwrap();
5613 assert!(
5614 ha_entry.is_some(),
5615 "HashedAccounts missing for block {block_num} acct {acct_idx}"
5616 );
5617
5618 for s in 1..=slots_per_account as u64 {
5619 let slot = U256::from(s + acct_idx as u64 * 100);
5620 let slot_key = B256::from(slot);
5621 let hashed_slot = keccak256(slot_key);
5622
5623 let hs_entry = provider
5624 .tx_ref()
5625 .cursor_dup_read::<tables::HashedStorages>()
5626 .unwrap()
5627 .seek_by_key_subkey(hashed_address, hashed_slot)
5628 .unwrap();
5629 assert!(
5630 hs_entry.is_some(),
5631 "HashedStorages missing for block {block_num} acct {acct_idx} slot {s}"
5632 );
5633 let entry = hs_entry.unwrap();
5634 assert_eq!(entry.key, hashed_slot);
5635 assert_eq!(entry.value, U256::from(block_num * 1000 + s));
5636 }
5637 }
5638 }
5639
5640 for block_num in 1..=num_blocks {
5641 let header = provider.header_by_number(block_num).unwrap();
5642 assert!(header.is_some(), "Header missing for block {block_num}");
5643
5644 let indices = provider.block_body_indices(block_num).unwrap();
5645 assert!(indices.is_some(), "BlockBodyIndices missing for block {block_num}");
5646 }
5647
5648 let plain_accounts = provider.tx_ref().entries::<tables::PlainAccountState>().unwrap();
5649 let plain_storage = provider.tx_ref().entries::<tables::PlainStorageState>().unwrap();
5650
5651 if mode == StorageMode::V2 {
5652 assert_eq!(plain_accounts, 0, "v2: PlainAccountState should be empty");
5653 assert_eq!(plain_storage, 0, "v2: PlainStorageState should be empty");
5654
5655 let mdbx_account_cs = provider.tx_ref().entries::<tables::AccountChangeSets>().unwrap();
5656 assert_eq!(mdbx_account_cs, 0, "v2: AccountChangeSets in MDBX should be empty");
5657
5658 let mdbx_storage_cs = provider.tx_ref().entries::<tables::StorageChangeSets>().unwrap();
5659 assert_eq!(mdbx_storage_cs, 0, "v2: StorageChangeSets in MDBX should be empty");
5660
5661 provider.static_file_provider().commit().unwrap();
5662 let sf = factory.static_file_provider();
5663
5664 for block_num in 1..=num_blocks {
5665 let account_cs = sf.account_block_changeset(block_num).unwrap();
5666 assert!(
5667 !account_cs.is_empty(),
5668 "v2: static file AccountChangeSets should exist for block {block_num}"
5669 );
5670
5671 let storage_cs = sf.storage_changeset(block_num).unwrap();
5672 assert!(
5673 !storage_cs.is_empty(),
5674 "v2: static file StorageChangeSets should exist for block {block_num}"
5675 );
5676
5677 for (_, entry) in &storage_cs {
5678 assert!(
5679 entry.key != keccak256(entry.key),
5680 "v2: static file storage changeset should have plain slot keys"
5681 );
5682 }
5683 }
5684
5685 let rocksdb = factory.rocksdb_provider();
5686 for block_num in 1..=num_blocks {
5687 for acct_idx in 0..accounts_per_block {
5688 let address = Address::with_last_byte((block_num * 10 + acct_idx as u64) as u8);
5689 let shards = rocksdb.account_history_shards(address).unwrap();
5690 assert!(
5691 !shards.is_empty(),
5692 "v2: RocksDB AccountsHistory missing for block {block_num} acct {acct_idx}"
5693 );
5694
5695 for s in 1..=slots_per_account as u64 {
5696 let slot = U256::from(s + acct_idx as u64 * 100);
5697 let slot_key = B256::from(slot);
5698 let shards = rocksdb.storage_history_shards(address, slot_key).unwrap();
5699 assert!(
5700 !shards.is_empty(),
5701 "v2: RocksDB StoragesHistory missing for block {block_num} acct {acct_idx} slot {s}"
5702 );
5703 }
5704 }
5705 }
5706 } else {
5707 assert!(plain_accounts > 0, "v1: PlainAccountState should not be empty");
5708 assert!(plain_storage > 0, "v1: PlainStorageState should not be empty");
5709
5710 let mdbx_account_cs = provider.tx_ref().entries::<tables::AccountChangeSets>().unwrap();
5711 assert!(mdbx_account_cs > 0, "v1: AccountChangeSets in MDBX should not be empty");
5712
5713 let mdbx_storage_cs = provider.tx_ref().entries::<tables::StorageChangeSets>().unwrap();
5714 assert!(mdbx_storage_cs > 0, "v1: StorageChangeSets in MDBX should not be empty");
5715
5716 for block_num in 1..=num_blocks {
5717 let storage_entries: Vec<_> = provider
5718 .tx_ref()
5719 .cursor_dup_read::<tables::StorageChangeSets>()
5720 .unwrap()
5721 .walk_range(BlockNumberAddress::range(block_num..=block_num))
5722 .unwrap()
5723 .collect::<Result<Vec<_>, _>>()
5724 .unwrap();
5725 assert!(
5726 !storage_entries.is_empty(),
5727 "v1: MDBX StorageChangeSets should have entries for block {block_num}"
5728 );
5729
5730 for (_, entry) in &storage_entries {
5731 let slot_key = B256::from(entry.key);
5732 assert!(
5733 slot_key != keccak256(slot_key),
5734 "v1: storage changeset keys should be plain (not hashed)"
5735 );
5736 }
5737 }
5738
5739 let mdbx_account_history =
5740 provider.tx_ref().entries::<tables::AccountsHistory>().unwrap();
5741 assert!(mdbx_account_history > 0, "v1: AccountsHistory in MDBX should not be empty");
5742
5743 let mdbx_storage_history =
5744 provider.tx_ref().entries::<tables::StoragesHistory>().unwrap();
5745 assert!(mdbx_storage_history > 0, "v1: StoragesHistory in MDBX should not be empty");
5746 }
5747 }
5748
5749 #[test]
5750 fn test_save_blocks_v1_table_assertions() {
5751 run_save_blocks_and_verify(StorageMode::V1);
5752 }
5753
5754 #[test]
5755 fn test_save_blocks_v2_table_assertions() {
5756 run_save_blocks_and_verify(StorageMode::V2);
5757 }
5758
5759 #[test]
5760 fn test_write_and_remove_state_roundtrip_v2() {
5761 let factory = create_test_provider_factory();
5762 let storage_settings = StorageSettings::v2();
5763 assert!(storage_settings.use_hashed_state());
5764 factory.set_storage_settings_cache(storage_settings);
5765
5766 let address = Address::with_last_byte(1);
5767 let hashed_address = keccak256(address);
5768 let slot = U256::from(5);
5769 let slot_key = B256::from(slot);
5770 let hashed_slot = keccak256(slot_key);
5771
5772 {
5773 let sf = factory.static_file_provider();
5774 let mut hw = sf.latest_writer(StaticFileSegment::Headers).unwrap();
5775 let h0 = alloy_consensus::Header { number: 0, ..Default::default() };
5776 hw.append_header(&h0, &B256::ZERO).unwrap();
5777 let h1 = alloy_consensus::Header { number: 1, ..Default::default() };
5778 hw.append_header(&h1, &B256::ZERO).unwrap();
5779 hw.commit().unwrap();
5780
5781 let mut aw = sf.latest_writer(StaticFileSegment::AccountChangeSets).unwrap();
5782 aw.append_account_changeset(vec![], 0).unwrap();
5783 aw.commit().unwrap();
5784
5785 let mut sw = sf.latest_writer(StaticFileSegment::StorageChangeSets).unwrap();
5786 sw.append_storage_changeset(vec![], 0).unwrap();
5787 sw.commit().unwrap();
5788 }
5789
5790 {
5791 let provider_rw = factory.provider_rw().unwrap();
5792 provider_rw
5793 .tx
5794 .put::<tables::BlockBodyIndices>(
5795 0,
5796 StoredBlockBodyIndices { first_tx_num: 0, tx_count: 0 },
5797 )
5798 .unwrap();
5799 provider_rw
5800 .tx
5801 .put::<tables::BlockBodyIndices>(
5802 1,
5803 StoredBlockBodyIndices { first_tx_num: 0, tx_count: 0 },
5804 )
5805 .unwrap();
5806 provider_rw
5807 .tx
5808 .cursor_write::<tables::HashedAccounts>()
5809 .unwrap()
5810 .upsert(hashed_address, &Account::default())
5811 .unwrap();
5812 provider_rw.commit().unwrap();
5813 }
5814
5815 let provider_rw = factory.provider_rw().unwrap();
5816
5817 let bundle = BundleState::builder(1..=1)
5818 .state_present_account_info(
5819 address,
5820 AccountInfo { nonce: 1, balance: U256::from(10), ..Default::default() },
5821 )
5822 .state_storage(address, HashMap::from_iter([(slot, (U256::ZERO, U256::from(10)))]))
5823 .revert_account_info(1, address, Some(None))
5824 .revert_storage(1, address, vec![(slot, U256::ZERO)])
5825 .build();
5826
5827 let execution_outcome = ExecutionOutcome::new(bundle.clone(), vec![vec![]], 1, Vec::new());
5828
5829 provider_rw
5830 .write_state(
5831 &execution_outcome,
5832 OriginalValuesKnown::Yes,
5833 StateWriteConfig {
5834 write_receipts: false,
5835 write_account_changesets: true,
5836 write_storage_changesets: true,
5837 },
5838 )
5839 .unwrap();
5840
5841 let hashed_state =
5842 HashedPostState::from_bundle_state::<KeccakKeyHasher>(bundle.state()).into_sorted();
5843 provider_rw.write_hashed_state(&hashed_state).unwrap();
5844
5845 let hashed_account = provider_rw
5846 .tx
5847 .cursor_read::<tables::HashedAccounts>()
5848 .unwrap()
5849 .seek_exact(hashed_address)
5850 .unwrap()
5851 .unwrap()
5852 .1;
5853 assert_eq!(hashed_account.nonce, 1);
5854
5855 let hashed_entry = provider_rw
5856 .tx
5857 .cursor_dup_read::<tables::HashedStorages>()
5858 .unwrap()
5859 .seek_by_key_subkey(hashed_address, hashed_slot)
5860 .unwrap()
5861 .unwrap();
5862 assert_eq!(hashed_entry.key, hashed_slot);
5863 assert_eq!(hashed_entry.value, U256::from(10));
5864
5865 let plain_accounts = provider_rw.tx.entries::<tables::PlainAccountState>().unwrap();
5866 assert_eq!(plain_accounts, 0, "v2: PlainAccountState should be empty");
5867
5868 let plain_storage = provider_rw.tx.entries::<tables::PlainStorageState>().unwrap();
5869 assert_eq!(plain_storage, 0, "v2: PlainStorageState should be empty");
5870
5871 provider_rw.static_file_provider().commit().unwrap();
5872
5873 let sf = factory.static_file_provider();
5874 let storage_cs = sf.storage_changeset(1).unwrap();
5875 assert!(!storage_cs.is_empty(), "v2: storage changesets should be in static files");
5876 assert_eq!(storage_cs[0].1.key, slot_key, "v2: changeset key should be plain");
5877
5878 provider_rw.remove_state_above(0).unwrap();
5879
5880 let restored_account = provider_rw
5881 .tx
5882 .cursor_read::<tables::HashedAccounts>()
5883 .unwrap()
5884 .seek_exact(hashed_address)
5885 .unwrap();
5886 assert!(
5887 restored_account.is_none(),
5888 "v2: account should be removed (didn't exist before block 1)"
5889 );
5890
5891 let storage_gone = provider_rw
5892 .tx
5893 .cursor_dup_read::<tables::HashedStorages>()
5894 .unwrap()
5895 .seek_by_key_subkey(hashed_address, hashed_slot)
5896 .unwrap();
5897 assert!(
5898 storage_gone.is_none() || storage_gone.unwrap().key != hashed_slot,
5899 "v2: storage should be reverted (removed or different key)"
5900 );
5901
5902 let mdbx_storage_cs = provider_rw.tx.entries::<tables::StorageChangeSets>().unwrap();
5903 assert_eq!(mdbx_storage_cs, 0, "v2: MDBX StorageChangeSets should remain empty");
5904
5905 let mdbx_account_cs = provider_rw.tx.entries::<tables::AccountChangeSets>().unwrap();
5906 assert_eq!(mdbx_account_cs, 0, "v2: MDBX AccountChangeSets should remain empty");
5907 }
5908
5909 #[test]
5910 fn test_unwind_storage_history_indices_v2() {
5911 let factory = create_test_provider_factory();
5912 factory.set_storage_settings_cache(StorageSettings::v2());
5913
5914 let address = Address::with_last_byte(1);
5915 let slot_key = B256::from(U256::from(42));
5916
5917 {
5918 let rocksdb = factory.rocksdb_provider();
5919 let mut batch = rocksdb.batch();
5920 batch.append_storage_history_shard(address, slot_key, vec![3u64, 7, 10]).unwrap();
5921 batch.commit().unwrap();
5922
5923 let shards = rocksdb.storage_history_shards(address, slot_key).unwrap();
5924 assert!(!shards.is_empty(), "history should be written to rocksdb");
5925 }
5926
5927 let provider_rw = factory.provider_rw().unwrap();
5928
5929 let changesets = vec![
5930 (
5931 BlockNumberAddress((7, address)),
5932 StorageEntry { key: slot_key, value: U256::from(5) },
5933 ),
5934 (
5935 BlockNumberAddress((10, address)),
5936 StorageEntry { key: slot_key, value: U256::from(8) },
5937 ),
5938 ];
5939
5940 let count = provider_rw.unwind_storage_history_indices(changesets.into_iter()).unwrap();
5941 assert_eq!(count, 2);
5942
5943 provider_rw.commit().unwrap();
5944
5945 let rocksdb = factory.rocksdb_provider();
5946 let shards = rocksdb.storage_history_shards(address, slot_key).unwrap();
5947
5948 assert!(
5949 !shards.is_empty(),
5950 "history shards should still exist with block 3 after partial unwind"
5951 );
5952
5953 let all_blocks: Vec<u64> = shards.iter().flat_map(|(_, list)| list.iter()).collect();
5954 assert!(all_blocks.contains(&3), "block 3 should remain");
5955 assert!(!all_blocks.contains(&7), "block 7 should be unwound");
5956 assert!(!all_blocks.contains(&10), "block 10 should be unwound");
5957 }
5958}