1use crate::authority::authority_per_epoch_store::AuthorityPerEpochStore;
4use crate::consensus_adapter::{BlockStatusReceiver, ConsensusClient};
5use crate::consensus_handler::{ConsensusHandlerInitializer, MysticetiConsensusHandler};
6use crate::consensus_transaction_pool::{
7 ConsensusTransactionPool, TransactionPoolClient, TransactionPoolContext,
8};
9use crate::consensus_validator::SuiTxValidator;
10use crate::mysticeti_adapter::LazyMysticetiClient;
11use arc_swap::ArcSwapOption;
12use async_trait::async_trait;
13use consensus_config::{
14 ChainType, Committee, ConsensusProtocolConfig, NetworkKeyPair,
15 NetworkPublicKey as ConsensusNetworkPublicKey, Parameters, ProtocolKeyPair, Stake,
16};
17use consensus_core::{
18 Clock, CommitConsumerArgs, CommitConsumerMonitor, CommitIndex, ConsensusAuthority, NetworkType,
19 RandomnessSignatureHandler, TransactionPool, storage::rocksdb_store::RocksDBStore,
20};
21use core::panic;
22use fastcrypto::encoding::{Encoding, Hex};
23use fastcrypto::traits::KeyPair as _;
24use mysten_common::debug_fatal;
25use mysten_metrics::{RegistryID, RegistryService};
26use mysten_network::Multiaddr;
27use prometheus::{
28 IntGauge, IntGaugeVec, Registry, register_int_gauge_vec_with_registry,
29 register_int_gauge_with_registry,
30};
31use std::collections::BTreeMap;
32use std::path::PathBuf;
33use std::sync::Arc;
34use std::time::{Duration, Instant};
35use sui_config::{ConsensusConfig, NodeConfig};
36use sui_network::endpoint_manager::{AddressSource, ConsensusAddressUpdater};
37use sui_protocol_config::{Chain, ProtocolConfig, ProtocolVersion};
38use sui_types::crypto::NetworkPublicKey;
39use sui_types::error::{SuiErrorKind, SuiResult};
40use sui_types::messages_consensus::{ConsensusPosition, ConsensusTransaction};
41use sui_types::node_role::NodeRole;
42use sui_types::{
43 committee::EpochId, sui_system_state::epoch_start_sui_system_state::EpochStartSystemStateTrait,
44};
45use tokio::sync::{Mutex, broadcast};
46use tokio::time::{sleep, timeout};
47use tracing::{error, info};
48
49#[cfg(test)]
50#[path = "../unit_tests/consensus_manager_tests.rs"]
51pub mod consensus_manager_tests;
52
53#[derive(PartialEq)]
54enum Running {
55 True(EpochId, ProtocolVersion),
56 False,
57}
58
59struct AddressOverridesMap {
62 map: BTreeMap<
64 ConsensusNetworkPublicKey,
65 BTreeMap<sui_network::endpoint_manager::AddressSource, Vec<Multiaddr>>,
66 >,
67}
68
69impl AddressOverridesMap {
70 pub fn new() -> Self {
71 Self {
72 map: BTreeMap::new(),
73 }
74 }
75
76 pub fn insert(
77 &mut self,
78 network_pubkey: ConsensusNetworkPublicKey,
79 source: sui_network::endpoint_manager::AddressSource,
80 addresses: Vec<Multiaddr>,
81 ) {
82 self.map
83 .entry(network_pubkey)
84 .or_default()
85 .insert(source, addresses);
86 }
87
88 pub fn remove(
89 &mut self,
90 network_pubkey: ConsensusNetworkPublicKey,
91 source: sui_network::endpoint_manager::AddressSource,
92 ) {
93 self.map
94 .entry(network_pubkey.clone())
95 .or_default()
96 .remove(&source);
97
98 if self.map.get(&network_pubkey.clone()).unwrap().is_empty() {
100 self.map.remove(&network_pubkey);
101 }
102 }
103
104 pub fn get_highest_priority_source_and_address(
108 &self,
109 network_pubkey: ConsensusNetworkPublicKey,
110 ) -> Option<(sui_network::endpoint_manager::AddressSource, Multiaddr)> {
111 self.map
112 .get(&network_pubkey)
113 .and_then(|sources| sources.first_key_value())
114 .and_then(|(source, addresses)| {
115 addresses.first().cloned().map(|address| (*source, address))
116 })
117 }
118
119 pub fn get_all_highest_priority_addresses(
120 &self,
121 ) -> Vec<(ConsensusNetworkPublicKey, Multiaddr)> {
122 let mut result = Vec::new();
123
124 for (network_pubkey, sources) in self.map.iter() {
125 if let Some((_source, addresses)) = sources.first_key_value()
126 && let Some(address) = addresses.first()
127 {
128 result.push((network_pubkey.clone(), address.clone()));
129 }
130 }
131 result
132 }
133}
134
135fn apply_v3_threshold_overrides(committee: Committee) -> Committee {
141 let malicious_stake: Stake = std::env::var("SUI_CONSENSUS_V3_MALICIOUS_STAKE")
142 .ok()
143 .and_then(|s| s.parse().ok())
144 .unwrap_or(1_250);
145 let crash_stake: Stake = std::env::var("SUI_CONSENSUS_V3_CRASH_STAKE")
146 .ok()
147 .and_then(|s| s.parse().ok())
148 .unwrap_or(1_250);
149 info!(
150 "consensus_manager: applying v3 committee thresholds \
151 (malicious_stake={malicious_stake}, crash_stake={crash_stake})"
152 );
153 Committee::new_v3(
154 committee.epoch(),
155 committee.authorities_slice().to_vec(),
156 malicious_stake,
157 crash_stake,
158 )
159}
160
161fn to_consensus_protocol_config(config: &ProtocolConfig) -> ConsensusProtocolConfig {
162 let chain_type = match config.chain() {
163 Chain::Mainnet => ChainType::Mainnet,
164 Chain::Testnet => ChainType::Testnet,
165 Chain::Unknown => ChainType::Unknown,
166 };
167 ConsensusProtocolConfig::new(
168 config.version.as_u64(),
169 chain_type,
170 config.max_transaction_size_bytes(),
171 config.max_transactions_in_block_bytes(),
172 config.max_num_transactions_in_block(),
173 config.gc_depth(),
174 config.consensus_slim_block_propagation(),
175 true,
176 config.mysticeti_num_leaders_per_round(),
177 config.consensus_bad_nodes_stake_threshold(),
178 false,
179 300,
180 12,
181 )
182}
183
184pub struct ConsensusManager {
187 consensus_config: ConsensusConfig,
188 protocol_keypair: Option<ProtocolKeyPair>,
189 network_keypair: NetworkKeyPair,
190 storage_base_path: PathBuf,
191 metrics: Arc<ConsensusManagerMetrics>,
192 registry_service: RegistryService,
193 authority: ArcSwapOption<(ConsensusAuthority, RegistryID)>,
194
195 client: Arc<LazyMysticetiClient>,
198 consensus_client: Arc<UpdatableConsensusClient>,
199 transaction_pool_context: Option<Arc<TransactionPoolContext>>,
200 transaction_pool: ArcSwapOption<ConsensusTransactionPool>,
201
202 consensus_handler: Mutex<Option<MysticetiConsensusHandler>>,
203
204 #[cfg(test)]
205 pub(crate) consumer_monitor: ArcSwapOption<CommitConsumerMonitor>,
206 #[cfg(not(test))]
207 consumer_monitor: ArcSwapOption<CommitConsumerMonitor>,
208 consumer_monitor_sender: broadcast::Sender<Arc<CommitConsumerMonitor>>,
209
210 running: Mutex<Running>,
211
212 #[cfg(test)]
213 pub(crate) boot_counter: Mutex<u64>,
214 #[cfg(not(test))]
215 boot_counter: Mutex<u64>,
216
217 address_overrides: parking_lot::Mutex<AddressOverridesMap>,
220}
221
222impl ConsensusManager {
223 pub fn new(
224 node_config: &NodeConfig,
225 consensus_config: &ConsensusConfig,
226 registry_service: &RegistryService,
227 consensus_client: Arc<UpdatableConsensusClient>,
228 transaction_pool_context: Option<Arc<TransactionPoolContext>>,
229 node_role: NodeRole,
230 ) -> Self {
231 let metrics = Arc::new(ConsensusManagerMetrics::new(
232 ®istry_service.default_registry(),
233 ));
234 let client = Arc::new(LazyMysticetiClient::new());
235 let (consumer_monitor_sender, _) = broadcast::channel(1);
236 let protocol_keypair = if node_role.is_validator() {
237 Some(ProtocolKeyPair::new(node_config.worker_key_pair().copy()))
238 } else {
239 None
240 };
241 Self {
242 consensus_config: consensus_config.clone(),
243 protocol_keypair,
244 network_keypair: NetworkKeyPair::new(node_config.network_key_pair().copy()),
245 storage_base_path: consensus_config.db_path().to_path_buf(),
246 metrics,
247 registry_service: registry_service.clone(),
248 authority: ArcSwapOption::empty(),
249 client,
250 consensus_client,
251 transaction_pool_context,
252 transaction_pool: ArcSwapOption::empty(),
253 consensus_handler: Mutex::new(None),
254 consumer_monitor: ArcSwapOption::empty(),
255 consumer_monitor_sender,
256 running: Mutex::new(Running::False),
257 boot_counter: Mutex::new(0),
258 address_overrides: parking_lot::Mutex::new(AddressOverridesMap::new()),
259 }
260 }
261
262 pub async fn start(
263 &self,
264 node_config: &NodeConfig,
265 epoch_store: Arc<AuthorityPerEpochStore>,
266 consensus_handler_initializer: ConsensusHandlerInitializer,
267 tx_validator: SuiTxValidator,
268 randomness_signature_handler: Option<Arc<dyn RandomnessSignatureHandler>>,
269 ) {
270 let epoch = epoch_store.epoch();
271 let protocol_config = epoch_store.protocol_config();
272 let consensus_protocol_config = to_consensus_protocol_config(protocol_config);
273 let system_state = epoch_store.epoch_start_state();
274 let committee = if consensus_protocol_config.enable_v3() {
275 apply_v3_threshold_overrides(system_state.get_consensus_committee())
276 } else {
277 system_state.get_consensus_committee()
278 };
279
280 let start_time = Instant::now();
282 let mut running = self.running.lock().await;
283 if let Running::True(running_epoch, running_version) = *running {
284 error!(
285 "Consensus is already Running for epoch {running_epoch:?} & protocol version {running_version:?} - shutdown first before starting",
286 );
287 return;
288 }
289 *running = Running::True(epoch, protocol_config.version);
290
291 info!(
292 "Starting up consensus for epoch {epoch:?} & protocol version {:?}",
293 protocol_config.version
294 );
295
296 let is_validator = epoch_store.is_validator();
297 if is_validator && self.protocol_keypair.is_none() {
298 debug_fatal!("validator epoch {epoch} started without a protocol keypair");
301 }
302 let pool_context = self
303 .transaction_pool_context
304 .as_ref()
305 .filter(|_| is_validator);
306 let transaction_pool: Option<Arc<dyn TransactionPool>> = if let Some(context) = pool_context
307 {
308 let config = &node_config.consensus_transaction_pool;
309 let pool = Arc::new(ConsensusTransactionPool::new(
310 epoch_store.clone(),
311 config.max_pending_transactions(&self.consensus_config),
312 context.metrics().clone(),
313 context.adapter_metrics().clone(),
314 ));
315 context.set_active(epoch, pool.clone());
316 self.transaction_pool.store(Some(pool.clone()));
317 self.consensus_client
318 .set(Arc::new(TransactionPoolClient::new(context.clone())));
319 Some(pool)
320 } else {
321 if let Some(context) = &self.transaction_pool_context {
322 context.set_unavailable(epoch);
323 }
324 self.consensus_client.set(self.client.clone());
325 None
326 };
327
328 let consensus_config = node_config
329 .consensus_config()
330 .expect("consensus_config should exist");
331
332 let parameters = Parameters {
333 db_path: self.get_store_path(epoch),
334 listen_address_override: consensus_config.listen_address.clone(),
335 ..consensus_config.parameters.clone().unwrap_or_default()
336 };
337
338 let registry = Registry::new_custom(Some("consensus".to_string()), None).unwrap();
339
340 let consensus_handler = consensus_handler_initializer.new_consensus_handler();
341
342 let num_prior_commits = protocol_config.consensus_num_requested_prior_commits_at_startup();
343 let last_processed_commit_index =
344 consensus_handler.last_processed_subdag_index() as CommitIndex;
345 let replay_after_commit_index =
346 last_processed_commit_index.saturating_sub(num_prior_commits);
347
348 let (commit_consumer, commit_receiver) =
349 CommitConsumerArgs::new(replay_after_commit_index, last_processed_commit_index);
350 let monitor = commit_consumer.monitor();
351
352 let handler = MysticetiConsensusHandler::new(
354 last_processed_commit_index,
355 consensus_handler,
356 commit_receiver,
357 monitor.clone(),
358 );
359 let mut consensus_handler = self.consensus_handler.lock().await;
360 *consensus_handler = Some(handler);
361
362 let participated_on_previous_run =
366 if let Some(previous_monitor) = self.consumer_monitor.swap(Some(monitor.clone())) {
367 previous_monitor.highest_handled_commit() > 0
368 } else {
369 false
370 };
371
372 let mut boot_counter = self.boot_counter.lock().await;
377 if participated_on_previous_run {
378 *boot_counter += 1;
379 } else {
380 info!(
381 "Node has not participated in previous epoch consensus. Boot counter ({}) will not increment.",
382 *boot_counter
383 );
384 }
385
386 let authority = ConsensusAuthority::start(
387 NetworkType::Tonic,
388 epoch_store.epoch_start_config().epoch_start_timestamp_ms(),
389 committee.clone(),
390 parameters.clone(),
391 consensus_protocol_config,
392 self.protocol_keypair.clone(),
393 self.network_keypair.clone(),
394 Arc::new(Clock::default()),
395 Arc::new(tx_validator.clone()),
396 transaction_pool,
397 commit_consumer,
398 registry.clone(),
399 *boot_counter,
400 randomness_signature_handler,
401 )
402 .await;
403 let client = pool_context
404 .is_none()
405 .then(|| authority.transaction_client());
406
407 let registry_id = self.registry_service.add(registry.clone());
408
409 let registered_authority = Arc::new((authority, registry_id));
410 self.authority.swap(Some(registered_authority.clone()));
411
412 let highest_priority_addresses = self
414 .address_overrides
415 .lock()
416 .get_all_highest_priority_addresses();
417 for (network_pubkey, address) in highest_priority_addresses {
418 registered_authority
419 .0
420 .update_peer_address(network_pubkey, Some(address.clone()));
421 }
422
423 if let Some(client) = client {
425 self.client.set(client);
426 }
427
428 let _ = self.consumer_monitor_sender.send(monitor);
430
431 let elapsed = start_time.elapsed().as_secs_f64();
432 self.metrics.start_latency.set(elapsed as i64);
433
434 tracing::info!(
435 "Started consensus for epoch {} & protocol version {:?} completed - took {} seconds",
436 epoch,
437 protocol_config.version,
438 elapsed
439 );
440 }
441
442 pub async fn shutdown(&self) {
443 info!("Shutting down consensus ...");
444
445 let start_time = Instant::now();
447 let mut running = self.running.lock().await;
448 let (shutdown_epoch, shutdown_version) = match *running {
449 Running::True(epoch, version) => {
450 tracing::info!(
451 "Shutting down consensus for epoch {epoch:?} & protocol version {version:?}"
452 );
453 *running = Running::False;
454 (epoch, version)
455 }
456 Running::False => {
457 error!("Consensus shutdown was called but consensus is not running");
458 return;
459 }
460 };
461
462 let pool = self.transaction_pool.swap(None);
464 if let Some(pool) = &pool {
465 pool.close();
466 }
467 self.client.clear();
468
469 let r = self.authority.swap(None).unwrap();
471 let Ok((authority, registry_id)) = Arc::try_unwrap(r) else {
472 panic!("Failed to retrieve the Mysticeti authority");
473 };
474
475 authority.stop().await;
477
478 let mut consensus_handler = self.consensus_handler.lock().await;
480 if let Some(mut handler) = consensus_handler.take() {
481 handler.abort().await;
482 }
483
484 self.registry_service.remove(registry_id);
486
487 if pool.is_none() {
488 self.consensus_client.clear();
489 }
490
491 let elapsed = start_time.elapsed().as_secs_f64();
492 self.metrics.shutdown_latency.set(elapsed as i64);
493
494 tracing::info!(
495 "Consensus stopped for epoch {shutdown_epoch:?} & protocol version {shutdown_version:?} is complete - took {} seconds",
496 elapsed
497 );
498 }
499
500 pub async fn is_running(&self) -> bool {
501 let running = self.running.lock().await;
502 matches!(*running, Running::True(_, _))
503 }
504
505 pub fn replay_waiter(&self) -> ReplayWaiter {
506 let consumer_monitor_receiver = self.consumer_monitor_sender.subscribe();
507 ReplayWaiter::new(consumer_monitor_receiver)
508 }
509
510 pub fn get_storage_base_path(&self) -> PathBuf {
511 self.consensus_config.db_path().to_path_buf()
512 }
513
514 pub fn consensus_store(&self) -> Option<Arc<RocksDBStore>> {
515 self.authority.load().as_ref().map(|a| a.0.store())
516 }
517
518 pub fn address_overrides_snapshot(
519 &self,
520 ) -> BTreeMap<
521 ConsensusNetworkPublicKey,
522 BTreeMap<sui_network::endpoint_manager::AddressSource, Vec<Multiaddr>>,
523 > {
524 self.address_overrides.lock().map.clone()
525 }
526
527 fn get_store_path(&self, epoch: EpochId) -> PathBuf {
528 let mut store_path = self.storage_base_path.clone();
529 store_path.push(format!("{}", epoch));
530 store_path
531 }
532}
533
534impl Drop for ConsensusManager {
535 fn drop(&mut self) {
536 if let Some(pool) = self.transaction_pool.swap(None) {
539 pool.close();
540 }
541 }
542}
543
544impl ConsensusAddressUpdater for ConsensusManager {
546 fn update_address(
547 &self,
548 network_pubkey: NetworkPublicKey,
549 source: sui_network::endpoint_manager::AddressSource,
550 addresses: Vec<Multiaddr>,
551 ) -> SuiResult<()> {
552 let network_pubkey = ConsensusNetworkPublicKey::new(network_pubkey.clone());
554
555 let highest_priority = {
557 let mut address_overrides = self.address_overrides.lock();
558
559 if addresses.is_empty() {
560 address_overrides.remove(network_pubkey.clone(), source);
561 } else {
562 address_overrides.insert(network_pubkey.clone(), source, addresses.clone());
563 }
564
565 address_overrides.get_highest_priority_source_and_address(network_pubkey.clone())
566 };
567 self.metrics.set_active_address_source(
568 &Hex::encode(network_pubkey.to_bytes()),
569 highest_priority.as_ref().map(|(source, _)| *source),
570 );
571
572 let address_to_apply = highest_priority.map(|(_, address)| address);
574 if let Some(authority) = self.authority.load_full() {
575 authority
576 .0
577 .update_peer_address(network_pubkey, address_to_apply);
578 Ok(())
579 } else {
580 info!(
581 "Consensus authority node is not running, address update persisted for peer {:?} from source {:?} and will be applied on next start",
582 network_pubkey, source
583 );
584 Err(SuiErrorKind::GenericAuthorityError {
585 error: "Consensus authority node is not running. Can not apply address update"
586 .to_string(),
587 }
588 .into())
589 }
590 }
591}
592
593#[derive(Default)]
596pub struct UpdatableConsensusClient {
597 client: ArcSwapOption<Arc<dyn ConsensusClient>>,
599}
600
601impl UpdatableConsensusClient {
602 pub fn new() -> Self {
603 Self {
604 client: ArcSwapOption::empty(),
605 }
606 }
607
608 async fn get(&self) -> Arc<Arc<dyn ConsensusClient>> {
609 const START_TIMEOUT: Duration = Duration::from_secs(300);
610 const RETRY_INTERVAL: Duration = Duration::from_millis(100);
611 if let Ok(client) = timeout(START_TIMEOUT, async {
612 loop {
613 let Some(client) = self.client.load_full() else {
614 sleep(RETRY_INTERVAL).await;
615 continue;
616 };
617 return client;
618 }
619 })
620 .await
621 {
622 return client;
623 }
624
625 panic!(
626 "Timed out after {:?} waiting for Consensus to start!",
627 START_TIMEOUT,
628 );
629 }
630
631 pub fn set(&self, client: Arc<dyn ConsensusClient>) {
632 self.client.store(Some(Arc::new(client)));
633 }
634
635 pub fn clear(&self) {
636 self.client.store(None);
637 }
638}
639
640#[async_trait]
641impl ConsensusClient for UpdatableConsensusClient {
642 async fn submit(
643 &self,
644 transactions: &[ConsensusTransaction],
645 epoch_store: &Arc<AuthorityPerEpochStore>,
646 ) -> SuiResult<(Vec<ConsensusPosition>, BlockStatusReceiver)> {
647 let client = self.get().await;
648 client.submit(transactions, epoch_store).await
649 }
650}
651
652pub struct ReplayWaiter {
654 consumer_monitor_receiver: broadcast::Receiver<Arc<CommitConsumerMonitor>>,
655}
656
657impl ReplayWaiter {
658 pub(crate) fn new(
659 consumer_monitor_receiver: broadcast::Receiver<Arc<CommitConsumerMonitor>>,
660 ) -> Self {
661 Self {
662 consumer_monitor_receiver,
663 }
664 }
665
666 pub(crate) async fn wait_for_replay(mut self) {
667 loop {
668 info!("Waiting for consensus to start replaying ...");
669 let Ok(monitor) = self.consumer_monitor_receiver.recv().await else {
670 continue;
671 };
672 info!("Waiting for consensus handler to finish replaying ...");
673 monitor
674 .replay_to_consumer_last_processed_commit_complete()
675 .await;
676 break;
677 }
678 }
679}
680
681impl Clone for ReplayWaiter {
682 fn clone(&self) -> Self {
683 Self {
684 consumer_monitor_receiver: self.consumer_monitor_receiver.resubscribe(),
685 }
686 }
687}
688
689pub struct ConsensusManagerMetrics {
690 start_latency: IntGauge,
691 shutdown_latency: IntGauge,
692 active_address_source: IntGaugeVec,
693}
694
695impl ConsensusManagerMetrics {
696 pub fn new(registry: &Registry) -> Self {
697 Self {
698 start_latency: register_int_gauge_with_registry!(
699 "consensus_manager_start_latency",
700 "The latency of starting up consensus nodes",
701 registry,
702 )
703 .unwrap(),
704 shutdown_latency: register_int_gauge_with_registry!(
705 "consensus_manager_shutdown_latency",
706 "The latency of shutting down consensus nodes",
707 registry,
708 )
709 .unwrap(),
710 active_address_source: register_int_gauge_vec_with_registry!(
711 "consensus_active_address_source",
712 "Active consensus address source per committee peer, encoded as the gauge \
713 value: 0=committee (no override active; the on-chain committee address is in \
714 use), 1=admin, 2=config, 3=discovery, 4=seed, 5=chain (override priority \
715 highest to lowest). One series per peer; `peer_id` is the full hex consensus \
716 network public key.",
717 &["peer_id"],
718 registry,
719 )
720 .unwrap(),
721 }
722 }
723
724 fn set_active_address_source(
728 &self,
729 peer_id: &str,
730 active: Option<sui_network::endpoint_manager::AddressSource>,
731 ) {
732 let code = active.map_or(
733 AddressSource::DEFAULT_ADDRESS_SOURCE_CODE,
734 sui_network::endpoint_manager::AddressSource::metric_code,
735 );
736 self.active_address_source
737 .with_label_values(&[peer_id])
738 .set(code);
739 }
740}