Skip to main content

consensus_core/
authority_node.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4use std::{
5    sync::{Arc, Weak},
6    time::{Duration, Instant},
7};
8
9use consensus_config::{
10    Committee, ConsensusProtocolConfig, NetworkKeyPair, NetworkPublicKey, Parameters,
11    ProtocolKeyPair,
12};
13use consensus_types::block::Round;
14use itertools::Itertools;
15use mysten_common::debug_fatal;
16use mysten_network::Multiaddr;
17use parking_lot::RwLock;
18use prometheus::Registry;
19use tracing::{info, warn};
20
21use crate::{
22    BlockAPI as _, CommitConsumerArgs, RandomnessSignatureHandler,
23    authority_service::AuthorityService,
24    block_manager::BlockManager,
25    block_sync_service::BlockSyncService,
26    block_verifier::SignedBlockVerifier,
27    commit_observer::CommitObserver,
28    commit_syncer::{CommitSyncer, CommitSyncerHandle},
29    commit_vote_monitor::CommitVoteMonitor,
30    context::{Clock, Context},
31    core::{Core, CoreSignals},
32    core_thread::{ChannelCoreThreadDispatcher, CoreThreadHandle},
33    dag_state::DagState,
34    leader_schedule::LeaderSchedule,
35    leader_timeout::{LeaderTimeoutTask, LeaderTimeoutTaskHandle},
36    metrics::initialise_metrics,
37    network::{
38        CommitSyncerClient, NetworkManager, PeerId, SynchronizerClient, tonic_network::TonicManager,
39    },
40    observer_service::ObserverService,
41    observer_subscriber::ObserverSubscriber,
42    peers_pool::PeersPool,
43    round_prober::{RoundProber, RoundProberHandle},
44    round_tracker::RoundTracker,
45    storage::rocksdb_store::RocksDBStore,
46    subscriber::Subscriber,
47    synchronizer::{Synchronizer, SynchronizerHandle},
48    transaction::{
49        TransactionClient, TransactionConsumer, TransactionConsumerPool, TransactionPool,
50        TransactionVerifier,
51    },
52    transaction_vote_tracker::TransactionVoteTracker,
53};
54
55/// ConsensusAuthority is used by Sui to manage the lifetime of AuthorityNode.
56/// It hides the details of the implementation from the caller, MysticetiManager.
57#[allow(private_interfaces)]
58pub enum ConsensusAuthority {
59    WithTonic(AuthorityNode<TonicManager>),
60}
61
62impl ConsensusAuthority {
63    pub async fn start(
64        network_type: NetworkType,
65        epoch_start_timestamp_ms: u64,
66        committee: Committee,
67        parameters: Parameters,
68        protocol_config: ConsensusProtocolConfig,
69        // Only required for validator nodes. Observer nodes don't have a protocol keypair.
70        protocol_keypair: Option<ProtocolKeyPair>,
71        network_keypair: NetworkKeyPair,
72        clock: Arc<Clock>,
73        transaction_verifier: Arc<dyn TransactionVerifier>,
74        // When provided, the proposer takes transactions from this pool and the
75        // `TransactionClient` submission path is unused. Only relevant for validator nodes.
76        transaction_pool: Option<Arc<dyn TransactionPool>>,
77        commit_consumer: CommitConsumerArgs,
78        registry: Registry,
79        // A counter that keeps track of how many times the consensus authority has been booted while the process
80        // has been running. It's useful for making decisions on whether amnesia recovery should run.
81        // When `boot_counter` is 0, `ConsensusAuthority` will initiate the process of amnesia recovery if that's enabled in the parameters.
82        boot_counter: u64,
83        randomness_signature_handler: Option<Arc<dyn RandomnessSignatureHandler>>,
84    ) -> Self {
85        match network_type {
86            NetworkType::Tonic => {
87                let authority = AuthorityNode::start(
88                    epoch_start_timestamp_ms,
89                    committee,
90                    parameters,
91                    protocol_config,
92                    protocol_keypair,
93                    network_keypair,
94                    clock,
95                    transaction_verifier,
96                    transaction_pool,
97                    commit_consumer,
98                    registry,
99                    boot_counter,
100                    randomness_signature_handler,
101                )
102                .await;
103                Self::WithTonic(authority)
104            }
105        }
106    }
107
108    pub async fn stop(self) {
109        match self {
110            Self::WithTonic(authority) => authority.stop().await,
111        }
112    }
113
114    pub fn update_peer_address(
115        &self,
116        network_pubkey: NetworkPublicKey,
117        address: Option<Multiaddr>,
118    ) {
119        match self {
120            Self::WithTonic(authority) => authority.update_peer_address(network_pubkey, address),
121        }
122    }
123
124    pub fn transaction_client(&self) -> Arc<TransactionClient> {
125        match self {
126            Self::WithTonic(authority) => authority.transaction_client(),
127        }
128    }
129
130    pub fn store(&self) -> Arc<RocksDBStore> {
131        match self {
132            Self::WithTonic(authority) => authority.store(),
133        }
134    }
135
136    #[cfg(test)]
137    fn context(&self) -> &Arc<Context> {
138        match self {
139            Self::WithTonic(authority) => &authority.context,
140        }
141    }
142}
143
144#[derive(Clone, Copy, PartialEq, Eq, Debug)]
145pub enum NetworkType {
146    Tonic,
147}
148
149/// Enum to handle different subscriber types based on whether the node is a validator or observer
150enum SubscriberType<N: NetworkManager> {
151    Validator(Subscriber<N::ValidatorClient, AuthorityService<ChannelCoreThreadDispatcher>>),
152    Observer(ObserverSubscriber<N::ObserverClient, ObserverService>),
153}
154
155impl<N: NetworkManager> SubscriberType<N> {
156    async fn stop(&self) {
157        match self {
158            SubscriberType::Validator(subscriber) => subscriber.stop().await,
159            SubscriberType::Observer(subscriber) => subscriber.stop().await,
160        }
161    }
162}
163
164pub(crate) struct AuthorityNode<N>
165where
166    N: NetworkManager,
167{
168    context: Arc<Context>,
169    start_time: Instant,
170    transaction_client: Arc<TransactionClient>,
171    synchronizer: Arc<SynchronizerHandle>,
172    store: Arc<RocksDBStore>,
173    // Only use for verification and logging during shutdown.
174    // To avoid keeping the DagState alive at the end of shutdown, this is only a weak reference.
175    dag_state: Weak<RwLock<DagState>>,
176
177    commit_syncer_handle: CommitSyncerHandle,
178    round_prober_handle: Option<RoundProberHandle>,
179    leader_timeout_handle: LeaderTimeoutTaskHandle,
180    core_thread_handle: CoreThreadHandle,
181    subscriber: SubscriberType<N>,
182    // Network proxies hold ObserverService weakly, so AuthorityNode keeps the service alive until
183    // the network servers have stopped.
184    observer_service: Option<Arc<ObserverService>>,
185    network_manager: N,
186}
187
188impl<N> AuthorityNode<N>
189where
190    N: NetworkManager,
191{
192    // See comments above ConsensusAuthority::start() for details on the input.
193    pub(crate) async fn start(
194        epoch_start_timestamp_ms: u64,
195        committee: Committee,
196        parameters: Parameters,
197        protocol_config: ConsensusProtocolConfig,
198        protocol_keypair: Option<ProtocolKeyPair>,
199        network_keypair: NetworkKeyPair,
200        clock: Arc<Clock>,
201        transaction_verifier: Arc<dyn TransactionVerifier>,
202        transaction_pool: Option<Arc<dyn TransactionPool>>,
203        commit_consumer: CommitConsumerArgs,
204        registry: Registry,
205        boot_counter: u64,
206        randomness_signature_handler: Option<Arc<dyn RandomnessSignatureHandler>>,
207    ) -> Self {
208        let metrics = initialise_metrics(registry);
209
210        // If a protocol key pair is provided, then this is a validator node.
211        let own_index = if let Some(protocol_keypair) = &protocol_keypair {
212            let (own_index, _) = committee
213                .authorities()
214                .find(|(_, a)| a.protocol_key == protocol_keypair.public())
215                .expect("Own authority should be among the consensus authorities!");
216
217            let own_hostname = committee.authority(own_index).hostname.clone();
218            info!(
219                "Starting consensus validator authority {} {}, {:?}, epoch start timestamp {}, boot counter {}, replay mode {:?}",
220                own_index,
221                own_hostname,
222                protocol_config.protocol_version(),
223                epoch_start_timestamp_ms,
224                boot_counter,
225                commit_consumer.replay_mode
226            );
227
228            metrics
229                .node_metrics
230                .authority_index
231                .with_label_values(&[&own_hostname])
232                .set(own_index.value() as i64);
233            Some(own_index)
234        } else {
235            // Otherwise this is an observer node and no index exists for it.
236            info!(
237                "Starting consensus observer authority, {:?}, epoch start timestamp {}, boot counter {}, replay mode {:?}",
238                protocol_config.protocol_version(),
239                epoch_start_timestamp_ms,
240                boot_counter,
241                commit_consumer.replay_mode
242            );
243            None
244        };
245
246        info!(
247            "Consensus authorities: {}",
248            committee
249                .authorities()
250                .map(|(i, a)| format!("{}: {}", i, a.hostname))
251                .join(", ")
252        );
253        info!("Consensus parameters: {:?}", parameters);
254        info!("Consensus committee: {:?}", committee);
255        let context = Arc::new(Context::new(
256            epoch_start_timestamp_ms,
257            own_index,
258            committee,
259            parameters,
260            protocol_config,
261            metrics,
262            clock,
263        ));
264        let start_time = Instant::now();
265
266        context
267            .metrics
268            .node_metrics
269            .protocol_version
270            .set(context.protocol_config.protocol_version() as i64);
271
272        let (tx_client, tx_receiver, priority_tx_receiver) =
273            TransactionClient::new(context.clone());
274        let transaction_pool: Arc<dyn TransactionPool> = match transaction_pool {
275            // With an external pool, the TransactionClient path is unused. The channel
276            // receiver is dropped so accidental client submissions fail fast instead of
277            // hanging.
278            Some(pool) => {
279                drop(tx_receiver);
280                drop(priority_tx_receiver);
281                pool
282            }
283            None => Arc::new(TransactionConsumerPool::new(TransactionConsumer::new(
284                tx_receiver,
285                priority_tx_receiver,
286                context.clone(),
287            ))),
288        };
289
290        let (core_signals, signals_receivers) = CoreSignals::new(context.clone());
291
292        let mut network_manager = N::new(context.clone(), network_keypair);
293        let validator_client = network_manager.validator_client();
294        let observer_client = network_manager.observer_client();
295
296        let synchronizer_client = Arc::new(SynchronizerClient::<
297            N::ValidatorClient,
298            N::ObserverClient,
299        >::new(
300            context.clone(),
301            Some(validator_client.clone()),
302            Some(observer_client.clone()),
303        ));
304        let commit_syncer_client = Arc::new(CommitSyncerClient::<
305            N::ValidatorClient,
306            N::ObserverClient,
307        >::new(
308            context.clone(),
309            Some(validator_client.clone()),
310            Some(observer_client.clone()),
311        ));
312
313        let store_path = context.parameters.db_path.as_path().to_str().unwrap();
314        let store = Arc::new(RocksDBStore::new(store_path));
315        let dag_state = Arc::new(RwLock::new(DagState::new(context.clone(), store.clone())));
316
317        let block_verifier = Arc::new(SignedBlockVerifier::new(
318            context.clone(),
319            transaction_verifier,
320        ));
321
322        let transaction_vote_tracker =
323            TransactionVoteTracker::new(context.clone(), block_verifier.clone(), dag_state.clone());
324
325        // Only sync last known own block if we are a validator and it's the first boot.
326        let sync_last_known_own_block = boot_counter == 0
327            && !context
328                .parameters
329                .sync_last_known_own_block_timeout
330                .is_zero()
331            && context.is_validator();
332        info!(
333            "Sync last known own block: {}. Boot count: {}. Timeout: {:?}.",
334            sync_last_known_own_block,
335            boot_counter,
336            context.parameters.sync_last_known_own_block_timeout
337        );
338
339        let block_manager = BlockManager::new(context.clone(), dag_state.clone());
340
341        let leader_schedule = Arc::new(LeaderSchedule::from_store(
342            context.clone(),
343            dag_state.clone(),
344        ));
345
346        let commit_consumer_monitor = commit_consumer.monitor();
347        let commit_observer = CommitObserver::new(
348            context.clone(),
349            commit_consumer,
350            dag_state.clone(),
351            transaction_vote_tracker.clone(),
352        )
353        .await;
354
355        let initial_received_rounds = dag_state
356            .read()
357            .get_last_cached_block_per_authority(Round::MAX)
358            .into_iter()
359            .map(|(block, _)| block.round())
360            .collect::<Vec<_>>();
361        let round_tracker = Arc::new(RwLock::new(RoundTracker::new(
362            context.clone(),
363            initial_received_rounds,
364        )));
365
366        // To avoid accidentally leaking the private key, the protocol key pair should only be
367        // kept in Core.
368        let core = if context.is_validator() {
369            Core::new_validator(
370                context.clone(),
371                leader_schedule,
372                transaction_pool,
373                transaction_vote_tracker.clone(),
374                block_manager,
375                commit_observer,
376                core_signals,
377                protocol_keypair.expect("protocol keypair is required when running as validator"),
378                dag_state.clone(),
379                sync_last_known_own_block,
380                round_tracker.clone(),
381            )
382        } else {
383            Core::new_observer(
384                context.clone(),
385                leader_schedule,
386                block_manager,
387                commit_observer,
388                core_signals,
389                dag_state.clone(),
390            )
391        };
392
393        let (core_dispatcher, core_thread_handle) =
394            ChannelCoreThreadDispatcher::start(context.clone(), &dag_state, core);
395        let core_dispatcher = Arc::new(core_dispatcher);
396        let leader_timeout_handle =
397            LeaderTimeoutTask::start(core_dispatcher.clone(), &signals_receivers, context.clone());
398
399        let commit_vote_monitor = Arc::new(CommitVoteMonitor::new(context.clone()));
400
401        // Create the PeersPool
402        let peers_pool = Arc::new(PeersPool::new(context.clone()));
403
404        let synchronizer = Synchronizer::start(
405            synchronizer_client.clone(),
406            context.clone(),
407            core_dispatcher.clone(),
408            commit_vote_monitor.clone(),
409            block_verifier.clone(),
410            transaction_vote_tracker.clone(),
411            round_tracker.clone(),
412            dag_state.clone(),
413            peers_pool.clone(),
414            sync_last_known_own_block,
415        );
416
417        let commit_syncer_handle = CommitSyncer::new(
418            context.clone(),
419            core_dispatcher.clone(),
420            commit_vote_monitor.clone(),
421            commit_consumer_monitor.clone(),
422            block_verifier.clone(),
423            transaction_vote_tracker.clone(),
424            round_tracker.clone(),
425            commit_syncer_client.clone(),
426            dag_state.clone(),
427            peers_pool.clone(),
428        )
429        .start();
430
431        // Create BlockSyncService that will be shared by both AuthorityService and ObserverService
432        let block_sync_service = Arc::new(BlockSyncService::new(
433            context.clone(),
434            dag_state.clone(),
435            store.clone(),
436        ));
437
438        let (subscriber, round_prober_handle, observer_service) = if context.is_validator() {
439            let authority_service = Arc::new(AuthorityService::new(
440                context.clone(),
441                block_verifier.clone(),
442                commit_vote_monitor.clone(),
443                round_tracker.clone(),
444                synchronizer.clone(),
445                core_dispatcher.clone(),
446                signals_receivers.block_broadcast_receiver(),
447                transaction_vote_tracker.clone(),
448                dag_state.clone(),
449                block_sync_service.clone(),
450            ));
451
452            // Start the validator server if this is a validator node.
453            network_manager
454                .start_validator_server(authority_service.clone())
455                .await;
456
457            // Validator node: subscribe to all other validators
458            let s = Subscriber::new(
459                context.clone(),
460                validator_client.clone(),
461                authority_service.clone(),
462                dag_state.clone(),
463            );
464            for (peer, _) in context.committee.authorities() {
465                if peer != context.own_index {
466                    s.subscribe(peer);
467                }
468            }
469
470            // Start the round prober
471            let round_prober_handle = Some(
472                RoundProber::new(
473                    context.clone(),
474                    core_dispatcher.clone(),
475                    round_tracker.clone(),
476                    dag_state.clone(),
477                    validator_client,
478                )
479                .start(),
480            );
481
482            // Start the observer server if the observer server is enabled in the parameters.
483            let observer_service = if context.parameters.observer.is_server_enabled() {
484                let observer_service = Arc::new(ObserverService::new(
485                    context.clone(),
486                    core_dispatcher.clone(),
487                    dag_state.clone(),
488                    signals_receivers.accepted_block_broadcast_receiver(),
489                    block_verifier,
490                    commit_vote_monitor.clone(),
491                    transaction_vote_tracker.clone(),
492                    synchronizer.clone(),
493                    block_sync_service.clone(),
494                    randomness_signature_handler.clone(),
495                ));
496                network_manager
497                    .start_observer_server(observer_service.clone())
498                    .await;
499                Some(observer_service)
500            } else {
501                None
502            };
503
504            (
505                SubscriberType::Validator(s),
506                round_prober_handle,
507                observer_service,
508            )
509        } else {
510            // Observer node: subscribe to specified peer(s) using ObserverSubscriber
511            let observer_client = network_manager.observer_client();
512            let observer_service = Arc::new(ObserverService::new(
513                context.clone(),
514                core_dispatcher.clone(),
515                dag_state.clone(),
516                signals_receivers.accepted_block_broadcast_receiver(),
517                block_verifier,
518                commit_vote_monitor.clone(),
519                transaction_vote_tracker.clone(),
520                synchronizer.clone(),
521                block_sync_service.clone(),
522                randomness_signature_handler.clone(),
523            ));
524
525            let observer_subscriber = ObserverSubscriber::new(
526                context.clone(),
527                observer_client,
528                observer_service.clone(),
529                commit_vote_monitor.clone(),
530                dag_state.clone(),
531                randomness_signature_handler,
532            );
533
534            network_manager
535                .start_observer_server(observer_service.clone())
536                .await;
537
538            // Subscribe to peers specified in the configuration
539            // For now get the first peer from the list to connect to.
540            // TODO: support multiple peers - as in choose/detect which one to connect to.
541            for peer_record in context.parameters.observer.peers.iter().take(1) {
542                let peer_id = if let Some((index, _)) = context
543                    .committee
544                    .authorities()
545                    .find(|(_, authority)| authority.network_key == peer_record.public_key)
546                {
547                    PeerId::Validator(index)
548                } else {
549                    PeerId::Observer(Box::new(peer_record.public_key.clone()))
550                };
551
552                info!("Observer subscribing to peer: {:?}", peer_id);
553                observer_subscriber.subscribe(peer_id);
554            }
555
556            (
557                SubscriberType::Observer(observer_subscriber),
558                None,
559                Some(observer_service),
560            )
561        };
562
563        info!(
564            "Consensus authority started, took {:?}",
565            start_time.elapsed()
566        );
567
568        Self {
569            context,
570            start_time,
571            transaction_client: Arc::new(tx_client),
572            synchronizer,
573            store,
574            dag_state: Arc::downgrade(&dag_state),
575            commit_syncer_handle,
576            round_prober_handle,
577            leader_timeout_handle,
578            core_thread_handle,
579            subscriber,
580            observer_service,
581            network_manager,
582        }
583    }
584
585    pub(crate) async fn stop(self) {
586        let Self {
587            context,
588            start_time,
589            transaction_client,
590            synchronizer,
591            store,
592            dag_state,
593            commit_syncer_handle,
594            round_prober_handle,
595            leader_timeout_handle,
596            core_thread_handle,
597            subscriber,
598            observer_service,
599            mut network_manager,
600        } = self;
601
602        info!(
603            "Stopping authority. Total run time: {:?}",
604            start_time.elapsed()
605        );
606
607        // First shutdown components calling into Core.
608        synchronizer.stop().await;
609        commit_syncer_handle.stop().await;
610        if let Some(round_prober_handle) = round_prober_handle {
611            round_prober_handle.stop().await;
612        }
613        leader_timeout_handle.stop().await;
614        // Shutdown Core to stop block productions and broadcast.
615        core_thread_handle.stop().await;
616        // Stop block subscriptions before stopping network server.
617        subscriber.stop().await;
618        network_manager.stop().await;
619
620        context
621            .metrics
622            .node_metrics
623            .uptime
624            .observe(start_time.elapsed().as_secs_f64());
625
626        drop((
627            context,
628            transaction_client,
629            synchronizer,
630            store,
631            subscriber,
632            observer_service,
633            network_manager,
634        ));
635
636        // Canceled tasks should wait to report cancellation on next poll, but they report from
637        // their JoinHandles immediately under msim. Their futures and captured references are
638        // only dropped when the executor next processes the canceled tasks. Sleeping (instead
639        // of yielding) ensures this task resumes only after the pending drops have run.
640        // In production tokio, a canceled task's future is guaranteed to have been dropped when
641        // its JoinHandle resolves, so the loop exits on the first check.
642        let mut dag_state_owners = dag_state.strong_count();
643        for _ in 0..5 {
644            if dag_state_owners == 0 {
645                break;
646            }
647            tokio::time::sleep(Duration::from_millis(1)).await;
648            dag_state_owners = dag_state.strong_count();
649        }
650        if dag_state_owners != 0 {
651            debug_fatal!(
652                "DagState still has {} owner(s) after stopping ConsensusAuthority",
653                dag_state_owners
654            );
655        }
656    }
657
658    pub(crate) fn transaction_client(&self) -> Arc<TransactionClient> {
659        self.transaction_client.clone()
660    }
661
662    pub(crate) fn store(&self) -> Arc<RocksDBStore> {
663        self.store.clone()
664    }
665
666    pub(crate) fn update_peer_address(
667        &self,
668        network_pubkey: NetworkPublicKey,
669        address: Option<Multiaddr>,
670    ) {
671        // Find the peer index for this network key
672        let Some(peer) = self
673            .context
674            .committee
675            .authorities()
676            .find(|(_, authority)| authority.network_key == network_pubkey)
677            .map(|(index, _)| index)
678        else {
679            warn!(
680                "Network public key {:?} not found in committee, ignoring address update",
681                network_pubkey
682            );
683            return;
684        };
685
686        // Update the address in the network manager
687        self.network_manager.update_peer_address(peer, address);
688
689        // Re-subscribe to the peer to force reconnection with new address
690        if peer != self.context.own_index {
691            info!("Re-subscribing to peer {} after address update", peer);
692            match &self.subscriber {
693                SubscriberType::Validator(s) => s.subscribe(peer),
694                SubscriberType::Observer(s) => {
695                    // For observer, create a PeerId for the validator
696                    s.subscribe(PeerId::Validator(peer));
697                }
698            }
699        }
700    }
701}
702
703#[cfg(test)]
704mod tests {
705    #![allow(non_snake_case)]
706
707    use std::{
708        collections::{BTreeMap, BTreeSet},
709        sync::Arc,
710        time::Duration,
711    };
712
713    use consensus_config::{
714        AuthorityIndex, ObserverParameters, Parameters, PeerRecord, local_committee_and_keys,
715    };
716    use mysten_metrics::RegistryService;
717    use mysten_metrics::monitored_mpsc::UnboundedReceiver;
718    use prometheus::Registry;
719    use rand::{SeedableRng, rngs::StdRng};
720    use rstest::rstest;
721    use tempfile::TempDir;
722    use tokio::time::{sleep, timeout};
723    use typed_store::DBMetrics;
724
725    use super::*;
726    use crate::{
727        CommittedSubDag,
728        block::{BlockAPI as _, GENESIS_ROUND},
729        transaction::{NoopTransactionVerifier, Priority},
730    };
731
732    #[rstest]
733    #[tokio::test]
734    async fn test_authority_start_and_stop(
735        #[values(NetworkType::Tonic)] network_type: NetworkType,
736    ) {
737        let (committee, keypairs) = local_committee_and_keys(0, vec![1]);
738        let registry = Registry::new();
739
740        let temp_dir = TempDir::new().unwrap();
741        let parameters = Parameters {
742            db_path: temp_dir.keep(),
743            ..Default::default()
744        };
745        let txn_verifier = NoopTransactionVerifier {};
746
747        let own_index = committee.to_authority_index(0).unwrap();
748        let protocol_keypair = keypairs[own_index].1.clone();
749        let network_keypair = keypairs[own_index].0.clone();
750
751        let (commit_consumer, _) = CommitConsumerArgs::new(0, 0);
752
753        let authority = ConsensusAuthority::start(
754            network_type,
755            0,
756            committee,
757            parameters,
758            ConsensusProtocolConfig::for_testing(),
759            Some(protocol_keypair),
760            network_keypair,
761            Arc::new(Clock::default()),
762            Arc::new(txn_verifier),
763            None,
764            commit_consumer,
765            registry,
766            0,
767            None,
768        )
769        .await;
770
771        assert_eq!(authority.context().own_index, own_index);
772        assert_eq!(authority.context().committee.epoch(), 0);
773        assert_eq!(authority.context().committee.size(), 1);
774
775        authority.stop().await;
776    }
777
778    #[rstest]
779    #[tokio::test]
780    async fn test_observer_start_and_stop(#[values(NetworkType::Tonic)] network_type: NetworkType) {
781        let (committee, keypairs) = local_committee_and_keys(0, vec![1]);
782        let registry = Registry::new();
783
784        let temp_dir = TempDir::new().unwrap();
785        let parameters = Parameters {
786            db_path: temp_dir.keep(),
787            ..Default::default()
788        };
789        let txn_verifier = NoopTransactionVerifier {};
790
791        // Use any network keypair for the observer, it doesn't need to match a committee member
792        let network_keypair = keypairs[0].0.clone();
793
794        let (commit_consumer, _) = CommitConsumerArgs::new(0, 0);
795
796        let observer = ConsensusAuthority::start(
797            network_type,
798            0,
799            committee.clone(),
800            parameters,
801            ConsensusProtocolConfig::for_testing(),
802            None, // No protocol keypair for observer node
803            network_keypair,
804            Arc::new(Clock::default()),
805            Arc::new(txn_verifier),
806            None,
807            commit_consumer,
808            registry,
809            0,
810            None,
811        )
812        .await;
813
814        sleep(Duration::from_secs(2)).await;
815
816        // Observer nodes have own_index set to MAX as a special value
817        assert_eq!(observer.context().own_index, AuthorityIndex::MAX);
818        assert_eq!(observer.context().committee.epoch(), 0);
819        assert_eq!(observer.context().committee.size(), 1);
820        assert!(!observer.context().is_validator());
821
822        observer.stop().await;
823    }
824
825    // TODO: build AuthorityFixture.
826    // Spins up a committee of authorities and an observer node that connects to authority 0.
827    // Verifies that the network is progressing, advancing rounds and commits. It also verifies
828    // that the Observer node is receiving blocks from the network.
829    #[rstest]
830    #[tokio::test(flavor = "current_thread")]
831    async fn test_authority_committee(
832        #[values(NetworkType::Tonic)] network_type: NetworkType,
833        #[values(5, 10)] gc_depth: u32,
834    ) {
835        telemetry_subscribers::init_for_testing();
836        let db_registry = Registry::new();
837        DBMetrics::init(RegistryService::new(db_registry));
838
839        const NUM_OF_AUTHORITIES: usize = 4;
840        let (committee, keypairs) = local_committee_and_keys(0, [1; NUM_OF_AUTHORITIES].to_vec());
841        let mut protocol_config = ConsensusProtocolConfig::for_testing();
842        protocol_config.set_gc_depth_for_testing(gc_depth);
843
844        let temp_dirs = (0..NUM_OF_AUTHORITIES)
845            .map(|_| TempDir::new().unwrap())
846            .collect::<Vec<_>>();
847
848        let mut commit_receivers = Vec::with_capacity(committee.size());
849        let mut authorities = Vec::with_capacity(committee.size());
850        let mut boot_counters = [0; NUM_OF_AUTHORITIES];
851
852        // Use a unique port based on gc_depth to avoid conflicts between parallel tests
853        let observer_server_port = 8900 + gc_depth as u16;
854
855        // Create authorities with observer server enabled for authority 0
856        let mut authority_0_network_key = None;
857        for (index, authority_info) in committee.authorities() {
858            let (authority, commit_receiver) = if index.value() == 0 {
859                // Save authority 0's network key for Observer connection
860                authority_0_network_key = Some(authority_info.network_key.clone());
861                // Enable observer server for authority 0
862                make_authority_with_observer_server(
863                    index,
864                    &temp_dirs[index.value()],
865                    committee.clone(),
866                    keypairs.clone(),
867                    network_type,
868                    boot_counters[index],
869                    protocol_config.clone(),
870                    Some(observer_server_port),
871                )
872                .await
873            } else {
874                make_authority(
875                    index,
876                    &temp_dirs[index.value()],
877                    committee.clone(),
878                    keypairs.clone(),
879                    network_type,
880                    boot_counters[index],
881                    protocol_config.clone(),
882                )
883                .await
884            };
885            boot_counters[index] += 1;
886            commit_receivers.push(commit_receiver);
887            authorities.push(authority);
888        }
889
890        // Create an Observer node that connects to authority 0
891        let observer_temp_dir = TempDir::new().unwrap();
892        let mut rng = StdRng::from_seed([99; 32]);
893        let observer_network_keypair = consensus_config::NetworkKeyPair::generate(&mut rng);
894
895        let observer_parameters = Parameters {
896            db_path: observer_temp_dir.path().to_path_buf(),
897            observer: ObserverParameters {
898                // Configure Observer to connect to authority 0
899                peers: vec![PeerRecord {
900                    public_key: authority_0_network_key
901                        .clone()
902                        .expect("Authority 0 network key should be set"),
903                    address: format!("/ip4/127.0.0.1/udp/{}", observer_server_port)
904                        .parse()
905                        .unwrap(),
906                }],
907                ..Default::default()
908            },
909            ..Default::default()
910        };
911
912        let (observer_commit_consumer, observer_commit_receiver) = CommitConsumerArgs::new(0, 0);
913        let observer = ConsensusAuthority::start(
914            network_type,
915            0,
916            committee.clone(),
917            observer_parameters,
918            protocol_config.clone(),
919            None, // No protocol keypair for observer
920            observer_network_keypair,
921            Arc::new(Clock::default()),
922            Arc::new(NoopTransactionVerifier {}),
923            None,
924            observer_commit_consumer,
925            Registry::new(),
926            0,
927            None,
928        )
929        .await;
930        // The relevant endpoints are now implemented for the synchronizer and commit_syncer components, so the Observer node should be able to catch up and
931        // fetch blocks beyond the latest ones that are fetched from the stream.
932        commit_receivers.push(observer_commit_receiver);
933
934        // Give Observer more time to connect and sync
935        sleep(Duration::from_secs(5)).await;
936
937        const NUM_TRANSACTIONS: u8 = 15;
938        let mut submitted_transactions = BTreeSet::<Vec<u8>>::new();
939        for i in 0..NUM_TRANSACTIONS {
940            let txn = vec![i; 16];
941            submitted_transactions.insert(txn.clone());
942            authorities[i as usize % authorities.len()]
943                .transaction_client()
944                .submit(vec![txn], Priority::Normal)
945                .await
946                .unwrap();
947        }
948
949        for receiver in &mut commit_receivers {
950            let mut expected_transactions = submitted_transactions.clone();
951            loop {
952                let committed_subdag =
953                    tokio::time::timeout(Duration::from_secs(1), receiver.recv())
954                        .await
955                        .unwrap()
956                        .unwrap();
957                for b in committed_subdag.blocks {
958                    for txn in b.transactions().iter().map(|t| t.data().to_vec()) {
959                        assert!(
960                            expected_transactions.remove(&txn),
961                            "Transaction not submitted or already seen: {:?}",
962                            txn
963                        );
964                    }
965                }
966                if expected_transactions.is_empty() {
967                    break;
968                }
969            }
970        }
971
972        // Stop authority 1.
973        let index = committee.to_authority_index(1).unwrap();
974        authorities.remove(index.value()).stop().await;
975        sleep(Duration::from_secs(10)).await;
976
977        // Restart authority 1 and let it run.
978        let (authority, commit_receiver) = make_authority(
979            index,
980            &temp_dirs[index.value()],
981            committee.clone(),
982            keypairs.clone(),
983            network_type,
984            boot_counters[index],
985            protocol_config.clone(),
986        )
987        .await;
988        boot_counters[index] += 1;
989        commit_receivers[index] = commit_receiver;
990        authorities.insert(index.value(), authority);
991        sleep(Duration::from_secs(10)).await;
992
993        // Verify that the Observer node is running
994        // TODO: The actual block processing for observers is not fully implemented yet
995        // for now we just verify that blocks are received and the number of received blocks is not far from
996        // the number of blocks sent by authority 0.
997        let observer_context = observer.context();
998        assert!(
999            observer_context.is_observer(),
1000            "It should be an observer node"
1001        );
1002
1003        // Get the total verified_blocks from authority 0 (sum across all sending authorities)
1004        let authority_0 = &authorities[0];
1005        let authority_0_context = authority_0.context();
1006        let mut authority_0_total_verified_blocks = 0;
1007
1008        // Sum verified_blocks from all authorities as seen by authority 0
1009        for (_, authority_info) in committee.authorities() {
1010            if let Ok(metric) = authority_0_context
1011                .metrics
1012                .node_metrics
1013                .verified_blocks
1014                .get_metric_with_label_values(&[&authority_info.hostname])
1015            {
1016                authority_0_total_verified_blocks += metric.get();
1017                println!(
1018                    "authority_info.hostname: {}, metric: {:?}",
1019                    authority_info.hostname, authority_0_total_verified_blocks
1020                );
1021            }
1022        }
1023
1024        let mut authority_0_total_proposed_blocks = 0;
1025        for force in [true, false] {
1026            if let Ok(metric) = authority_0_context
1027                .metrics
1028                .node_metrics
1029                .proposed_blocks
1030                .get_metric_with_label_values(&[&force.to_string()])
1031            {
1032                authority_0_total_proposed_blocks += metric.get();
1033            }
1034        }
1035
1036        authority_0_total_verified_blocks += authority_0_total_proposed_blocks;
1037
1038        // Sum verified_blocks from all authorities as seen by the observer
1039        let mut observer_received_blocks = 0;
1040        for (_, authority_info) in committee.authorities() {
1041            if let Ok(metric) = observer_context
1042                .metrics
1043                .node_metrics
1044                .verified_blocks
1045                .get_metric_with_label_values(&[&authority_info.hostname])
1046            {
1047                observer_received_blocks += metric.get();
1048            }
1049        }
1050
1051        // Compare the values - they should be related but might not be exactly equal
1052        // due to timing and the observer connecting mid-stream
1053        assert!(
1054            observer_received_blocks > 0,
1055            "Observer should have received at least some blocks, got: {}",
1056            observer_received_blocks
1057        );
1058
1059        println!(
1060            "authority_0_total_verified_blocks: {}, observer_received_blocks: {}",
1061            authority_0_total_verified_blocks, observer_received_blocks
1062        );
1063
1064        const TOLERANCE: u64 = 20;
1065        assert!(
1066            authority_0_total_verified_blocks - observer_received_blocks <= TOLERANCE,
1067            "The number of blocks received by the observer ({}) should be close to the number of blocks verified by authority 0 ({})",
1068            observer_received_blocks,
1069            authority_0_total_verified_blocks,
1070        );
1071
1072        // Stop observer first
1073        observer.stop().await;
1074
1075        // Stop all authorities and exit.
1076        for authority in authorities {
1077            authority.stop().await;
1078        }
1079    }
1080
1081    #[rstest]
1082    #[tokio::test(flavor = "current_thread")]
1083    async fn test_small_committee(
1084        #[values(NetworkType::Tonic)] network_type: NetworkType,
1085        #[values(1, 2, 3)] num_authorities: usize,
1086    ) {
1087        telemetry_subscribers::init_for_testing();
1088        let db_registry = Registry::new();
1089        DBMetrics::init(RegistryService::new(db_registry));
1090
1091        let (committee, keypairs) = local_committee_and_keys(0, vec![1; num_authorities]);
1092        let protocol_config = ConsensusProtocolConfig::for_testing();
1093
1094        let temp_dirs = (0..num_authorities)
1095            .map(|_| TempDir::new().unwrap())
1096            .collect::<Vec<_>>();
1097
1098        let mut output_receivers = Vec::with_capacity(committee.size());
1099        let mut authorities: Vec<ConsensusAuthority> = Vec::with_capacity(committee.size());
1100        let mut boot_counters = vec![0; num_authorities];
1101
1102        for (index, _authority_info) in committee.authorities() {
1103            let (authority, commit_receiver) = make_authority(
1104                index,
1105                &temp_dirs[index.value()],
1106                committee.clone(),
1107                keypairs.clone(),
1108                network_type,
1109                boot_counters[index],
1110                protocol_config.clone(),
1111            )
1112            .await;
1113            boot_counters[index] += 1;
1114            output_receivers.push(commit_receiver);
1115            authorities.push(authority);
1116        }
1117
1118        const NUM_TRANSACTIONS: u8 = 15;
1119        let mut submitted_transactions = BTreeSet::<Vec<u8>>::new();
1120        for i in 0..NUM_TRANSACTIONS {
1121            let txn = vec![i; 16];
1122            submitted_transactions.insert(txn.clone());
1123            authorities[i as usize % authorities.len()]
1124                .transaction_client()
1125                .submit(vec![txn], Priority::Normal)
1126                .await
1127                .unwrap();
1128        }
1129
1130        for receiver in &mut output_receivers {
1131            let mut expected_transactions = submitted_transactions.clone();
1132            loop {
1133                let committed_subdag =
1134                    tokio::time::timeout(Duration::from_secs(1), receiver.recv())
1135                        .await
1136                        .unwrap()
1137                        .unwrap();
1138                for b in committed_subdag.blocks {
1139                    for txn in b.transactions().iter().map(|t| t.data().to_vec()) {
1140                        assert!(
1141                            expected_transactions.remove(&txn),
1142                            "Transaction not submitted or already seen: {:?}",
1143                            txn
1144                        );
1145                    }
1146                }
1147                if expected_transactions.is_empty() {
1148                    break;
1149                }
1150            }
1151        }
1152
1153        // Stop authority 0.
1154        let index = committee.to_authority_index(0).unwrap();
1155        authorities.remove(index.value()).stop().await;
1156        sleep(Duration::from_secs(10)).await;
1157
1158        // Restart authority 0 and let it run.
1159        let (authority, commit_receiver) = make_authority(
1160            index,
1161            &temp_dirs[index.value()],
1162            committee.clone(),
1163            keypairs.clone(),
1164            network_type,
1165            boot_counters[index],
1166            protocol_config.clone(),
1167        )
1168        .await;
1169        boot_counters[index] += 1;
1170        output_receivers[index] = commit_receiver;
1171        authorities.insert(index.value(), authority);
1172        sleep(Duration::from_secs(10)).await;
1173
1174        // Stop all authorities and exit.
1175        for authority in authorities {
1176            authority.stop().await;
1177        }
1178    }
1179
1180    #[rstest]
1181    #[tokio::test(flavor = "current_thread")]
1182    async fn test_amnesia_recovery_success(#[values(5, 10)] gc_depth: u32) {
1183        telemetry_subscribers::init_for_testing();
1184        let db_registry = Registry::new();
1185        DBMetrics::init(RegistryService::new(db_registry));
1186
1187        const NUM_OF_AUTHORITIES: usize = 4;
1188        let (committee, keypairs) = local_committee_and_keys(0, [1; NUM_OF_AUTHORITIES].to_vec());
1189        let mut commit_receivers = vec![];
1190        let mut authorities = BTreeMap::new();
1191        let mut temp_dirs = BTreeMap::new();
1192        let mut boot_counters = [0; NUM_OF_AUTHORITIES];
1193
1194        let mut protocol_config = ConsensusProtocolConfig::for_testing();
1195        protocol_config.set_gc_depth_for_testing(gc_depth);
1196
1197        for (index, _authority_info) in committee.authorities() {
1198            let dir = TempDir::new().unwrap();
1199            let (authority, commit_receiver) = make_authority(
1200                index,
1201                &dir,
1202                committee.clone(),
1203                keypairs.clone(),
1204                NetworkType::Tonic,
1205                boot_counters[index],
1206                protocol_config.clone(),
1207            )
1208            .await;
1209            boot_counters[index] += 1;
1210            commit_receivers.push(commit_receiver);
1211            authorities.insert(index, authority);
1212            temp_dirs.insert(index, dir);
1213        }
1214
1215        // Now we take the receiver of authority 1 and we wait until we see at least one block committed from this authority
1216        // We wait until we see at least one committed block authored from this authority. That way we'll be 100% sure that
1217        // at least one block has been proposed and successfully received by a quorum of nodes.
1218        let index_1 = committee.to_authority_index(1).unwrap();
1219        'outer: while let Some(result) =
1220            timeout(Duration::from_secs(10), commit_receivers[index_1].recv())
1221                .await
1222                .expect("Timed out while waiting for at least one committed block from authority 1")
1223        {
1224            for block in result.blocks {
1225                if block.round() > GENESIS_ROUND && block.author() == index_1 {
1226                    break 'outer;
1227                }
1228            }
1229        }
1230
1231        // Stop authority 1 & 2.
1232        // * Authority 1 will be used to wipe out their DB and practically "force" the amnesia recovery.
1233        // * Authority 2 is stopped in order to simulate less than f+1 availability which will
1234        // make authority 1 retry during amnesia recovery until it has finally managed to successfully get back f+1 responses.
1235        // once authority 2 is up and running again.
1236        authorities.remove(&index_1).unwrap().stop().await;
1237        let index_2 = committee.to_authority_index(2).unwrap();
1238        authorities.remove(&index_2).unwrap().stop().await;
1239        sleep(Duration::from_secs(5)).await;
1240
1241        // Authority 1: create a new directory to simulate amnesia. The node will start having participated previously
1242        // to consensus but now will attempt to synchronize the last own block and recover from there. It won't be able
1243        // to do that successfully as authority 2 is still down.
1244        let dir = TempDir::new().unwrap();
1245        // We do reset the boot counter for this one to simulate a "binary" restart
1246        boot_counters[index_1] = 0;
1247        let (authority, mut commit_receiver) = make_authority(
1248            index_1,
1249            &dir,
1250            committee.clone(),
1251            keypairs.clone(),
1252            NetworkType::Tonic,
1253            boot_counters[index_1],
1254            protocol_config.clone(),
1255        )
1256        .await;
1257        boot_counters[index_1] += 1;
1258        authorities.insert(index_1, authority);
1259        temp_dirs.insert(index_1, dir);
1260        sleep(Duration::from_secs(5)).await;
1261
1262        // Now spin up authority 2 using its earlier directly - so no amnesia recovery should be forced here.
1263        // Authority 1 should be able to recover from amnesia successfully.
1264        let (authority, _commit_receiver) = make_authority(
1265            index_2,
1266            &temp_dirs[&index_2],
1267            committee.clone(),
1268            keypairs,
1269            NetworkType::Tonic,
1270            boot_counters[index_2],
1271            protocol_config.clone(),
1272        )
1273        .await;
1274        boot_counters[index_2] += 1;
1275        authorities.insert(index_2, authority);
1276        sleep(Duration::from_secs(5)).await;
1277
1278        // We wait until we see at least one committed block authored from this authority
1279        'outer: while let Some(result) = commit_receiver.recv().await {
1280            for block in result.blocks {
1281                if block.round() > GENESIS_ROUND && block.author() == index_1 {
1282                    break 'outer;
1283                }
1284            }
1285        }
1286
1287        // Stop all authorities and exit.
1288        for (_, authority) in authorities {
1289            authority.stop().await;
1290        }
1291    }
1292
1293    // TODO: create a fixture
1294    async fn make_authority(
1295        index: AuthorityIndex,
1296        db_dir: &TempDir,
1297        committee: Committee,
1298        keypairs: Vec<(NetworkKeyPair, ProtocolKeyPair)>,
1299        network_type: NetworkType,
1300        boot_counter: u64,
1301        protocol_config: ConsensusProtocolConfig,
1302    ) -> (ConsensusAuthority, UnboundedReceiver<CommittedSubDag>) {
1303        make_authority_with_observer_server(
1304            index,
1305            db_dir,
1306            committee,
1307            keypairs,
1308            network_type,
1309            boot_counter,
1310            protocol_config,
1311            None, // No observer server port
1312        )
1313        .await
1314    }
1315
1316    async fn make_authority_with_observer_server(
1317        index: AuthorityIndex,
1318        db_dir: &TempDir,
1319        committee: Committee,
1320        keypairs: Vec<(NetworkKeyPair, ProtocolKeyPair)>,
1321        network_type: NetworkType,
1322        boot_counter: u64,
1323        protocol_config: ConsensusProtocolConfig,
1324        observer_server_port: Option<u16>,
1325    ) -> (ConsensusAuthority, UnboundedReceiver<CommittedSubDag>) {
1326        let registry = Registry::new();
1327
1328        // Cache less blocks to exercise commit sync.
1329        let mut parameters = Parameters {
1330            db_path: db_dir.path().to_path_buf(),
1331            dag_state_cached_rounds: 5,
1332            commit_sync_parallel_fetches: 2,
1333            commit_sync_batch_size: 3,
1334            sync_last_known_own_block_timeout: Duration::from_millis(2_000),
1335            ..Default::default()
1336        };
1337
1338        // Enable observer server if port is provided
1339        if let Some(port) = observer_server_port {
1340            parameters.observer.server_port = Some(port);
1341        }
1342
1343        let txn_verifier = NoopTransactionVerifier {};
1344
1345        let protocol_keypair = keypairs[index].1.clone();
1346        let network_keypair = keypairs[index].0.clone();
1347
1348        let (commit_consumer, commit_receiver) = CommitConsumerArgs::new(0, 0);
1349
1350        let authority = ConsensusAuthority::start(
1351            network_type,
1352            0,
1353            committee,
1354            parameters,
1355            protocol_config,
1356            Some(protocol_keypair),
1357            network_keypair,
1358            Arc::new(Clock::default()),
1359            Arc::new(txn_verifier),
1360            None,
1361            commit_consumer,
1362            registry,
1363            boot_counter,
1364            None,
1365        )
1366        .await;
1367
1368        (authority, commit_receiver)
1369    }
1370}