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