1use super::{DatabaseProviderRO, ProviderFactory, ProviderNodeTypes};
2use crate::{
3 providers::{StaticFileProvider, StaticFileProviderRWRefMut},
4 to_range, BlockHashReader, BlockIdReader, BlockNumReader, BlockReader, BlockReaderIdExt,
5 BlockSource, ChainSpecProvider, ChangeSetReader, HeaderProvider, ProviderError,
6 PruneCheckpointReader, ReceiptProvider, ReceiptProviderIdExt, StageCheckpointReader,
7 StaticFileProviderFactory, TransactionVariant, TransactionsProvider,
8};
9use alloy_consensus::{
10 transaction::{TransactionMeta, TxHashRef},
11 BlockHeader,
12};
13use alloy_eips::{BlockHashOrNumber, BlockId, BlockNumHash, BlockNumberOrTag, HashOrNumber};
14use alloy_primitives::{Address, BlockHash, BlockNumber, TxHash, TxNumber, B256};
15use reth_chain_state::{BlockState, CanonicalInMemoryState};
16use reth_chainspec::ChainInfo;
17use reth_db_api::models::{AccountBeforeTx, BlockNumberAddress, StoredBlockBodyIndices};
18use reth_node_types::{BlockTy, HeaderTy, ReceiptTy, TxTy};
19use reth_primitives_traits::{
20 BlockBody, RecoveredBlock, SealedHeader, SealedOrRecoveredBlock, StorageEntry,
21};
22use reth_prune_types::{PruneCheckpoint, PruneSegment};
23use reth_stages_types::{StageCheckpoint, StageId};
24use reth_static_file_types::StaticFileSegment;
25use reth_storage_api::{
26 BlockBodyIndicesProvider, DatabaseProviderFactory, NodePrimitivesProvider,
27 StorageChangeSetReader,
28};
29use reth_storage_errors::provider::ProviderResult;
30use revm::database::states::PlainStorageRevert;
31use std::{
32 ops::{Add, Bound, RangeBounds, RangeInclusive, Sub},
33 sync::Arc,
34};
35
36#[derive(Debug)]
43#[doc(hidden)] pub struct ConsistentProvider<N: ProviderNodeTypes> {
45 storage_provider: <ProviderFactory<N> as DatabaseProviderFactory>::Provider,
47 head_block: Option<Arc<BlockState<N::Primitives>>>,
49 canonical_in_memory_state: CanonicalInMemoryState<N::Primitives>,
51}
52
53impl<N: ProviderNodeTypes> ConsistentProvider<N> {
54 pub fn new(
60 storage_provider_factory: ProviderFactory<N>,
61 state: CanonicalInMemoryState<N::Primitives>,
62 ) -> ProviderResult<Self> {
63 let head_block = state.head_state();
71 let storage_provider = storage_provider_factory.database_provider_ro()?;
72 Ok(Self { storage_provider, head_block, canonical_in_memory_state: state })
73 }
74
75 pub(crate) fn into_database_provider(
77 self,
78 ) -> <ProviderFactory<N> as DatabaseProviderFactory>::Provider {
79 self.storage_provider
80 }
81
82 fn convert_range_bounds<T>(
84 &self,
85 range: impl RangeBounds<T>,
86 end_unbounded: impl FnOnce() -> T,
87 ) -> (T, T)
88 where
89 T: Copy + Add<Output = T> + Sub<Output = T> + From<u8>,
90 {
91 let start = match range.start_bound() {
92 Bound::Included(&n) => n,
93 Bound::Excluded(&n) => n + T::from(1u8),
94 Bound::Unbounded => T::from(0u8),
95 };
96
97 let end = match range.end_bound() {
98 Bound::Included(&n) => n,
99 Bound::Excluded(&n) => n - T::from(1u8),
100 Bound::Unbounded => end_unbounded(),
101 };
102
103 (start, end)
104 }
105
106 fn get_in_memory_or_storage_by_block_range_while<T, F, G, P>(
118 &self,
119 range: impl RangeBounds<BlockNumber>,
120 fetch_db_range: F,
121 map_block_state_item: G,
122 mut predicate: P,
123 ) -> ProviderResult<Vec<T>>
124 where
125 F: FnOnce(
126 &DatabaseProviderRO<N::DB, N>,
127 RangeInclusive<BlockNumber>,
128 &mut P,
129 ) -> ProviderResult<Vec<T>>,
130 G: Fn(&BlockState<N::Primitives>, &mut P) -> Option<T>,
131 P: FnMut(&T) -> bool,
132 {
133 let mut in_memory_chain =
141 self.head_block.as_ref().map(|b| b.chain().collect::<Vec<_>>()).unwrap_or_default();
142 let db_provider = &self.storage_provider;
143
144 let (start, end) = self.convert_range_bounds(range, || {
145 in_memory_chain
147 .first()
148 .map(|b| b.number())
149 .unwrap_or_else(|| db_provider.last_block_number().unwrap_or_default())
150 });
151
152 if start > end {
153 return Ok(vec![])
154 }
155
156 let (in_memory, storage_range) = match in_memory_chain.last().as_ref().map(|b| b.number()) {
161 Some(lowest_memory_block) if lowest_memory_block <= end => {
162 let highest_memory_block =
163 in_memory_chain.first().as_ref().map(|b| b.number()).expect("qed");
164
165 let in_memory_range =
169 lowest_memory_block.max(start)..=end.min(highest_memory_block);
170
171 in_memory_chain.truncate(
174 in_memory_chain
175 .len()
176 .saturating_sub(start.saturating_sub(lowest_memory_block) as usize),
177 );
178
179 let storage_range =
180 (lowest_memory_block > start).then(|| start..=lowest_memory_block - 1);
181
182 (Some((in_memory_chain, in_memory_range)), storage_range)
183 }
184 _ => {
185 drop(in_memory_chain);
187
188 (None, Some(start..=end))
189 }
190 };
191
192 let mut items = Vec::with_capacity((end - start + 1) as usize);
193
194 if let Some(storage_range) = storage_range {
195 let mut db_items = fetch_db_range(db_provider, storage_range.clone(), &mut predicate)?;
196 items.append(&mut db_items);
197
198 if items.len() as u64 != storage_range.end() - storage_range.start() + 1 {
201 return Ok(items)
202 }
203 }
204
205 if let Some((in_memory_chain, in_memory_range)) = in_memory {
206 for (num, block) in in_memory_range.zip(in_memory_chain.into_iter().rev()) {
207 debug_assert!(num == block.number());
208 if let Some(item) = map_block_state_item(block, &mut predicate) {
209 items.push(item);
210 } else {
211 break
212 }
213 }
214 }
215
216 Ok(items)
217 }
218
219 fn get_in_memory_or_storage_by_tx_range<S, M, R>(
225 &self,
226 range: impl RangeBounds<BlockNumber>,
227 fetch_from_db: S,
228 fetch_from_block_state: M,
229 ) -> ProviderResult<Vec<R>>
230 where
231 S: FnOnce(
232 &DatabaseProviderRO<N::DB, N>,
233 RangeInclusive<TxNumber>,
234 ) -> ProviderResult<Vec<R>>,
235 M: Fn(RangeInclusive<usize>, &BlockState<N::Primitives>) -> ProviderResult<Vec<R>>,
236 {
237 let in_mem_chain = self.head_block.iter().flat_map(|b| b.chain()).collect::<Vec<_>>();
238 let provider = &self.storage_provider;
239
240 let last_database_block_number = in_mem_chain
243 .last()
244 .map(|b| Ok(b.anchor().number))
245 .unwrap_or_else(|| provider.last_block_number())?;
246
247 let last_block_body_index = provider
250 .block_body_indices(last_database_block_number)?
251 .ok_or(ProviderError::BlockBodyIndicesNotFound(last_database_block_number))?;
252 let mut in_memory_tx_num = last_block_body_index.next_tx_num();
253
254 let (start, end) = self.convert_range_bounds(range, || {
255 in_mem_chain
256 .iter()
257 .map(|b| b.block_ref().recovered_block().body().transactions().len() as u64)
258 .sum::<u64>() +
259 last_block_body_index.last_tx_num()
260 });
261
262 if start > end {
263 return Ok(vec![])
264 }
265
266 let mut tx_range = start..=end;
267
268 if *tx_range.end() < in_memory_tx_num {
271 return fetch_from_db(provider, tx_range);
272 }
273
274 let mut items = Vec::with_capacity((tx_range.end() - tx_range.start() + 1) as usize);
275
276 if *tx_range.start() < in_memory_tx_num {
278 let db_range = *tx_range.start()..=in_memory_tx_num.saturating_sub(1);
280
281 tx_range = in_memory_tx_num..=*tx_range.end();
283
284 items.extend(fetch_from_db(provider, db_range)?);
285 }
286
287 for block_state in in_mem_chain.iter().rev() {
289 let block_tx_count =
290 block_state.block_ref().recovered_block().body().transactions().len();
291 let remaining = (tx_range.end() - tx_range.start() + 1) as usize;
292
293 if *tx_range.start() >= in_memory_tx_num + block_tx_count as u64 {
296 in_memory_tx_num += block_tx_count as u64;
297 continue
298 }
299
300 let skip = (tx_range.start() - in_memory_tx_num) as usize;
302
303 items.extend(fetch_from_block_state(
304 skip..=skip + (remaining.min(block_tx_count - skip) - 1),
305 block_state,
306 )?);
307
308 in_memory_tx_num += block_tx_count as u64;
309
310 if in_memory_tx_num > *tx_range.end() {
312 break
313 }
314
315 tx_range = in_memory_tx_num..=*tx_range.end();
317 }
318
319 Ok(items)
320 }
321
322 fn get_in_memory_or_storage_by_tx<S, M, R>(
325 &self,
326 id: HashOrNumber,
327 fetch_from_db: S,
328 fetch_from_block_state: M,
329 ) -> ProviderResult<Option<R>>
330 where
331 S: FnOnce(&DatabaseProviderRO<N::DB, N>) -> ProviderResult<Option<R>>,
332 M: Fn(usize, TxNumber, &BlockState<N::Primitives>) -> ProviderResult<Option<R>>,
333 {
334 let in_mem_chain = self.head_block.iter().flat_map(|b| b.chain()).collect::<Vec<_>>();
335 let provider = &self.storage_provider;
336
337 let last_database_block_number = in_mem_chain
340 .last()
341 .map(|b| Ok(b.anchor().number))
342 .unwrap_or_else(|| provider.last_block_number())?;
343
344 let last_block_body_index = provider
347 .block_body_indices(last_database_block_number)?
348 .ok_or(ProviderError::BlockBodyIndicesNotFound(last_database_block_number))?;
349 let mut in_memory_tx_num = last_block_body_index.next_tx_num();
350
351 if let HashOrNumber::Number(id) = id &&
354 id < in_memory_tx_num
355 {
356 return fetch_from_db(provider)
357 }
358
359 for block_state in in_mem_chain.iter().rev() {
361 let executed_block = block_state.block_ref();
362 let block = executed_block.recovered_block();
363
364 for tx_index in 0..block.body().transactions().len() {
365 match id {
366 HashOrNumber::Hash(tx_hash) => {
367 if tx_hash == *block.body().transactions()[tx_index].tx_hash() {
368 return fetch_from_block_state(tx_index, in_memory_tx_num, block_state)
369 }
370 }
371 HashOrNumber::Number(id) => {
372 if id == in_memory_tx_num {
373 return fetch_from_block_state(tx_index, in_memory_tx_num, block_state)
374 }
375 }
376 }
377
378 in_memory_tx_num += 1;
379 }
380 }
381
382 if let HashOrNumber::Hash(_) = id {
384 return fetch_from_db(provider)
385 }
386
387 Ok(None)
388 }
389
390 pub(crate) fn get_in_memory_or_storage_by_block<S, M, R>(
392 &self,
393 id: BlockHashOrNumber,
394 fetch_from_db: S,
395 fetch_from_block_state: M,
396 ) -> ProviderResult<R>
397 where
398 S: FnOnce(&DatabaseProviderRO<N::DB, N>) -> ProviderResult<R>,
399 M: Fn(&BlockState<N::Primitives>) -> ProviderResult<R>,
400 {
401 if let Some(Some(block_state)) = self.head_block.as_ref().map(|b| b.block_on_chain(id)) {
402 return fetch_from_block_state(block_state)
403 }
404 fetch_from_db(&self.storage_provider)
405 }
406}
407
408impl<N: ProviderNodeTypes> NodePrimitivesProvider for ConsistentProvider<N> {
409 type Primitives = N::Primitives;
410}
411
412impl<N: ProviderNodeTypes> StaticFileProviderFactory for ConsistentProvider<N> {
413 fn static_file_provider(&self) -> StaticFileProvider<N::Primitives> {
414 self.storage_provider.static_file_provider()
415 }
416
417 fn get_static_file_writer(
418 &self,
419 block: BlockNumber,
420 segment: StaticFileSegment,
421 ) -> ProviderResult<StaticFileProviderRWRefMut<'_, Self::Primitives>> {
422 self.storage_provider.get_static_file_writer(block, segment)
423 }
424}
425
426impl<N: ProviderNodeTypes> HeaderProvider for ConsistentProvider<N> {
427 type Header = HeaderTy<N>;
428
429 fn header(&self, block_hash: BlockHash) -> ProviderResult<Option<Self::Header>> {
430 self.get_in_memory_or_storage_by_block(
431 block_hash.into(),
432 |db_provider| db_provider.header(block_hash),
433 |block_state| Ok(Some(block_state.block_ref().recovered_block().clone_header())),
434 )
435 }
436
437 fn header_by_number(&self, num: BlockNumber) -> ProviderResult<Option<Self::Header>> {
438 self.get_in_memory_or_storage_by_block(
439 num.into(),
440 |db_provider| db_provider.header_by_number(num),
441 |block_state| Ok(Some(block_state.block_ref().recovered_block().clone_header())),
442 )
443 }
444
445 fn headers_range(
446 &self,
447 range: impl RangeBounds<BlockNumber>,
448 ) -> ProviderResult<Vec<Self::Header>> {
449 self.get_in_memory_or_storage_by_block_range_while(
450 range,
451 |db_provider, range, _| db_provider.headers_range(range),
452 |block_state, _| Some(block_state.block_ref().recovered_block().header().clone()),
453 |_| true,
454 )
455 }
456
457 fn sealed_header(
458 &self,
459 number: BlockNumber,
460 ) -> ProviderResult<Option<SealedHeader<Self::Header>>> {
461 self.get_in_memory_or_storage_by_block(
462 number.into(),
463 |db_provider| db_provider.sealed_header(number),
464 |block_state| Ok(Some(block_state.block_ref().recovered_block().clone_sealed_header())),
465 )
466 }
467
468 fn sealed_headers_range(
469 &self,
470 range: impl RangeBounds<BlockNumber>,
471 ) -> ProviderResult<Vec<SealedHeader<Self::Header>>> {
472 self.get_in_memory_or_storage_by_block_range_while(
473 range,
474 |db_provider, range, _| db_provider.sealed_headers_range(range),
475 |block_state, _| Some(block_state.block_ref().recovered_block().clone_sealed_header()),
476 |_| true,
477 )
478 }
479
480 fn sealed_headers_while(
481 &self,
482 range: impl RangeBounds<BlockNumber>,
483 predicate: impl FnMut(&SealedHeader<Self::Header>) -> bool,
484 ) -> ProviderResult<Vec<SealedHeader<Self::Header>>> {
485 self.get_in_memory_or_storage_by_block_range_while(
486 range,
487 |db_provider, range, predicate| db_provider.sealed_headers_while(range, predicate),
488 |block_state, predicate| {
489 let header = block_state.block_ref().recovered_block().sealed_header();
490 predicate(header).then(|| header.clone())
491 },
492 predicate,
493 )
494 }
495}
496
497impl<N: ProviderNodeTypes> BlockHashReader for ConsistentProvider<N> {
498 fn block_hash(&self, number: u64) -> ProviderResult<Option<B256>> {
499 self.get_in_memory_or_storage_by_block(
500 number.into(),
501 |db_provider| db_provider.block_hash(number),
502 |block_state| Ok(Some(block_state.hash())),
503 )
504 }
505
506 fn canonical_hashes_range(
507 &self,
508 start: BlockNumber,
509 end: BlockNumber,
510 ) -> ProviderResult<Vec<B256>> {
511 self.get_in_memory_or_storage_by_block_range_while(
512 start..end,
513 |db_provider, inclusive_range, _| {
514 db_provider
515 .canonical_hashes_range(*inclusive_range.start(), *inclusive_range.end() + 1)
516 },
517 |block_state, _| Some(block_state.hash()),
518 |_| true,
519 )
520 }
521}
522
523impl<N: ProviderNodeTypes> BlockNumReader for ConsistentProvider<N> {
524 fn chain_info(&self) -> ProviderResult<ChainInfo> {
525 let best_number = self.best_block_number()?;
526 Ok(ChainInfo { best_hash: self.block_hash(best_number)?.unwrap_or_default(), best_number })
527 }
528
529 fn best_block_number(&self) -> ProviderResult<BlockNumber> {
530 self.head_block.as_ref().map(|b| Ok(b.number())).unwrap_or_else(|| self.last_block_number())
531 }
532
533 fn last_block_number(&self) -> ProviderResult<BlockNumber> {
534 self.storage_provider.last_block_number()
535 }
536
537 fn block_number(&self, hash: B256) -> ProviderResult<Option<BlockNumber>> {
538 self.get_in_memory_or_storage_by_block(
539 hash.into(),
540 |db_provider| db_provider.block_number(hash),
541 |block_state| Ok(Some(block_state.number())),
542 )
543 }
544}
545
546impl<N: ProviderNodeTypes> BlockIdReader for ConsistentProvider<N> {
547 fn pending_block_num_hash(&self) -> ProviderResult<Option<BlockNumHash>> {
548 Ok(self.canonical_in_memory_state.pending_block_num_hash())
549 }
550
551 fn safe_block_num_hash(&self) -> ProviderResult<Option<BlockNumHash>> {
552 Ok(self.canonical_in_memory_state.get_safe_num_hash())
553 }
554
555 fn finalized_block_num_hash(&self) -> ProviderResult<Option<BlockNumHash>> {
556 Ok(self.canonical_in_memory_state.get_finalized_num_hash())
557 }
558}
559
560impl<N: ProviderNodeTypes> BlockReader for ConsistentProvider<N> {
561 type Block = BlockTy<N>;
562
563 fn find_block_by_hash(
564 &self,
565 hash: B256,
566 source: BlockSource,
567 ) -> ProviderResult<Option<Self::Block>> {
568 if matches!(source, BlockSource::Canonical | BlockSource::Any) &&
569 let Some(block) = self.get_in_memory_or_storage_by_block(
570 hash.into(),
571 |db_provider| db_provider.find_block_by_hash(hash, BlockSource::Canonical),
572 |block_state| Ok(Some(block_state.block_ref().recovered_block().clone_block())),
573 )?
574 {
575 return Ok(Some(block))
576 }
577
578 if matches!(source, BlockSource::Pending | BlockSource::Any) {
579 return Ok(self
580 .canonical_in_memory_state
581 .pending_block()
582 .filter(|b| b.hash() == hash)
583 .map(|b| b.into_block()))
584 }
585
586 Ok(None)
587 }
588
589 fn find_sealed_or_recovered_block(
590 &self,
591 hash: B256,
592 source: BlockSource,
593 ) -> ProviderResult<Option<SealedOrRecoveredBlock<Self::Block>>> {
594 if matches!(source, BlockSource::Canonical | BlockSource::Any) &&
595 let Some(block) = self.get_in_memory_or_storage_by_block(
596 hash.into(),
597 |db_provider| {
598 db_provider.find_sealed_or_recovered_block(hash, BlockSource::Canonical)
599 },
600 |block_state| {
601 Ok(Some(SealedOrRecoveredBlock::recovered_arc(Arc::clone(
602 &block_state.block_ref().recovered_block,
603 ))))
604 },
605 )?
606 {
607 return Ok(Some(block))
608 }
609
610 if matches!(source, BlockSource::Pending | BlockSource::Any) &&
611 let Some(block_state) = self.canonical_in_memory_state.pending_state()
612 {
613 let recovered_block = Arc::clone(&block_state.block_ref().recovered_block);
614 if recovered_block.hash() == hash {
615 return Ok(Some(SealedOrRecoveredBlock::recovered_arc(recovered_block)))
616 }
617 }
618
619 Ok(None)
620 }
621
622 fn block(&self, id: BlockHashOrNumber) -> ProviderResult<Option<Self::Block>> {
623 self.get_in_memory_or_storage_by_block(
624 id,
625 |db_provider| db_provider.block(id),
626 |block_state| Ok(Some(block_state.block_ref().recovered_block().clone_block())),
627 )
628 }
629
630 fn pending_block(&self) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
631 Ok(self.canonical_in_memory_state.pending_recovered_block())
632 }
633
634 fn pending_block_and_receipts(
635 &self,
636 ) -> ProviderResult<Option<(RecoveredBlock<Self::Block>, Vec<Self::Receipt>)>> {
637 Ok(self.canonical_in_memory_state.pending_block_and_receipts())
638 }
639
640 fn recovered_block(
647 &self,
648 id: BlockHashOrNumber,
649 transaction_kind: TransactionVariant,
650 ) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
651 self.get_in_memory_or_storage_by_block(
652 id,
653 |db_provider| db_provider.recovered_block(id, transaction_kind),
654 |block_state| Ok(Some(block_state.block().recovered_block().clone())),
655 )
656 }
657
658 fn sealed_block_with_senders(
659 &self,
660 id: BlockHashOrNumber,
661 transaction_kind: TransactionVariant,
662 ) -> ProviderResult<Option<RecoveredBlock<Self::Block>>> {
663 self.get_in_memory_or_storage_by_block(
664 id,
665 |db_provider| db_provider.sealed_block_with_senders(id, transaction_kind),
666 |block_state| Ok(Some(block_state.block().recovered_block().clone())),
667 )
668 }
669
670 fn block_range(&self, range: RangeInclusive<BlockNumber>) -> ProviderResult<Vec<Self::Block>> {
671 self.get_in_memory_or_storage_by_block_range_while(
672 range,
673 |db_provider, range, _| db_provider.block_range(range),
674 |block_state, _| Some(block_state.block_ref().recovered_block().clone_block()),
675 |_| true,
676 )
677 }
678
679 fn block_with_senders_range(
680 &self,
681 range: RangeInclusive<BlockNumber>,
682 ) -> ProviderResult<Vec<RecoveredBlock<Self::Block>>> {
683 self.get_in_memory_or_storage_by_block_range_while(
684 range,
685 |db_provider, range, _| db_provider.block_with_senders_range(range),
686 |block_state, _| Some(block_state.block().recovered_block().clone()),
687 |_| true,
688 )
689 }
690
691 fn recovered_block_range(
692 &self,
693 range: RangeInclusive<BlockNumber>,
694 ) -> ProviderResult<Vec<RecoveredBlock<Self::Block>>> {
695 self.get_in_memory_or_storage_by_block_range_while(
696 range,
697 |db_provider, range, _| db_provider.recovered_block_range(range),
698 |block_state, _| Some(block_state.block().recovered_block().clone()),
699 |_| true,
700 )
701 }
702
703 fn block_by_transaction_id(&self, id: TxNumber) -> ProviderResult<Option<BlockNumber>> {
704 self.get_in_memory_or_storage_by_tx(
705 id.into(),
706 |db_provider| db_provider.block_by_transaction_id(id),
707 |_, _, block_state| Ok(Some(block_state.number())),
708 )
709 }
710}
711
712impl<N: ProviderNodeTypes> TransactionsProvider for ConsistentProvider<N> {
713 type Transaction = TxTy<N>;
714
715 fn transaction_id(&self, tx_hash: TxHash) -> ProviderResult<Option<TxNumber>> {
716 self.get_in_memory_or_storage_by_tx(
717 tx_hash.into(),
718 |db_provider| db_provider.transaction_id(tx_hash),
719 |_, tx_number, _| Ok(Some(tx_number)),
720 )
721 }
722
723 fn transaction_by_id(&self, id: TxNumber) -> ProviderResult<Option<Self::Transaction>> {
724 self.get_in_memory_or_storage_by_tx(
725 id.into(),
726 |provider| provider.transaction_by_id(id),
727 |tx_index, _, block_state| {
728 Ok(block_state
729 .block_ref()
730 .recovered_block()
731 .body()
732 .transactions()
733 .get(tx_index)
734 .cloned())
735 },
736 )
737 }
738
739 fn transaction_by_id_unhashed(
740 &self,
741 id: TxNumber,
742 ) -> ProviderResult<Option<Self::Transaction>> {
743 self.get_in_memory_or_storage_by_tx(
744 id.into(),
745 |provider| provider.transaction_by_id_unhashed(id),
746 |tx_index, _, block_state| {
747 Ok(block_state
748 .block_ref()
749 .recovered_block()
750 .body()
751 .transactions()
752 .get(tx_index)
753 .cloned())
754 },
755 )
756 }
757
758 fn transaction_by_hash(&self, hash: TxHash) -> ProviderResult<Option<Self::Transaction>> {
759 if let Some(tx) = self.head_block.as_ref().and_then(|b| b.transaction_on_chain(hash)) {
760 return Ok(Some(tx))
761 }
762
763 self.storage_provider.transaction_by_hash(hash)
764 }
765
766 fn transaction_by_hash_with_meta(
767 &self,
768 tx_hash: TxHash,
769 ) -> ProviderResult<Option<(Self::Transaction, TransactionMeta)>> {
770 if let Some((tx, meta)) =
771 self.head_block.as_ref().and_then(|b| b.transaction_meta_on_chain(tx_hash))
772 {
773 return Ok(Some((tx, meta)))
774 }
775
776 self.storage_provider.transaction_by_hash_with_meta(tx_hash)
777 }
778
779 fn transactions_by_block(
780 &self,
781 id: BlockHashOrNumber,
782 ) -> ProviderResult<Option<Vec<Self::Transaction>>> {
783 self.get_in_memory_or_storage_by_block(
784 id,
785 |provider| provider.transactions_by_block(id),
786 |block_state| {
787 Ok(Some(block_state.block_ref().recovered_block().body().transactions().to_vec()))
788 },
789 )
790 }
791
792 fn transactions_by_block_range(
793 &self,
794 range: impl RangeBounds<BlockNumber>,
795 ) -> ProviderResult<Vec<Vec<Self::Transaction>>> {
796 self.get_in_memory_or_storage_by_block_range_while(
797 range,
798 |db_provider, range, _| db_provider.transactions_by_block_range(range),
799 |block_state, _| {
800 Some(block_state.block_ref().recovered_block().body().transactions().to_vec())
801 },
802 |_| true,
803 )
804 }
805
806 fn transactions_by_tx_range(
807 &self,
808 range: impl RangeBounds<TxNumber>,
809 ) -> ProviderResult<Vec<Self::Transaction>> {
810 self.get_in_memory_or_storage_by_tx_range(
811 range,
812 |db_provider, db_range| db_provider.transactions_by_tx_range(db_range),
813 |index_range, block_state| {
814 Ok(block_state.block_ref().recovered_block().body().transactions()[index_range]
815 .to_vec())
816 },
817 )
818 }
819
820 fn senders_by_tx_range(
821 &self,
822 range: impl RangeBounds<TxNumber>,
823 ) -> ProviderResult<Vec<Address>> {
824 self.get_in_memory_or_storage_by_tx_range(
825 range,
826 |db_provider, db_range| db_provider.senders_by_tx_range(db_range),
827 |index_range, block_state| {
828 Ok(block_state.block_ref().recovered_block.senders()[index_range].to_vec())
829 },
830 )
831 }
832
833 fn transaction_sender(&self, id: TxNumber) -> ProviderResult<Option<Address>> {
834 self.get_in_memory_or_storage_by_tx(
835 id.into(),
836 |provider| provider.transaction_sender(id),
837 |tx_index, _, block_state| {
838 Ok(block_state.block_ref().recovered_block.senders().get(tx_index).copied())
839 },
840 )
841 }
842}
843
844impl<N: ProviderNodeTypes> ReceiptProvider for ConsistentProvider<N> {
845 type Receipt = ReceiptTy<N>;
846
847 fn receipt(&self, id: TxNumber) -> ProviderResult<Option<Self::Receipt>> {
848 self.get_in_memory_or_storage_by_tx(
849 id.into(),
850 |provider| provider.receipt(id),
851 |tx_index, _, block_state| {
852 Ok(block_state.executed_block_receipts_ref().get(tx_index).cloned())
853 },
854 )
855 }
856
857 fn receipt_by_hash(&self, hash: TxHash) -> ProviderResult<Option<Self::Receipt>> {
858 for block_state in self.head_block.iter().flat_map(|b| b.chain()) {
859 let executed_block = block_state.block_ref();
860 let block = executed_block.recovered_block();
861 let receipts = block_state.executed_block_receipts_ref();
862
863 debug_assert_eq!(
865 block.body().transactions().len(),
866 receipts.len(),
867 "Mismatch between transaction and receipt count"
868 );
869
870 if let Some(tx_index) =
871 block.body().transactions_iter().position(|tx| *tx.tx_hash() == hash)
872 {
873 return Ok(receipts.get(tx_index).cloned());
875 }
876 }
877
878 self.storage_provider.receipt_by_hash(hash)
879 }
880
881 fn receipts_by_block(
882 &self,
883 block: BlockHashOrNumber,
884 ) -> ProviderResult<Option<Vec<Self::Receipt>>> {
885 self.get_in_memory_or_storage_by_block(
886 block,
887 |db_provider| db_provider.receipts_by_block(block),
888 |block_state| Ok(Some(block_state.executed_block_receipts())),
889 )
890 }
891
892 fn receipts_by_tx_range(
893 &self,
894 range: impl RangeBounds<TxNumber>,
895 ) -> ProviderResult<Vec<Self::Receipt>> {
896 self.get_in_memory_or_storage_by_tx_range(
897 range,
898 |db_provider, db_range| db_provider.receipts_by_tx_range(db_range),
899 |index_range, block_state| {
900 Ok(block_state.executed_block_receipts_ref()[index_range].to_vec())
901 },
902 )
903 }
904
905 fn receipts_by_block_range(
906 &self,
907 block_range: RangeInclusive<BlockNumber>,
908 ) -> ProviderResult<Vec<Vec<Self::Receipt>>> {
909 self.storage_provider.receipts_by_block_range(block_range)
910 }
911}
912
913impl<N: ProviderNodeTypes> ReceiptProviderIdExt for ConsistentProvider<N> {
914 fn receipts_by_block_id(&self, block: BlockId) -> ProviderResult<Option<Vec<Self::Receipt>>> {
915 match block {
916 BlockId::Hash(rpc_block_hash) => {
917 let mut receipts = self.receipts_by_block(rpc_block_hash.block_hash.into())?;
918 if receipts.is_none() &&
919 !rpc_block_hash.require_canonical.unwrap_or(false) &&
920 let Some(state) = self
921 .head_block
922 .as_ref()
923 .and_then(|b| b.block_on_chain(rpc_block_hash.block_hash.into()))
924 {
925 receipts = Some(state.executed_block_receipts());
926 }
927 Ok(receipts)
928 }
929 BlockId::Number(num_tag) => match num_tag {
930 BlockNumberOrTag::Pending => Ok(self
931 .canonical_in_memory_state
932 .pending_state()
933 .map(|block_state| block_state.executed_block_receipts())),
934 _ => {
935 if let Some(num) = self.convert_block_number(num_tag)? {
936 self.receipts_by_block(num.into())
937 } else {
938 Ok(None)
939 }
940 }
941 },
942 }
943 }
944}
945
946impl<N: ProviderNodeTypes> BlockBodyIndicesProvider for ConsistentProvider<N> {
947 fn block_body_indices(
948 &self,
949 number: BlockNumber,
950 ) -> ProviderResult<Option<StoredBlockBodyIndices>> {
951 self.get_in_memory_or_storage_by_block(
952 number.into(),
953 |db_provider| db_provider.block_body_indices(number),
954 |block_state| {
955 let last_storage_block_number = block_state.anchor().number;
957 let mut stored_indices = self
958 .storage_provider
959 .block_body_indices(last_storage_block_number)?
960 .ok_or(ProviderError::BlockBodyIndicesNotFound(last_storage_block_number))?;
961
962 stored_indices.first_tx_num = stored_indices.next_tx_num();
964 stored_indices.tx_count = 0;
965
966 for state in block_state.chain().collect::<Vec<_>>().into_iter().rev() {
968 let block_tx_count =
969 state.block_ref().recovered_block().body().transactions().len() as u64;
970 if state.block_ref().recovered_block().number() == number {
971 stored_indices.tx_count = block_tx_count;
972 } else {
973 stored_indices.first_tx_num += block_tx_count;
974 }
975 }
976
977 Ok(Some(stored_indices))
978 },
979 )
980 }
981
982 fn block_body_indices_range(
983 &self,
984 range: RangeInclusive<BlockNumber>,
985 ) -> ProviderResult<Vec<StoredBlockBodyIndices>> {
986 range.map_while(|b| self.block_body_indices(b).transpose()).collect()
987 }
988}
989
990impl<N: ProviderNodeTypes> StageCheckpointReader for ConsistentProvider<N> {
991 fn get_stage_checkpoint(&self, id: StageId) -> ProviderResult<Option<StageCheckpoint>> {
992 self.storage_provider.get_stage_checkpoint(id)
993 }
994
995 fn get_stage_checkpoint_progress(&self, id: StageId) -> ProviderResult<Option<Vec<u8>>> {
996 self.storage_provider.get_stage_checkpoint_progress(id)
997 }
998
999 fn get_all_checkpoints(&self) -> ProviderResult<Vec<(String, StageCheckpoint)>> {
1000 self.storage_provider.get_all_checkpoints()
1001 }
1002}
1003
1004impl<N: ProviderNodeTypes> PruneCheckpointReader for ConsistentProvider<N> {
1005 fn get_prune_checkpoint(
1006 &self,
1007 segment: PruneSegment,
1008 ) -> ProviderResult<Option<PruneCheckpoint>> {
1009 self.storage_provider.get_prune_checkpoint(segment)
1010 }
1011
1012 fn get_prune_checkpoints(&self) -> ProviderResult<Vec<(PruneSegment, PruneCheckpoint)>> {
1013 self.storage_provider.get_prune_checkpoints()
1014 }
1015}
1016
1017impl<N: ProviderNodeTypes> ChainSpecProvider for ConsistentProvider<N> {
1018 type ChainSpec = N::ChainSpec;
1019
1020 fn chain_spec(&self) -> Arc<N::ChainSpec> {
1021 ChainSpecProvider::chain_spec(&self.storage_provider)
1022 }
1023}
1024
1025impl<N: ProviderNodeTypes> BlockReaderIdExt for ConsistentProvider<N> {
1026 fn block_by_id(&self, id: BlockId) -> ProviderResult<Option<Self::Block>> {
1027 match id {
1028 BlockId::Number(num) => self.block_by_number_or_tag(num),
1029 BlockId::Hash(hash) => {
1030 if Some(true) == hash.require_canonical {
1035 self.find_block_by_hash(hash.block_hash, BlockSource::Canonical)
1037 } else {
1038 self.block_by_hash(hash.block_hash)
1039 }
1040 }
1041 }
1042 }
1043
1044 fn header_by_number_or_tag(&self, id: BlockNumberOrTag) -> ProviderResult<Option<HeaderTy<N>>> {
1045 Ok(match id {
1046 BlockNumberOrTag::Latest => {
1047 Some(self.canonical_in_memory_state.get_canonical_head().unseal())
1048 }
1049 BlockNumberOrTag::Finalized => {
1050 self.canonical_in_memory_state.get_finalized_header().map(|h| h.unseal())
1051 }
1052 BlockNumberOrTag::Safe => {
1053 self.canonical_in_memory_state.get_safe_header().map(|h| h.unseal())
1054 }
1055 BlockNumberOrTag::Earliest => self.header_by_number(self.earliest_block_number()?)?,
1056 BlockNumberOrTag::Pending => self.canonical_in_memory_state.pending_header(),
1057
1058 BlockNumberOrTag::Number(num) => self.header_by_number(num)?,
1059 })
1060 }
1061
1062 fn sealed_header_by_number_or_tag(
1063 &self,
1064 id: BlockNumberOrTag,
1065 ) -> ProviderResult<Option<SealedHeader<HeaderTy<N>>>> {
1066 match id {
1067 BlockNumberOrTag::Latest => {
1068 Ok(Some(self.canonical_in_memory_state.get_canonical_head()))
1069 }
1070 BlockNumberOrTag::Finalized => {
1071 Ok(self.canonical_in_memory_state.get_finalized_header())
1072 }
1073 BlockNumberOrTag::Safe => Ok(self.canonical_in_memory_state.get_safe_header()),
1074 BlockNumberOrTag::Earliest => self
1075 .header_by_number(self.earliest_block_number()?)?
1076 .map_or_else(|| Ok(None), |h| Ok(Some(SealedHeader::seal_slow(h)))),
1077 BlockNumberOrTag::Pending => Ok(self.canonical_in_memory_state.pending_sealed_header()),
1078 BlockNumberOrTag::Number(num) => self
1079 .header_by_number(num)?
1080 .map_or_else(|| Ok(None), |h| Ok(Some(SealedHeader::seal_slow(h)))),
1081 }
1082 }
1083
1084 fn sealed_header_by_id(
1085 &self,
1086 id: BlockId,
1087 ) -> ProviderResult<Option<SealedHeader<HeaderTy<N>>>> {
1088 Ok(match id {
1089 BlockId::Number(num) => self.sealed_header_by_number_or_tag(num)?,
1090 BlockId::Hash(hash) => self
1091 .header(hash.block_hash)?
1092 .map(|header| SealedHeader::new(header, hash.block_hash)),
1093 })
1094 }
1095
1096 fn header_by_id(&self, id: BlockId) -> ProviderResult<Option<HeaderTy<N>>> {
1097 Ok(match id {
1098 BlockId::Number(num) => self.header_by_number_or_tag(num)?,
1099 BlockId::Hash(hash) => self.header(hash.block_hash)?,
1100 })
1101 }
1102}
1103
1104impl<N: ProviderNodeTypes> StorageChangeSetReader for ConsistentProvider<N> {
1105 fn storage_changeset(
1106 &self,
1107 block_number: BlockNumber,
1108 ) -> ProviderResult<Vec<(BlockNumberAddress, StorageEntry)>> {
1109 if let Some(state) =
1110 self.head_block.as_ref().and_then(|b| b.block_on_chain(block_number.into()))
1111 {
1112 let changesets = state
1113 .block()
1114 .execution_output
1115 .state
1116 .reverts
1117 .to_plain_state_reverts()
1118 .storage
1119 .into_iter()
1120 .flatten()
1121 .flat_map(|revert: PlainStorageRevert| {
1122 revert.storage_revert.into_iter().map(move |(key, value)| {
1123 let plain_key = B256::from(key.to_be_bytes());
1124 (
1125 BlockNumberAddress((block_number, revert.address)),
1126 StorageEntry { key: plain_key, value: value.to_previous_value() },
1127 )
1128 })
1129 })
1130 .collect();
1131 Ok(changesets)
1132 } else {
1133 let storage_history_exists = self
1137 .storage_provider
1138 .get_prune_checkpoint(PruneSegment::StorageHistory)?
1139 .and_then(|checkpoint| {
1140 checkpoint.block_number.map(|checkpoint| block_number > checkpoint)
1145 })
1146 .unwrap_or(true);
1147
1148 if !storage_history_exists {
1149 return Err(ProviderError::StateAtBlockPruned(block_number))
1150 }
1151
1152 self.storage_provider.storage_changeset(block_number)
1153 }
1154 }
1155
1156 fn get_storage_before_block(
1157 &self,
1158 block_number: BlockNumber,
1159 address: Address,
1160 storage_key: B256,
1161 ) -> ProviderResult<Option<StorageEntry>> {
1162 if let Some(state) =
1163 self.head_block.as_ref().and_then(|b| b.block_on_chain(block_number.into()))
1164 {
1165 let changeset = state
1166 .block_ref()
1167 .execution_output
1168 .state
1169 .reverts
1170 .to_plain_state_reverts()
1171 .storage
1172 .into_iter()
1173 .flatten()
1174 .find_map(|revert: PlainStorageRevert| {
1175 if revert.address != address {
1176 return None
1177 }
1178 revert.storage_revert.into_iter().find_map(|(key, value)| {
1179 let plain_key = B256::from(key.to_be_bytes());
1180 (plain_key == storage_key).then(|| StorageEntry {
1181 key: plain_key,
1182 value: value.to_previous_value(),
1183 })
1184 })
1185 });
1186 Ok(changeset)
1187 } else {
1188 let storage_history_exists = self
1189 .storage_provider
1190 .get_prune_checkpoint(PruneSegment::StorageHistory)?
1191 .and_then(|checkpoint| {
1192 checkpoint.block_number.map(|checkpoint| block_number > checkpoint)
1193 })
1194 .unwrap_or(true);
1195
1196 if !storage_history_exists {
1197 return Err(ProviderError::StateAtBlockPruned(block_number))
1198 }
1199
1200 self.storage_provider.get_storage_before_block(block_number, address, storage_key)
1201 }
1202 }
1203
1204 fn storage_changesets_range(
1205 &self,
1206 range: impl RangeBounds<BlockNumber>,
1207 ) -> ProviderResult<Vec<(BlockNumberAddress, StorageEntry)>> {
1208 let range = to_range(range);
1209 let mut changesets = Vec::new();
1210 let database_start = range.start;
1211 let mut database_end = range.end;
1212
1213 if let Some(head_block) = &self.head_block {
1214 database_end = head_block.anchor().number;
1215
1216 for state in head_block.chain() {
1217 let block_changesets = state
1218 .block_ref()
1219 .execution_output
1220 .state
1221 .reverts
1222 .to_plain_state_reverts()
1223 .storage
1224 .into_iter()
1225 .flatten()
1226 .flat_map(|revert: PlainStorageRevert| {
1227 revert.storage_revert.into_iter().map(move |(key, value)| {
1228 let plain_key = B256::from(key.to_be_bytes());
1229 (
1230 BlockNumberAddress((state.number(), revert.address)),
1231 StorageEntry { key: plain_key, value: value.to_previous_value() },
1232 )
1233 })
1234 });
1235
1236 changesets.extend(block_changesets);
1237 }
1238 }
1239
1240 if database_start < database_end {
1241 let storage_history_exists = self
1242 .storage_provider
1243 .get_prune_checkpoint(PruneSegment::StorageHistory)?
1244 .and_then(|checkpoint| {
1245 checkpoint.block_number.map(|checkpoint| database_start > checkpoint)
1246 })
1247 .unwrap_or(true);
1248
1249 if !storage_history_exists {
1250 return Err(ProviderError::StateAtBlockPruned(database_start))
1251 }
1252
1253 let db_changesets = self
1254 .storage_provider
1255 .storage_changesets_range(database_start..=database_end - 1)?;
1256 changesets.extend(db_changesets);
1257 }
1258
1259 changesets.sort_by_key(|(block_address, _)| block_address.block_number());
1260
1261 Ok(changesets)
1262 }
1263}
1264
1265impl<N: ProviderNodeTypes> ChangeSetReader for ConsistentProvider<N> {
1266 fn account_block_changeset(
1267 &self,
1268 block_number: BlockNumber,
1269 ) -> ProviderResult<Vec<AccountBeforeTx>> {
1270 if let Some(state) =
1271 self.head_block.as_ref().and_then(|b| b.block_on_chain(block_number.into()))
1272 {
1273 let changesets = state
1274 .block_ref()
1275 .execution_output
1276 .state
1277 .reverts
1278 .to_plain_state_reverts()
1279 .accounts
1280 .into_iter()
1281 .flatten()
1282 .map(|(address, info)| AccountBeforeTx { address, info: info.map(Into::into) })
1283 .collect();
1284 Ok(changesets)
1285 } else {
1286 let account_history_exists = self
1290 .storage_provider
1291 .get_prune_checkpoint(PruneSegment::AccountHistory)?
1292 .and_then(|checkpoint| {
1293 checkpoint.block_number.map(|checkpoint| block_number > checkpoint)
1298 })
1299 .unwrap_or(true);
1300
1301 if !account_history_exists {
1302 return Err(ProviderError::StateAtBlockPruned(block_number))
1303 }
1304
1305 self.storage_provider.account_block_changeset(block_number)
1306 }
1307 }
1308
1309 fn get_account_before_block(
1310 &self,
1311 block_number: BlockNumber,
1312 address: Address,
1313 ) -> ProviderResult<Option<AccountBeforeTx>> {
1314 if let Some(state) =
1315 self.head_block.as_ref().and_then(|b| b.block_on_chain(block_number.into()))
1316 {
1317 let changeset = state
1319 .block_ref()
1320 .execution_output
1321 .state
1322 .reverts
1323 .to_plain_state_reverts()
1324 .accounts
1325 .into_iter()
1326 .flatten()
1327 .find(|(addr, _)| addr == &address)
1328 .map(|(address, info)| AccountBeforeTx { address, info: info.map(Into::into) });
1329 Ok(changeset)
1330 } else {
1331 let account_history_exists = self
1334 .storage_provider
1335 .get_prune_checkpoint(PruneSegment::AccountHistory)?
1336 .and_then(|checkpoint| {
1337 checkpoint.block_number.map(|checkpoint| block_number > checkpoint)
1342 })
1343 .unwrap_or(true);
1344
1345 if !account_history_exists {
1346 return Err(ProviderError::StateAtBlockPruned(block_number))
1347 }
1348
1349 self.storage_provider.get_account_before_block(block_number, address)
1351 }
1352 }
1353
1354 fn account_changesets_range(
1355 &self,
1356 range: impl core::ops::RangeBounds<BlockNumber>,
1357 ) -> ProviderResult<Vec<(BlockNumber, AccountBeforeTx)>> {
1358 let range = to_range(range);
1359 let mut changesets = Vec::new();
1360 let database_start = range.start;
1361 let mut database_end = range.end;
1362
1363 if let Some(head_block) = &self.head_block {
1365 database_end = head_block.anchor().number;
1367
1368 for state in head_block.chain() {
1369 let block_changesets = state
1371 .block_ref()
1372 .execution_output
1373 .state
1374 .reverts
1375 .to_plain_state_reverts()
1376 .accounts
1377 .into_iter()
1378 .flatten()
1379 .map(|(address, info)| AccountBeforeTx { address, info: info.map(Into::into) });
1380
1381 for changeset in block_changesets {
1382 changesets.push((state.number(), changeset));
1383 }
1384 }
1385 }
1386
1387 if database_start < database_end {
1389 let account_history_exists = self
1391 .storage_provider
1392 .get_prune_checkpoint(PruneSegment::AccountHistory)?
1393 .and_then(|checkpoint| {
1394 checkpoint.block_number.map(|checkpoint| database_start > checkpoint)
1395 })
1396 .unwrap_or(true);
1397
1398 if !account_history_exists {
1399 return Err(ProviderError::StateAtBlockPruned(database_start))
1400 }
1401
1402 let db_changesets =
1403 self.storage_provider.account_changesets_range(database_start..database_end)?;
1404 changesets.extend(db_changesets);
1405 }
1406
1407 changesets.sort_by_key(|(block_num, _)| *block_num);
1408
1409 Ok(changesets)
1410 }
1411}
1412
1413#[cfg(test)]
1414mod tests {
1415 use crate::{
1416 providers::blockchain_provider::BlockchainProvider,
1417 test_utils::create_test_provider_factory, BlockWriter,
1418 };
1419 use alloy_eips::BlockHashOrNumber;
1420 use alloy_primitives::B256;
1421 use itertools::Itertools;
1422 use rand::Rng;
1423 use reth_chain_state::{ExecutedBlock, NewCanonicalChain};
1424 use reth_db_api::models::AccountBeforeTx;
1425 use reth_ethereum_primitives::Block;
1426 use reth_execution_types::{BlockExecutionOutput, BlockExecutionResult, ExecutionOutcome};
1427 use reth_primitives_traits::{RecoveredBlock, SealedBlock};
1428 use reth_storage_api::{BlockReader, BlockSource, ChangeSetReader, StateReader};
1429 use reth_testing_utils::generators::{
1430 self, random_block_range, random_changeset_range, random_eoa_accounts, BlockRangeParams,
1431 };
1432 use revm::database::BundleState;
1433 use std::{
1434 ops::{Bound, Range, RangeBounds},
1435 sync::Arc,
1436 };
1437
1438 const TEST_BLOCKS_COUNT: usize = 5;
1439
1440 fn random_blocks(
1441 rng: &mut impl Rng,
1442 database_blocks: usize,
1443 in_memory_blocks: usize,
1444 requests_count: Option<Range<u8>>,
1445 withdrawals_count: Option<Range<u8>>,
1446 tx_count: impl RangeBounds<u8>,
1447 ) -> (Vec<SealedBlock<Block>>, Vec<SealedBlock<Block>>) {
1448 let block_range = (database_blocks + in_memory_blocks - 1) as u64;
1449
1450 let tx_start = match tx_count.start_bound() {
1451 Bound::Included(&n) | Bound::Excluded(&n) => n,
1452 Bound::Unbounded => u8::MIN,
1453 };
1454 let tx_end = match tx_count.end_bound() {
1455 Bound::Included(&n) | Bound::Excluded(&n) => n + 1,
1456 Bound::Unbounded => u8::MAX,
1457 };
1458
1459 let blocks = random_block_range(
1460 rng,
1461 0..=block_range,
1462 BlockRangeParams {
1463 parent: Some(B256::ZERO),
1464 tx_count: tx_start..tx_end,
1465 requests_count,
1466 withdrawals_count,
1467 },
1468 );
1469 let (database_blocks, in_memory_blocks) = blocks.split_at(database_blocks);
1470 (database_blocks.to_vec(), in_memory_blocks.to_vec())
1471 }
1472
1473 #[test]
1474 fn test_block_reader_find_block_by_hash() -> eyre::Result<()> {
1475 let mut rng = generators::rng();
1477 let factory = create_test_provider_factory();
1478
1479 let blocks = random_block_range(
1481 &mut rng,
1482 0..=10,
1483 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
1484 );
1485 let (database_blocks, in_memory_blocks) = blocks.split_at(5);
1486
1487 let provider_rw = factory.provider_rw()?;
1489 for block in database_blocks {
1490 provider_rw.insert_block(
1491 &block.clone().try_recover().expect("failed to seal block with senders"),
1492 )?;
1493 }
1494 provider_rw.commit()?;
1495
1496 let provider = BlockchainProvider::new(factory)?;
1498 let consistent_provider = provider.consistent_provider()?;
1499
1500 let first_db_block = database_blocks.first().unwrap();
1502 let first_in_mem_block = in_memory_blocks.first().unwrap();
1503 let last_in_mem_block = in_memory_blocks.last().unwrap();
1504
1505 assert_eq!(
1507 consistent_provider.find_block_by_hash(first_in_mem_block.hash(), BlockSource::Any)?,
1508 None
1509 );
1510 assert!(consistent_provider
1511 .find_sealed_or_recovered_block(first_in_mem_block.hash(), BlockSource::Any)?
1512 .is_none());
1513 assert_eq!(
1514 consistent_provider
1515 .find_block_by_hash(first_in_mem_block.hash(), BlockSource::Canonical)?,
1516 None
1517 );
1518 assert_eq!(
1520 consistent_provider
1521 .find_block_by_hash(first_in_mem_block.hash(), BlockSource::Pending)?,
1522 None
1523 );
1524
1525 let in_memory_block_senders =
1527 first_in_mem_block.senders().expect("failed to recover senders");
1528 let chain = NewCanonicalChain::Commit {
1529 new: vec![ExecutedBlock {
1530 recovered_block: Arc::new(RecoveredBlock::new_sealed(
1531 first_in_mem_block.clone(),
1532 in_memory_block_senders,
1533 )),
1534 ..Default::default()
1535 }],
1536 };
1537 consistent_provider.canonical_in_memory_state.update_chain(chain);
1538 let consistent_provider = provider.consistent_provider()?;
1539
1540 assert_eq!(
1542 consistent_provider.find_block_by_hash(first_in_mem_block.hash(), BlockSource::Any)?,
1543 Some(first_in_mem_block.clone().into_block())
1544 );
1545 let block = consistent_provider
1546 .find_sealed_or_recovered_block(first_in_mem_block.hash(), BlockSource::Any)?
1547 .expect("in-memory block should be found");
1548 assert_eq!(block.sealed_block(), first_in_mem_block);
1549 assert!(block.recovered_block().is_some());
1550
1551 assert_eq!(
1552 consistent_provider
1553 .find_block_by_hash(first_in_mem_block.hash(), BlockSource::Canonical)?,
1554 Some(first_in_mem_block.clone().into_block())
1555 );
1556 let block = consistent_provider
1557 .find_sealed_or_recovered_block(first_in_mem_block.hash(), BlockSource::Canonical)?
1558 .expect("canonical in-memory block should be found");
1559 assert_eq!(block.sealed_block(), first_in_mem_block);
1560 assert!(block.recovered_block().is_some());
1561
1562 assert_eq!(
1564 consistent_provider.find_block_by_hash(first_db_block.hash(), BlockSource::Any)?,
1565 Some(first_db_block.clone().into_block())
1566 );
1567 let block = consistent_provider
1568 .find_sealed_or_recovered_block(first_db_block.hash(), BlockSource::Any)?
1569 .expect("database block should be found");
1570 assert_eq!(block.sealed_block(), first_db_block);
1571 assert!(block.recovered_block().is_none());
1572
1573 assert_eq!(
1574 consistent_provider
1575 .find_block_by_hash(first_db_block.hash(), BlockSource::Canonical)?,
1576 Some(first_db_block.clone().into_block())
1577 );
1578 let block = consistent_provider
1579 .find_sealed_or_recovered_block(first_db_block.hash(), BlockSource::Canonical)?
1580 .expect("canonical database block should be found");
1581 assert_eq!(block.sealed_block(), first_db_block);
1582 assert!(block.recovered_block().is_none());
1583
1584 assert_eq!(
1586 consistent_provider.find_block_by_hash(first_db_block.hash(), BlockSource::Pending)?,
1587 None
1588 );
1589 assert!(consistent_provider
1590 .find_sealed_or_recovered_block(first_db_block.hash(), BlockSource::Pending)?
1591 .is_none());
1592
1593 provider.canonical_in_memory_state.set_pending_block(ExecutedBlock {
1595 recovered_block: Arc::new(RecoveredBlock::new_sealed(
1596 last_in_mem_block.clone(),
1597 Default::default(),
1598 )),
1599 ..Default::default()
1600 });
1601
1602 assert_eq!(
1604 consistent_provider
1605 .find_block_by_hash(last_in_mem_block.hash(), BlockSource::Pending)?,
1606 Some(last_in_mem_block.clone_block())
1607 );
1608 let block = consistent_provider
1609 .find_sealed_or_recovered_block(last_in_mem_block.hash(), BlockSource::Pending)?
1610 .expect("pending block should be found");
1611 assert_eq!(block.sealed_block(), last_in_mem_block);
1612 assert!(block.recovered_block().is_some());
1613
1614 Ok(())
1615 }
1616
1617 #[test]
1618 fn test_block_reader_block() -> eyre::Result<()> {
1619 let mut rng = generators::rng();
1621 let factory = create_test_provider_factory();
1622
1623 let blocks = random_block_range(
1625 &mut rng,
1626 0..=10,
1627 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
1628 );
1629 let (database_blocks, in_memory_blocks) = blocks.split_at(5);
1630
1631 let provider_rw = factory.provider_rw()?;
1633 for block in database_blocks {
1634 provider_rw.insert_block(
1635 &block.clone().try_recover().expect("failed to seal block with senders"),
1636 )?;
1637 }
1638 provider_rw.commit()?;
1639
1640 let provider = BlockchainProvider::new(factory)?;
1642 let consistent_provider = provider.consistent_provider()?;
1643
1644 let first_in_mem_block = in_memory_blocks.first().unwrap();
1646 let first_db_block = database_blocks.first().unwrap();
1648
1649 assert_eq!(
1651 consistent_provider.block(BlockHashOrNumber::Hash(first_in_mem_block.hash()))?,
1652 None
1653 );
1654 assert_eq!(
1655 consistent_provider.block(BlockHashOrNumber::Number(first_in_mem_block.number))?,
1656 None
1657 );
1658
1659 let in_memory_block_senders =
1661 first_in_mem_block.senders().expect("failed to recover senders");
1662 let chain = NewCanonicalChain::Commit {
1663 new: vec![ExecutedBlock {
1664 recovered_block: Arc::new(RecoveredBlock::new_sealed(
1665 first_in_mem_block.clone(),
1666 in_memory_block_senders,
1667 )),
1668 ..Default::default()
1669 }],
1670 };
1671 consistent_provider.canonical_in_memory_state.update_chain(chain);
1672
1673 let consistent_provider = provider.consistent_provider()?;
1674
1675 assert_eq!(
1677 consistent_provider.block(BlockHashOrNumber::Hash(first_in_mem_block.hash()))?,
1678 Some(first_in_mem_block.clone().into_block())
1679 );
1680 assert_eq!(
1681 consistent_provider.block(BlockHashOrNumber::Number(first_in_mem_block.number))?,
1682 Some(first_in_mem_block.clone().into_block())
1683 );
1684
1685 assert_eq!(
1687 consistent_provider.block(BlockHashOrNumber::Hash(first_db_block.hash()))?,
1688 Some(first_db_block.clone().into_block())
1689 );
1690 assert_eq!(
1691 consistent_provider.block(BlockHashOrNumber::Number(first_db_block.number))?,
1692 Some(first_db_block.clone().into_block())
1693 );
1694
1695 Ok(())
1696 }
1697
1698 #[test]
1699 fn test_changeset_reader() -> eyre::Result<()> {
1700 let mut rng = generators::rng();
1701
1702 let (database_blocks, in_memory_blocks) =
1703 random_blocks(&mut rng, TEST_BLOCKS_COUNT, 1, None, None, 0..1);
1704
1705 let first_database_block = database_blocks.first().map(|block| block.number).unwrap();
1706 let last_database_block = database_blocks.last().map(|block| block.number).unwrap();
1707 let first_in_memory_block = in_memory_blocks.first().map(|block| block.number).unwrap();
1708
1709 let accounts = random_eoa_accounts(&mut rng, 2);
1710
1711 let (database_changesets, database_state) = random_changeset_range(
1712 &mut rng,
1713 &database_blocks,
1714 accounts.into_iter().map(|(address, account)| (address, (account, Vec::new()))),
1715 0..0,
1716 0..0,
1717 );
1718 let (in_memory_changesets, in_memory_state) = random_changeset_range(
1719 &mut rng,
1720 &in_memory_blocks,
1721 database_state
1722 .iter()
1723 .map(|(address, (account, storage))| (*address, (*account, storage.clone()))),
1724 0..0,
1725 0..0,
1726 );
1727
1728 let factory = create_test_provider_factory();
1729
1730 let provider_rw = factory.provider_rw()?;
1731 provider_rw.append_blocks_with_state(
1732 database_blocks
1733 .into_iter()
1734 .map(|b| b.try_recover().expect("failed to seal block with senders"))
1735 .collect(),
1736 &ExecutionOutcome {
1737 bundle: BundleState::new(
1738 database_state.into_iter().map(|(address, (account, _))| {
1739 (address, None, Some(account.into()), Default::default())
1740 }),
1741 database_changesets.iter().map(|block_changesets| {
1742 block_changesets.iter().map(|(address, account, _)| {
1743 (*address, Some(Some((*account).into())), [])
1744 })
1745 }),
1746 Vec::new(),
1747 ),
1748 first_block: first_database_block,
1749 ..Default::default()
1750 },
1751 Default::default(),
1752 )?;
1753 provider_rw.commit()?;
1754
1755 let provider = BlockchainProvider::new(factory)?;
1756
1757 let in_memory_changesets = in_memory_changesets.into_iter().next().unwrap();
1758 let chain = NewCanonicalChain::Commit {
1759 new: vec![in_memory_blocks
1760 .first()
1761 .map(|block| {
1762 let senders = block.senders().expect("failed to recover senders");
1763 ExecutedBlock {
1764 recovered_block: Arc::new(RecoveredBlock::new_sealed(
1765 block.clone(),
1766 senders,
1767 )),
1768 execution_output: Arc::new(BlockExecutionOutput {
1769 state: BundleState::new(
1770 in_memory_state.into_iter().map(|(address, (account, _))| {
1771 (address, None, Some(account.into()), Default::default())
1772 }),
1773 [in_memory_changesets.iter().map(|(address, account, _)| {
1774 (*address, Some(Some((*account).into())), Vec::new())
1775 })],
1776 [],
1777 ),
1778 result: BlockExecutionResult {
1779 receipts: Default::default(),
1780 requests: Default::default(),
1781 gas_used: 0,
1782 blob_gas_used: 0,
1783 },
1784 }),
1785 ..Default::default()
1786 }
1787 })
1788 .unwrap()],
1789 };
1790 provider.canonical_in_memory_state.update_chain(chain);
1791
1792 let consistent_provider = provider.consistent_provider()?;
1793
1794 assert_eq!(
1795 consistent_provider.account_block_changeset(last_database_block).unwrap(),
1796 database_changesets
1797 .into_iter()
1798 .next_back()
1799 .unwrap()
1800 .into_iter()
1801 .sorted_by_key(|(address, _, _)| *address)
1802 .map(|(address, account, _)| AccountBeforeTx { address, info: Some(account) })
1803 .collect::<Vec<_>>()
1804 );
1805 assert_eq!(
1806 consistent_provider.account_block_changeset(first_in_memory_block).unwrap(),
1807 in_memory_changesets
1808 .into_iter()
1809 .sorted_by_key(|(address, _, _)| *address)
1810 .map(|(address, account, _)| AccountBeforeTx { address, info: Some(account) })
1811 .collect::<Vec<_>>()
1812 );
1813
1814 Ok(())
1815 }
1816 #[test]
1817 fn test_get_state_storage_value_plain_state() -> eyre::Result<()> {
1818 use alloy_primitives::U256;
1819 use reth_db_api::{models::StorageSettings, tables, transaction::DbTxMut};
1820 use reth_primitives_traits::StorageEntry;
1821 use reth_storage_api::StorageSettingsCache;
1822 use std::collections::HashMap;
1823
1824 let address = alloy_primitives::Address::with_last_byte(1);
1825 let account = reth_primitives_traits::Account {
1826 nonce: 1,
1827 balance: U256::from(1000),
1828 bytecode_hash: None,
1829 };
1830 let slot = U256::from(0x42);
1831 let slot_b256 = B256::from(slot);
1832
1833 let mut rng = generators::rng();
1834 let factory = create_test_provider_factory();
1835 factory.set_storage_settings_cache(StorageSettings::v1());
1836
1837 let blocks = random_block_range(
1838 &mut rng,
1839 0..=1,
1840 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
1841 );
1842
1843 let provider_rw = factory.provider_rw()?;
1844 provider_rw.append_blocks_with_state(
1845 blocks
1846 .into_iter()
1847 .map(|b| b.try_recover().expect("failed to seal block with senders"))
1848 .collect(),
1849 &ExecutionOutcome {
1850 bundle: BundleState::new(
1851 [(address, None, Some(account.into()), {
1852 let mut s = HashMap::default();
1853 s.insert(slot, (U256::ZERO, U256::from(100)));
1854 s
1855 })],
1856 [
1857 Vec::new(),
1858 vec![(address, Some(Some(account.into())), vec![(slot, U256::ZERO)])],
1859 ],
1860 [],
1861 ),
1862 first_block: 0,
1863 ..Default::default()
1864 },
1865 Default::default(),
1866 )?;
1867
1868 provider_rw.tx_ref().put::<tables::PlainStorageState>(
1869 address,
1870 StorageEntry { key: slot_b256, value: U256::from(100) },
1871 )?;
1872 provider_rw.tx_ref().put::<tables::PlainAccountState>(address, account)?;
1873
1874 provider_rw.commit()?;
1875
1876 let provider = BlockchainProvider::new(factory)?;
1877 let outcome = provider.get_state(1)?.expect("should return execution outcome");
1878
1879 let state = &outcome.bundle.state;
1880 let account_state = state.get(&address).expect("should have account in bundle state");
1881 let storage = &account_state.storage;
1882
1883 let storage_slot = storage.get(&slot).expect("should have the slot in storage");
1884
1885 assert_eq!(
1886 storage_slot.present_value,
1887 U256::from(100),
1888 "present_value should be 100 (the actual value in PlainStorageState)"
1889 );
1890
1891 Ok(())
1892 }
1893
1894 #[test]
1895 fn test_storage_changeset_consistent_keys_plain_state() -> eyre::Result<()> {
1896 use alloy_primitives::U256;
1897 use reth_db_api::models::StorageSettings;
1898 use reth_storage_api::{StorageChangeSetReader, StorageSettingsCache};
1899 use std::collections::HashMap;
1900
1901 let mut rng = generators::rng();
1902 let factory = create_test_provider_factory();
1903 factory.set_storage_settings_cache(StorageSettings::v1());
1904
1905 let (database_blocks, in_memory_blocks) = random_blocks(&mut rng, 1, 1, None, None, 0..1);
1906
1907 let address = alloy_primitives::Address::with_last_byte(1);
1908 let account = reth_primitives_traits::Account {
1909 nonce: 1,
1910 balance: U256::from(1000),
1911 bytecode_hash: None,
1912 };
1913 let slot = U256::from(0x42);
1914
1915 let provider_rw = factory.provider_rw()?;
1916 provider_rw.append_blocks_with_state(
1917 database_blocks
1918 .into_iter()
1919 .map(|b| b.try_recover().expect("failed to seal block with senders"))
1920 .collect(),
1921 &ExecutionOutcome {
1922 bundle: BundleState::new(
1923 [(address, None, Some(account.into()), {
1924 let mut s = HashMap::default();
1925 s.insert(slot, (U256::ZERO, U256::from(100)));
1926 s
1927 })],
1928 [[(address, Some(Some(account.into())), vec![(slot, U256::ZERO)])]],
1929 [],
1930 ),
1931 first_block: 0,
1932 ..Default::default()
1933 },
1934 Default::default(),
1935 )?;
1936 provider_rw.commit()?;
1937
1938 let provider = BlockchainProvider::new(factory)?;
1939
1940 let in_mem_block = in_memory_blocks.first().unwrap();
1941 let senders = in_mem_block.senders().expect("failed to recover senders");
1942 let chain = NewCanonicalChain::Commit {
1943 new: vec![ExecutedBlock {
1944 recovered_block: Arc::new(RecoveredBlock::new_sealed(
1945 in_mem_block.clone(),
1946 senders,
1947 )),
1948 execution_output: Arc::new(BlockExecutionOutput {
1949 state: BundleState::new(
1950 [(address, None, Some(account.into()), {
1951 let mut s = HashMap::default();
1952 s.insert(slot, (U256::from(100), U256::from(200)));
1953 s
1954 })],
1955 [[(address, Some(Some(account.into())), vec![(slot, U256::from(100))])]],
1956 [],
1957 ),
1958 result: BlockExecutionResult {
1959 receipts: Default::default(),
1960 requests: Default::default(),
1961 gas_used: 0,
1962 blob_gas_used: 0,
1963 },
1964 }),
1965 ..Default::default()
1966 }],
1967 };
1968 provider.canonical_in_memory_state.update_chain(chain);
1969
1970 let consistent_provider = provider.consistent_provider()?;
1971
1972 let db_changeset = consistent_provider.storage_changeset(0)?;
1973 let mem_changeset = consistent_provider.storage_changeset(1)?;
1974
1975 let slot_b256 = B256::from(slot);
1976
1977 assert_eq!(db_changeset.len(), 1);
1978 assert_eq!(mem_changeset.len(), 1);
1979
1980 let db_key = db_changeset[0].1.key;
1981 let mem_key = mem_changeset[0].1.key;
1982
1983 assert_eq!(db_key, slot_b256, "DB changeset should use plain (unhashed) key");
1984 assert_eq!(mem_key, slot_b256, "In-memory changeset should use plain (unhashed) key");
1985 assert_eq!(
1986 db_key, mem_key,
1987 "DB and in-memory changesets should return the same key format (plain) for the same logical slot"
1988 );
1989
1990 Ok(())
1991 }
1992
1993 #[test]
1994 fn test_storage_changesets_range_consistent_keys_plain_state() -> eyre::Result<()> {
1995 use alloy_primitives::U256;
1996 use reth_db_api::models::StorageSettings;
1997 use reth_storage_api::{StorageChangeSetReader, StorageSettingsCache};
1998 use std::collections::HashMap;
1999
2000 let mut rng = generators::rng();
2001 let factory = create_test_provider_factory();
2002 factory.set_storage_settings_cache(StorageSettings::v1());
2003
2004 let (database_blocks, in_memory_blocks) = random_blocks(&mut rng, 2, 1, None, None, 0..1);
2005
2006 let address = alloy_primitives::Address::with_last_byte(1);
2007 let account = reth_primitives_traits::Account {
2008 nonce: 1,
2009 balance: U256::from(1000),
2010 bytecode_hash: None,
2011 };
2012 let slot = U256::from(0x42);
2013
2014 let provider_rw = factory.provider_rw()?;
2015 provider_rw.append_blocks_with_state(
2016 database_blocks
2017 .into_iter()
2018 .map(|b| b.try_recover().expect("failed to seal block with senders"))
2019 .collect(),
2020 &ExecutionOutcome {
2021 bundle: BundleState::new(
2022 [(address, None, Some(account.into()), {
2023 let mut s = HashMap::default();
2024 s.insert(slot, (U256::ZERO, U256::from(100)));
2025 s
2026 })],
2027 vec![
2028 vec![(address, Some(Some(account.into())), vec![(slot, U256::ZERO)])],
2029 vec![],
2030 ],
2031 [],
2032 ),
2033 first_block: 0,
2034 ..Default::default()
2035 },
2036 Default::default(),
2037 )?;
2038 provider_rw.commit()?;
2039
2040 let provider = BlockchainProvider::new(factory)?;
2041
2042 let in_mem_block = in_memory_blocks.first().unwrap();
2043 let senders = in_mem_block.senders().expect("failed to recover senders");
2044 let chain = NewCanonicalChain::Commit {
2045 new: vec![ExecutedBlock {
2046 recovered_block: Arc::new(RecoveredBlock::new_sealed(
2047 in_mem_block.clone(),
2048 senders,
2049 )),
2050 execution_output: Arc::new(BlockExecutionOutput {
2051 state: BundleState::new(
2052 [(address, None, Some(account.into()), {
2053 let mut s = HashMap::default();
2054 s.insert(slot, (U256::from(100), U256::from(200)));
2055 s
2056 })],
2057 [[(address, Some(Some(account.into())), vec![(slot, U256::from(100))])]],
2058 [],
2059 ),
2060 result: BlockExecutionResult {
2061 receipts: Default::default(),
2062 requests: Default::default(),
2063 gas_used: 0,
2064 blob_gas_used: 0,
2065 },
2066 }),
2067 ..Default::default()
2068 }],
2069 };
2070 provider.canonical_in_memory_state.update_chain(chain);
2071
2072 let consistent_provider = provider.consistent_provider()?;
2073
2074 let all_changesets = consistent_provider.storage_changesets_range(0..=2)?;
2075
2076 assert_eq!(all_changesets.len(), 2, "should have one changeset entry per block");
2077
2078 let slot_b256 = B256::from(slot);
2079 let keys: Vec<B256> = all_changesets.iter().map(|(_, entry)| entry.key).collect();
2080
2081 assert_eq!(
2082 keys[0], keys[1],
2083 "same logical slot should produce identical keys whether from DB or memory"
2084 );
2085 assert_eq!(
2086 keys[0], slot_b256,
2087 "keys should be plain/unhashed when use_hashed_state is false"
2088 );
2089
2090 Ok(())
2091 }
2092}