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