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