1use 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
19pub type CanonStateNotifications<N = reth_ethereum_primitives::EthPrimitives> =
21 broadcast::Receiver<CanonStateNotification<N>>;
22
23pub type CanonStateNotificationSender<N = reth_ethereum_primitives::EthPrimitives> =
25 broadcast::Sender<CanonStateNotification<N>>;
26
27pub trait CanonStateSubscriptions: Send + Sync {
29 type Primitives: NodePrimitives;
31
32 fn subscribe_to_canonical_state(&self) -> CanonStateNotifications<Self::Primitives>;
36
37 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#[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#[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 Commit {
93 new: Arc<Chain<N>>,
95 },
96 Reorg {
103 old: Arc<Chain<N>>,
105 new: Arc<Chain<N>>,
110 },
111}
112
113impl<N: NodePrimitives> CanonStateNotification<N> {
114 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 pub fn committed(&self) -> Arc<Chain<N>> {
124 match self {
125 Self::Commit { new } | Self::Reorg { new, .. } => new.clone(),
126 }
127 }
128
129 pub fn tip(&self) -> &RecoveredBlock<N::Block> {
138 match self {
139 Self::Commit { new } | Self::Reorg { new, .. } => new.tip(),
140 }
141 }
142
143 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 pub fn block_receipts(&self) -> Vec<(BlockReceipts<N::Receipt>, bool)> {
165 let mut receipts = Vec::new();
166
167 if let Some(old) = self.reverted() {
169 receipts
170 .extend(old.receipts_with_attachment().into_iter().map(|receipt| (receipt, true)));
171 }
172 receipts.extend(
174 self.committed().receipts_with_attachment().into_iter().map(|receipt| (receipt, false)),
175 );
176 receipts
177 }
178}
179
180#[derive(Debug, Deref, DerefMut)]
182pub struct ForkChoiceNotifications<T = alloy_consensus::Header>(
183 pub watch::Receiver<Option<SealedHeader<T>>>,
184);
185
186pub trait ForkChoiceSubscriptions: Send + Sync {
189 type Header: Clone + Send + Sync + 'static;
191
192 fn subscribe_safe_block(&self) -> ForkChoiceNotifications<Self::Header>;
194
195 fn subscribe_finalized_block(&self) -> ForkChoiceNotifications<Self::Header>;
197
198 fn safe_block_stream(&self) -> ForkChoiceStream<SealedHeader<Self::Header>> {
200 ForkChoiceStream::new(self.subscribe_safe_block().0)
201 }
202
203 fn finalized_block_stream(&self) -> ForkChoiceStream<SealedHeader<Self::Header>> {
205 ForkChoiceStream::new(self.subscribe_finalized_block().0)
206 }
207}
208
209#[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 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
238pub type ForkChoiceStream<T> = WatchValueStream<T>;
240
241#[derive(Debug, Deref, DerefMut)]
243pub struct PersistedBlockNotifications(pub watch::Receiver<Option<BlockNumHash>>);
244
245pub trait PersistedBlockSubscriptions: Send + Sync {
247 fn subscribe_persisted_block(&self) -> PersistedBlockNotifications;
249
250 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 let notification = CanonStateNotification::Commit { new: chain.clone() };
288
289 assert_eq!(notification.committed(), chain);
291
292 assert!(notification.reverted().is_none());
294
295 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 let notification =
331 CanonStateNotification::Reorg { old: old_chain.clone(), new: new_chain.clone() };
332
333 assert_eq!(notification.reverted(), Some(old_chain));
335
336 assert_eq!(notification.committed(), new_chain);
338
339 assert_eq!(*notification.tip(), block3);
341 }
342
343 #[test]
344 fn test_block_receipts_commit() {
345 let mut body = BlockBody::<TransactionSigned>::default();
347
348 let block1_hash = B256::new([0x01; 32]);
350 let block2_hash = B256::new([0x02; 32]);
351
352 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 let mut block1 = block.clone();
365 block1.set_block_number(1);
366 block1.set_hash(block1_hash);
367
368 let mut block2 = block;
370 block2.set_block_number(2);
371 block2.set_hash(block2_hash);
372
373 let receipt1 = Receipt {
375 tx_type: TxType::Legacy,
376 cumulative_gas_used: 12345,
377 logs: vec![],
378 success: true,
379 };
380
381 let receipts = vec![vec![receipt1.clone()]];
383
384 let execution_outcome = ExecutionOutcome { receipts, ..Default::default() };
386
387 let new_chain: Arc<Chain> = Arc::new(Chain::new(
389 vec![block1.clone(), block2.clone()],
390 execution_outcome,
391 BTreeMap::new(),
392 ));
393
394 let notification = CanonStateNotification::Commit { new: new_chain };
396
397 let block_receipts = notification.block_receipts();
399
400 assert_eq!(block_receipts.len(), 1);
402
403 assert_eq!(
405 block_receipts[0].0,
406 BlockReceipts {
407 block: block1.num_hash(),
408 timestamp: block1.timestamp,
409 tx_receipts: vec![(
410 b256!("0x20b5378c6fe992c118b557d2f8e8bbe0b7567f6fe5483a8f0f1c51e93a9d91ab"),
412 receipt1
413 )]
414 }
415 );
416
417 assert!(!block_receipts[0].1);
419 }
420
421 #[test]
422 fn test_block_receipts_reorg() {
423 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 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 let old_chain: Arc<Chain> =
450 Arc::new(Chain::new(vec![old_block1.clone()], old_execution_outcome, BTreeMap::new()));
451
452 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 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 let new_chain =
479 Arc::new(Chain::new(vec![new_block1.clone()], new_execution_outcome, BTreeMap::new()));
480
481 let notification = CanonStateNotification::Reorg { old: old_chain, new: new_chain };
483
484 let block_receipts = notification.block_receipts();
486
487 assert_eq!(block_receipts.len(), 2);
489
490 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 b256!("0x20b5378c6fe992c118b557d2f8e8bbe0b7567f6fe5483a8f0f1c51e93a9d91ab"),
499 old_receipt
500 )]
501 }
502 );
503 assert!(block_receipts[0].1);
505
506 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 b256!("0x20b5378c6fe992c118b557d2f8e8bbe0b7567f6fe5483a8f0f1c51e93a9d91ab"),
516 new_receipt
517 )]
518 }
519 );
520 assert!(!block_receipts[1].1);
522 }
523}