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