1use std::{collections::BTreeMap, pin::Pin, sync::Arc, time::Duration};
5
6use async_trait::async_trait;
7use bytes::Bytes;
8use consensus_config::AuthorityIndex;
9use consensus_types::block::{BlockRef, Round};
10use futures::{Stream, StreamExt, ready, stream, task};
11use mysten_metrics::spawn_monitored_task;
12use parking_lot::RwLock;
13use sui_macros::fail_point_async;
14use tap::TapFallible;
15use tokio::sync::broadcast;
16use tokio_util::sync::ReusableBoxFuture;
17use tracing::{debug, info, warn};
18
19use crate::{
20 block::{BlockAPI as _, ExtendedBlock, GENESIS_ROUND, SignedBlock, VerifiedBlock},
21 block_sync_service::BlockSyncService,
22 block_verifier::BlockVerifier,
23 commit::{CommitRange, TrustedCommit},
24 commit_vote_monitor::{CommitVoteMonitor, is_commit_lagging},
25 context::Context,
26 core_thread::CoreThreadDispatcher,
27 dag_state::DagState,
28 error::{ConsensusError, ConsensusResult},
29 network::{
30 BlockStream, ExtendedSerializedBlock, PeerId, SerializedBlockForm, ValidatorNetworkService,
31 },
32 round_tracker::RoundTracker,
33 synchronizer::SynchronizerHandle,
34 task::spawn_blocking,
35 transaction_vote_tracker::TransactionVoteTracker,
36};
37
38pub(crate) struct AuthorityService<C: CoreThreadDispatcher> {
40 context: Arc<Context>,
41 commit_vote_monitor: Arc<CommitVoteMonitor>,
42 block_verifier: Arc<dyn BlockVerifier>,
43 synchronizer: Arc<SynchronizerHandle>,
44 core_dispatcher: Arc<C>,
45 rx_block_broadcast: broadcast::Receiver<ExtendedBlock>,
46 subscription_counter: Arc<SubscriptionCounter>,
47 transaction_vote_tracker: TransactionVoteTracker,
48 dag_state: Arc<RwLock<DagState>>,
49 round_tracker: Arc<RwLock<RoundTracker>>,
50 block_sync_service: Arc<BlockSyncService>,
51}
52
53impl<C: CoreThreadDispatcher> AuthorityService<C> {
54 pub(crate) fn new(
55 context: Arc<Context>,
56 block_verifier: Arc<dyn BlockVerifier>,
57 commit_vote_monitor: Arc<CommitVoteMonitor>,
58 round_tracker: Arc<RwLock<RoundTracker>>,
59 synchronizer: Arc<SynchronizerHandle>,
60 core_dispatcher: Arc<C>,
61 rx_block_broadcast: broadcast::Receiver<ExtendedBlock>,
62 transaction_vote_tracker: TransactionVoteTracker,
63 dag_state: Arc<RwLock<DagState>>,
64 block_sync_service: Arc<BlockSyncService>,
65 ) -> Self {
66 let subscription_counter = Arc::new(SubscriptionCounter::new(context.clone()));
67 Self {
68 context,
69 block_verifier,
70 commit_vote_monitor,
71 synchronizer,
72 core_dispatcher,
73 rx_block_broadcast,
74 subscription_counter,
75 transaction_vote_tracker,
76 dag_state,
77 round_tracker,
78 block_sync_service,
79 }
80 }
81
82 fn parse_excluded_ancestors(
84 &self,
85 peer: AuthorityIndex,
86 block: &VerifiedBlock,
87 mut excluded_ancestors: Vec<Vec<u8>>,
88 ) -> ConsensusResult<Vec<BlockRef>> {
89 let peer_hostname = &self.context.committee.authority(peer).hostname;
90
91 let excluded_ancestors_limit = self.context.committee.size() * 2;
92 if excluded_ancestors.len() > excluded_ancestors_limit {
93 debug!(
94 "Dropping {} excluded ancestor(s) from {} {} due to size limit",
95 excluded_ancestors.len() - excluded_ancestors_limit,
96 peer,
97 peer_hostname,
98 );
99 excluded_ancestors.truncate(excluded_ancestors_limit);
100 }
101
102 let excluded_ancestors = excluded_ancestors
103 .into_iter()
104 .map(|serialized| {
105 let block_ref: BlockRef =
106 bcs::from_bytes(&serialized).map_err(ConsensusError::MalformedBlock)?;
107 if !self.context.committee.is_valid_index(block_ref.author) {
108 return Err(ConsensusError::InvalidAuthorityIndex {
109 loc: format!("excluded ancestor {}", block_ref),
110 index: block_ref.author,
111 max: self.context.committee.size() - 1,
112 });
113 }
114 if block_ref.round >= block.round() {
115 return Err(ConsensusError::InvalidAncestorRound {
116 ancestor: block_ref.round,
117 block: block.round(),
118 });
119 }
120 Ok(block_ref)
121 })
122 .collect::<ConsensusResult<Vec<BlockRef>>>()?;
123
124 for excluded_ancestor in &excluded_ancestors {
125 let excluded_ancestor_hostname = &self
126 .context
127 .committee
128 .authority(excluded_ancestor.author)
129 .hostname;
130 self.context
131 .metrics
132 .node_metrics
133 .network_excluded_ancestors_count_by_authority
134 .with_label_values(&[excluded_ancestor_hostname])
135 .inc();
136 }
137 self.context
138 .metrics
139 .node_metrics
140 .network_received_excluded_ancestors_from_authority
141 .with_label_values(&[peer_hostname])
142 .inc_by(excluded_ancestors.len() as u64);
143
144 Ok(excluded_ancestors)
145 }
146}
147
148#[async_trait]
149impl<C: CoreThreadDispatcher> ValidatorNetworkService for AuthorityService<C> {
150 async fn handle_send_block(
151 &self,
152 peer: AuthorityIndex,
153 serialized_block: ExtendedSerializedBlock,
154 ) -> ConsensusResult<()> {
155 fail_point_async!("consensus-rpc-response");
156
157 let peer_hostname = &self.context.committee.authority(peer).hostname;
158
159 let SerializedBlockForm::Full(serialized_bytes) = serialized_block.block else {
163 return Err(ConsensusError::UnexpectedBlockForm);
164 };
165 let signed_block: SignedBlock =
166 bcs::from_bytes(&serialized_bytes).map_err(ConsensusError::MalformedBlock)?;
167
168 if peer != signed_block.author() {
170 self.context
171 .metrics
172 .node_metrics
173 .invalid_blocks
174 .with_label_values(&[
175 peer_hostname.as_str(),
176 "handle_send_block",
177 "UnexpectedAuthority",
178 ])
179 .inc();
180 let e = ConsensusError::UnexpectedAuthority(signed_block.author(), peer);
181 info!("Block with wrong authority from {}: {}", peer, e);
182 return Err(e);
183 }
184
185 let block_verifier = self.block_verifier.clone();
187 let serialized = serialized_bytes.clone();
188 let (verified_block, reject_txn_votes) =
189 spawn_blocking(move || block_verifier.verify_and_vote(signed_block, serialized))
190 .await?
191 .tap_err(|e| {
192 self.context
193 .metrics
194 .node_metrics
195 .invalid_blocks
196 .with_label_values(&[peer_hostname.as_str(), "handle_send_block", e.name()])
197 .inc();
198 info!("Invalid block from {}: {}", peer, e);
199 })?;
200 let excluded_ancestors = self
201 .parse_excluded_ancestors(peer, &verified_block, serialized_block.excluded_ancestors)
202 .tap_err(|e| {
203 debug!("Failed to parse excluded ancestors from {peer} {peer_hostname}: {e}");
204 self.context
205 .metrics
206 .node_metrics
207 .invalid_blocks
208 .with_label_values(&[peer_hostname.as_str(), "handle_send_block", e.name()])
209 .inc();
210 })?;
211
212 let block_ref = verified_block.reference();
213 debug!("Received block {} via send block.", block_ref);
214
215 self.context
216 .metrics
217 .node_metrics
218 .verified_blocks
219 .with_label_values(&[peer_hostname])
220 .inc();
221
222 let now = self.context.clock.timestamp_utc_ms();
223 let forward_time_drift =
224 Duration::from_millis(verified_block.timestamp_ms().saturating_sub(now));
225
226 self.context
227 .metrics
228 .node_metrics
229 .block_timestamp_drift_ms
230 .with_label_values(&[peer_hostname.as_str(), "handle_send_block"])
231 .inc_by(forward_time_drift.as_millis() as u64);
232
233 self.commit_vote_monitor.observe_block(&verified_block);
236
237 self.round_tracker
239 .write()
240 .update_from_verified_block(&ExtendedBlock {
241 block: verified_block.clone(),
242 excluded_ancestors: excluded_ancestors.clone(),
243 });
244
245 let last_commit_index = self.dag_state.read().last_commit_index();
254 let quorum_commit_index = self.commit_vote_monitor.quorum_commit_index();
255 if is_commit_lagging(
256 self.context.as_ref(),
257 last_commit_index,
258 quorum_commit_index,
259 ) {
260 self.context
261 .metrics
262 .node_metrics
263 .rejected_blocks
264 .with_label_values(&["commit_lagging"])
265 .inc();
266 debug!(
267 "Block {:?} is rejected because last commit index is lagging quorum commit index too much ({} < {})",
268 block_ref, last_commit_index, quorum_commit_index,
269 );
270 return Err(ConsensusError::BlockRejected {
271 block_ref,
272 reason: format!(
273 "Last commit index is lagging quorum commit index too much ({} < {})",
274 last_commit_index, quorum_commit_index,
275 ),
276 });
277 }
278
279 if self.context.protocol_config.transaction_voting_enabled() {
282 self.transaction_vote_tracker
283 .add_voted_blocks(vec![(verified_block.clone(), reject_txn_votes)]);
284 }
285
286 let missing_ancestors = self
288 .core_dispatcher
289 .add_blocks(vec![verified_block.clone()])
290 .await
291 .map_err(|_| ConsensusError::Shutdown)?;
292
293 if !missing_ancestors.is_empty() {
295 self.context
296 .metrics
297 .node_metrics
298 .handler_received_block_missing_ancestors
299 .with_label_values(&[peer_hostname])
300 .inc_by(missing_ancestors.len() as u64);
301 let synchronizer = self.synchronizer.clone();
302 spawn_monitored_task!(async move {
303 if let Err(err) = synchronizer
308 .fetch_blocks(missing_ancestors, PeerId::Validator(peer))
309 .await
310 {
311 debug!("Failed to fetch missing ancestors via synchronizer: {err}");
312 }
313 });
314 }
315
316 let missing_excluded_ancestors = self
318 .core_dispatcher
319 .check_block_refs(excluded_ancestors)
320 .await
321 .map_err(|_| ConsensusError::Shutdown)?;
322 if !missing_excluded_ancestors.is_empty() {
323 self.context
324 .metrics
325 .node_metrics
326 .network_excluded_ancestors_sent_to_fetch
327 .with_label_values(&[peer_hostname])
328 .inc_by(missing_excluded_ancestors.len() as u64);
329
330 let synchronizer = self.synchronizer.clone();
331 spawn_monitored_task!(async move {
332 if let Err(err) = synchronizer
333 .fetch_blocks(missing_excluded_ancestors, PeerId::Validator(peer))
334 .await
335 {
336 debug!("Failed to fetch excluded ancestors via synchronizer: {err}");
337 }
338 });
339 }
340
341 Ok(())
342 }
343
344 async fn handle_subscribe_blocks(
345 &self,
346 peer: AuthorityIndex,
347 last_received: Round,
348 ) -> ConsensusResult<BlockStream> {
349 fail_point_async!("consensus-rpc-response");
350
351 let broadcast_rx = self.rx_block_broadcast.resubscribe();
356
357 let past_proposed_blocks = {
366 let dag_state = self.dag_state.read();
367
368 let mut proposed_blocks = dag_state
371 .get_cached_blocks(self.context.own_index, last_received.saturating_add(1));
372 if proposed_blocks.is_empty() {
373 let last_proposed_block = dag_state
374 .get_last_proposed_block()
375 .expect("Last proposed block should be returned on validators");
376 proposed_blocks = if last_proposed_block.round() > GENESIS_ROUND {
377 vec![last_proposed_block]
378 } else {
379 vec![]
380 };
381 }
382 stream::iter(
383 proposed_blocks
384 .into_iter()
385 .map(|block| ExtendedSerializedBlock {
386 block: SerializedBlockForm::Full(block.serialized().clone()),
387 excluded_ancestors: vec![],
388 }),
389 )
390 };
391
392 const MAX_BLOCKS_PER_POLL: usize = 1;
394 let broadcasted_blocks = BroadcastedBlockStream::new(
395 PeerId::Validator(peer),
396 broadcast_rx,
397 MAX_BLOCKS_PER_POLL,
398 self.subscription_counter.clone(),
399 );
400
401 Ok(Box::pin(past_proposed_blocks.chain(
403 broadcasted_blocks.flat_map(|items| {
404 debug_assert!(
405 items.len() <= MAX_BLOCKS_PER_POLL,
406 "Too many blocks received from broadcast"
407 );
408 stream::iter(items.into_iter().map(ExtendedSerializedBlock::from))
409 }),
410 )))
411 }
412
413 async fn handle_fetch_blocks(
429 &self,
430 _peer: AuthorityIndex,
431 block_refs: Vec<BlockRef>,
432 fetch_after_rounds: Vec<Round>,
433 fetch_missing_ancestors: bool,
434 ) -> ConsensusResult<Vec<Bytes>> {
435 fail_point_async!("consensus-rpc-response");
436
437 self.block_sync_service
439 .fetch_blocks(block_refs, fetch_after_rounds, fetch_missing_ancestors)
440 .await
441 }
442
443 async fn handle_fetch_commits(
444 &self,
445 _peer: AuthorityIndex,
446 commit_range: CommitRange,
447 ) -> ConsensusResult<(Vec<TrustedCommit>, Vec<VerifiedBlock>)> {
448 fail_point_async!("consensus-rpc-response");
449
450 self.block_sync_service.fetch_commits(commit_range).await
452 }
453
454 async fn handle_fetch_latest_blocks(
455 &self,
456 peer: AuthorityIndex,
457 authorities: Vec<AuthorityIndex>,
458 ) -> ConsensusResult<Vec<Bytes>> {
459 fail_point_async!("consensus-rpc-response");
460
461 self.block_sync_service
463 .fetch_latest_blocks(peer, authorities)
464 .await
465 }
466
467 async fn handle_get_latest_rounds(
468 &self,
469 _peer: AuthorityIndex,
470 ) -> ConsensusResult<(Vec<Round>, Vec<Round>)> {
471 fail_point_async!("consensus-rpc-response");
472
473 let highest_received_rounds = self.round_tracker.read().local_highest_received_rounds();
474
475 let blocks = self
476 .dag_state
477 .read()
478 .get_last_cached_block_per_authority(Round::MAX);
479 let highest_accepted_rounds = blocks
480 .into_iter()
481 .map(|(block, _)| block.round())
482 .collect::<Vec<_>>();
483
484 Ok((highest_received_rounds, highest_accepted_rounds))
485 }
486}
487
488struct Counter {
489 count: usize,
490 subscriptions_by_peer: BTreeMap<PeerId, usize>,
491}
492
493pub(crate) struct SubscriptionCounter {
495 context: Arc<Context>,
496 counter: parking_lot::Mutex<Counter>,
497}
498
499impl SubscriptionCounter {
500 pub(crate) fn new(context: Arc<Context>) -> Self {
501 for (_, authority) in context.committee.authorities() {
503 context
504 .metrics
505 .node_metrics
506 .subscribed_by
507 .with_label_values(&[authority.hostname.as_str()])
508 .set(0);
509 }
510
511 Self {
512 counter: parking_lot::Mutex::new(Counter {
513 count: 0,
514 subscriptions_by_peer: BTreeMap::new(),
515 }),
516 context,
517 }
518 }
519
520 fn increment(&self, peer: &PeerId) {
521 let mut counter = self.counter.lock();
522 counter.count += 1;
523 let peer_count = {
524 let count = counter
525 .subscriptions_by_peer
526 .entry(peer.clone())
527 .or_default();
528 *count += 1;
529 *count
530 };
531
532 match peer {
533 PeerId::Validator(authority) => {
534 let peer_hostname = &self.context.committee.authority(*authority).hostname;
535 self.context
536 .metrics
537 .node_metrics
538 .subscribed_by
539 .with_label_values(&[peer_hostname])
540 .set(1);
541 }
542 PeerId::Observer(_) => {
543 if peer_count == 1 {
545 self.context
546 .metrics
547 .node_metrics
548 .subscribed_by
549 .with_label_values(&["observer"])
550 .inc();
551 }
552 }
553 }
554 }
555
556 fn decrement(&self, peer: &PeerId) {
557 let mut counter = self.counter.lock();
558 counter.count = counter.count.saturating_sub(1);
559 let peer_count = counter
560 .subscriptions_by_peer
561 .entry(peer.clone())
562 .or_default();
563 *peer_count = peer_count.saturating_sub(1);
564
565 if *peer_count == 0 {
566 match peer {
567 PeerId::Validator(authority) => {
568 let peer_hostname = &self.context.committee.authority(*authority).hostname;
569 self.context
570 .metrics
571 .node_metrics
572 .subscribed_by
573 .with_label_values(&[peer_hostname])
574 .set(0);
575 }
576 PeerId::Observer(_) => {
577 self.context
578 .metrics
579 .node_metrics
580 .subscribed_by
581 .with_label_values(&["observer"])
582 .dec();
583 }
584 }
585 }
586 }
587}
588
589type BroadcastedBlockStream = BroadcastStream<ExtendedBlock>;
592
593pub(crate) struct BroadcastStream<T> {
596 peer: Option<PeerId>,
597 inner: ReusableBoxFuture<
599 'static,
600 (
601 Result<T, broadcast::error::RecvError>,
602 broadcast::Receiver<T>,
603 ),
604 >,
605 max_items_per_poll: usize,
607 subscription_counter: Option<Arc<SubscriptionCounter>>,
609}
610
611impl<T: 'static + Clone + Send> BroadcastStream<T> {
612 pub fn new(
613 peer: PeerId,
614 rx: broadcast::Receiver<T>,
615 max_items_per_poll: usize,
616 subscription_counter: Arc<SubscriptionCounter>,
617 ) -> Self {
618 assert!(max_items_per_poll > 0, "max_items_per_poll must be > 0");
619 subscription_counter.increment(&peer);
620 Self {
621 peer: Some(peer),
622 inner: ReusableBoxFuture::new(make_recv_future(rx)),
623 max_items_per_poll,
624 subscription_counter: Some(subscription_counter),
625 }
626 }
627
628 pub fn new_untracked(rx: broadcast::Receiver<T>, max_items_per_poll: usize) -> Self {
630 assert!(max_items_per_poll > 0, "max_items_per_poll must be > 0");
631 Self {
632 peer: None,
633 inner: ReusableBoxFuture::new(make_recv_future(rx)),
634 max_items_per_poll,
635 subscription_counter: None,
636 }
637 }
638}
639
640impl<T: 'static + Clone + Send> Stream for BroadcastStream<T> {
641 type Item = Vec<T>;
642
643 fn poll_next(
644 mut self: Pin<&mut Self>,
645 cx: &mut task::Context<'_>,
646 ) -> task::Poll<Option<Self::Item>> {
647 loop {
648 let (result, mut rx) = ready!(self.inner.poll(cx));
649
650 match result {
651 Ok(item) => {
652 let mut items = Vec::new();
653 items.push(item);
654
655 while items.len() < self.max_items_per_poll {
657 match rx.try_recv() {
658 Ok(item) => items.push(item),
659 Err(broadcast::error::TryRecvError::Empty) => break,
660 Err(broadcast::error::TryRecvError::Closed) => break,
661 Err(broadcast::error::TryRecvError::Lagged(n)) => {
662 warn!("BroadcastStream {:?} lagged by {} messages", self.peer, n);
663 break;
664 }
665 }
666 }
667
668 self.inner.set(make_recv_future(rx));
669 return task::Poll::Ready(Some(items));
670 }
671 Err(broadcast::error::RecvError::Closed) => {
672 info!("BroadcastStream {:?} closed", self.peer);
673 return task::Poll::Ready(None);
674 }
675 Err(broadcast::error::RecvError::Lagged(n)) => {
676 warn!("BroadcastStream {:?} lagged by {} messages", self.peer, n);
677 self.inner.set(make_recv_future(rx));
679 continue;
680 }
681 }
682 }
683 }
684}
685
686impl<T> Drop for BroadcastStream<T> {
687 fn drop(&mut self) {
688 if let (Some(counter), Some(peer)) = (&self.subscription_counter, &self.peer) {
689 counter.decrement(peer);
690 }
691 }
692}
693
694async fn make_recv_future<T: Clone>(
695 mut rx: broadcast::Receiver<T>,
696) -> (
697 Result<T, broadcast::error::RecvError>,
698 broadcast::Receiver<T>,
699) {
700 let result = rx.recv().await;
701 (result, rx)
702}
703
704#[cfg(test)]
705mod tests {
706 use std::{
707 collections::{BTreeMap, BTreeSet},
708 sync::Arc,
709 time::Duration,
710 };
711
712 use async_trait::async_trait;
713 use bytes::Bytes;
714 use consensus_config::AuthorityIndex;
715 use consensus_types::block::{BlockDigest, BlockRef, Round};
716 use parking_lot::{Mutex, RwLock};
717 use tokio::{sync::broadcast, time::sleep};
718
719 use futures::StreamExt as _;
720
721 fn expect_full(form: &SerializedBlockForm) -> &[u8] {
723 match form {
724 SerializedBlockForm::Full(bytes) => bytes,
725 SerializedBlockForm::Slim(_) => panic!("expected a full block"),
726 }
727 }
728
729 use crate::{
730 authority_service::AuthorityService,
731 block::{BlockAPI, SignedBlock, TestBlock, VerifiedBlock},
732 block_sync_service::BlockSyncService,
733 commit::{CertifiedCommits, CommitRange},
734 commit_vote_monitor::CommitVoteMonitor,
735 context::Context,
736 core_thread::{CoreError, CoreThreadDispatcher},
737 dag_state::DagState,
738 error::ConsensusResult,
739 network::{
740 BlockStream, ExtendedSerializedBlock, ObserverNetworkClient, SerializedBlockForm,
741 SynchronizerClient, ValidatorNetworkClient, ValidatorNetworkService,
742 },
743 peers_pool::PeersPool,
744 round_tracker::RoundTracker,
745 storage::mem_store::MemStore,
746 synchronizer::Synchronizer,
747 test_dag_builder::DagBuilder,
748 transaction_vote_tracker::TransactionVoteTracker,
749 };
750 struct FakeCoreThreadDispatcher {
751 blocks: Mutex<Vec<VerifiedBlock>>,
752 }
753
754 impl FakeCoreThreadDispatcher {
755 fn new() -> Self {
756 Self {
757 blocks: Mutex::new(vec![]),
758 }
759 }
760
761 fn get_blocks(&self) -> Vec<VerifiedBlock> {
762 self.blocks.lock().clone()
763 }
764 }
765
766 #[async_trait]
767 impl CoreThreadDispatcher for FakeCoreThreadDispatcher {
768 async fn add_blocks(
769 &self,
770 blocks: Vec<VerifiedBlock>,
771 ) -> Result<BTreeSet<BlockRef>, CoreError> {
772 let block_refs = blocks.iter().map(|b| b.reference()).collect();
773 self.blocks.lock().extend(blocks);
774 Ok(block_refs)
775 }
776
777 async fn check_block_refs(
778 &self,
779 _block_refs: Vec<BlockRef>,
780 ) -> Result<BTreeSet<BlockRef>, CoreError> {
781 Ok(BTreeSet::new())
782 }
783
784 async fn add_certified_commits(
785 &self,
786 _commits: CertifiedCommits,
787 ) -> Result<BTreeSet<BlockRef>, CoreError> {
788 todo!()
789 }
790
791 async fn new_block(&self, _round: Round, _force: bool) -> Result<(), CoreError> {
792 Ok(())
793 }
794
795 async fn get_missing_blocks(&self) -> Result<BTreeSet<BlockRef>, CoreError> {
796 Ok(Default::default())
797 }
798
799 fn set_propagation_delay(&self, _propagation_delay: Round) -> Result<(), CoreError> {
800 todo!()
801 }
802
803 fn set_last_known_proposed_round(&self, _round: Round) -> Result<(), CoreError> {
804 todo!()
805 }
806 }
807
808 #[derive(Default)]
809 struct FakeNetworkClient {}
810
811 #[async_trait]
812 impl ValidatorNetworkClient for FakeNetworkClient {
813 async fn send_block(
814 &self,
815 _peer: AuthorityIndex,
816 _block: &VerifiedBlock,
817 _timeout: Duration,
818 ) -> ConsensusResult<()> {
819 unimplemented!("Unimplemented")
820 }
821
822 async fn subscribe_blocks(
823 &self,
824 _peer: AuthorityIndex,
825 _last_received: Round,
826 _timeout: Duration,
827 ) -> ConsensusResult<BlockStream> {
828 unimplemented!("Unimplemented")
829 }
830
831 async fn fetch_blocks(
832 &self,
833 _peer: AuthorityIndex,
834 _block_refs: Vec<BlockRef>,
835 _fetch_after_rounds: Vec<Round>,
836 _fetch_missing_ancestors: bool,
837 _timeout: Duration,
838 ) -> ConsensusResult<Vec<Bytes>> {
839 unimplemented!("Unimplemented")
840 }
841
842 async fn fetch_commits(
843 &self,
844 _peer: AuthorityIndex,
845 _commit_range: CommitRange,
846 _timeout: Duration,
847 ) -> ConsensusResult<(Vec<Bytes>, Vec<Bytes>)> {
848 unimplemented!("Unimplemented")
849 }
850
851 async fn fetch_latest_blocks(
852 &self,
853 _peer: AuthorityIndex,
854 _authorities: Vec<AuthorityIndex>,
855 _timeout: Duration,
856 ) -> ConsensusResult<Vec<Bytes>> {
857 unimplemented!("Unimplemented")
858 }
859
860 async fn get_latest_rounds(
861 &self,
862 _peer: AuthorityIndex,
863 _timeout: Duration,
864 ) -> ConsensusResult<(Vec<Round>, Vec<Round>)> {
865 unimplemented!("Unimplemented")
866 }
867 }
868
869 #[async_trait]
870 impl ObserverNetworkClient for FakeNetworkClient {
871 async fn stream_blocks(
872 &self,
873 _peer: crate::network::PeerId,
874 _highest_round_per_authority: Vec<Round>,
875 _timeout: Duration,
876 ) -> ConsensusResult<crate::network::ObserverBlockStream> {
877 unimplemented!("Unimplemented")
878 }
879
880 async fn fetch_blocks(
881 &self,
882 _peer: crate::network::PeerId,
883 _block_refs: Vec<BlockRef>,
884 _fetch_after_rounds: Vec<Round>,
885 _fetch_missing_ancestors: bool,
886 _timeout: Duration,
887 ) -> ConsensusResult<Vec<Bytes>> {
888 unimplemented!("Unimplemented")
889 }
890
891 async fn fetch_commits(
892 &self,
893 _peer: crate::network::PeerId,
894 _commit_range: CommitRange,
895 _timeout: Duration,
896 ) -> ConsensusResult<(Vec<Bytes>, Vec<Bytes>)> {
897 unimplemented!("Unimplemented")
898 }
899 }
900
901 #[tokio::test(flavor = "current_thread", start_paused = true)]
902 async fn test_handle_send_block() {
903 let (context, _keys) = Context::new_for_test(4);
904 let context = Arc::new(context);
905 let block_verifier = Arc::new(crate::block_verifier::NoopBlockVerifier {});
906 let commit_vote_monitor = Arc::new(CommitVoteMonitor::new(context.clone()));
907 let core_dispatcher = Arc::new(FakeCoreThreadDispatcher::new());
908 let (_tx_block_broadcast, rx_block_broadcast) = broadcast::channel(100);
909 let fake_client = Arc::new(FakeNetworkClient::default());
910 let network_client = Arc::new(SynchronizerClient::new(
911 context.clone(),
912 Some(fake_client.clone()),
913 Some(fake_client.clone()),
914 ));
915 let store = Arc::new(MemStore::new());
916 let dag_state = Arc::new(RwLock::new(DagState::new(context.clone(), store.clone())));
917 let transaction_vote_tracker =
918 TransactionVoteTracker::new(context.clone(), block_verifier.clone(), dag_state.clone());
919 let round_tracker = Arc::new(RwLock::new(RoundTracker::new(context.clone(), vec![])));
920 let peers_pool = Arc::new(PeersPool::new(context.clone()));
921 let synchronizer = Synchronizer::start(
922 network_client,
923 context.clone(),
924 core_dispatcher.clone(),
925 commit_vote_monitor.clone(),
926 block_verifier.clone(),
927 transaction_vote_tracker.clone(),
928 round_tracker.clone(),
929 dag_state.clone(),
930 peers_pool.clone(),
931 false,
932 );
933 let block_sync_service = Arc::new(BlockSyncService::new(
934 context.clone(),
935 dag_state.clone(),
936 store.clone(),
937 ));
938 let authority_service = Arc::new(AuthorityService::new(
939 context.clone(),
940 block_verifier,
941 commit_vote_monitor,
942 round_tracker,
943 synchronizer,
944 core_dispatcher.clone(),
945 rx_block_broadcast,
946 transaction_vote_tracker,
947 dag_state,
948 block_sync_service,
949 ));
950
951 let now = context.clock.timestamp_utc_ms();
953 let max_drift = context.parameters.max_forward_time_drift;
954 let input_block = VerifiedBlock::new_for_test(
955 TestBlock::new(9, 0)
956 .set_timestamp_ms(now + max_drift.as_millis() as u64)
957 .build(),
958 );
959
960 let service = authority_service.clone();
961 let serialized = ExtendedSerializedBlock {
962 block: SerializedBlockForm::Full(input_block.serialized().clone()),
963 excluded_ancestors: vec![],
964 };
965
966 tokio::spawn({
967 let service = service.clone();
968 let context = context.clone();
969 async move {
970 service
971 .handle_send_block(context.committee.to_authority_index(0).unwrap(), serialized)
972 .await
973 .unwrap();
974 }
975 });
976
977 sleep(max_drift / 2).await;
978
979 let blocks = core_dispatcher.get_blocks();
980 assert_eq!(blocks.len(), 1);
981 assert_eq!(blocks[0], input_block);
982
983 let invalid_block =
985 VerifiedBlock::new_for_test(TestBlock::new(10, 1000).set_timestamp_ms(10).build());
986 let extended_block = ExtendedSerializedBlock {
987 block: SerializedBlockForm::Full(invalid_block.serialized().clone()),
988 excluded_ancestors: vec![],
989 };
990 service
991 .handle_send_block(
992 context.committee.to_authority_index(0).unwrap(),
993 extended_block,
994 )
995 .await
996 .unwrap_err();
997
998 let slim_block = ExtendedSerializedBlock {
1001 block: SerializedBlockForm::Slim(Bytes::from_static(b"slim")),
1002 excluded_ancestors: vec![],
1003 };
1004 let result = service
1005 .handle_send_block(context.committee.to_authority_index(0).unwrap(), slim_block)
1006 .await;
1007 assert!(matches!(
1008 result,
1009 Err(crate::error::ConsensusError::UnexpectedBlockForm)
1010 ));
1011
1012 let invalid_excluded_ancestors = vec![
1014 bcs::to_bytes(&BlockRef::new(
1015 10,
1016 AuthorityIndex::new_for_test(1000),
1017 BlockDigest::MIN,
1018 ))
1019 .unwrap(),
1020 vec![3u8; 40],
1021 bcs::to_bytes(&invalid_block.reference()).unwrap(),
1022 ];
1023 let extended_block = ExtendedSerializedBlock {
1024 block: SerializedBlockForm::Full(input_block.serialized().clone()),
1025 excluded_ancestors: invalid_excluded_ancestors,
1026 };
1027 service
1028 .handle_send_block(
1029 context.committee.to_authority_index(0).unwrap(),
1030 extended_block,
1031 )
1032 .await
1033 .unwrap_err();
1034 }
1035
1036 #[tokio::test(flavor = "current_thread", start_paused = true)]
1037 async fn test_handle_fetch_blocks() {
1038 const NUM_AUTHORITIES: usize = 40;
1041 const NUM_ROUNDS: usize = 40;
1042 let (mut context, _keys) = Context::new_for_test(NUM_AUTHORITIES);
1043 context.parameters.max_blocks_per_fetch = 50;
1044 let context = Arc::new(context);
1045 let block_verifier = Arc::new(crate::block_verifier::NoopBlockVerifier {});
1046 let commit_vote_monitor = Arc::new(CommitVoteMonitor::new(context.clone()));
1047 let core_dispatcher = Arc::new(FakeCoreThreadDispatcher::new());
1048 let (_tx_block_broadcast, rx_block_broadcast) = broadcast::channel(100);
1049 let fake_client = Arc::new(FakeNetworkClient::default());
1050 let network_client = Arc::new(SynchronizerClient::new(
1051 context.clone(),
1052 Some(fake_client.clone()),
1053 Some(fake_client.clone()),
1054 ));
1055 let store = Arc::new(MemStore::new());
1056 let dag_state = Arc::new(RwLock::new(DagState::new(context.clone(), store.clone())));
1057 let transaction_vote_tracker =
1058 TransactionVoteTracker::new(context.clone(), block_verifier.clone(), dag_state.clone());
1059 let round_tracker = Arc::new(RwLock::new(RoundTracker::new(context.clone(), vec![])));
1060 let peers_pool = Arc::new(PeersPool::new(context.clone()));
1061 let synchronizer = Synchronizer::start(
1062 network_client,
1063 context.clone(),
1064 core_dispatcher.clone(),
1065 commit_vote_monitor.clone(),
1066 block_verifier.clone(),
1067 transaction_vote_tracker.clone(),
1068 round_tracker.clone(),
1069 dag_state.clone(),
1070 peers_pool.clone(),
1071 false,
1072 );
1073 let block_sync_service = Arc::new(BlockSyncService::new(
1074 context.clone(),
1075 dag_state.clone(),
1076 store.clone(),
1077 ));
1078 let authority_service = Arc::new(AuthorityService::new(
1079 context.clone(),
1080 block_verifier,
1081 commit_vote_monitor,
1082 round_tracker,
1083 synchronizer,
1084 core_dispatcher.clone(),
1085 rx_block_broadcast,
1086 transaction_vote_tracker,
1087 dag_state.clone(),
1088 block_sync_service,
1089 ));
1090
1091 let mut dag_builder = DagBuilder::new(context.clone());
1093 dag_builder
1094 .layers(1..=(NUM_ROUNDS as u32))
1095 .build()
1096 .persist_layers(dag_state.clone());
1097 dag_state.write().flush();
1098 let all_blocks = dag_builder.all_blocks();
1099
1100 let missing_block_refs: Vec<BlockRef> = all_blocks
1102 .iter()
1103 .rev()
1104 .take(2)
1105 .map(|b| b.reference())
1106 .collect();
1107 let highest_accepted_rounds: Vec<Round> = vec![1; NUM_AUTHORITIES];
1108 let results = authority_service
1109 .handle_fetch_blocks(
1110 AuthorityIndex::new_for_test(0),
1111 missing_block_refs.clone(),
1112 highest_accepted_rounds,
1113 true,
1114 )
1115 .await
1116 .unwrap();
1117
1118 let blocks: BTreeMap<BlockRef, VerifiedBlock> = results
1120 .iter()
1121 .map(|b| {
1122 let signed = bcs::from_bytes(b).unwrap();
1123 let block = VerifiedBlock::new_verified(signed, b.clone());
1124 (block.reference(), block)
1125 })
1126 .collect();
1127 assert_eq!(blocks.len(), context.parameters.max_blocks_per_sync);
1128 for b in &missing_block_refs {
1130 assert!(blocks.contains_key(b));
1131 }
1132 let num_missing_ancestors = blocks
1133 .keys()
1134 .filter(|b| b.round == NUM_ROUNDS as Round - 1)
1135 .count();
1136 assert_eq!(
1137 num_missing_ancestors,
1138 context.parameters.max_blocks_per_sync - missing_block_refs.len()
1139 );
1140
1141 let missing_round = NUM_ROUNDS as Round - 3;
1143 let missing_block_refs: Vec<BlockRef> = all_blocks
1144 .iter()
1145 .filter(|b| b.reference().round == missing_round)
1146 .map(|b| b.reference())
1147 .take(2)
1148 .collect();
1149 let mut highest_accepted_rounds: Vec<Round> = vec![1; NUM_AUTHORITIES];
1150 highest_accepted_rounds[missing_block_refs[0].author] = missing_round - 5;
1152 let results = authority_service
1153 .handle_fetch_blocks(
1154 AuthorityIndex::new_for_test(0),
1155 missing_block_refs.clone(),
1156 highest_accepted_rounds,
1157 false,
1158 )
1159 .await
1160 .unwrap();
1161
1162 let blocks: BTreeMap<BlockRef, VerifiedBlock> = results
1164 .iter()
1165 .map(|b| {
1166 let signed = bcs::from_bytes(b).unwrap();
1167 let block = VerifiedBlock::new_verified(signed, b.clone());
1168 (block.reference(), block)
1169 })
1170 .collect();
1171 assert_eq!(blocks.len(), context.parameters.max_blocks_per_sync);
1172 for b in &missing_block_refs {
1174 assert!(blocks.contains_key(b));
1175 }
1176 let expected_authors = [missing_block_refs[0].author, missing_block_refs[1].author];
1178 for b in blocks.keys() {
1179 assert!(b.round <= missing_round);
1180 assert!(expected_authors.contains(&b.author));
1181 }
1182
1183 let mut highest_accepted_rounds: Vec<Round> = vec![1; NUM_AUTHORITIES];
1185 highest_accepted_rounds[0] = (NUM_ROUNDS as Round) - 5;
1187 highest_accepted_rounds[1] = (NUM_ROUNDS as Round) - 3;
1188 let results = authority_service
1189 .handle_fetch_blocks(
1190 AuthorityIndex::new_for_test(0),
1191 vec![],
1192 highest_accepted_rounds.clone(),
1193 false,
1194 )
1195 .await
1196 .unwrap();
1197
1198 let blocks: BTreeMap<BlockRef, VerifiedBlock> = results
1200 .iter()
1201 .map(|b| {
1202 let signed = bcs::from_bytes(b).unwrap();
1203 let block = VerifiedBlock::new_verified(signed, b.clone());
1204 (block.reference(), block)
1205 })
1206 .collect();
1207 assert_eq!(blocks.len(), context.parameters.max_blocks_per_fetch);
1208 for block_ref in blocks.keys() {
1210 let accepted = highest_accepted_rounds[block_ref.author];
1211 assert!(block_ref.round > accepted);
1212 }
1213 let max_round_in_result = blocks.keys().map(|b| b.round).max().unwrap();
1216 assert!(
1219 max_round_in_result <= 4,
1220 "Expected low rounds from fair round-order fetching, got max round {}",
1221 max_round_in_result
1222 );
1223
1224 let missing_block_refs: Vec<BlockRef> = all_blocks
1226 .iter()
1227 .filter(|b| b.reference().round == NUM_ROUNDS as Round - 10)
1228 .map(|b| b.reference())
1229 .take(5)
1230 .collect();
1231 let results = authority_service
1232 .handle_fetch_blocks(
1233 AuthorityIndex::new_for_test(0),
1234 missing_block_refs.clone(),
1235 vec![],
1236 false,
1237 )
1238 .await
1239 .unwrap();
1240
1241 let blocks: BTreeMap<BlockRef, VerifiedBlock> = results
1243 .iter()
1244 .map(|b| {
1245 let signed = bcs::from_bytes(b).unwrap();
1246 let block = VerifiedBlock::new_verified(signed, b.clone());
1247 (block.reference(), block)
1248 })
1249 .collect();
1250 assert_eq!(blocks.len(), 5);
1251 for b in &missing_block_refs {
1252 assert!(blocks.contains_key(b));
1253 }
1254 }
1255
1256 #[tokio::test(flavor = "current_thread", start_paused = true)]
1257 async fn test_handle_fetch_latest_blocks() {
1258 let (context, _keys) = Context::new_for_test(4);
1260 let context = Arc::new(context);
1261 let block_verifier = Arc::new(crate::block_verifier::NoopBlockVerifier {});
1262 let commit_vote_monitor = Arc::new(CommitVoteMonitor::new(context.clone()));
1263 let core_dispatcher = Arc::new(FakeCoreThreadDispatcher::new());
1264 let (_tx_block_broadcast, rx_block_broadcast) = broadcast::channel(100);
1265 let fake_client = Arc::new(FakeNetworkClient::default());
1266 let network_client = Arc::new(SynchronizerClient::new(
1267 context.clone(),
1268 Some(fake_client.clone()),
1269 Some(fake_client.clone()),
1270 ));
1271 let store = Arc::new(MemStore::new());
1272 let dag_state = Arc::new(RwLock::new(DagState::new(context.clone(), store.clone())));
1273 let transaction_vote_tracker =
1274 TransactionVoteTracker::new(context.clone(), block_verifier.clone(), dag_state.clone());
1275 let round_tracker = Arc::new(RwLock::new(RoundTracker::new(context.clone(), vec![])));
1276 let peers_pool = Arc::new(PeersPool::new(context.clone()));
1277 let synchronizer = Synchronizer::start(
1278 network_client,
1279 context.clone(),
1280 core_dispatcher.clone(),
1281 commit_vote_monitor.clone(),
1282 block_verifier.clone(),
1283 transaction_vote_tracker.clone(),
1284 round_tracker.clone(),
1285 dag_state.clone(),
1286 peers_pool.clone(),
1287 true,
1288 );
1289 let block_sync_service = Arc::new(BlockSyncService::new(
1290 context.clone(),
1291 dag_state.clone(),
1292 store.clone(),
1293 ));
1294 let authority_service = Arc::new(AuthorityService::new(
1295 context.clone(),
1296 block_verifier,
1297 commit_vote_monitor,
1298 round_tracker,
1299 synchronizer,
1300 core_dispatcher.clone(),
1301 rx_block_broadcast,
1302 transaction_vote_tracker,
1303 dag_state.clone(),
1304 block_sync_service,
1305 ));
1306
1307 let mut dag_builder = DagBuilder::new(context.clone());
1309 dag_builder
1310 .layers(1..=10)
1311 .authorities(vec![AuthorityIndex::new_for_test(2)])
1312 .equivocate(1)
1313 .build()
1314 .persist_layers(dag_state);
1315
1316 let authorities_to_request = vec![
1318 AuthorityIndex::new_for_test(1),
1319 AuthorityIndex::new_for_test(2),
1320 ];
1321 let results = authority_service
1322 .handle_fetch_latest_blocks(AuthorityIndex::new_for_test(1), authorities_to_request)
1323 .await;
1324
1325 let serialised_blocks = results.unwrap();
1327 for serialised_block in serialised_blocks {
1328 let signed_block: SignedBlock =
1329 bcs::from_bytes(&serialised_block).expect("Error while deserialising block");
1330 let verified_block = VerifiedBlock::new_verified(signed_block, serialised_block);
1331
1332 assert_eq!(verified_block.round(), 10);
1333 }
1334 }
1335
1336 #[tokio::test(flavor = "current_thread", start_paused = true)]
1337 async fn test_handle_subscribe_blocks() {
1338 let (context, _keys) = Context::new_for_test(4);
1339 let context = Arc::new(context);
1340 let block_verifier = Arc::new(crate::block_verifier::NoopBlockVerifier {});
1341 let commit_vote_monitor = Arc::new(CommitVoteMonitor::new(context.clone()));
1342 let core_dispatcher = Arc::new(FakeCoreThreadDispatcher::new());
1343 let (_tx_block_broadcast, rx_block_broadcast) = broadcast::channel(100);
1344 let fake_client = Arc::new(FakeNetworkClient::default());
1345 let network_client = Arc::new(SynchronizerClient::new(
1346 context.clone(),
1347 Some(fake_client.clone()),
1348 Some(fake_client.clone()),
1349 ));
1350 let store = Arc::new(MemStore::new());
1351 let dag_state = Arc::new(RwLock::new(DagState::new(context.clone(), store.clone())));
1352 let transaction_vote_tracker =
1353 TransactionVoteTracker::new(context.clone(), block_verifier.clone(), dag_state.clone());
1354 let round_tracker = Arc::new(RwLock::new(RoundTracker::new(context.clone(), vec![])));
1355 let peers_pool = Arc::new(PeersPool::new(context.clone()));
1356 let synchronizer = Synchronizer::start(
1357 network_client,
1358 context.clone(),
1359 core_dispatcher.clone(),
1360 commit_vote_monitor.clone(),
1361 block_verifier.clone(),
1362 transaction_vote_tracker.clone(),
1363 round_tracker.clone(),
1364 dag_state.clone(),
1365 peers_pool.clone(),
1366 false,
1367 );
1368 let block_sync_service = Arc::new(BlockSyncService::new(
1369 context.clone(),
1370 dag_state.clone(),
1371 store.clone(),
1372 ));
1373
1374 dag_state
1376 .write()
1377 .accept_block(VerifiedBlock::new_for_test(TestBlock::new(5, 0).build()));
1378 dag_state
1379 .write()
1380 .accept_block(VerifiedBlock::new_for_test(TestBlock::new(10, 0).build()));
1381 dag_state
1382 .write()
1383 .accept_block(VerifiedBlock::new_for_test(TestBlock::new(15, 0).build()));
1384
1385 let authority_service = Arc::new(AuthorityService::new(
1386 context.clone(),
1387 block_verifier,
1388 commit_vote_monitor,
1389 round_tracker,
1390 synchronizer,
1391 core_dispatcher.clone(),
1392 rx_block_broadcast,
1393 transaction_vote_tracker,
1394 dag_state.clone(),
1395 block_sync_service,
1396 ));
1397
1398 let peer = context.committee.to_authority_index(1).unwrap();
1399
1400 {
1403 let mut stream = authority_service
1404 .handle_subscribe_blocks(peer, 100)
1405 .await
1406 .unwrap();
1407 let block: SignedBlock =
1408 bcs::from_bytes(expect_full(&stream.next().await.unwrap().block)).unwrap();
1409 assert_eq!(
1410 block.round(),
1411 15,
1412 "Should return last proposed block as fallback"
1413 );
1414 assert_eq!(block.author().value(), 0);
1415 }
1416
1417 {
1420 let mut stream = authority_service
1421 .handle_subscribe_blocks(peer, 7)
1422 .await
1423 .unwrap();
1424
1425 let block1: SignedBlock =
1426 bcs::from_bytes(expect_full(&stream.next().await.unwrap().block)).unwrap();
1427 assert_eq!(block1.round(), 10, "Should return block at round 10");
1428
1429 let block2: SignedBlock =
1430 bcs::from_bytes(expect_full(&stream.next().await.unwrap().block)).unwrap();
1431 assert_eq!(block2.round(), 15, "Should return block at round 15");
1432 }
1433 }
1434
1435 #[tokio::test(flavor = "current_thread", start_paused = true)]
1436 async fn test_handle_subscribe_blocks_not_proposed() {
1437 let (context, _keys) = Context::new_for_test(4);
1438 let context = Arc::new(context);
1439 let block_verifier = Arc::new(crate::block_verifier::NoopBlockVerifier {});
1440 let commit_vote_monitor = Arc::new(CommitVoteMonitor::new(context.clone()));
1441 let core_dispatcher = Arc::new(FakeCoreThreadDispatcher::new());
1442 let (_tx_block_broadcast, rx_block_broadcast) = broadcast::channel(100);
1443 let fake_client = Arc::new(FakeNetworkClient::default());
1444 let network_client = Arc::new(SynchronizerClient::new(
1445 context.clone(),
1446 Some(fake_client.clone()),
1447 Some(fake_client.clone()),
1448 ));
1449 let store = Arc::new(MemStore::new());
1450 let dag_state = Arc::new(RwLock::new(DagState::new(context.clone(), store.clone())));
1451 let transaction_vote_tracker =
1452 TransactionVoteTracker::new(context.clone(), block_verifier.clone(), dag_state.clone());
1453 let round_tracker = Arc::new(RwLock::new(RoundTracker::new(context.clone(), vec![])));
1454 let peers_pool = Arc::new(PeersPool::new(context.clone()));
1455 let synchronizer = Synchronizer::start(
1456 network_client,
1457 context.clone(),
1458 core_dispatcher.clone(),
1459 commit_vote_monitor.clone(),
1460 block_verifier.clone(),
1461 transaction_vote_tracker.clone(),
1462 round_tracker.clone(),
1463 dag_state.clone(),
1464 peers_pool.clone(),
1465 false,
1466 );
1467 let block_sync_service = Arc::new(BlockSyncService::new(
1468 context.clone(),
1469 dag_state.clone(),
1470 store.clone(),
1471 ));
1472
1473 let authority_service = Arc::new(AuthorityService::new(
1476 context.clone(),
1477 block_verifier,
1478 commit_vote_monitor,
1479 round_tracker,
1480 synchronizer,
1481 core_dispatcher.clone(),
1482 rx_block_broadcast,
1483 transaction_vote_tracker,
1484 dag_state.clone(),
1485 block_sync_service,
1486 ));
1487
1488 let peer = context.committee.to_authority_index(1).unwrap();
1489
1490 let mut stream = authority_service
1492 .handle_subscribe_blocks(peer, 0)
1493 .await
1494 .unwrap();
1495
1496 use futures::poll;
1498 use std::task::Poll;
1499 let poll_result = poll!(stream.next());
1500 assert!(
1501 matches!(poll_result, Poll::Pending),
1502 "Should not receive genesis block on subscription stream"
1503 );
1504 }
1505}