Skip to main content

sui_core/consensus_manager/
mod.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3use 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
59/// Stores address updates that should be persisted across epoch changes.
60/// We store the consensus NetworkPublicKey to avoid repeated conversions.
61struct AddressOverridesMap {
62    // We store the AddressSource on a BTreeMap as it helps accessing the keys in priority order.
63    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 no sources remain for this peer, remove the peer entry entirely
99        if self.map.get(&network_pubkey.clone()).unwrap().is_empty() {
100            self.map.remove(&network_pubkey);
101        }
102    }
103
104    /// Returns the highest-priority active override `(source, address)` for the
105    /// peer, or `None` when no override is installed (the on-chain committee
106    /// address is in use).
107    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
135/// Rebuilds the consensus `Committee` with Mysticeti v3 threshold parameters.
136/// `malicious_stake` and `crash_stake` come from env vars with reference-budget
137/// defaults (`f = c = 1250`); the nominal `threshold_total_stake = 5f + 3c + 1`
138/// is derived inside `Committee::new_v3`. This is a temporary iteration knob
139/// until the thresholds are promoted into `ProtocolConfig`.
140fn 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        /* transaction_voting_enabled */ true,
176        config.mysticeti_num_leaders_per_round(),
177        config.consensus_bad_nodes_stake_threshold(),
178        /* enable_v3 */ false,
179        /* leader_schedule_window_size */ 300,
180        /* leader_schedule_update_interval */ 12,
181    )
182}
183
184/// Used by Sui to start consensus protocol for each epoch.
185/// Supports both validator mode (with protocol keypair) and observer mode (without).
186pub 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    // Use a shared lazy Mysticeti client so we can update the internal Mysticeti
196    // client that gets created for every new epoch.
197    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    // Persistent storage for address updates across epoch changes.
218    // Keyed by NetworkPublicKey and then by AddressSource.
219    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            &registry_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        // Ensure start() is not called twice.
281        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            // The manager was built for a non-validator role and reused across a
299            // promotion to validator; consensus cannot sign proposals in this state.
300            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        // Spin up the new Mysticeti consensus handler to listen for committed sub dags, before starting authority.
353        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        // If there is a previous consumer monitor, it indicates that the consensus engine has been restarted, due to an epoch change. However, that on its
363        // own doesn't tell us much whether it participated on an active epoch or an old one. We need to check if it has handled any commits to determine this.
364        // If indeed any commits did happen, then we assume that node did participate on previous run.
365        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        // Increment the boot counter only if the consensus successfully participated in the previous run.
373        // This is typical during normal epoch changes, where the node restarts as expected, and the boot counter is incremented to prevent amnesia recovery on the next start.
374        // If the node is recovering from a restore process and catching up across multiple epochs, it won't handle any commits until it reaches the last active epoch.
375        // In this scenario, we do not increment the boot counter, as we need amnesia recovery to run.
376        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        // Reapply all stored address updates to the new consensus instance.
413        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        // Initialize the client to send transactions to this Mysticeti instance.
424        if let Some(client) = client {
425            self.client.set(client);
426        }
427
428        // Send the consumer monitor to the replay waiter.
429        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        // Ensure shutdown() is called on a running consensus and get the epoch/version info.
446        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        // Stop consensus submissions.
463        let pool = self.transaction_pool.swap(None);
464        if let Some(pool) = &pool {
465            pool.close();
466        }
467        self.client.clear();
468
469        // swap with empty to ensure there is no other reference to authority and we can safely do Arc unwrap
470        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        // shutdown the authority and wait for it
476        authority.stop().await;
477
478        // drop the old consensus handler to force stop any underlying task running.
479        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        // unregister the registry id
485        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        // Abrupt node teardown can bypass async shutdown; explicitly disarm any
537        // pending pool acknowledgements before the manager's fields are dropped.
538        if let Some(pool) = self.transaction_pool.swap(None) {
539            pool.close();
540        }
541    }
542}
543
544// Implementing the interface so we can update the consensus peer addresses when requested.
545impl 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        // Convert to consensus network public key once
553        let network_pubkey = ConsensusNetworkPublicKey::new(network_pubkey.clone());
554
555        // Determine which override (if any) should be used after this update.
556        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        // Apply the update to running consensus if it exists
573        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/// A ConsensusClient that can be updated internally at any time. This usually happening during epoch
594/// change where a client is set after the new consensus is started for the new epoch.
595#[derive(Default)]
596pub struct UpdatableConsensusClient {
597    // An extra layer of Arc<> is needed as required by ArcSwapAny.
598    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
652/// Waits for consensus to finish replaying at consensus handler.
653pub 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    /// Records the active consensus address source for a committee peer as a single
725    /// per-peer gauge whose value is the source's `metric_code`. `active =
726    /// Some(source)` means an override is installed; `None` means no override.
727    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}