Skip to main content

reth_provider/providers/
consistent.rs

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