1use crate::{
4 CanonStateNotification, CanonStateNotificationSender, CanonStateNotifications, ChainInfoTracker,
5};
6use alloy_consensus::{transaction::TransactionMeta, BlockHeader};
7use alloy_eips::{BlockHashOrNumber, BlockNumHash};
8use alloy_primitives::{map::B256Map, BlockNumber, TxHash, B256};
9use parking_lot::RwLock;
10use reth_chainspec::ChainInfo;
11use reth_ethereum_primitives::EthPrimitives;
12use reth_execution_types::{
13 BlockExecutionOutput, BlockExecutionResult, Chain, DecodedRevmBal, ExecutionOutcome,
14 RecoveredBlockAndExecutionOutput,
15};
16use reth_metrics::{metrics::Gauge, Metrics};
17use reth_primitives_traits::{
18 BlockBody as _, IndexedTx, NodePrimitives, RecoveredBlock, SealedBlock, SealedHeader,
19 SignedTransaction,
20};
21use reth_trie::{
22 updates::TrieUpdatesSorted, BlockTrieData, HashedPostStateSorted, LazyHashedPostStateSorted,
23};
24use std::{collections::BTreeMap, sync::Arc, time::Instant};
25use tokio::sync::{broadcast, watch};
26
27const CANON_STATE_NOTIFICATION_CHANNEL_SIZE: usize = 256;
29
30#[derive(Metrics)]
32#[metrics(scope = "blockchain_tree.in_mem_state")]
33pub(crate) struct InMemoryStateMetrics {
34 pub(crate) earliest_block: Gauge,
36 pub(crate) latest_block: Gauge,
38 pub(crate) num_blocks: Gauge,
40}
41
42#[derive(Debug, Default)]
58pub(crate) struct InMemoryState<N: NodePrimitives = EthPrimitives> {
59 blocks: RwLock<B256Map<Arc<BlockState<N>>>>,
61 numbers: RwLock<BTreeMap<u64, B256>>,
63 pending: watch::Sender<Option<BlockState<N>>>,
65 metrics: InMemoryStateMetrics,
67}
68
69impl<N: NodePrimitives> InMemoryState<N> {
70 pub(crate) fn new(
71 blocks: B256Map<Arc<BlockState<N>>>,
72 numbers: BTreeMap<u64, B256>,
73 pending: Option<BlockState<N>>,
74 ) -> Self {
75 let (pending, _) = watch::channel(pending);
76 let this = Self {
77 blocks: RwLock::new(blocks),
78 numbers: RwLock::new(numbers),
79 pending,
80 metrics: Default::default(),
81 };
82 this.update_metrics();
83 this
84 }
85
86 pub(crate) fn update_metrics(&self) {
92 let (count, earliest, latest) = {
93 let numbers = self.numbers.read();
94 let count = numbers.len();
95 let earliest = numbers.first_key_value().map(|(number, _)| *number);
96 let latest = numbers.last_key_value().map(|(number, _)| *number);
97 (count, earliest, latest)
98 };
99 if let Some(earliest_block_number) = earliest {
100 self.metrics.earliest_block.set(earliest_block_number as f64);
101 }
102 if let Some(latest_block_number) = latest {
103 self.metrics.latest_block.set(latest_block_number as f64);
104 }
105 self.metrics.num_blocks.set(count as f64);
106 }
107
108 pub(crate) fn state_by_hash(&self, hash: B256) -> Option<Arc<BlockState<N>>> {
110 self.blocks.read().get(&hash).cloned()
111 }
112
113 pub(crate) fn state_by_number(&self, number: u64) -> Option<Arc<BlockState<N>>> {
115 let hash = self.hash_by_number(number)?;
116 self.state_by_hash(hash)
117 }
118
119 pub(crate) fn hash_by_number(&self, number: u64) -> Option<B256> {
121 self.numbers.read().get(&number).copied()
122 }
123
124 pub(crate) fn head_state(&self) -> Option<Arc<BlockState<N>>> {
126 let hash = *self.numbers.read().last_key_value()?.1;
127 self.state_by_hash(hash)
128 }
129
130 pub(crate) fn pending_state(&self) -> Option<BlockState<N>> {
133 self.pending.borrow().clone()
134 }
135
136 #[cfg(test)]
137 fn block_count(&self) -> usize {
138 self.blocks.read().len()
139 }
140}
141
142#[derive(Debug)]
145pub(crate) struct CanonicalInMemoryStateInner<N: NodePrimitives> {
146 pub(crate) chain_info_tracker: ChainInfoTracker<N>,
149 pub(crate) in_memory_state: InMemoryState<N>,
151 pub(crate) canon_state_notification_sender: CanonStateNotificationSender<N>,
153}
154
155impl<N: NodePrimitives> CanonicalInMemoryStateInner<N> {
156 fn clear(&self) {
158 {
159 let mut numbers = self.in_memory_state.numbers.write();
161 let mut blocks = self.in_memory_state.blocks.write();
162 numbers.clear();
163 blocks.clear();
164 self.in_memory_state.pending.send_modify(|p| {
165 p.take();
166 });
167 }
168 self.in_memory_state.update_metrics();
169 }
170}
171
172#[derive(Debug, Clone)]
176pub struct CanonicalInMemoryState<N: NodePrimitives = EthPrimitives> {
177 pub(crate) inner: Arc<CanonicalInMemoryStateInner<N>>,
178}
179
180impl<N: NodePrimitives> CanonicalInMemoryState<N> {
181 pub fn new(
184 blocks: B256Map<Arc<BlockState<N>>>,
185 numbers: BTreeMap<u64, B256>,
186 pending: Option<BlockState<N>>,
187 finalized: Option<SealedHeader<N::BlockHeader>>,
188 safe: Option<SealedHeader<N::BlockHeader>>,
189 ) -> Self {
190 let in_memory_state = InMemoryState::new(blocks, numbers, pending);
191 let header = in_memory_state.head_state().map_or_else(SealedHeader::default, |state| {
192 state.block_ref().recovered_block().clone_sealed_header()
193 });
194 let chain_info_tracker = ChainInfoTracker::new(header, finalized, safe);
195 let (canon_state_notification_sender, _) =
196 broadcast::channel(CANON_STATE_NOTIFICATION_CHANNEL_SIZE);
197
198 Self {
199 inner: Arc::new(CanonicalInMemoryStateInner {
200 chain_info_tracker,
201 in_memory_state,
202 canon_state_notification_sender,
203 }),
204 }
205 }
206
207 pub fn empty() -> Self {
209 Self::new(B256Map::default(), BTreeMap::new(), None, None, None)
210 }
211
212 pub fn with_head(
215 head: SealedHeader<N::BlockHeader>,
216 finalized: Option<SealedHeader<N::BlockHeader>>,
217 safe: Option<SealedHeader<N::BlockHeader>>,
218 ) -> Self {
219 let chain_info_tracker = ChainInfoTracker::new(head, finalized, safe);
220 let in_memory_state = InMemoryState::default();
221 let (canon_state_notification_sender, _) =
222 broadcast::channel(CANON_STATE_NOTIFICATION_CHANNEL_SIZE);
223 let inner = CanonicalInMemoryStateInner {
224 chain_info_tracker,
225 in_memory_state,
226 canon_state_notification_sender,
227 };
228
229 Self { inner: Arc::new(inner) }
230 }
231
232 pub fn hash_by_number(&self, number: u64) -> Option<B256> {
234 self.inner.in_memory_state.hash_by_number(number)
235 }
236
237 pub fn header_by_hash(&self, hash: B256) -> Option<SealedHeader<N::BlockHeader>> {
239 self.state_by_hash(hash)
240 .map(|block| block.block_ref().recovered_block().clone_sealed_header())
241 }
242
243 pub fn clear_state(&self) {
245 self.inner.clear()
246 }
247
248 pub fn set_pending_block(&self, pending: ExecutedBlock<N>) {
252 let parent = self.state_by_hash(pending.recovered_block().parent_hash());
254 let pending = BlockState::with_parent(pending, parent);
255 self.inner.in_memory_state.pending.send_modify(|p| {
256 p.replace(pending);
257 });
258 self.inner.in_memory_state.update_metrics();
259 }
260
261 fn update_blocks<I, R>(&self, new_blocks: I, reorged: R)
266 where
267 I: IntoIterator<Item = ExecutedBlock<N>>,
268 R: IntoIterator<Item = ExecutedBlock<N>>,
269 {
270 {
271 let mut numbers = self.inner.in_memory_state.numbers.write();
273 let mut blocks = self.inner.in_memory_state.blocks.write();
274
275 for block in reorged {
277 let hash = block.recovered_block().hash();
278 let number = block.recovered_block().number();
279 blocks.remove(&hash);
280 numbers.remove(&number);
281 }
282
283 for block in new_blocks {
285 let parent = blocks.get(&block.recovered_block().parent_hash()).cloned();
286 let block_state = BlockState::with_parent(block, parent);
287 let hash = block_state.hash();
288 let number = block_state.number();
289
290 blocks.insert(hash, Arc::new(block_state));
292 numbers.insert(number, hash);
293 }
294
295 self.inner.in_memory_state.pending.send_modify(|p| {
297 p.take();
298 });
299 }
300 self.inner.in_memory_state.update_metrics();
301 }
302
303 pub fn update_chain(&self, new_chain: NewCanonicalChain<N>) {
305 match new_chain {
306 NewCanonicalChain::Commit { new } => {
307 self.update_blocks(new, vec![]);
308 }
309 NewCanonicalChain::Reorg { new, old } => {
310 self.update_blocks(new, old);
311 }
312 }
313 }
314
315 pub fn remove_persisted_blocks(&self, persisted_num_hash: BlockNumHash) {
320 self.remove_persisted_blocks_until(persisted_num_hash, persisted_num_hash.number);
321 }
322
323 pub fn remove_persisted_blocks_until(
326 &self,
327 persisted_num_hash: BlockNumHash,
328 remove_until: BlockNumber,
329 ) {
330 self.set_persisted(persisted_num_hash);
331 {
336 if self.inner.in_memory_state.blocks.read().get(&persisted_num_hash.hash).is_none() {
337 return
339 }
340 }
341
342 {
343 let mut numbers = self.inner.in_memory_state.numbers.write();
345 let mut blocks = self.inner.in_memory_state.blocks.write();
346
347 let remove_until = remove_until.min(persisted_num_hash.number);
348
349 numbers.clear();
351
352 let mut old_blocks = blocks
354 .drain()
355 .filter(|(_, b)| b.block_ref().recovered_block().number() > remove_until)
356 .map(|(_, b)| b.block.clone())
357 .collect::<Vec<_>>();
358
359 old_blocks.sort_unstable_by_key(|block| block.recovered_block().number());
361
362 for block in old_blocks {
364 let parent = blocks.get(&block.recovered_block().parent_hash()).cloned();
365 let block_state = BlockState::with_parent(block, parent);
366 let hash = block_state.hash();
367 let number = block_state.number();
368
369 blocks.insert(hash, Arc::new(block_state));
371 numbers.insert(number, hash);
372 }
373
374 self.inner.in_memory_state.pending.send_modify(|p| {
376 if let Some(p) = p.as_mut() {
377 p.parent = blocks.get(&p.block_ref().recovered_block().parent_hash()).cloned();
378 }
379 });
380 }
381 self.inner.in_memory_state.update_metrics();
382 }
383
384 pub fn state_by_hash(&self, hash: B256) -> Option<Arc<BlockState<N>>> {
386 self.inner.in_memory_state.state_by_hash(hash)
387 }
388
389 pub fn state_by_number(&self, number: u64) -> Option<Arc<BlockState<N>>> {
391 self.inner.in_memory_state.state_by_number(number)
392 }
393
394 pub fn head_state(&self) -> Option<Arc<BlockState<N>>> {
396 self.inner.in_memory_state.head_state()
397 }
398
399 pub fn pending_state(&self) -> Option<BlockState<N>> {
401 self.inner.in_memory_state.pending_state()
402 }
403
404 pub fn pending_block_num_hash(&self) -> Option<BlockNumHash> {
406 self.inner
407 .in_memory_state
408 .pending_state()
409 .map(|state| BlockNumHash { number: state.number(), hash: state.hash() })
410 }
411
412 pub fn chain_info(&self) -> ChainInfo {
414 self.inner.chain_info_tracker.chain_info()
415 }
416
417 pub fn get_canonical_block_number(&self) -> u64 {
419 self.inner.chain_info_tracker.get_canonical_block_number()
420 }
421
422 pub fn get_safe_num_hash(&self) -> Option<BlockNumHash> {
424 self.inner.chain_info_tracker.get_safe_num_hash()
425 }
426
427 pub fn get_finalized_num_hash(&self) -> Option<BlockNumHash> {
429 self.inner.chain_info_tracker.get_finalized_num_hash()
430 }
431
432 pub fn on_forkchoice_update_received(&self) {
434 self.inner.chain_info_tracker.on_forkchoice_update_received();
435 }
436
437 pub fn last_received_update_timestamp(&self) -> Option<Instant> {
439 self.inner.chain_info_tracker.last_forkchoice_update_received_at()
440 }
441
442 pub fn set_canonical_head(&self, header: SealedHeader<N::BlockHeader>) {
444 self.inner.chain_info_tracker.set_canonical_head(header);
445 }
446
447 pub fn set_safe(&self, header: SealedHeader<N::BlockHeader>) {
449 self.inner.chain_info_tracker.set_safe(header);
450 }
451
452 pub fn set_finalized(&self, header: SealedHeader<N::BlockHeader>) {
454 self.inner.chain_info_tracker.set_finalized(header);
455 }
456
457 pub fn set_persisted(&self, num_hash: BlockNumHash) {
459 self.inner.chain_info_tracker.set_persisted(num_hash);
460 }
461
462 pub fn get_canonical_head(&self) -> SealedHeader<N::BlockHeader> {
464 self.inner.chain_info_tracker.get_canonical_head()
465 }
466
467 pub fn get_finalized_header(&self) -> Option<SealedHeader<N::BlockHeader>> {
469 self.inner.chain_info_tracker.get_finalized_header()
470 }
471
472 pub fn get_safe_header(&self) -> Option<SealedHeader<N::BlockHeader>> {
474 self.inner.chain_info_tracker.get_safe_header()
475 }
476
477 pub fn get_persisted_num_hash(&self) -> Option<BlockNumHash> {
479 self.inner.chain_info_tracker.get_persisted_num_hash()
480 }
481
482 pub fn pending_sealed_header(&self) -> Option<SealedHeader<N::BlockHeader>> {
484 self.pending_state().map(|h| h.block_ref().recovered_block().clone_sealed_header())
485 }
486
487 pub fn pending_header(&self) -> Option<N::BlockHeader> {
489 self.pending_sealed_header().map(|sealed_header| sealed_header.unseal())
490 }
491
492 pub fn pending_block(&self) -> Option<SealedBlock<N::Block>> {
494 self.pending_state()
495 .map(|block_state| block_state.block_ref().recovered_block().sealed_block().clone())
496 }
497
498 pub fn pending_recovered_block(&self) -> Option<Arc<RecoveredBlock<N::Block>>> {
500 self.pending_state().map(|block_state| Arc::clone(&block_state.block_ref().recovered_block))
501 }
502
503 pub fn pending_block_and_receipts(
505 &self,
506 ) -> Option<RecoveredBlockAndExecutionOutput<N::Block, N::Receipt>> {
507 self.pending_state().map(|block_state| {
508 RecoveredBlockAndExecutionOutput::new(
509 Arc::clone(&block_state.block_ref().recovered_block),
510 Arc::clone(&block_state.block_ref().execution_output),
511 )
512 })
513 }
514
515 pub fn subscribe_canon_state(&self) -> CanonStateNotifications<N> {
517 self.inner.canon_state_notification_sender.subscribe()
518 }
519
520 pub fn subscribe_safe_block(&self) -> watch::Receiver<Option<SealedHeader<N::BlockHeader>>> {
522 self.inner.chain_info_tracker.subscribe_safe_block()
523 }
524
525 pub fn subscribe_finalized_block(
527 &self,
528 ) -> watch::Receiver<Option<SealedHeader<N::BlockHeader>>> {
529 self.inner.chain_info_tracker.subscribe_finalized_block()
530 }
531
532 pub fn subscribe_persisted_block(&self) -> watch::Receiver<Option<BlockNumHash>> {
534 self.inner.chain_info_tracker.subscribe_persisted_block()
535 }
536
537 pub fn notify_canon_state(&self, event: CanonStateNotification<N>) {
539 self.inner.canon_state_notification_sender.send(event).ok();
540 }
541
542 pub fn canonical_chain(&self) -> impl Iterator<Item = Arc<BlockState<N>>> {
547 self.inner.in_memory_state.head_state().into_iter().flat_map(|head| head.iter())
548 }
549
550 pub fn transaction_by_hash(&self, hash: TxHash) -> Option<N::SignedTx> {
552 for block_state in self.canonical_chain() {
553 if let Some(tx) =
554 block_state.block_ref().recovered_block().body().transaction_by_hash(&hash)
555 {
556 return Some(tx.clone())
557 }
558 }
559 None
560 }
561
562 pub fn transaction_by_hash_with_meta(
565 &self,
566 tx_hash: TxHash,
567 ) -> Option<(N::SignedTx, TransactionMeta)> {
568 for block_state in self.canonical_chain() {
569 if let Some(indexed) = block_state.find_indexed(tx_hash) {
570 return Some((indexed.tx().clone(), indexed.meta()));
571 }
572 }
573 None
574 }
575}
576
577#[derive(Debug, Clone)]
580pub struct BlockState<N: NodePrimitives = EthPrimitives> {
581 block: ExecutedBlock<N>,
583 parent: Option<Arc<Self>>,
585}
586
587impl<N: NodePrimitives> PartialEq for BlockState<N> {
588 fn eq(&self, other: &Self) -> bool {
589 self.block == other.block && self.parent == other.parent
590 }
591}
592
593impl<N: NodePrimitives> BlockState<N> {
594 pub const fn new(block: ExecutedBlock<N>) -> Self {
596 Self { block, parent: None }
597 }
598
599 pub const fn with_parent(block: ExecutedBlock<N>, parent: Option<Arc<Self>>) -> Self {
601 Self { block, parent }
602 }
603
604 pub fn anchor(&self) -> BlockNumHash {
606 let mut current = self;
607 while let Some(parent) = ¤t.parent {
608 current = parent;
609 }
610 current.block.recovered_block().parent_num_hash()
611 }
612
613 pub fn block(&self) -> ExecutedBlock<N> {
615 self.block.clone()
616 }
617
618 pub const fn block_ref(&self) -> &ExecutedBlock<N> {
620 &self.block
621 }
622
623 pub fn hash(&self) -> B256 {
625 self.block.recovered_block().hash()
626 }
627
628 pub fn number(&self) -> u64 {
630 self.block.recovered_block().number()
631 }
632
633 pub fn state_root(&self) -> B256 {
636 self.block.recovered_block().state_root()
637 }
638
639 pub fn receipts(&self) -> &Vec<N::Receipt> {
641 &self.block.execution_outcome().receipts
642 }
643
644 pub fn executed_block_receipts(&self) -> Vec<N::Receipt> {
651 self.receipts().clone()
652 }
653
654 pub fn executed_block_receipts_ref(&self) -> &[N::Receipt] {
659 self.receipts()
660 }
661
662 pub fn parent_state_chain(&self) -> impl Iterator<Item = &Self> + '_ {
669 std::iter::successors(self.parent.as_deref(), |state| state.parent.as_deref())
670 }
671
672 pub fn chain(&self) -> impl Iterator<Item = &Self> {
676 std::iter::successors(Some(self), |state| state.parent.as_deref())
677 }
678
679 pub fn append_parent_chain<'a>(&'a self, chain: &mut Vec<&'a Self>) {
686 chain.extend(self.parent_state_chain());
687 }
688
689 pub fn iter(self: Arc<Self>) -> impl Iterator<Item = Arc<Self>> {
693 std::iter::successors(Some(self), |state| state.parent.clone())
694 }
695
696 pub fn block_on_chain(&self, hash_or_num: BlockHashOrNumber) -> Option<&Self> {
698 self.chain().find(|block| match hash_or_num {
699 BlockHashOrNumber::Hash(hash) => block.hash() == hash,
700 BlockHashOrNumber::Number(number) => block.number() == number,
701 })
702 }
703
704 pub fn transaction_on_chain(&self, hash: TxHash) -> Option<N::SignedTx> {
706 self.chain().find_map(|block_state| {
707 block_state.block_ref().recovered_block().body().transaction_by_hash(&hash).cloned()
708 })
709 }
710
711 pub fn transaction_meta_on_chain(
713 &self,
714 tx_hash: TxHash,
715 ) -> Option<(N::SignedTx, TransactionMeta)> {
716 self.chain().find_map(|block_state| {
717 block_state.find_indexed(tx_hash).map(|indexed| (indexed.tx().clone(), indexed.meta()))
718 })
719 }
720
721 pub fn find_indexed(&self, tx_hash: TxHash) -> Option<IndexedTx<'_, N::Block>> {
723 self.block_ref().recovered_block().find_indexed(tx_hash)
724 }
725}
726
727#[derive(Clone, Debug)]
729pub struct ExecutedBlock<N: NodePrimitives = EthPrimitives> {
730 pub recovered_block: Arc<RecoveredBlock<N::Block>>,
732 pub execution_output: Arc<BlockExecutionOutput<N::Receipt>>,
734 pub hashed_state: LazyHashedPostStateSorted,
736 pub trie_updates: Arc<TrieUpdatesSorted>,
738 pub bal: Option<Arc<DecodedRevmBal>>,
750}
751
752impl<N: NodePrimitives> Default for ExecutedBlock<N> {
753 fn default() -> Self {
754 Self {
755 recovered_block: Default::default(),
756 execution_output: Arc::new(BlockExecutionOutput {
757 result: BlockExecutionResult {
758 receipts: Default::default(),
759 requests: Default::default(),
760 gas_used: 0,
761 blob_gas_used: 0,
762 },
763 state: Default::default(),
764 }),
765 hashed_state: LazyHashedPostStateSorted::ready(Default::default()),
766 trie_updates: Default::default(),
767 bal: None,
768 }
769 }
770}
771
772impl<N: NodePrimitives> PartialEq for ExecutedBlock<N> {
773 fn eq(&self, other: &Self) -> bool {
774 self.recovered_block == other.recovered_block &&
776 self.execution_output == other.execution_output
777 }
778}
779
780impl<N: NodePrimitives> ExecutedBlock<N> {
781 pub fn new(
783 recovered_block: Arc<RecoveredBlock<N::Block>>,
784 execution_output: Arc<BlockExecutionOutput<N::Receipt>>,
785 hashed_state: Arc<HashedPostStateSorted>,
786 trie_updates: Arc<TrieUpdatesSorted>,
787 ) -> Self {
788 Self::with_deferred_hashed_state(
789 recovered_block,
790 execution_output,
791 LazyHashedPostStateSorted::ready(hashed_state),
792 trie_updates,
793 )
794 }
795
796 pub const fn with_deferred_hashed_state(
799 recovered_block: Arc<RecoveredBlock<N::Block>>,
800 execution_output: Arc<BlockExecutionOutput<N::Receipt>>,
801 hashed_state: LazyHashedPostStateSorted,
802 trie_updates: Arc<TrieUpdatesSorted>,
803 ) -> Self {
804 Self { recovered_block, execution_output, hashed_state, trie_updates, bal: None }
805 }
806
807 pub fn with_bal(mut self, bal: Option<Arc<DecodedRevmBal>>) -> Self {
809 self.bal = bal;
810 self
811 }
812
813 #[inline]
817 pub const fn bal(&self) -> Option<&Arc<DecodedRevmBal>> {
818 self.bal.as_ref()
819 }
820
821 #[inline]
823 pub fn sealed_block(&self) -> &SealedBlock<N::Block> {
824 self.recovered_block.sealed_block()
825 }
826
827 #[inline]
829 pub fn recovered_block(&self) -> &RecoveredBlock<N::Block> {
830 &self.recovered_block
831 }
832
833 #[inline]
835 pub fn execution_outcome(&self) -> &BlockExecutionOutput<N::Receipt> {
836 &self.execution_output
837 }
838
839 pub fn trie_data(&self) -> BlockTrieData {
841 BlockTrieData {
842 hashed_state: self.hashed_state.clone(),
843 trie_updates: Arc::clone(&self.trie_updates),
844 }
845 }
846
847 #[inline]
851 pub fn hashed_state(&self) -> Arc<HashedPostStateSorted> {
852 self.hashed_state.get().clone()
853 }
854
855 #[inline]
859 pub fn hashed_state_ref(&self) -> &HashedPostStateSorted {
860 self.hashed_state.get()
861 }
862
863 pub fn hashed_state_refs(blocks: &[Self]) -> Vec<&HashedPostStateSorted> {
867 blocks.iter().map(Self::hashed_state_ref).collect()
868 }
869
870 #[inline]
872 pub fn trie_updates(&self) -> Arc<TrieUpdatesSorted> {
873 Arc::clone(&self.trie_updates)
874 }
875
876 #[inline]
878 pub fn trie_updates_ref(&self) -> &TrieUpdatesSorted {
879 &self.trie_updates
880 }
881
882 pub fn trie_updates_refs(blocks: &[Self]) -> Vec<&TrieUpdatesSorted> {
884 blocks.iter().map(Self::trie_updates_ref).collect()
885 }
886
887 #[inline]
889 pub fn block_number(&self) -> BlockNumber {
890 self.recovered_block.header().number()
891 }
892}
893
894#[derive(Debug)]
896pub enum NewCanonicalChain<N: NodePrimitives = EthPrimitives> {
897 Commit {
899 new: Vec<ExecutedBlock<N>>,
901 },
902 Reorg {
905 new: Vec<ExecutedBlock<N>>,
907 old: Vec<ExecutedBlock<N>>,
909 },
910}
911
912impl<N: NodePrimitives<SignedTx: SignedTransaction>> NewCanonicalChain<N> {
913 pub const fn new_block_count(&self) -> usize {
915 match self {
916 Self::Commit { new } | Self::Reorg { new, .. } => new.len(),
917 }
918 }
919
920 pub const fn reorged_block_count(&self) -> usize {
922 match self {
923 Self::Commit { .. } => 0,
924 Self::Reorg { old, .. } => old.len(),
925 }
926 }
927
928 pub fn to_chain_notification(&self) -> CanonStateNotification<N> {
930 match self {
931 Self::Commit { new } => {
932 CanonStateNotification::Commit { new: Arc::new(Self::blocks_to_chain(new)) }
933 }
934 Self::Reorg { new, old } => CanonStateNotification::Reorg {
935 new: Arc::new(Self::blocks_to_chain(new)),
936 old: Arc::new(Self::blocks_to_chain(old)),
937 },
938 }
939 }
940
941 fn blocks_to_chain(blocks: &[ExecutedBlock<N>]) -> Chain<N> {
943 let mut chain = match blocks {
944 [] => Chain::default(),
945 [first, rest @ ..] => {
946 let mut chain = Chain::from_block(
947 Arc::clone(&first.recovered_block),
948 ExecutionOutcome::from((
949 first.execution_outcome().clone(),
950 first.block_number(),
951 )),
952 first.trie_data(),
953 );
954 for exec in rest {
955 chain.append_block(
956 Arc::clone(&exec.recovered_block),
957 ExecutionOutcome::from((
958 exec.execution_outcome().clone(),
959 exec.block_number(),
960 )),
961 exec.trie_data(),
962 );
963 }
964 chain
965 }
966 };
967 for exec in blocks {
968 if let Some(bal) = exec.bal() {
969 chain.insert_bal(exec.block_number(), Arc::clone(bal));
970 }
971 }
972 chain
973 }
974
975 pub fn tip(&self) -> &RecoveredBlock<N::Block> {
980 self.new_blocks().last().expect("non empty blocks").recovered_block()
981 }
982
983 pub fn new_blocks(&self) -> &[ExecutedBlock<N>] {
985 match self {
986 Self::Commit { new } | Self::Reorg { new, .. } => new,
987 }
988 }
989
990 pub fn contains(&self, hash: B256) -> bool {
994 self.new_blocks().iter().any(|block| block.recovered_block().hash() == hash)
995 }
996}
997
998#[cfg(test)]
999mod tests {
1000 use super::*;
1001 use crate::test_utils::TestBlockBuilder;
1002 use alloy_eips::eip7685::Requests;
1003 use alloy_primitives::Bytes;
1004 use rand::Rng;
1005 use reth_ethereum_primitives::{EthPrimitives, Receipt};
1006
1007 #[test]
1008 fn trie_updates_are_available_before_hashed_state_is_published() {
1009 let updates = Arc::new(TrieUpdatesSorted::new(
1010 vec![(reth_trie::Nibbles::from_nibbles([1]), None)],
1011 Default::default(),
1012 ));
1013 let (hashed_state, producer) = LazyHashedPostStateSorted::pending(Arc::default());
1014 let block = ExecutedBlock::<EthPrimitives>::with_deferred_hashed_state(
1015 Default::default(),
1016 Default::default(),
1017 hashed_state,
1018 Arc::clone(&updates),
1019 );
1020 let (tx, rx) = std::sync::mpsc::channel();
1021 let reader = std::thread::spawn(move || {
1022 let count = block.trie_data().trie_updates.total_len();
1023 tx.send((count, block.trie_updates())).unwrap();
1024 block
1025 });
1026 let result = rx.recv_timeout(std::time::Duration::from_secs(1));
1027 let sorted = producer.compute_and_publish();
1028 let block = reader.join().unwrap();
1029
1030 let (count, actual) = result.expect("trie updates must not wait for hashed state");
1031 assert_eq!(count, 1);
1032 assert!(Arc::ptr_eq(&updates, &actual));
1033 assert!(Arc::ptr_eq(&sorted, &block.hashed_state()));
1034 }
1035
1036 fn create_mock_state(
1037 test_block_builder: &mut TestBlockBuilder<EthPrimitives>,
1038 block_number: u64,
1039 parent_hash: B256,
1040 ) -> BlockState {
1041 BlockState::new(
1042 test_block_builder.get_executed_block_with_number(block_number, parent_hash),
1043 )
1044 }
1045
1046 fn create_mock_state_chain(
1047 test_block_builder: &mut TestBlockBuilder<EthPrimitives>,
1048 num_blocks: u64,
1049 ) -> Vec<BlockState> {
1050 let mut chain = Vec::with_capacity(num_blocks as usize);
1051 let mut parent_hash = B256::random();
1052 let mut parent_state: Option<BlockState> = None;
1053
1054 for i in 1..=num_blocks {
1055 let mut state = create_mock_state(test_block_builder, i, parent_hash);
1056 if let Some(parent) = parent_state {
1057 state.parent = Some(Arc::new(parent));
1058 }
1059 parent_hash = state.hash();
1060 parent_state = Some(state.clone());
1061 chain.push(state);
1062 }
1063
1064 chain
1065 }
1066
1067 #[test]
1068 fn test_in_memory_state_impl_state_by_hash() {
1069 let mut state_by_hash = B256Map::default();
1070 let number = rand::rng().random::<u64>();
1071 let mut test_block_builder: TestBlockBuilder = TestBlockBuilder::default();
1072 let state = Arc::new(create_mock_state(&mut test_block_builder, number, B256::random()));
1073 state_by_hash.insert(state.hash(), state.clone());
1074
1075 let in_memory_state = InMemoryState::new(state_by_hash, BTreeMap::new(), None);
1076
1077 assert_eq!(in_memory_state.state_by_hash(state.hash()), Some(state));
1078 assert_eq!(in_memory_state.state_by_hash(B256::random()), None);
1079 }
1080
1081 #[test]
1082 fn test_in_memory_state_impl_state_by_number() {
1083 let mut state_by_hash = B256Map::default();
1084 let mut hash_by_number = BTreeMap::new();
1085
1086 let number = rand::rng().random::<u64>();
1087 let mut test_block_builder: TestBlockBuilder = TestBlockBuilder::default();
1088 let state = Arc::new(create_mock_state(&mut test_block_builder, number, B256::random()));
1089 let hash = state.hash();
1090
1091 state_by_hash.insert(hash, state.clone());
1092 hash_by_number.insert(number, hash);
1093
1094 let in_memory_state = InMemoryState::new(state_by_hash, hash_by_number, None);
1095
1096 assert_eq!(in_memory_state.state_by_number(number), Some(state));
1097 assert_eq!(in_memory_state.state_by_number(number + 1), None);
1098 }
1099
1100 #[test]
1101 fn test_in_memory_state_impl_head_state() {
1102 let mut state_by_hash = B256Map::default();
1103 let mut hash_by_number = BTreeMap::new();
1104 let mut test_block_builder: TestBlockBuilder = TestBlockBuilder::default();
1105 let state1 = Arc::new(create_mock_state(&mut test_block_builder, 1, B256::random()));
1106 let hash1 = state1.hash();
1107 let state2 = Arc::new(create_mock_state(&mut test_block_builder, 2, hash1));
1108 let hash2 = state2.hash();
1109 hash_by_number.insert(1, hash1);
1110 hash_by_number.insert(2, hash2);
1111 state_by_hash.insert(hash1, state1);
1112 state_by_hash.insert(hash2, state2);
1113
1114 let in_memory_state = InMemoryState::new(state_by_hash, hash_by_number, None);
1115 let head_state = in_memory_state.head_state().unwrap();
1116
1117 assert_eq!(head_state.hash(), hash2);
1118 assert_eq!(head_state.number(), 2);
1119 }
1120
1121 #[test]
1122 fn test_in_memory_state_impl_pending_state() {
1123 let pending_number = rand::rng().random::<u64>();
1124 let mut test_block_builder: TestBlockBuilder = TestBlockBuilder::default();
1125 let pending_state =
1126 create_mock_state(&mut test_block_builder, pending_number, B256::random());
1127 let pending_hash = pending_state.hash();
1128
1129 let in_memory_state =
1130 InMemoryState::new(B256Map::default(), BTreeMap::new(), Some(pending_state));
1131
1132 let result = in_memory_state.pending_state();
1133 assert!(result.is_some());
1134 let actual_pending_state = result.unwrap();
1135 assert_eq!(actual_pending_state.block.recovered_block().hash(), pending_hash);
1136 assert_eq!(actual_pending_state.block.recovered_block().number, pending_number);
1137 }
1138
1139 #[test]
1140 fn test_in_memory_state_impl_no_pending_state() {
1141 let in_memory_state: InMemoryState =
1142 InMemoryState::new(B256Map::default(), BTreeMap::new(), None);
1143
1144 assert_eq!(in_memory_state.pending_state(), None);
1145 }
1146
1147 #[test]
1148 fn test_state() {
1149 let number = rand::rng().random::<u64>();
1150 let mut test_block_builder: TestBlockBuilder = TestBlockBuilder::default();
1151 let block = test_block_builder.get_executed_block_with_number(number, B256::random());
1152
1153 let state = BlockState::new(block.clone());
1154
1155 assert_eq!(state.block(), block);
1156 assert_eq!(state.hash(), block.recovered_block().hash());
1157 assert_eq!(state.number(), number);
1158 assert_eq!(state.state_root(), block.recovered_block().state_root);
1159 }
1160
1161 #[test]
1162 fn test_state_receipts() {
1163 let receipts = vec![vec![Receipt::default()]];
1164 let mut test_block_builder: TestBlockBuilder = TestBlockBuilder::default();
1165 let block =
1166 test_block_builder.get_executed_block_with_receipts(receipts.clone(), B256::random());
1167
1168 let state = BlockState::new(block);
1169
1170 assert_eq!(state.receipts(), receipts.first().unwrap());
1171 }
1172
1173 #[test]
1174 fn test_in_memory_state_chain_update() {
1175 let state: CanonicalInMemoryState = CanonicalInMemoryState::empty();
1176 let mut test_block_builder: TestBlockBuilder = TestBlockBuilder::default();
1177 let block1 = test_block_builder.get_executed_block_with_number(0, B256::random());
1178 let block2 = test_block_builder.get_executed_block_with_number(0, B256::random());
1179 let chain = NewCanonicalChain::Commit { new: vec![block1.clone()] };
1180 state.update_chain(chain);
1181 assert_eq!(
1182 state.head_state().unwrap().block_ref().recovered_block().hash(),
1183 block1.recovered_block().hash()
1184 );
1185 assert_eq!(
1186 state.state_by_number(0).unwrap().block_ref().recovered_block().hash(),
1187 block1.recovered_block().hash()
1188 );
1189
1190 let chain = NewCanonicalChain::Reorg { new: vec![block2.clone()], old: vec![block1] };
1191 state.update_chain(chain);
1192 assert_eq!(
1193 state.head_state().unwrap().block_ref().recovered_block().hash(),
1194 block2.recovered_block().hash()
1195 );
1196 assert_eq!(
1197 state.state_by_number(0).unwrap().block_ref().recovered_block().hash(),
1198 block2.recovered_block().hash()
1199 );
1200
1201 assert_eq!(state.inner.in_memory_state.block_count(), 1);
1202 }
1203
1204 #[test]
1205 fn test_in_memory_state_set_pending_block() {
1206 let state: CanonicalInMemoryState = CanonicalInMemoryState::empty();
1207 let mut test_block_builder: TestBlockBuilder = TestBlockBuilder::default();
1208
1209 let block1 = test_block_builder.get_executed_block_with_number(0, B256::random());
1211
1212 let block2 =
1214 test_block_builder.get_executed_block_with_number(1, block1.recovered_block().hash());
1215
1216 let chain = NewCanonicalChain::Commit { new: vec![block1.clone(), block2.clone()] };
1218 state.update_chain(chain);
1219
1220 assert!(state.pending_state().is_none());
1222
1223 state.set_pending_block(block2.clone());
1225
1226 assert_eq!(
1228 state.pending_state().unwrap(),
1229 BlockState::with_parent(block2.clone(), Some(Arc::new(BlockState::new(block1))))
1230 );
1231
1232 assert_eq!(state.pending_block().unwrap(), block2.recovered_block().sealed_block().clone());
1234
1235 assert_eq!(
1237 state.pending_block_num_hash().unwrap(),
1238 BlockNumHash { number: 1, hash: block2.recovered_block().hash() }
1239 );
1240
1241 assert_eq!(state.pending_header().unwrap(), block2.recovered_block().header().clone());
1243
1244 assert_eq!(
1246 state.pending_sealed_header().unwrap(),
1247 block2.recovered_block().clone_sealed_header()
1248 );
1249
1250 let pending_block = state.pending_recovered_block().unwrap();
1252 assert!(Arc::ptr_eq(&pending_block, &block2.recovered_block));
1253
1254 let pending = state.pending_block_and_receipts().unwrap();
1256 assert!(Arc::ptr_eq(pending.block(), &block2.recovered_block));
1257 assert!(Arc::ptr_eq(pending.execution_output(), &block2.execution_output));
1258 }
1259
1260 #[test]
1261 fn test_canonical_in_memory_state_canonical_chain_empty() {
1262 let state: CanonicalInMemoryState = CanonicalInMemoryState::empty();
1263 assert!(state.canonical_chain().next().is_none());
1264 }
1265
1266 #[test]
1267 fn test_canonical_in_memory_state_canonical_chain_single_block() {
1268 let block = TestBlockBuilder::eth().get_executed_block_with_number(1, B256::random());
1269 let hash = block.recovered_block().hash();
1270 let mut blocks = B256Map::default();
1271 blocks.insert(hash, Arc::new(BlockState::new(block)));
1272 let mut numbers = BTreeMap::new();
1273 numbers.insert(1, hash);
1274
1275 let state = CanonicalInMemoryState::new(blocks, numbers, None, None, None);
1276 let chain: Vec<_> = state.canonical_chain().collect();
1277
1278 assert_eq!(chain.len(), 1);
1279 assert_eq!(chain[0].number(), 1);
1280 assert_eq!(chain[0].hash(), hash);
1281 }
1282
1283 #[test]
1284 fn test_canonical_in_memory_state_canonical_chain_multiple_blocks() {
1285 let mut parent_hash = B256::random();
1286 let mut block_builder = TestBlockBuilder::eth();
1287 let state: CanonicalInMemoryState = CanonicalInMemoryState::empty();
1288
1289 for i in 1..=3 {
1290 let block = block_builder.get_executed_block_with_number(i, parent_hash);
1291 let hash = block.recovered_block().hash();
1292 state.update_blocks(Some(block), None);
1293 parent_hash = hash;
1294 }
1295
1296 let chain: Vec<_> = state.canonical_chain().collect();
1297
1298 assert_eq!(chain.len(), 3);
1299 assert_eq!(chain[0].number(), 3);
1300 assert_eq!(chain[1].number(), 2);
1301 assert_eq!(chain[2].number(), 1);
1302 }
1303
1304 #[test]
1306 fn test_canonical_in_memory_state_canonical_chain_with_pending_block() {
1307 let mut parent_hash = B256::random();
1308 let mut block_builder = TestBlockBuilder::<EthPrimitives>::eth();
1309 let state: CanonicalInMemoryState = CanonicalInMemoryState::empty();
1310
1311 for i in 1..=2 {
1312 let block = block_builder.get_executed_block_with_number(i, parent_hash);
1313 let hash = block.recovered_block().hash();
1314 state.update_blocks(Some(block), None);
1315 parent_hash = hash;
1316 }
1317
1318 let pending_block = block_builder.get_executed_block_with_number(3, parent_hash);
1319 state.set_pending_block(pending_block);
1320 let chain: Vec<_> = state.canonical_chain().collect();
1321
1322 assert_eq!(chain.len(), 2);
1323 assert_eq!(chain[0].number(), 2);
1324 assert_eq!(chain[1].number(), 1);
1325 }
1326
1327 #[test]
1328 fn test_block_state_parent_blocks() {
1329 let mut test_block_builder: TestBlockBuilder = TestBlockBuilder::default();
1330 let chain = create_mock_state_chain(&mut test_block_builder, 4);
1331
1332 let parents: Vec<_> = chain[3].parent_state_chain().collect();
1333 assert_eq!(parents.len(), 3);
1334 assert_eq!(parents[0].block().recovered_block().number, 3);
1335 assert_eq!(parents[1].block().recovered_block().number, 2);
1336 assert_eq!(parents[2].block().recovered_block().number, 1);
1337
1338 let parents: Vec<_> = chain[2].parent_state_chain().collect();
1339 assert_eq!(parents.len(), 2);
1340 assert_eq!(parents[0].block().recovered_block().number, 2);
1341 assert_eq!(parents[1].block().recovered_block().number, 1);
1342
1343 assert_eq!(chain[0].parent_state_chain().count(), 0);
1344 }
1345
1346 #[test]
1347 fn test_block_state_single_block_state_chain() {
1348 let single_block_number = 1;
1349 let mut test_block_builder: TestBlockBuilder = TestBlockBuilder::default();
1350 let single_block =
1351 create_mock_state(&mut test_block_builder, single_block_number, B256::random());
1352 let single_block_hash = single_block.block().recovered_block().hash();
1353
1354 assert_eq!(single_block.parent_state_chain().count(), 0);
1355
1356 let block_state_chain = single_block.chain().collect::<Vec<_>>();
1357 assert_eq!(block_state_chain.len(), 1);
1358 assert_eq!(block_state_chain[0].block().recovered_block().number, single_block_number);
1359 assert_eq!(block_state_chain[0].block().recovered_block().hash(), single_block_hash);
1360 }
1361
1362 #[test]
1363 fn test_block_state_chain() {
1364 let mut test_block_builder: TestBlockBuilder = TestBlockBuilder::default();
1365 let chain = create_mock_state_chain(&mut test_block_builder, 3);
1366
1367 let block_state_chain = chain[2].chain().collect::<Vec<_>>();
1368 assert_eq!(block_state_chain.len(), 3);
1369 assert_eq!(block_state_chain[0].block().recovered_block().number, 3);
1370 assert_eq!(block_state_chain[1].block().recovered_block().number, 2);
1371 assert_eq!(block_state_chain[2].block().recovered_block().number, 1);
1372
1373 let block_state_chain = chain[1].chain().collect::<Vec<_>>();
1374 assert_eq!(block_state_chain.len(), 2);
1375 assert_eq!(block_state_chain[0].block().recovered_block().number, 2);
1376 assert_eq!(block_state_chain[1].block().recovered_block().number, 1);
1377
1378 let block_state_chain = chain[0].chain().collect::<Vec<_>>();
1379 assert_eq!(block_state_chain.len(), 1);
1380 assert_eq!(block_state_chain[0].block().recovered_block().number, 1);
1381 }
1382
1383 #[test]
1384 fn test_to_chain_notification() {
1385 let mut test_block_builder: TestBlockBuilder = TestBlockBuilder::default();
1387 let block0 = test_block_builder.get_executed_block_with_number(0, B256::random());
1388 let block1 =
1389 test_block_builder.get_executed_block_with_number(1, block0.recovered_block.hash());
1390 let block1a =
1391 test_block_builder.get_executed_block_with_number(1, block0.recovered_block.hash());
1392 let block2 =
1393 test_block_builder.get_executed_block_with_number(2, block1.recovered_block.hash());
1394 let block2a =
1395 test_block_builder.get_executed_block_with_number(2, block1.recovered_block.hash());
1396
1397 let chain_commit = NewCanonicalChain::Commit { new: vec![block0.clone(), block1.clone()] };
1399
1400 let mut expected_trie_data = BTreeMap::new();
1402 expected_trie_data.insert(0, block0.trie_data());
1403 expected_trie_data.insert(1, block1.trie_data());
1404
1405 let commit_execution_outcome = ExecutionOutcome {
1407 receipts: vec![vec![], vec![]],
1408 requests: vec![Requests::default(), Requests::default()],
1409 first_block: 0,
1410 ..Default::default()
1411 };
1412
1413 assert_eq!(
1414 chain_commit.to_chain_notification(),
1415 CanonStateNotification::Commit {
1416 new: Arc::new(Chain::new(
1417 vec![block0.recovered_block().clone(), block1.recovered_block().clone()],
1418 commit_execution_outcome,
1419 expected_trie_data,
1420 ))
1421 }
1422 );
1423
1424 let chain_reorg = NewCanonicalChain::Reorg {
1426 new: vec![block1a.clone(), block2a.clone()],
1427 old: vec![block1.clone(), block2.clone()],
1428 };
1429
1430 let mut old_trie_data = BTreeMap::new();
1432 old_trie_data.insert(1, block1.trie_data());
1433 old_trie_data.insert(2, block2.trie_data());
1434
1435 let mut new_trie_data = BTreeMap::new();
1437 new_trie_data.insert(1, block1a.trie_data());
1438 new_trie_data.insert(2, block2a.trie_data());
1439
1440 let reorg_execution_outcome = ExecutionOutcome {
1443 receipts: vec![vec![], vec![]],
1444 requests: vec![Requests::default(), Requests::default()],
1445 first_block: 1,
1446 ..Default::default()
1447 };
1448
1449 assert_eq!(
1450 chain_reorg.to_chain_notification(),
1451 CanonStateNotification::Reorg {
1452 old: Arc::new(Chain::new(
1453 vec![block1.recovered_block().clone(), block2.recovered_block().clone()],
1454 reorg_execution_outcome.clone(),
1455 old_trie_data,
1456 )),
1457 new: Arc::new(Chain::new(
1458 vec![block1a.recovered_block().clone(), block2a.recovered_block().clone()],
1459 reorg_execution_outcome,
1460 new_trie_data,
1461 ))
1462 }
1463 );
1464 }
1465
1466 #[test]
1467 fn test_to_chain_notification_carries_prepared_bal() {
1468 let mut test_block_builder: TestBlockBuilder = TestBlockBuilder::default();
1469 let block0 = test_block_builder.get_executed_block_with_number(0, B256::random());
1470 let block1 = test_block_builder
1471 .get_executed_block_with_number(1, block0.recovered_block.hash())
1472 .with_bal(Some(Arc::new(DecodedRevmBal::new(
1473 Arc::new(revm::state::bal::Bal::default()),
1474 Bytes::from_static(&[0xc0]),
1475 ))));
1476
1477 let chain = NewCanonicalChain::Commit { new: vec![block0, block1.clone()] };
1478 let CanonStateNotification::Commit { new } = chain.to_chain_notification() else {
1479 panic!("expected a commit notification")
1480 };
1481
1482 assert_eq!(new.bals().len(), 1);
1484 assert_eq!(new.bal_at(1), block1.bal());
1485
1486 let mut blocks_and_bals = new.blocks_and_bals();
1487 let (block, bal) = blocks_and_bals.next().expect("block with BAL");
1488 assert_eq!(block.hash(), block1.recovered_block.hash());
1489 assert_eq!(Some(bal), block1.bal());
1490 assert!(blocks_and_bals.next().is_none());
1491 }
1492}