Skip to main content

consensus_core/
authority_service.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4use 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
38/// Authority's network service implementation, agnostic to the actual networking stack used.
39pub(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    // Parses and validates serialized excluded ancestors.
83    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        // TODO: dedup block verifications, here and with fetched blocks.
160        // Only the full form reaches this service: a slim payload is decoded back to
161        // full upstream (or dropped, before the decoder exists).
162        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        // Reject blocks not produced by the peer.
169        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        // Reject blocks failing parsing and validations.
186        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        // Observe the block for the commit votes. When local commit is lagging too much,
234        // commit sync loop will trigger fetching.
235        self.commit_vote_monitor.observe_block(&verified_block);
236
237        // Update own received rounds and peer accepted rounds from this verified block.
238        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        // Reject blocks when local commit index is lagging too far from quorum commit index,
246        // to avoid the memory overhead from suspended blocks.
247        //
248        // IMPORTANT: this must be done after observing votes from the block, otherwise
249        // observed quorum commit will no longer progress.
250        //
251        // Since the main issue with too many suspended blocks is memory usage not CPU,
252        // it is ok to reject after block verifications instead of before.
253        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        // The block is verified and current, so record own votes on the block
280        // before sending the block to Core.
281        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        // Send the block to Core to try accepting it into the DAG.
287        let missing_ancestors = self
288            .core_dispatcher
289            .add_blocks(vec![verified_block.clone()])
290            .await
291            .map_err(|_| ConsensusError::Shutdown)?;
292
293        // Schedule fetching missing ancestors from this peer in the background.
294        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                // This does not wait for the fetch request to complete.
304                // It only waits for synchronizer to queue the request to a peer.
305                // When this fails, it usually means the queue is full.
306                // The fetch will retry from other peers via live and periodic syncs.
307                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        // Schedule fetching missing soft links from this peer in the background.
317        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        // Subscribe before snapshotting past blocks below. This can duplicate
352        // a block in both the subscription stream and snapshot, which is fine.
353        // Otherwise, it is possible to miss a block if it is broadcasted after snapshotting
354        // but before subscribing.
355        let broadcast_rx = self.rx_block_broadcast.resubscribe();
356
357        // Find past proposed blocks as the initial blocks to send to the peer.
358        //
359        // If there are cached blocks in the range which the peer requested, send all of them.
360        // The size is bounded by the local GC round and DagState cache size.
361        //
362        // Otherwise if there is no cached block in the range which the peer requested,
363        // and this node has proposed blocks before, at least one block should be sent to the peer
364        // to help with liveness.
365        let past_proposed_blocks = {
366            let dag_state = self.dag_state.read();
367
368            // Saturate so an out-of-range round from the peer cannot wrap to 0 and
369            // replay the entire block cache.
370            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        // Ok to not batch own proposed blocks, which is < 20/s.
393        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        // Return a stream of blocks that first yields missed blocks as requested, then new blocks.
402        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    // Handles 3 types of requests:
414    // 1. Live sync:
415    //    - Both missing block refs and highest accepted rounds are specified.
416    //    - fetch_missing_ancestors is true.
417    //    - response returns max_blocks_per_sync blocks.
418    // 2. Periodic sync:
419    //    - Highest accepted rounds must be specified.
420    //    - Missing block refs are optional.
421    //    - fetch_missing_ancestors is false (default).
422    //    - response returns max_blocks_per_fetch blocks.
423    // 3. Commit sync:
424    //    - Missing block refs are specified.
425    //    - Highest accepted rounds are empty.
426    //    - fetch_missing_ancestors is false (default).
427    //    - response returns max_blocks_per_fetch blocks.
428    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        // Delegate to BlockSyncService
438        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        // Delegate to BlockSyncService
451        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        // Delegate to BlockSyncService
462        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
493/// Atomically counts the number of active subscriptions to the block broadcast stream.
494pub(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        // Set the subscribed peers by default to 0
502        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                // Only count the first subscription from each peer.
544                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
589/// Each broadcasted block stream wraps a broadcast receiver for blocks.
590/// It yields blocks that are broadcasted after the stream is created.
591type BroadcastedBlockStream = BroadcastStream<ExtendedBlock>;
592
593/// Adapted from `tokio_stream::wrappers::BroadcastStream`. The main difference is that
594/// this tolerates lags with only logging, without yielding errors.
595pub(crate) struct BroadcastStream<T> {
596    peer: Option<PeerId>,
597    // Stores the receiver across poll_next() calls.
598    inner: ReusableBoxFuture<
599        'static,
600        (
601            Result<T, broadcast::error::RecvError>,
602            broadcast::Receiver<T>,
603        ),
604    >,
605    // Maximum number of items to return per poll.
606    max_items_per_poll: usize,
607    // Counts total subscriptions / active BroadcastStreams.
608    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    /// Creates a stream without subscription tracking.
629    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                    // Drain any additional items that are already available, up to the cap.
656                    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                    // Re-arm the future and loop to await the next item.
678                    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    /// Every wire payload in these tests is the full form.
722    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        // Test delaying blocks with time drift.
952        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        // Test invalid block.
984        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        // A slim payload that reaches the service is a bug upstream (the subscriber
999        // drops them until the codec lands); it must be rejected, not parsed.
1000        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        // Test invalid excluded ancestors.
1013        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        // GIVEN
1039        // Use NUM_AUTHORITIES and NUM_ROUNDS higher than max_blocks_per_sync to test limits.
1040        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        // GIVEN: 40 rounds of blocks in the dag state.
1092        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        // WHEN: Request 2 blocks from round 40, fetch missing ancestors enabled.
1101        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        // THEN: the expected number of unique blocks are returned.
1119        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        // All missing blocks are returned.
1129        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        // WHEN: Request 2 blocks from round 37, fetch missing ancestors disabled.
1142        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        // Try to fill up the blocks from the 1st authority in missing_block_refs.
1151        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        // THEN: the expected number of unique blocks are returned.
1163        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        // All missing blocks are returned.
1173        for b in &missing_block_refs {
1174            assert!(blocks.contains_key(b));
1175        }
1176        // Ancestor blocks are from the expected rounds and authorities.
1177        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        // WHEN: Request with empty block_refs, fetch missing ancestors disabled.
1184        let mut highest_accepted_rounds: Vec<Round> = vec![1; NUM_AUTHORITIES];
1185        // Set a few authorities to higher accepted rounds.
1186        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        // THEN: the expected number of unique blocks are returned.
1199        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        // Blocks should be from all authorities, within the expected round range.
1209        for block_ref in blocks.keys() {
1210            let accepted = highest_accepted_rounds[block_ref.author];
1211            assert!(block_ref.round > accepted);
1212        }
1213        // Blocks should be fetched in ascending round order across authorities,
1214        // so blocks should have low rounds near the accepted rounds.
1215        let max_round_in_result = blocks.keys().map(|b| b.round).max().unwrap();
1216        // With 40 authorities mostly at accepted round 1 and max_blocks_per_fetch=50,
1217        // the min-heap fills ~1-2 rounds per authority.
1218        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        // WHEN: Request 5 block from round 40, not getting ancestors.
1225        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        // THEN: the expected number of unique blocks are returned.
1242        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        // GIVEN
1259        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        // Create some blocks for a few authorities. Create some equivocations as well and store in dag state.
1308        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        // WHEN
1317        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        // THEN
1326        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        // Create 3 proposed blocks at rounds 5, 10, 15 for own authority (index 0)
1375        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        // Case A: Subscribe with last_received = 100 (after all proposed blocks)
1401        // Should return last proposed block (round 15) as fallback
1402        {
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        // Case B: Subscribe with last_received = 7 (includes rounds 10, 15)
1418        // Should return cached blocks from round 8+
1419        {
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        // No blocks added to DagState - only genesis exists
1474
1475        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        // Subscribe - no blocks have been proposed yet (only genesis exists)
1491        let mut stream = authority_service
1492            .handle_subscribe_blocks(peer, 0)
1493            .await
1494            .unwrap();
1495
1496        // Should NOT receive any block (genesis must not be returned)
1497        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}