Skip to main content

reth_chain_state/
notifications.rs

1//! Canonical chain state notification trait and types.
2
3use alloy_eips::BlockNumHash;
4use derive_more::{Deref, DerefMut};
5use reth_execution_types::{BlockReceipts, Chain};
6use reth_primitives_traits::{NodePrimitives, RecoveredBlock, SealedHeader};
7use std::{
8    pin::Pin,
9    sync::Arc,
10    task::{ready, Context, Poll},
11};
12use tokio::sync::{broadcast, watch};
13use tokio_stream::{
14    wrappers::{BroadcastStream, WatchStream},
15    Stream,
16};
17use tracing::debug;
18
19/// Type alias for a receiver that receives [`CanonStateNotification`]
20pub type CanonStateNotifications<N = reth_ethereum_primitives::EthPrimitives> =
21    broadcast::Receiver<CanonStateNotification<N>>;
22
23/// Type alias for a sender that sends [`CanonStateNotification`]
24pub type CanonStateNotificationSender<N = reth_ethereum_primitives::EthPrimitives> =
25    broadcast::Sender<CanonStateNotification<N>>;
26
27/// A type that allows to register chain related event subscriptions.
28pub trait CanonStateSubscriptions: Send + Sync {
29    /// The node primitive types.
30    type Primitives: NodePrimitives;
31
32    /// Get notified when a new canonical chain was imported.
33    ///
34    /// A canonical chain be one or more blocks, a reorg or a revert.
35    fn subscribe_to_canonical_state(&self) -> CanonStateNotifications<Self::Primitives>;
36
37    /// Convenience method to get a stream of [`CanonStateNotification`].
38    fn canonical_state_stream(&self) -> CanonStateNotificationStream<Self::Primitives> {
39        CanonStateNotificationStream {
40            st: BroadcastStream::new(self.subscribe_to_canonical_state()),
41        }
42    }
43}
44
45impl<T: CanonStateSubscriptions> CanonStateSubscriptions for &T {
46    type Primitives = T::Primitives;
47
48    fn subscribe_to_canonical_state(&self) -> CanonStateNotifications<Self::Primitives> {
49        (*self).subscribe_to_canonical_state()
50    }
51
52    fn canonical_state_stream(&self) -> CanonStateNotificationStream<Self::Primitives> {
53        (*self).canonical_state_stream()
54    }
55}
56
57/// A Stream of [`CanonStateNotification`].
58#[derive(Debug)]
59#[pin_project::pin_project]
60pub struct CanonStateNotificationStream<N: NodePrimitives = reth_ethereum_primitives::EthPrimitives>
61{
62    #[pin]
63    st: BroadcastStream<CanonStateNotification<N>>,
64}
65
66impl<N: NodePrimitives> Stream for CanonStateNotificationStream<N> {
67    type Item = CanonStateNotification<N>;
68
69    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
70        loop {
71            return match ready!(self.as_mut().project().st.poll_next(cx)) {
72                Some(Ok(notification)) => Poll::Ready(Some(notification)),
73                Some(Err(err)) => {
74                    debug!(%err, "canonical state notification stream lagging behind");
75                    continue
76                }
77                None => Poll::Ready(None),
78            }
79        }
80    }
81}
82
83/// A notification that is sent when a new block is imported, or an old block is reverted.
84///
85/// The notification contains at least one [`Chain`] with the imported segment. If some blocks were
86/// reverted (e.g. during a reorg), the old chain is also returned.
87#[derive(Clone, Debug, PartialEq, Eq)]
88#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
89#[cfg_attr(feature = "serde", serde(bound = ""))]
90pub enum CanonStateNotification<N: NodePrimitives = reth_ethereum_primitives::EthPrimitives> {
91    /// The canonical chain was extended.
92    Commit {
93        /// The newly added chain segment.
94        new: Arc<Chain<N>>,
95    },
96    /// A chain segment was reverted or reorged.
97    ///
98    /// - In the case of a reorg, the reverted blocks are present in `old`, and the new blocks are
99    ///   present in `new`.
100    /// - In the case of a revert, the reverted blocks are present in `old`, and `new` is an empty
101    ///   chain segment.
102    Reorg {
103        /// The chain segment that was reverted.
104        old: Arc<Chain<N>>,
105        /// The chain segment that was added on top of the canonical chain, minus the reverted
106        /// blocks.
107        ///
108        /// In the case of a revert, not a reorg, this chain segment is empty.
109        new: Arc<Chain<N>>,
110    },
111}
112
113impl<N: NodePrimitives> CanonStateNotification<N> {
114    /// Get the chain segment that was reverted, if any.
115    pub fn reverted(&self) -> Option<Arc<Chain<N>>> {
116        match self {
117            Self::Commit { .. } => None,
118            Self::Reorg { old, .. } => Some(old.clone()),
119        }
120    }
121
122    /// Get the newly imported chain segment, if any.
123    pub fn committed(&self) -> Arc<Chain<N>> {
124        match self {
125            Self::Commit { new } | Self::Reorg { new, .. } => new.clone(),
126        }
127    }
128
129    /// Gets the new tip of the chain.
130    ///
131    /// Returns the new tip for [`Self::Reorg`] and [`Self::Commit`] variants which commit at least
132    /// 1 new block.
133    ///
134    /// # Panics
135    ///
136    /// If chain doesn't have any blocks.
137    pub fn tip(&self) -> &RecoveredBlock<N::Block> {
138        match self {
139            Self::Commit { new } | Self::Reorg { new, .. } => new.tip(),
140        }
141    }
142
143    /// Gets the new tip of the chain.
144    ///
145    /// If the chain has no blocks, it returns `None`. Otherwise, it returns the new tip for
146    /// [`Self::Reorg`] and [`Self::Commit`] variants.
147    pub fn tip_checked(&self) -> Option<&RecoveredBlock<N::Block>> {
148        match self {
149            Self::Commit { new } | Self::Reorg { new, .. } => {
150                if new.is_empty() {
151                    None
152                } else {
153                    Some(new.tip())
154                }
155            }
156        }
157    }
158
159    /// Get receipts in the reverted and newly imported chain segments with their corresponding
160    /// block numbers and transaction hashes.
161    ///
162    /// The boolean in the tuple (2nd element) denotes whether the receipt was from the reverted
163    /// chain segment.
164    pub fn block_receipts(&self) -> Vec<(BlockReceipts<N::Receipt>, bool)> {
165        let mut receipts = Vec::new();
166
167        // get old receipts
168        if let Some(old) = self.reverted() {
169            receipts
170                .extend(old.receipts_with_attachment().into_iter().map(|receipt| (receipt, true)));
171        }
172        // get new receipts
173        receipts.extend(
174            self.committed().receipts_with_attachment().into_iter().map(|receipt| (receipt, false)),
175        );
176        receipts
177    }
178}
179
180/// Wrapper around a broadcast receiver that receives fork choice notifications.
181#[derive(Debug, Deref, DerefMut)]
182pub struct ForkChoiceNotifications<T = alloy_consensus::Header>(
183    pub watch::Receiver<Option<SealedHeader<T>>>,
184);
185
186/// A trait that allows to register to fork choice related events
187/// and get notified when a new fork choice is available.
188pub trait ForkChoiceSubscriptions: Send + Sync {
189    /// Block Header type.
190    type Header: Clone + Send + Sync + 'static;
191
192    /// Get notified when a new safe block of the chain is selected.
193    fn subscribe_safe_block(&self) -> ForkChoiceNotifications<Self::Header>;
194
195    /// Get notified when a new finalized block of the chain is selected.
196    fn subscribe_finalized_block(&self) -> ForkChoiceNotifications<Self::Header>;
197
198    /// Convenience method to get a stream of the new safe blocks of the chain.
199    fn safe_block_stream(&self) -> ForkChoiceStream<SealedHeader<Self::Header>> {
200        ForkChoiceStream::new(self.subscribe_safe_block().0)
201    }
202
203    /// Convenience method to get a stream of the new finalized blocks of the chain.
204    fn finalized_block_stream(&self) -> ForkChoiceStream<SealedHeader<Self::Header>> {
205        ForkChoiceStream::new(self.subscribe_finalized_block().0)
206    }
207}
208
209/// A stream that yields values from a `watch::Receiver<Option<T>>`, filtering out `None` values.
210#[derive(Debug)]
211#[pin_project::pin_project]
212pub struct WatchValueStream<T> {
213    #[pin]
214    st: WatchStream<Option<T>>,
215}
216
217impl<T: Clone + Sync + Send + 'static> WatchValueStream<T> {
218    /// Creates a new [`WatchValueStream`]
219    pub fn new(rx: watch::Receiver<Option<T>>) -> Self {
220        Self { st: WatchStream::from_changes(rx) }
221    }
222}
223
224impl<T: Clone + Sync + Send + 'static> Stream for WatchValueStream<T> {
225    type Item = T;
226
227    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
228        loop {
229            match ready!(self.as_mut().project().st.poll_next(cx)) {
230                Some(Some(notification)) => return Poll::Ready(Some(notification)),
231                Some(None) => {}
232                None => return Poll::Ready(None),
233            }
234        }
235    }
236}
237
238/// Alias for [`WatchValueStream`] for fork choice watch channels.
239pub type ForkChoiceStream<T> = WatchValueStream<T>;
240
241/// Wrapper around a watch receiver that receives persisted block notifications.
242#[derive(Debug, Deref, DerefMut)]
243pub struct PersistedBlockNotifications(pub watch::Receiver<Option<BlockNumHash>>);
244
245/// A trait that allows subscribing to persisted block events.
246pub trait PersistedBlockSubscriptions: Send + Sync {
247    /// Get notified when a new block is persisted to disk.
248    fn subscribe_persisted_block(&self) -> PersistedBlockNotifications;
249
250    /// Convenience method to get a stream of the persisted blocks.
251    fn persisted_block_stream(&self) -> WatchValueStream<BlockNumHash> {
252        WatchValueStream::new(self.subscribe_persisted_block().0)
253    }
254}
255
256#[cfg(test)]
257mod tests {
258    use super::*;
259    use alloy_consensus::{BlockBody, SignableTransaction, TxLegacy};
260    use alloy_primitives::{b256, Signature, B256};
261    use reth_ethereum_primitives::{Receipt, TransactionSigned, TxType};
262    use reth_execution_types::ExecutionOutcome;
263    use reth_primitives_traits::SealedBlock;
264    use std::collections::BTreeMap;
265
266    #[test]
267    fn test_commit_notification() {
268        let block: RecoveredBlock<reth_ethereum_primitives::Block> = Default::default();
269        let block1_hash = B256::new([0x01; 32]);
270        let block2_hash = B256::new([0x02; 32]);
271
272        let mut block1 = block.clone();
273        block1.set_block_number(1);
274        block1.set_hash(block1_hash);
275
276        let mut block2 = block;
277        block2.set_block_number(2);
278        block2.set_hash(block2_hash);
279
280        let chain: Arc<Chain> = Arc::new(Chain::new(
281            vec![block1.clone(), block2.clone()],
282            ExecutionOutcome::default(),
283            BTreeMap::new(),
284        ));
285
286        // Create a commit notification
287        let notification = CanonStateNotification::Commit { new: chain.clone() };
288
289        // Test that `committed` returns the correct chain
290        assert_eq!(notification.committed(), chain);
291
292        // Test that `reverted` returns None for `Commit`
293        assert!(notification.reverted().is_none());
294
295        // Test that `tip` returns the correct block
296        assert_eq!(*notification.tip(), block2);
297    }
298
299    #[test]
300    fn test_reorg_notification() {
301        let block: RecoveredBlock<reth_ethereum_primitives::Block> = Default::default();
302        let block1_hash = B256::new([0x01; 32]);
303        let block2_hash = B256::new([0x02; 32]);
304        let block3_hash = B256::new([0x03; 32]);
305
306        let mut block1 = block.clone();
307        block1.set_block_number(1);
308        block1.set_hash(block1_hash);
309
310        let mut block2 = block.clone();
311        block2.set_block_number(2);
312        block2.set_hash(block2_hash);
313
314        let mut block3 = block;
315        block3.set_block_number(3);
316        block3.set_hash(block3_hash);
317
318        let old_chain: Arc<Chain> = Arc::new(Chain::new(
319            vec![block1.clone()],
320            ExecutionOutcome::default(),
321            BTreeMap::new(),
322        ));
323        let new_chain = Arc::new(Chain::new(
324            vec![block2.clone(), block3.clone()],
325            ExecutionOutcome::default(),
326            BTreeMap::new(),
327        ));
328
329        // Create a reorg notification
330        let notification =
331            CanonStateNotification::Reorg { old: old_chain.clone(), new: new_chain.clone() };
332
333        // Test that `reverted` returns the old chain
334        assert_eq!(notification.reverted(), Some(old_chain));
335
336        // Test that `committed` returns the new chain
337        assert_eq!(notification.committed(), new_chain);
338
339        // Test that `tip` returns the tip of the new chain (last block in the new chain)
340        assert_eq!(*notification.tip(), block3);
341    }
342
343    #[test]
344    fn test_block_receipts_commit() {
345        // Create a default block instance for use in block definitions.
346        let mut body = BlockBody::<TransactionSigned>::default();
347
348        // Define unique hashes for two blocks to differentiate them in the chain.
349        let block1_hash = B256::new([0x01; 32]);
350        let block2_hash = B256::new([0x02; 32]);
351
352        // Create a default transaction to include in block1's transactions.
353        let tx = TxLegacy::default().into_signed(Signature::test_signature()).into();
354        body.transactions.push(tx);
355
356        let block = SealedBlock::<alloy_consensus::Block<TransactionSigned>>::from_sealed_parts(
357            SealedHeader::seal_slow(alloy_consensus::Header::default()),
358            body,
359        )
360        .try_recover()
361        .unwrap();
362
363        // Create a clone of the default block and customize it to act as block1.
364        let mut block1 = block.clone();
365        block1.set_block_number(1);
366        block1.set_hash(block1_hash);
367
368        // Clone the default block and customize it to act as block2.
369        let mut block2 = block;
370        block2.set_block_number(2);
371        block2.set_hash(block2_hash);
372
373        // Create a receipt for the transaction in block1.
374        let receipt1 = Receipt {
375            tx_type: TxType::Legacy,
376            cumulative_gas_used: 12345,
377            logs: vec![],
378            success: true,
379        };
380
381        // Wrap the receipt in a `Receipts` structure, as expected in the `ExecutionOutcome`.
382        let receipts = vec![vec![receipt1.clone()]];
383
384        // Define an `ExecutionOutcome` with the created receipts.
385        let execution_outcome = ExecutionOutcome { receipts, ..Default::default() };
386
387        // Create a new chain segment with `block1` and `block2` and the execution outcome.
388        let new_chain: Arc<Chain> = Arc::new(Chain::new(
389            vec![block1.clone(), block2.clone()],
390            execution_outcome,
391            BTreeMap::new(),
392        ));
393
394        // Create a commit notification containing the new chain segment.
395        let notification = CanonStateNotification::Commit { new: new_chain };
396
397        // Call `block_receipts` on the commit notification to retrieve block receipts.
398        let block_receipts = notification.block_receipts();
399
400        // Assert that only one receipt entry exists in the `block_receipts` list.
401        assert_eq!(block_receipts.len(), 1);
402
403        // Verify that the first entry matches block1's hash and transaction receipt.
404        assert_eq!(
405            block_receipts[0].0,
406            BlockReceipts {
407                block: block1.num_hash(),
408                timestamp: block1.timestamp,
409                tx_receipts: vec![(
410                    // Transaction hash of a Transaction::default()
411                    b256!("0x20b5378c6fe992c118b557d2f8e8bbe0b7567f6fe5483a8f0f1c51e93a9d91ab"),
412                    receipt1
413                )]
414            }
415        );
416
417        // Assert that the receipt is from the committed segment (not reverted).
418        assert!(!block_receipts[0].1);
419    }
420
421    #[test]
422    fn test_block_receipts_reorg() {
423        // Define block1 for the old chain segment, which will be reverted.
424        let mut body = BlockBody::<TransactionSigned>::default();
425        body.transactions.push(TxLegacy::default().into_signed(Signature::test_signature()).into());
426        let mut old_block1 =
427            SealedBlock::<alloy_consensus::Block<TransactionSigned>>::from_sealed_parts(
428                SealedHeader::seal_slow(alloy_consensus::Header::default()),
429                body,
430            )
431            .try_recover()
432            .unwrap();
433        old_block1.set_block_number(1);
434        old_block1.set_hash(B256::new([0x01; 32]));
435
436        // Create a receipt for a transaction in the reverted block.
437        let old_receipt = Receipt {
438            tx_type: TxType::Legacy,
439            cumulative_gas_used: 54321,
440            logs: vec![],
441            success: false,
442        };
443        let old_receipts = vec![vec![old_receipt.clone()]];
444
445        let old_execution_outcome =
446            ExecutionOutcome { receipts: old_receipts, ..Default::default() };
447
448        // Create an old chain segment to be reverted, containing `old_block1`.
449        let old_chain: Arc<Chain> =
450            Arc::new(Chain::new(vec![old_block1.clone()], old_execution_outcome, BTreeMap::new()));
451
452        // Define block2 for the new chain segment, which will be committed.
453        let mut body = BlockBody::<TransactionSigned>::default();
454        body.transactions.push(TxLegacy::default().into_signed(Signature::test_signature()).into());
455        let mut new_block1 =
456            SealedBlock::<alloy_consensus::Block<TransactionSigned>>::from_sealed_parts(
457                SealedHeader::seal_slow(alloy_consensus::Header::default()),
458                body,
459            )
460            .try_recover()
461            .unwrap();
462        new_block1.set_block_number(2);
463        new_block1.set_hash(B256::new([0x02; 32]));
464
465        // Create a receipt for a transaction in the new committed block.
466        let new_receipt = Receipt {
467            tx_type: TxType::Legacy,
468            cumulative_gas_used: 12345,
469            logs: vec![],
470            success: true,
471        };
472        let new_receipts = vec![vec![new_receipt.clone()]];
473
474        let new_execution_outcome =
475            ExecutionOutcome { receipts: new_receipts, ..Default::default() };
476
477        // Create a new chain segment to be committed, containing `new_block1`.
478        let new_chain =
479            Arc::new(Chain::new(vec![new_block1.clone()], new_execution_outcome, BTreeMap::new()));
480
481        // Create a reorg notification with both reverted (old) and committed (new) chain segments.
482        let notification = CanonStateNotification::Reorg { old: old_chain, new: new_chain };
483
484        // Retrieve receipts from both old (reverted) and new (committed) segments.
485        let block_receipts = notification.block_receipts();
486
487        // Assert there are two receipt entries, one from each chain segment.
488        assert_eq!(block_receipts.len(), 2);
489
490        // Verify that the first entry matches old_block1 and its receipt from the reverted segment.
491        assert_eq!(
492            block_receipts[0].0,
493            BlockReceipts {
494                block: old_block1.num_hash(),
495                timestamp: old_block1.timestamp,
496                tx_receipts: vec![(
497                    // Transaction hash of a Transaction::default()
498                    b256!("0x20b5378c6fe992c118b557d2f8e8bbe0b7567f6fe5483a8f0f1c51e93a9d91ab"),
499                    old_receipt
500                )]
501            }
502        );
503        // Confirm this is from the reverted segment.
504        assert!(block_receipts[0].1);
505
506        // Verify that the second entry matches new_block1 and its receipt from the committed
507        // segment.
508        assert_eq!(
509            block_receipts[1].0,
510            BlockReceipts {
511                block: new_block1.num_hash(),
512                timestamp: new_block1.timestamp,
513                tx_receipts: vec![(
514                    // Transaction hash of a Transaction::default()
515                    b256!("0x20b5378c6fe992c118b557d2f8e8bbe0b7567f6fe5483a8f0f1c51e93a9d91ab"),
516                    new_receipt
517                )]
518            }
519        );
520        // Confirm this is from the committed segment.
521        assert!(!block_receipts[1].1);
522    }
523}