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