1use 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#[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 protocol_keypair: Option<ProtocolKeyPair>,
71 network_keypair: NetworkKeyPair,
72 clock: Arc<Clock>,
73 transaction_verifier: Arc<dyn TransactionVerifier>,
74 transaction_pool: Option<Arc<dyn TransactionPool>>,
77 commit_consumer: CommitConsumerArgs,
78 registry: Registry,
79 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
149enum 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 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 observer_service: Option<Arc<ObserverService>>,
185 network_manager: N,
186}
187
188impl<N> AuthorityNode<N>
189where
190 N: NetworkManager,
191{
192 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 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 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 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 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 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 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 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 network_manager
454 .start_validator_server(authority_service.clone())
455 .await;
456
457 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 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 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 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 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 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 core_thread_handle.stop().await;
616 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 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 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 self.network_manager.update_peer_address(peer, address);
688
689 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 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 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, 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 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 #[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 let observer_server_port = 8900 + gc_depth as u16;
854
855 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 authority_0_network_key = Some(authority_info.network_key.clone());
861 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 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 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, 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 commit_receivers.push(observer_commit_receiver);
933
934 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 let index = committee.to_authority_index(1).unwrap();
974 authorities.remove(index.value()).stop().await;
975 sleep(Duration::from_secs(10)).await;
976
977 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 let observer_context = observer.context();
998 assert!(
999 observer_context.is_observer(),
1000 "It should be an observer node"
1001 );
1002
1003 let authority_0 = &authorities[0];
1005 let authority_0_context = authority_0.context();
1006 let mut authority_0_total_verified_blocks = 0;
1007
1008 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 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 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 observer.stop().await;
1074
1075 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 let index = committee.to_authority_index(0).unwrap();
1155 authorities.remove(index.value()).stop().await;
1156 sleep(Duration::from_secs(10)).await;
1157
1158 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 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 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 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 let dir = TempDir::new().unwrap();
1245 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 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 '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 for (_, authority) in authorities {
1289 authority.stop().await;
1290 }
1291 }
1292
1293 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, )
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 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 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}