Skip to main content

sui_node/
lib.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4use anemo::Network;
5use anemo::PeerId;
6use anemo_tower::callback::CallbackLayer;
7use anemo_tower::trace::DefaultMakeSpan;
8use anemo_tower::trace::DefaultOnFailure;
9use anemo_tower::trace::TraceLayer;
10use anyhow::Context;
11use anyhow::Result;
12use anyhow::anyhow;
13use arc_swap::ArcSwap;
14use fastcrypto_zkp::bn254::zk_login::JwkId;
15use fastcrypto_zkp::bn254::zk_login::OIDCProvider;
16use futures::future::BoxFuture;
17use mysten_common::debug_fatal;
18use mysten_common::in_test_configuration;
19use prometheus::Registry;
20use std::collections::{BTreeSet, HashMap, HashSet};
21use std::fmt;
22use std::future::Future;
23use std::path::Path;
24use std::path::PathBuf;
25use std::str::FromStr;
26#[cfg(msim)]
27use std::sync::atomic::Ordering;
28use std::sync::{Arc, OnceLock, Weak};
29use std::time::Duration;
30use sui_core::admission_queue::{
31    AdmissionQueueContext, AdmissionQueueManager, AdmissionQueueMetrics,
32};
33use sui_core::authority::ExecutionEnv;
34use sui_core::authority::authority_store_tables::AuthorityPerpetualTablesOptions;
35use sui_core::authority::backpressure::BackpressureManager;
36use sui_core::authority::epoch_start_configuration::EpochFlag;
37use sui_core::authority::execution_time_estimator::ExecutionTimeObserver;
38use sui_core::consensus_adapter::ConsensusClient;
39use sui_core::consensus_manager::UpdatableConsensusClient;
40use sui_core::consensus_transaction_pool::TransactionPoolContext;
41use sui_core::epoch::randomness::RandomnessManager;
42use sui_core::execution_cache::build_execution_cache;
43use sui_core::randomness_round_receiver::{RandomnessRoundReceiver, RandomnessRoundReceiverHandle};
44use sui_network::endpoint_manager::{AddressSource, EndpointId};
45use sui_network::validator::server::SUI_TLS_SERVER_NAME;
46use sui_types::full_checkpoint_content::Checkpoint;
47use sui_types::node_role::NodeRole;
48
49use sui_core::global_state_hasher::GlobalStateHashMetrics;
50use sui_core::storage::RestReadStore;
51use sui_network::randomness;
52use sui_rpc_api::ServerVersion;
53use sui_rpc_api::subscription::SubscriptionService;
54use sui_types::base_types::ConciseableName;
55use sui_types::crypto::RandomnessRound;
56use sui_types::digests::{
57    ChainIdentifier, CheckpointDigest, TransactionDigest, TransactionEffectsDigest,
58};
59use sui_types::messages_consensus::AuthorityCapabilitiesV2;
60use sui_types::sui_system_state::SuiSystemState;
61use tap::tap::TapFallible;
62use tokio::sync::{Mutex, broadcast, mpsc};
63use tokio::task::JoinHandle;
64use tower::ServiceBuilder;
65use tracing::{Instrument, error_span, info};
66use tracing::{debug, error, warn};
67
68// Logs at debug level in test configuration, info level otherwise.
69// JWK logs cause significant volume in tests, but are insignificant in prod,
70// so we keep them at info
71macro_rules! jwk_log {
72    ($($arg:tt)+) => {
73        if in_test_configuration() {
74            debug!($($arg)+);
75        } else {
76            info!($($arg)+);
77        }
78    };
79}
80
81use fastcrypto_zkp::bn254::zk_login::JWK;
82pub use handle::SuiNodeHandle;
83use mysten_metrics::{RegistryService, spawn_monitored_task};
84use mysten_service::server_timing::server_timing_middleware;
85use sui_config::node::{DBCheckpointConfig, RunWithRange};
86use sui_config::node::{ForkCrashBehavior, ForkRecoveryConfig};
87use sui_config::transaction_deny_config::TransactionDenyRules;
88use sui_config::{ConsensusConfig, NodeConfig};
89use sui_core::authority::authority_per_epoch_store::AuthorityPerEpochStore;
90use sui_core::authority::authority_store_tables::AuthorityPerpetualTables;
91use sui_core::authority::epoch_start_configuration::EpochStartConfigTrait;
92use sui_core::authority::epoch_start_configuration::EpochStartConfiguration;
93use sui_core::authority::submitted_transaction_cache::SubmittedTransactionCacheMetrics;
94use sui_core::authority_aggregator::AuthorityAggregator;
95use sui_core::authority_server::{UserSubmissionPath, ValidatorService, ValidatorServiceMetrics};
96use sui_core::checkpoints::checkpoint_executor::metrics::CheckpointExecutorMetrics;
97use sui_core::checkpoints::checkpoint_executor::{CheckpointExecutor, StopReason};
98use sui_core::checkpoints::{
99    CheckpointMetrics, CheckpointOutput, CheckpointService, CheckpointStore, LogCheckpointOutput,
100    SendCheckpointToStateSync, SubmitCheckpointToConsensus,
101};
102use sui_core::consensus_adapter::{ConsensusAdapter, ConsensusAdapterMetrics};
103use sui_core::consensus_manager::ConsensusManager;
104use sui_core::consensus_throughput_calculator::ConsensusThroughputCalculator;
105use sui_core::consensus_validator::{SuiTxValidator, SuiTxValidatorMetrics};
106use sui_core::db_checkpoint_handler::DBCheckpointHandler;
107use sui_core::epoch::committee_store::CommitteeStore;
108use sui_core::epoch::consensus_store_pruner::ConsensusStorePruner;
109use sui_core::epoch::epoch_metrics::EpochMetrics;
110use sui_core::epoch::reconfiguration::ReconfigurationInitiator;
111use sui_core::global_state_hasher::GlobalStateHasher;
112use sui_core::module_cache_metrics::ResolverMetrics;
113use sui_core::overload_monitor::overload_monitor;
114use sui_core::rpc_store_embed::EmbeddedRpcStore;
115use sui_core::signature_verifier::SignatureVerifierMetrics;
116use sui_core::storage::RocksDbStore;
117use sui_core::storage::RpcStoreReadStore;
118use sui_core::transaction_orchestrator::TransactionOrchestrator;
119use sui_core::{
120    authority::{AuthorityState, AuthorityStore},
121    authority_client::NetworkAuthorityClient,
122};
123use sui_macros::fail_point;
124use sui_macros::{fail_point_arg, fail_point_async, replay_log};
125use sui_network::api::ValidatorServer;
126use sui_network::discovery;
127use sui_network::endpoint_manager::EndpointManager;
128use sui_network::state_sync;
129use sui_network::validator::server::ServerBuilder;
130use sui_protocol_config::{ProtocolConfig, ProtocolVersion};
131use sui_snapshot::uploader::StateSnapshotUploader;
132use sui_types::base_types::{AuthorityName, EpochId};
133use sui_types::committee::Committee;
134use sui_types::crypto::KeypairTraits;
135use sui_types::error::{SuiError, SuiResult};
136use sui_types::messages_consensus::{ConsensusTransaction, check_total_jwk_size};
137use sui_types::storage::RpcStateReader;
138use sui_types::sui_system_state::SuiSystemStateTrait;
139use sui_types::sui_system_state::epoch_start_sui_system_state::EpochStartSystemState;
140use sui_types::sui_system_state::epoch_start_sui_system_state::EpochStartSystemStateTrait;
141use sui_types::supported_protocol_versions::SupportedProtocolVersions;
142use typed_store::DBMetrics;
143use typed_store::rocks::default_db_options;
144
145use crate::metrics::{GrpcMetrics, SuiNodeMetrics};
146
147pub mod address_prober;
148pub mod admin;
149pub mod db_shell;
150mod handle;
151pub mod metrics;
152
153pub struct ValidatorComponents {
154    validator_server_handle: Option<ValidatorGrpcServer>,
155    validator_overload_monitor_handle: Option<JoinHandle<()>>,
156    consensus_manager: Arc<ConsensusManager>,
157    consensus_store_pruner: ConsensusStorePruner,
158    consensus_adapter: Arc<ConsensusAdapter>,
159    checkpoint_metrics: Arc<CheckpointMetrics>,
160    sui_tx_validator_metrics: Arc<SuiTxValidatorMetrics>,
161    admission_queue: Option<AdmissionQueueContext>,
162    transaction_pool_context: Option<Arc<TransactionPoolContext>>,
163}
164
165pub struct P2pComponents {
166    p2p_network: Network,
167    known_peers: HashMap<PeerId, String>,
168    discovery_handle: discovery::Handle,
169    state_sync_handle: state_sync::Handle,
170    randomness_handle: randomness::Handle,
171    endpoint_manager: EndpointManager,
172}
173
174#[cfg(msim)]
175mod simulator {
176    use std::sync::Mutex;
177    use std::sync::atomic::AtomicBool;
178    use sui_types::error::SuiErrorKind;
179
180    use super::*;
181    pub(super) struct SimState {
182        pub sim_node: sui_simulator::runtime::NodeHandle,
183        pub sim_safe_mode_expected: AtomicBool,
184        _leak_detector: sui_simulator::NodeLeakDetector,
185    }
186
187    impl Default for SimState {
188        fn default() -> Self {
189            Self {
190                sim_node: sui_simulator::runtime::NodeHandle::current(),
191                sim_safe_mode_expected: AtomicBool::new(false),
192                _leak_detector: sui_simulator::NodeLeakDetector::new(),
193            }
194        }
195    }
196
197    type JwkInjector = dyn Fn(AuthorityName, &OIDCProvider) -> SuiResult<Vec<(JwkId, JWK)>>
198        + Send
199        + Sync
200        + 'static;
201
202    fn default_fetch_jwks(
203        _authority: AuthorityName,
204        _provider: &OIDCProvider,
205    ) -> SuiResult<Vec<(JwkId, JWK)>> {
206        use fastcrypto_zkp::bn254::zk_login::parse_jwks;
207        // Just load a default Twitch jwk for testing.
208        parse_jwks(
209            sui_types::zk_login_util::DEFAULT_JWK_BYTES,
210            &OIDCProvider::Twitch,
211            true,
212        )
213        .map_err(|_| SuiErrorKind::JWKRetrievalError.into())
214    }
215
216    static JWK_INJECTOR: Mutex<Option<Arc<JwkInjector>>> = Mutex::new(None);
217
218    pub(super) fn get_jwk_injector() -> Arc<JwkInjector> {
219        JWK_INJECTOR
220            .lock()
221            .unwrap()
222            .clone()
223            .unwrap_or_else(|| Arc::new(default_fetch_jwks))
224    }
225
226    pub fn set_jwk_injector(injector: Arc<JwkInjector>) {
227        *JWK_INJECTOR.lock().unwrap() = Some(injector);
228    }
229}
230
231#[cfg(msim)]
232pub use simulator::set_jwk_injector;
233#[cfg(msim)]
234use simulator::*;
235use sui_core::authority::authority_store_pruner::PrunerWatermarks;
236use sui_core::{
237    consensus_handler::ConsensusHandlerInitializer, safe_client::SafeClientMetricsBase,
238};
239
240const DEFAULT_GRPC_CONNECT_TIMEOUT: Duration = Duration::from_secs(60);
241const VALIDATOR_GRPC_SERVER_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(30);
242
243/// How long after reconfiguration every reference to the previous epoch's store must be gone.
244/// The longest legitimate holder is an RPC handler on a fullnode waiting up to the local
245/// execution timeout (10s) for a transaction to be checkpointed. Simtests use a shorter period
246/// so that the check is exercised even by tests with short epochs.
247fn epoch_store_release_grace_period() -> Duration {
248    if cfg!(msim) {
249        Duration::from_secs(15)
250    } else {
251        Duration::from_secs(60)
252    }
253}
254
255pub struct SuiNode {
256    config: NodeConfig,
257    validator_components: Mutex<Option<ValidatorComponents>>,
258
259    /// The http servers responsible for serving RPC traffic.
260    #[allow(unused)]
261    http_servers: HttpServers,
262
263    state: Arc<AuthorityState>,
264    transaction_orchestrator: Option<Arc<TransactionOrchestrator<NetworkAuthorityClient>>>,
265    registry_service: RegistryService,
266    metrics: Arc<SuiNodeMetrics>,
267    checkpoint_metrics: Arc<CheckpointMetrics>,
268
269    _discovery: discovery::Handle,
270    _connection_monitor_handle: mysten_network::anemo_connection_monitor::ConnectionMonitorHandle,
271    state_sync_handle: state_sync::Handle,
272    randomness_handle: randomness::Handle,
273    checkpoint_store: Arc<CheckpointStore>,
274    global_state_hasher: Mutex<Option<Arc<GlobalStateHasher>>>,
275
276    /// Broadcast channel to send the starting system state for the next epoch.
277    end_of_epoch_channel: broadcast::Sender<SuiSystemState>,
278
279    /// EndpointManager for updating peer network addresses.
280    endpoint_manager: EndpointManager,
281
282    /// Handle to the discovery-shared address prober (`None` when disabled).
283    address_prober: Option<address_prober::Handle>,
284
285    backpressure_manager: Arc<BackpressureManager>,
286
287    _db_checkpoint_handle: Option<tokio::sync::broadcast::Sender<()>>,
288
289    #[cfg(msim)]
290    sim_state: SimState,
291
292    _state_snapshot_uploader_handle: Option<broadcast::Sender<()>>,
293    // Channel to allow signaling upstream to shutdown sui-node
294    shutdown_channel_tx: broadcast::Sender<Option<RunWithRange>>,
295
296    /// Handle shared with RandomnessManager and the consensus layer.
297    randomness_receiver_handle: Arc<RandomnessRoundReceiverHandle>,
298
299    /// Per-epoch consensus transaction pool handoff, shared between the RPC
300    /// server and ConsensusManager (`Some` only in pull-based submission mode).
301    transaction_pool_context: Option<Arc<TransactionPoolContext>>,
302
303    consensus_adapter_metrics: OnceLock<ConsensusAdapterMetrics>,
304
305    /// AuthorityAggregator of the network, created at start and beginning of each epoch.
306    /// Use ArcSwap so that we could mutate it without taking mut reference.
307    // TODO: Eventually we can make this auth aggregator a shared reference so that this
308    // update will automatically propagate to other uses.
309    auth_agg: Arc<ArcSwap<AuthorityAggregator<NetworkAuthorityClient>>>,
310
311    subscription_service_checkpoint_sender: Option<tokio::sync::broadcast::Sender<Arc<Checkpoint>>>,
312
313    /// The embedded `sui-rpc-store`, present when the node is a fullnode
314    /// with indexing enabled. Held for the node's lifetime so its tip
315    /// indexer keeps running (dropping it aborts the indexer). Exposed
316    /// through [`SuiNode::embedded_rpc_store`] for introspection.
317    embedded_rpc_store: Option<EmbeddedRpcStore>,
318}
319
320impl fmt::Debug for SuiNode {
321    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
322        f.debug_struct("SuiNode")
323            .field("name", &self.state.name.concise())
324            .finish()
325    }
326}
327
328static MAX_JWK_KEYS_PER_FETCH: usize = 100;
329
330impl SuiNode {
331    pub async fn start(
332        config: NodeConfig,
333        registry_service: RegistryService,
334    ) -> Result<Arc<SuiNode>> {
335        Self::start_async(
336            config,
337            registry_service,
338            ServerVersion::new("sui-node", "unknown"),
339        )
340        .await
341    }
342
343    fn start_jwk_updater(
344        config: &NodeConfig,
345        metrics: Arc<SuiNodeMetrics>,
346        authority: AuthorityName,
347        epoch_store: Arc<AuthorityPerEpochStore>,
348        consensus_adapter: Arc<ConsensusAdapter>,
349    ) {
350        let epoch = epoch_store.epoch();
351
352        let supported_providers = config
353            .zklogin_oauth_providers
354            .get(&epoch_store.get_chain_identifier().chain())
355            .unwrap_or(&BTreeSet::new())
356            .iter()
357            .map(|s| OIDCProvider::from_str(s).expect("Invalid provider string"))
358            .collect::<Vec<_>>();
359
360        let fetch_interval = Duration::from_secs(config.jwk_fetch_interval_seconds);
361
362        info!(
363            ?fetch_interval,
364            "Starting JWK updater tasks with supported providers: {:?}", supported_providers
365        );
366
367        fn validate_jwk(
368            metrics: &Arc<SuiNodeMetrics>,
369            provider: &OIDCProvider,
370            id: &JwkId,
371            jwk: &JWK,
372        ) -> bool {
373            let Ok(iss_provider) = OIDCProvider::from_iss(&id.iss) else {
374                warn!(
375                    "JWK iss {:?} (retrieved from {:?}) is not a valid provider",
376                    id.iss, provider
377                );
378                metrics
379                    .invalid_jwks
380                    .with_label_values(&[&provider.to_string()])
381                    .inc();
382                return false;
383            };
384
385            if iss_provider != *provider {
386                warn!(
387                    "JWK iss {:?} (retrieved from {:?}) does not match provider {:?}",
388                    id.iss, provider, iss_provider
389                );
390                metrics
391                    .invalid_jwks
392                    .with_label_values(&[&provider.to_string()])
393                    .inc();
394                return false;
395            }
396
397            if !check_total_jwk_size(id, jwk) {
398                warn!("JWK {:?} (retrieved from {:?}) is too large", id, provider);
399                metrics
400                    .invalid_jwks
401                    .with_label_values(&[&provider.to_string()])
402                    .inc();
403                return false;
404            }
405
406            true
407        }
408
409        // metrics is:
410        //  pub struct SuiNodeMetrics {
411        //      pub jwk_requests: IntCounterVec,
412        //      pub jwk_request_errors: IntCounterVec,
413        //      pub total_jwks: IntCounterVec,
414        //      pub unique_jwks: IntCounterVec,
415        //  }
416
417        for p in supported_providers.into_iter() {
418            let provider_str = p.to_string();
419            let epoch_store = epoch_store.clone();
420            let consensus_adapter = consensus_adapter.clone();
421            let metrics = metrics.clone();
422            spawn_monitored_task!(epoch_store.clone().within_alive_epoch(
423                async move {
424                    // note: restart-safe de-duplication happens after consensus, this is
425                    // just best-effort to reduce unneeded submissions.
426                    let mut seen = HashSet::new();
427                    loop {
428                        jwk_log!("fetching JWK for provider {:?}", p);
429                        metrics.jwk_requests.with_label_values(&[&provider_str]).inc();
430                        match Self::fetch_jwks(authority, &p).await {
431                            Err(e) => {
432                                metrics.jwk_request_errors.with_label_values(&[&provider_str]).inc();
433                                warn!("Error when fetching JWK for provider {:?} {:?}", p, e);
434                                // Retry in 30 seconds
435                                tokio::time::sleep(Duration::from_secs(30)).await;
436                                continue;
437                            }
438                            Ok(mut keys) => {
439                                metrics.total_jwks
440                                    .with_label_values(&[&provider_str])
441                                    .inc_by(keys.len() as u64);
442
443                                keys.retain(|(id, jwk)| {
444                                    validate_jwk(&metrics, &p, id, jwk) &&
445                                    !epoch_store.jwk_active_in_current_epoch(id, jwk) &&
446                                    seen.insert((id.clone(), jwk.clone()))
447                                });
448
449                                metrics.unique_jwks
450                                    .with_label_values(&[&provider_str])
451                                    .inc_by(keys.len() as u64);
452
453                                // prevent oauth providers from sending too many keys,
454                                // inadvertently or otherwise
455                                if keys.len() > MAX_JWK_KEYS_PER_FETCH {
456                                    warn!("Provider {:?} sent too many JWKs, only the first {} will be used", p, MAX_JWK_KEYS_PER_FETCH);
457                                    keys.truncate(MAX_JWK_KEYS_PER_FETCH);
458                                }
459
460                                for (id, jwk) in keys.into_iter() {
461                                    jwk_log!("Submitting JWK to consensus: {:?}", id);
462
463                                    let txn = ConsensusTransaction::new_jwk_fetched(authority, id, jwk);
464                                    consensus_adapter.submit(txn, None, &epoch_store, None, None)
465                                        .tap_err(|e| warn!("Error when submitting JWKs to consensus {:?}", e))
466                                        .ok();
467                                }
468                            }
469                        }
470                        tokio::time::sleep(fetch_interval).await;
471                    }
472                }
473                .instrument(error_span!("jwk_updater_task", epoch)),
474            ));
475        }
476    }
477
478    pub async fn start_async(
479        config: NodeConfig,
480        registry_service: RegistryService,
481        server_version: ServerVersion,
482    ) -> Result<Arc<SuiNode>> {
483        // Fail fast on config errors before starting any node components.
484        if let Some(prober_config) = &config.address_prober {
485            prober_config.validate()?;
486        }
487
488        let mut config = config.clone();
489        if config.supported_protocol_versions.is_none() {
490            info!(
491                "populating config.supported_protocol_versions with default {:?}",
492                SupportedProtocolVersions::SYSTEM_DEFAULT
493            );
494            config.supported_protocol_versions = Some(SupportedProtocolVersions::SYSTEM_DEFAULT);
495        }
496
497        let run_with_range = config.run_with_range;
498        let prometheus_registry = registry_service.default_registry();
499        let node_role = config.intended_node_role();
500
501        info!(node =? config.protocol_public_key(),
502            "Initializing sui-node listening on {} with role {:?}", config.network_address, node_role
503        );
504
505        // Initialize metrics to track db usage before creating any stores
506        DBMetrics::init(registry_service.clone());
507
508        // Build the Bulletproofs generators up front, so that the first range proof verification
509        // does not pay for it.
510        fastcrypto::bulletproofs::initialize_generators();
511
512        // Initialize db sync-to-disk setting from config (falls back to env var if not set)
513        typed_store::init_write_sync(config.enable_db_sync_to_disk);
514
515        // Initialize Mysten metrics.
516        mysten_metrics::init_metrics(&prometheus_registry);
517        // Unsupported (because of the use of static variable) and unnecessary in simtests.
518        #[cfg(not(msim))]
519        mysten_metrics::thread_stall_monitor::start_thread_stall_monitor();
520
521        let genesis = config.genesis()?.clone();
522
523        let secret = Arc::pin(config.protocol_key_pair().copy());
524        let genesis_committee = genesis.committee();
525        let committee_store = Arc::new(CommitteeStore::new(
526            config.db_path().join("epochs"),
527            &genesis_committee,
528            None,
529        ));
530
531        let pruner_watermarks = Arc::new(PrunerWatermarks::default());
532        let checkpoint_store = CheckpointStore::new(
533            &config.db_path().join("checkpoints"),
534            pruner_watermarks.clone(),
535        );
536        let checkpoint_metrics = CheckpointMetrics::new(&registry_service.default_registry());
537
538        #[allow(unused_mut)]
539        let mut build_version = server_version.version.to_string();
540        fail_point_arg!("override_binary_version", |version: std::sync::Arc<
541            std::sync::Mutex<String>,
542        >| {
543            #[cfg(msim)]
544            {
545                build_version = version.lock().unwrap().clone();
546            }
547        });
548        // Embedded in any fork marker this run records, so recovery can refuse to clear a fork
549        // under the same binary version that produced it.
550        checkpoint_store.set_binary_version(&build_version);
551
552        if node_role.runs_consensus() {
553            Self::check_and_recover_forks(
554                &checkpoint_store,
555                &checkpoint_metrics,
556                config.fork_recovery.as_ref(),
557                &build_version,
558            )
559            .await?;
560        }
561
562        // By default, only enable write stall on nodes that run consensus.
563        let enable_write_stall = config
564            .enable_db_write_stall
565            .unwrap_or(node_role.runs_consensus());
566        // The tidehunter objects compactor retains only the latest version per
567        // ObjectID and is mutually exclusive with the object pruner. Enable it
568        // for validators (which always disable the pruner), and also for any
569        // node configured with `num_epochs_to_retain = 0` — that aggressive
570        // setting is what the compactor replaces. The pruner is force-disabled
571        // in `AuthorityStorePruner::new` whenever this is true.
572        let enable_objects_compactor = node_role.is_validator()
573            || config.authority_store_pruning_config.num_epochs_to_retain == 0;
574        let perpetual_tables_options = AuthorityPerpetualTablesOptions {
575            enable_write_stall,
576            enable_objects_compactor,
577        };
578        let perpetual_tables = Arc::new(AuthorityPerpetualTables::open(
579            &config.db_store_path(),
580            Some(perpetual_tables_options),
581            Some(pruner_watermarks.epoch_id.clone()),
582        ));
583        let is_genesis = perpetual_tables
584            .database_is_empty()
585            .expect("Database read should not fail at init.");
586
587        let backpressure_manager =
588            BackpressureManager::new_from_checkpoint_store(&checkpoint_store);
589
590        let store =
591            AuthorityStore::open(perpetual_tables, &genesis, &config, &prometheus_registry).await?;
592
593        let cur_epoch = store.get_recovery_epoch_at_restart()?;
594        let committee = committee_store
595            .get_committee(&cur_epoch)?
596            .expect("Committee of the current epoch must exist");
597        let epoch_start_configuration = store
598            .get_epoch_start_configuration()?
599            .expect("EpochStartConfiguration of the current epoch must exist");
600        let cache_metrics = Arc::new(ResolverMetrics::new(&prometheus_registry));
601        let signature_verifier_metrics = SignatureVerifierMetrics::new(&prometheus_registry);
602
603        let cache_traits = build_execution_cache(
604            &config.execution_cache,
605            &prometheus_registry,
606            &store,
607            backpressure_manager.clone(),
608        );
609
610        let auth_agg = {
611            let safe_client_metrics_base = SafeClientMetricsBase::new(&prometheus_registry);
612            Arc::new(ArcSwap::new(Arc::new(
613                AuthorityAggregator::new_from_epoch_start_state(
614                    epoch_start_configuration.epoch_start_state(),
615                    &committee_store,
616                    safe_client_metrics_base,
617                ),
618            )))
619        };
620
621        let chain_id = ChainIdentifier::from(*genesis.checkpoint().digest());
622        let chain = match config.chain_override_for_testing {
623            Some(chain) => chain,
624            None => ChainIdentifier::from(*genesis.checkpoint().digest()).chain(),
625        };
626
627        let highest_executed_checkpoint = checkpoint_store
628            .get_highest_executed_checkpoint_seq_number()
629            .expect("checkpoint store read cannot fail")
630            .unwrap_or(0);
631
632        let previous_epoch_last_checkpoint = if cur_epoch == 0 {
633            0
634        } else {
635            checkpoint_store
636                .get_epoch_last_checkpoint_seq_number(cur_epoch - 1)
637                .expect("checkpoint store read cannot fail")
638                .unwrap_or(highest_executed_checkpoint)
639        };
640
641        let epoch_options = default_db_options().optimize_db_for_write_throughput(4, false);
642        let epoch_store = AuthorityPerEpochStore::new(
643            config.protocol_public_key(),
644            committee.clone(),
645            &config.db_store_path(),
646            Some(epoch_options.options),
647            EpochMetrics::new(&registry_service.default_registry()),
648            epoch_start_configuration,
649            cache_traits.backing_package_store.clone(),
650            cache_traits.object_store.clone(),
651            cache_metrics,
652            signature_verifier_metrics,
653            &config.expensive_safety_check_config,
654            (chain_id, chain),
655            highest_executed_checkpoint,
656            previous_epoch_last_checkpoint,
657            Arc::new(SubmittedTransactionCacheMetrics::new(
658                &registry_service.default_registry(),
659            )),
660            config.fullnode_sync_mode,
661        )?;
662
663        info!("created epoch store");
664
665        replay_log!(
666            "Beginning replay run. Epoch: {:?}, Protocol config: {:?}",
667            epoch_store.epoch(),
668            epoch_store.protocol_config()
669        );
670
671        // the database is empty at genesis time
672        if is_genesis {
673            info!("checking SUI conservation at genesis");
674            // When we are opening the db table, the only time when it's safe to
675            // check SUI conservation is at genesis. Otherwise we may be in the middle of
676            // an epoch and the SUI conservation check will fail. This also initialize
677            // the expected_network_sui_amount table.
678            cache_traits
679                .reconfig_api
680                .expensive_check_sui_conservation(&epoch_store)
681                .expect("SUI conservation check cannot fail at genesis");
682        }
683
684        let effective_buffer_stake = epoch_store.get_effective_buffer_stake_bps();
685        let default_buffer_stake = epoch_store
686            .protocol_config()
687            .buffer_stake_for_protocol_upgrade_bps();
688        if effective_buffer_stake != default_buffer_stake {
689            warn!(
690                ?effective_buffer_stake,
691                ?default_buffer_stake,
692                "buffer_stake_for_protocol_upgrade_bps is currently overridden"
693            );
694        }
695
696        checkpoint_store.insert_genesis_checkpoint(
697            genesis.checkpoint(),
698            genesis.checkpoint_contents().clone(),
699            &epoch_store,
700        );
701
702        info!("creating state sync store");
703        let state_sync_store = RocksDbStore::new(
704            cache_traits.clone(),
705            committee_store.clone(),
706            checkpoint_store.clone(),
707        );
708
709        if node_role.is_fullnode() {
710            // Fullnodes upgraded from a version that still ran the legacy
711            // index backends may have their now-dead on-disk directories
712            // lying around; remove them so they stop wasting disk.
713            remove_legacy_index_stores(&config.db_path());
714        }
715
716        let chain_identifier = epoch_store.get_chain_identifier();
717
718        // The embedded `sui-rpc-store` is the node's index backend: when
719        // indexing is enabled it builds the derived-index and ledger-history
720        // column families (indexed independently of the authority store) and
721        // serves the index read paths from the embedded store. Raw chain data
722        // is still served from the perpetual store.
723        let mut embedded_rpc_store =
724            if node_role.is_fullnode() && config.rpc().is_some_and(|rpc| rpc.enable_indexing()) {
725                info!("creating embedded rpc-store");
726                // The tip indexer pulls checkpoints from the node's local
727                // checkpoint / perpetual stores via a dedicated read handle.
728                let ingestion_source = RocksDbStore::new(
729                    cache_traits.clone(),
730                    committee_store.clone(),
731                    checkpoint_store.clone(),
732                );
733                let embedded_rpc_store = EmbeddedRpcStore::bootstrap(
734                    &config,
735                    &store,
736                    &checkpoint_store,
737                    ingestion_source,
738                    chain_identifier,
739                    &prometheus_registry,
740                )
741                .await?;
742                Some(embedded_rpc_store)
743            } else {
744                None
745            };
746
747        info!("creating archive reader");
748        // Create network
749        let (randomness_tx, randomness_rx) = mpsc::channel(
750            config
751                .p2p_config
752                .randomness
753                .clone()
754                .unwrap_or_default()
755                .mailbox_capacity(),
756        );
757        let P2pComponents {
758            p2p_network,
759            known_peers,
760            discovery_handle,
761            state_sync_handle,
762            randomness_handle,
763            endpoint_manager,
764        } = Self::create_p2p_network(
765            &config,
766            state_sync_store.clone(),
767            chain_identifier,
768            randomness_tx,
769            &prometheus_registry,
770        )?;
771
772        // Inject configured peer address overrides.
773        for peer in &config.p2p_config.peer_address_overrides {
774            endpoint_manager
775                .update_endpoint(
776                    EndpointId::P2p(peer.peer_id),
777                    AddressSource::Config,
778                    peer.addresses.clone(),
779                )
780                .expect("Updating peer address overrides should not fail");
781        }
782
783        // Send initial peer addresses to the p2p network.
784        update_peer_addresses(
785            &config,
786            &endpoint_manager,
787            epoch_store.epoch_start_state(),
788            None,
789        );
790
791        info!("start snapshot upload");
792        // Start uploading state snapshot to remote store
793        let state_snapshot_handle = Self::start_state_snapshot(
794            &config,
795            &prometheus_registry,
796            checkpoint_store.clone(),
797            chain_identifier,
798        )?;
799
800        // Start uploading db checkpoints to remote store
801        info!("start db checkpoint");
802        let (db_checkpoint_config, db_checkpoint_handle) = Self::start_db_checkpoint(
803            &config,
804            &prometheus_registry,
805            state_snapshot_handle.is_some(),
806        )?;
807
808        if !epoch_store
809            .protocol_config()
810            .simplified_unwrap_then_delete()
811        {
812            // We cannot prune tombstones if simplified_unwrap_then_delete is not enabled.
813            config
814                .authority_store_pruning_config
815                .set_killswitch_tombstone_pruning(true);
816        }
817
818        let authority_name = config.protocol_public_key();
819
820        info!("create authority state");
821        let state = AuthorityState::new(
822            authority_name,
823            secret,
824            config.supported_protocol_versions.unwrap(),
825            store.clone(),
826            cache_traits.clone(),
827            epoch_store.clone(),
828            committee_store.clone(),
829            embedded_rpc_store.as_ref().map(|embedded| embedded.store()),
830            checkpoint_store.clone(),
831            &prometheus_registry,
832            &db_checkpoint_config,
833            config.clone(),
834            chain_identifier,
835            config.policy_config.clone(),
836            config.firewall_config.clone(),
837            pruner_watermarks,
838        )
839        .await;
840        // ensure genesis txn was executed
841        if epoch_store.epoch() == 0 {
842            let txn = &genesis.transaction();
843            let span = error_span!("genesis_txn", tx_digest = ?txn.digest());
844            let transaction =
845                sui_types::executable_transaction::VerifiedExecutableTransaction::new_unchecked(
846                    sui_types::executable_transaction::ExecutableTransaction::new_from_data_and_sig(
847                        genesis.transaction().data().clone(),
848                        sui_types::executable_transaction::CertificateProof::Checkpoint(0, 0),
849                    ),
850                );
851            let _enter = span.enter();
852            state
853                .try_execute_immediately(&transaction, ExecutionEnv::new(), &epoch_store)
854                .unwrap();
855        }
856
857        // Start the loop that receives new randomness and generates transactions for it.
858        // The returned is long-lived (node lifetime).
859        let randomness_receiver_handle =
860            RandomnessRoundReceiver::spawn(state.clone(), randomness_rx);
861
862        let (end_of_epoch_channel, end_of_epoch_receiver) =
863            broadcast::channel(config.end_of_epoch_broadcast_channel_capacity);
864
865        let transaction_orchestrator = if node_role.is_fullnode() && run_with_range.is_none() {
866            Some(Arc::new(TransactionOrchestrator::new_with_auth_aggregator(
867                auth_agg.load_full(),
868                state.clone(),
869                end_of_epoch_receiver,
870                &config.db_path(),
871                &prometheus_registry,
872                &config,
873            )))
874        } else {
875            None
876        };
877
878        let (http_servers, subscription_service_checkpoint_sender) = build_http_servers(
879            state.clone(),
880            state_sync_store,
881            &transaction_orchestrator.clone(),
882            &config,
883            &prometheus_registry,
884            server_version,
885            node_role,
886            embedded_rpc_store.as_ref(),
887        )
888        .await?;
889
890        // Start the embedded rpc-store's tip indexer. It follows the tip
891        // via the checkpoint executor's broadcast stream and backfills
892        // any gap from the perpetual store. Spawned on a background task
893        // (see `spawn_indexer`) so node startup does not block on the
894        // first checkpoint, which the executor only produces after this
895        // function returns.
896        if let Some(embedded) = embedded_rpc_store.as_mut() {
897            embedded.spawn_indexer(
898                subscription_service_checkpoint_sender.clone(),
899                prometheus_registry.clone(),
900            );
901        }
902
903        let global_state_hasher = Arc::new(GlobalStateHasher::new(
904            cache_traits.global_state_hash_store.clone(),
905            GlobalStateHashMetrics::new(&prometheus_registry),
906        ));
907
908        let network_connection_metrics = mysten_network::quinn_metrics::QuinnConnectionMetrics::new(
909            "sui",
910            &registry_service.default_registry(),
911        );
912
913        let connection_monitor_handle =
914            mysten_network::anemo_connection_monitor::AnemoConnectionMonitor::spawn(
915                p2p_network.downgrade(),
916                Arc::new(network_connection_metrics),
917                known_peers,
918            );
919
920        let sui_node_metrics = Arc::new(SuiNodeMetrics::new(&registry_service.default_registry()));
921
922        sui_node_metrics
923            .binary_max_protocol_version
924            .set(ProtocolVersion::MAX.as_u64() as i64);
925        sui_node_metrics
926            .configured_max_protocol_version
927            .set(config.supported_protocol_versions.unwrap().max.as_u64() as i64);
928
929        let consensus_adapter_metrics = OnceLock::new();
930        let transaction_pool_context = config.consensus_transaction_pool.enabled.then(|| {
931            Arc::new(TransactionPoolContext::new(
932                Arc::new(AdmissionQueueMetrics::new(
933                    &registry_service.default_registry(),
934                )),
935                consensus_adapter_metrics
936                    .get_or_init(|| {
937                        ConsensusAdapterMetrics::new(&registry_service.default_registry())
938                    })
939                    .clone(),
940            ))
941        });
942        let node_role = epoch_store.node_role();
943        if !node_role.is_validator()
944            && let Some(context) = &transaction_pool_context
945        {
946            context.set_unavailable(epoch_store.epoch());
947        }
948        let validator_components = if node_role.runs_consensus() {
949            let mut components = Self::construct_validator_components(
950                config.clone(),
951                state.clone(),
952                committee,
953                epoch_store.clone(),
954                checkpoint_store.clone(),
955                state_sync_handle.clone(),
956                randomness_handle.clone(),
957                Arc::downgrade(&global_state_hasher),
958                backpressure_manager.clone(),
959                &registry_service,
960                transaction_pool_context.clone(),
961                &consensus_adapter_metrics,
962                sui_node_metrics.clone(),
963                checkpoint_metrics.clone(),
964                node_role,
965                randomness_receiver_handle.clone(),
966            )
967            .await?;
968
969            if node_role.is_validator() {
970                components
971                    .consensus_adapter
972                    .recover_end_of_publish(&epoch_store);
973
974                // Start the gRPC server
975                components.validator_server_handle = Some(
976                    components
977                        .validator_server_handle
978                        .take()
979                        .unwrap()
980                        .start()
981                        .await,
982                );
983
984                // Set the consensus address updater so that we can update the consensus peer addresses when requested.
985                endpoint_manager
986                    .set_consensus_address_updater(components.consensus_manager.clone());
987            } else {
988                info!("Starting node as Observer — connecting to configured peers");
989            }
990
991            Some(components)
992        } else {
993            None
994        };
995
996        let address_prober = if Self::address_prober_enabled(&config) {
997            let handle = address_prober::Builder::new()
998                .config(config.address_prober.clone().unwrap_or_default())
999                .with_metrics(&prometheus_registry)
1000                .build()
1001                .start(
1002                    p2p_network.clone(),
1003                    discovery_handle.sender(),
1004                    consensus_config::NetworkKeyPair::new(config.network_key_pair().copy()),
1005                );
1006            // Seed the current epoch if we are starting as a validator.
1007            if node_role.is_validator()
1008                && let Some(components) = &validator_components
1009            {
1010                handle.update_epoch(
1011                    epoch_store.epoch(),
1012                    epoch_store.epoch_start_state().get_consensus_committee(),
1013                    components.consensus_manager.clone(),
1014                );
1015            }
1016            Some(handle)
1017        } else {
1018            None
1019        };
1020
1021        // setup shutdown channel
1022        let (shutdown_channel, _) = broadcast::channel::<Option<RunWithRange>>(1);
1023
1024        let node = Self {
1025            config,
1026            validator_components: Mutex::new(validator_components),
1027            http_servers,
1028            state,
1029            transaction_orchestrator,
1030            registry_service,
1031            metrics: sui_node_metrics,
1032            checkpoint_metrics,
1033
1034            _discovery: discovery_handle,
1035            _connection_monitor_handle: connection_monitor_handle,
1036            state_sync_handle,
1037            randomness_handle,
1038            checkpoint_store,
1039            global_state_hasher: Mutex::new(Some(global_state_hasher)),
1040            end_of_epoch_channel,
1041            endpoint_manager,
1042            backpressure_manager,
1043            address_prober,
1044
1045            _db_checkpoint_handle: db_checkpoint_handle,
1046
1047            #[cfg(msim)]
1048            sim_state: Default::default(),
1049
1050            _state_snapshot_uploader_handle: state_snapshot_handle,
1051            shutdown_channel_tx: shutdown_channel,
1052            randomness_receiver_handle,
1053            transaction_pool_context,
1054            consensus_adapter_metrics,
1055
1056            auth_agg,
1057            subscription_service_checkpoint_sender,
1058            embedded_rpc_store,
1059        };
1060
1061        info!("SuiNode started!");
1062        let node = Arc::new(node);
1063        let node_copy = node.clone();
1064        spawn_monitored_task!(async move {
1065            let result = Self::monitor_reconfiguration(node_copy, epoch_store).await;
1066            if let Err(error) = result {
1067                warn!("Reconfiguration finished with error {:?}", error);
1068            }
1069        });
1070
1071        Ok(node)
1072    }
1073
1074    pub fn subscribe_to_epoch_change(&self) -> broadcast::Receiver<SuiSystemState> {
1075        self.end_of_epoch_channel.subscribe()
1076    }
1077
1078    pub fn subscribe_to_shutdown_channel(&self) -> broadcast::Receiver<Option<RunWithRange>> {
1079        self.shutdown_channel_tx.subscribe()
1080    }
1081
1082    pub fn current_epoch_for_testing(&self) -> EpochId {
1083        self.state.current_epoch_for_testing()
1084    }
1085
1086    pub fn db_checkpoint_path(&self) -> PathBuf {
1087        self.config.db_checkpoint_path()
1088    }
1089
1090    // Init reconfig process by starting to reject user certs
1091    pub async fn close_epoch(&self, epoch_store: &Arc<AuthorityPerEpochStore>) -> SuiResult {
1092        info!("close_epoch (current epoch = {})", epoch_store.epoch());
1093        self.validator_components
1094            .lock()
1095            .await
1096            .as_ref()
1097            .ok_or_else(|| SuiError::from("Node is not a validator"))?
1098            .consensus_adapter
1099            .close_epoch(epoch_store);
1100        Ok(())
1101    }
1102
1103    pub fn clear_override_protocol_upgrade_buffer_stake(&self, epoch: EpochId) -> SuiResult {
1104        self.state
1105            .clear_override_protocol_upgrade_buffer_stake(epoch)
1106    }
1107
1108    pub fn set_override_protocol_upgrade_buffer_stake(
1109        &self,
1110        epoch: EpochId,
1111        buffer_stake_bps: u64,
1112    ) -> SuiResult {
1113        self.state
1114            .set_override_protocol_upgrade_buffer_stake(epoch, buffer_stake_bps)
1115    }
1116
1117    // Testing-only API to start epoch close process.
1118    // For production code, please use the non-testing version.
1119    pub async fn close_epoch_for_testing(&self) -> SuiResult {
1120        let epoch_store = self.state.epoch_store_for_testing();
1121        self.close_epoch(&epoch_store).await
1122    }
1123
1124    fn start_state_snapshot(
1125        config: &NodeConfig,
1126        prometheus_registry: &Registry,
1127        checkpoint_store: Arc<CheckpointStore>,
1128        chain_identifier: ChainIdentifier,
1129    ) -> Result<Option<tokio::sync::broadcast::Sender<()>>> {
1130        if let Some(remote_store_config) = &config.state_snapshot_write_config.object_store_config {
1131            let snapshot_uploader = StateSnapshotUploader::new(
1132                &config.db_checkpoint_path(),
1133                &config.snapshot_path(),
1134                remote_store_config.clone(),
1135                60,
1136                prometheus_registry,
1137                checkpoint_store,
1138                chain_identifier,
1139                config.state_snapshot_write_config.archive_interval_epochs,
1140            )?;
1141            Ok(Some(snapshot_uploader.start()))
1142        } else {
1143            Ok(None)
1144        }
1145    }
1146
1147    fn start_db_checkpoint(
1148        config: &NodeConfig,
1149        prometheus_registry: &Registry,
1150        state_snapshot_enabled: bool,
1151    ) -> Result<(
1152        DBCheckpointConfig,
1153        Option<tokio::sync::broadcast::Sender<()>>,
1154    )> {
1155        let checkpoint_path = Some(
1156            config
1157                .db_checkpoint_config
1158                .checkpoint_path
1159                .clone()
1160                .unwrap_or_else(|| config.db_checkpoint_path()),
1161        );
1162        let db_checkpoint_config = if config.db_checkpoint_config.checkpoint_path.is_none() {
1163            DBCheckpointConfig {
1164                checkpoint_path,
1165                perform_db_checkpoints_at_epoch_end: if state_snapshot_enabled {
1166                    true
1167                } else {
1168                    config
1169                        .db_checkpoint_config
1170                        .perform_db_checkpoints_at_epoch_end
1171                },
1172                ..config.db_checkpoint_config.clone()
1173            }
1174        } else {
1175            config.db_checkpoint_config.clone()
1176        };
1177
1178        match (
1179            db_checkpoint_config.object_store_config.as_ref(),
1180            state_snapshot_enabled,
1181        ) {
1182            // If db checkpoint config object store not specified but
1183            // state snapshot object store is specified, create handler
1184            // anyway for marking db checkpoints as completed so that they
1185            // can be uploaded as state snapshots.
1186            (None, false) => Ok((db_checkpoint_config, None)),
1187            (_, _) => {
1188                let handler = DBCheckpointHandler::new(
1189                    &db_checkpoint_config.checkpoint_path.clone().unwrap(),
1190                    db_checkpoint_config.object_store_config.as_ref(),
1191                    60,
1192                    db_checkpoint_config
1193                        .prune_and_compact_before_upload
1194                        .unwrap_or(true),
1195                    config.authority_store_pruning_config.clone(),
1196                    prometheus_registry,
1197                    state_snapshot_enabled,
1198                )?;
1199                Ok((
1200                    db_checkpoint_config,
1201                    Some(DBCheckpointHandler::start(handler)),
1202                ))
1203            }
1204        }
1205    }
1206
1207    fn create_p2p_network(
1208        config: &NodeConfig,
1209        state_sync_store: RocksDbStore,
1210        chain_identifier: ChainIdentifier,
1211        randomness_tx: mpsc::Sender<(EpochId, RandomnessRound, Vec<u8>)>,
1212        prometheus_registry: &Registry,
1213    ) -> Result<P2pComponents> {
1214        let mut p2p_config = config.p2p_config.clone();
1215        {
1216            let disc = p2p_config.discovery.get_or_insert_with(Default::default);
1217            if disc.peer_addr_store_path.is_none() {
1218                disc.peer_addr_store_path =
1219                    Some(config.db_path().join("discovery_peer_cache.yaml"));
1220            }
1221        }
1222        let mut discovery_builder = discovery::Builder::new().config(p2p_config.clone());
1223        if let Some(consensus_config) = &config.consensus_config {
1224            let effective_addr = consensus_config
1225                .external_address
1226                .as_ref()
1227                .or(consensus_config.listen_address.as_ref());
1228            if let Some(addr) = effective_addr {
1229                discovery_builder = discovery_builder.consensus_external_address(addr.clone());
1230            }
1231        }
1232        let (discovery, discovery_server, endpoint_manager) = discovery_builder.build();
1233        let discovery_sender = discovery.sender();
1234
1235        let (state_sync, state_sync_router) = state_sync::Builder::new()
1236            .config(config.p2p_config.state_sync.clone().unwrap_or_default())
1237            .store(state_sync_store)
1238            .archive_config(config.archive_reader_config())
1239            .discovery_sender(discovery_sender)
1240            .with_metrics(prometheus_registry)
1241            .build();
1242
1243        let discovery_config = config.p2p_config.discovery.clone().unwrap_or_default();
1244        let known_peers: HashMap<PeerId, String> = discovery_config
1245            .allowlisted_peers
1246            .clone()
1247            .into_iter()
1248            .map(|ap| (ap.peer_id, "allowlisted_peer".to_string()))
1249            .chain(config.p2p_config.seed_peers.iter().filter_map(|peer| {
1250                peer.peer_id
1251                    .map(|peer_id| (peer_id, "seed_peer".to_string()))
1252            }))
1253            .collect();
1254
1255        let (randomness, randomness_router) =
1256            randomness::Builder::new(config.protocol_public_key(), randomness_tx)
1257                .config(config.p2p_config.randomness.clone().unwrap_or_default())
1258                .with_metrics(prometheus_registry)
1259                .build();
1260
1261        let p2p_network = {
1262            let routes = anemo::Router::new()
1263                .add_rpc_service(discovery_server)
1264                .merge(state_sync_router);
1265            let routes = routes.merge(randomness_router);
1266
1267            let inbound_network_metrics =
1268                mysten_network::metrics::NetworkMetrics::new("sui", "inbound", prometheus_registry);
1269            let outbound_network_metrics = mysten_network::metrics::NetworkMetrics::new(
1270                "sui",
1271                "outbound",
1272                prometheus_registry,
1273            );
1274
1275            let service = ServiceBuilder::new()
1276                .layer(
1277                    TraceLayer::new_for_server_errors()
1278                        .make_span_with(DefaultMakeSpan::new().level(tracing::Level::INFO))
1279                        .on_failure(DefaultOnFailure::new().level(tracing::Level::WARN)),
1280                )
1281                .layer(CallbackLayer::new(
1282                    mysten_network::metrics::MetricsMakeCallbackHandler::new(
1283                        Arc::new(inbound_network_metrics),
1284                        config.p2p_config.excessive_message_size(),
1285                    ),
1286                ))
1287                .service(routes);
1288
1289            let outbound_layer = ServiceBuilder::new()
1290                .layer(
1291                    TraceLayer::new_for_client_and_server_errors()
1292                        .make_span_with(DefaultMakeSpan::new().level(tracing::Level::INFO))
1293                        .on_failure(DefaultOnFailure::new().level(tracing::Level::WARN)),
1294                )
1295                .layer(CallbackLayer::new(
1296                    mysten_network::metrics::MetricsMakeCallbackHandler::new(
1297                        Arc::new(outbound_network_metrics),
1298                        config.p2p_config.excessive_message_size(),
1299                    ),
1300                ))
1301                .into_inner();
1302
1303            let mut anemo_config = config.p2p_config.anemo_config.clone().unwrap_or_default();
1304            // Inbound requests on this network are small (signatures, queries, summaries).
1305            // Cap request frames at 1 MiB.
1306            anemo_config.max_request_frame_size = Some(1 << 20);
1307            // Responses can be larger (checkpoint contents).
1308            // Cap response frames at 128 MiB.
1309            anemo_config.max_response_frame_size = Some(128 << 20);
1310
1311            // Set a higher default value for socket send/receive buffers if not already
1312            // configured.
1313            let mut quic_config = anemo_config.quic.unwrap_or_default();
1314            if quic_config.socket_send_buffer_size.is_none() {
1315                quic_config.socket_send_buffer_size = Some(20 << 20);
1316            }
1317            if quic_config.socket_receive_buffer_size.is_none() {
1318                quic_config.socket_receive_buffer_size = Some(20 << 20);
1319            }
1320            quic_config.allow_failed_socket_buffer_size_setting = true;
1321
1322            // Set high-performance defaults for quinn transport.
1323            // With 200MiB buffer size and ~500ms RTT, max throughput ~400MiB/s.
1324            if quic_config.max_concurrent_bidi_streams.is_none() {
1325                quic_config.max_concurrent_bidi_streams = Some(500);
1326            }
1327            if quic_config.max_concurrent_uni_streams.is_none() {
1328                quic_config.max_concurrent_uni_streams = Some(500);
1329            }
1330            if quic_config.stream_receive_window.is_none() {
1331                quic_config.stream_receive_window = Some(100 << 20);
1332            }
1333            if quic_config.receive_window.is_none() {
1334                quic_config.receive_window = Some(200 << 20);
1335            }
1336            if quic_config.send_window.is_none() {
1337                quic_config.send_window = Some(200 << 20);
1338            }
1339            if quic_config.crypto_buffer_size.is_none() {
1340                quic_config.crypto_buffer_size = Some(1 << 20);
1341            }
1342            if quic_config.max_idle_timeout_ms.is_none() {
1343                quic_config.max_idle_timeout_ms = Some(10_000);
1344            }
1345            if quic_config.keep_alive_interval_ms.is_none() {
1346                quic_config.keep_alive_interval_ms = Some(5_000);
1347            }
1348            anemo_config.quic = Some(quic_config);
1349
1350            let server_name = format!("sui-{}", chain_identifier);
1351            let network = Network::bind(config.p2p_config.listen_address)
1352                .server_name(&server_name)
1353                .private_key(config.network_key_pair().copy().private().0.to_bytes())
1354                .config(anemo_config)
1355                .outbound_request_layer(outbound_layer)
1356                .start(service)?;
1357            info!(
1358                server_name = server_name,
1359                "P2p network started on {}",
1360                network.local_addr()
1361            );
1362
1363            network
1364        };
1365
1366        let discovery_handle =
1367            discovery.start(p2p_network.clone(), config.network_key_pair().copy());
1368        let state_sync_handle = state_sync.start(p2p_network.clone());
1369        let randomness_handle = randomness.start(p2p_network.clone());
1370
1371        Ok(P2pComponents {
1372            p2p_network,
1373            known_peers,
1374            discovery_handle,
1375            state_sync_handle,
1376            randomness_handle,
1377            endpoint_manager,
1378        })
1379    }
1380
1381    async fn construct_validator_components(
1382        config: NodeConfig,
1383        state: Arc<AuthorityState>,
1384        committee: Arc<Committee>,
1385        epoch_store: Arc<AuthorityPerEpochStore>,
1386        checkpoint_store: Arc<CheckpointStore>,
1387        state_sync_handle: state_sync::Handle,
1388        randomness_handle: randomness::Handle,
1389        global_state_hasher: Weak<GlobalStateHasher>,
1390        backpressure_manager: Arc<BackpressureManager>,
1391        registry_service: &RegistryService,
1392        transaction_pool_context: Option<Arc<TransactionPoolContext>>,
1393        consensus_adapter_metrics: &OnceLock<ConsensusAdapterMetrics>,
1394        sui_node_metrics: Arc<SuiNodeMetrics>,
1395        checkpoint_metrics: Arc<CheckpointMetrics>,
1396        node_role: NodeRole,
1397        randomness_receiver_handle: Arc<RandomnessRoundReceiverHandle>,
1398    ) -> Result<ValidatorComponents> {
1399        let mut config_clone = config.clone();
1400        let consensus_config = config_clone
1401            .consensus_config
1402            .as_mut()
1403            .ok_or_else(|| anyhow!("Node is missing consensus config"))?;
1404
1405        let client = Arc::new(UpdatableConsensusClient::new());
1406        let inflight_slot_freed_notify = Arc::new(tokio::sync::Notify::new());
1407        let consensus_adapter = Arc::new(Self::construct_consensus_adapter(
1408            &committee,
1409            consensus_config,
1410            state.name,
1411            consensus_adapter_metrics
1412                .get_or_init(|| ConsensusAdapterMetrics::new(&registry_service.default_registry()))
1413                .clone(),
1414            client.clone(),
1415            checkpoint_store.clone(),
1416            inflight_slot_freed_notify.clone(),
1417        ));
1418
1419        let consensus_manager = Arc::new(ConsensusManager::new(
1420            &config,
1421            consensus_config,
1422            registry_service,
1423            client,
1424            transaction_pool_context.clone(),
1425            node_role,
1426        ));
1427
1428        // This only gets started up once, not on every epoch. (Make call to remove every epoch.)
1429        let consensus_store_pruner = ConsensusStorePruner::new(
1430            consensus_manager.get_storage_base_path(),
1431            consensus_config.db_retention_epochs(),
1432            consensus_config.db_pruner_period(),
1433            &registry_service.default_registry(),
1434        );
1435
1436        let sui_tx_validator_metrics =
1437            SuiTxValidatorMetrics::new(&registry_service.default_registry());
1438
1439        let (validator_server_handle, admission_queue) = if node_role.is_validator() {
1440            let (handle, queue) = Self::start_grpc_validator_service(
1441                &config,
1442                state.clone(),
1443                consensus_adapter.clone(),
1444                epoch_store.clone(),
1445                &registry_service.default_registry(),
1446                inflight_slot_freed_notify,
1447                transaction_pool_context.clone(),
1448            )
1449            .await?;
1450            (Some(handle), queue)
1451        } else {
1452            (None, None)
1453        };
1454
1455        // Starts an overload monitor that monitors the execution of the authority.
1456        // Don't start the overload monitor when max_load_shedding_percentage is 0.
1457        let validator_overload_monitor_handle = if node_role.is_validator()
1458            && config
1459                .authority_overload_config
1460                .max_load_shedding_percentage
1461                > 0
1462        {
1463            let authority_state = Arc::downgrade(&state);
1464            let overload_config = config.authority_overload_config.clone();
1465            fail_point!("starting_overload_monitor");
1466            Some(spawn_monitored_task!(overload_monitor(
1467                authority_state,
1468                overload_config,
1469            )))
1470        } else {
1471            None
1472        };
1473
1474        Self::start_epoch_specific_validator_components(
1475            &config,
1476            state.clone(),
1477            consensus_adapter,
1478            checkpoint_store,
1479            epoch_store,
1480            state_sync_handle,
1481            randomness_handle,
1482            randomness_receiver_handle,
1483            consensus_manager,
1484            consensus_store_pruner,
1485            global_state_hasher,
1486            backpressure_manager,
1487            validator_server_handle,
1488            validator_overload_monitor_handle,
1489            checkpoint_metrics,
1490            sui_node_metrics,
1491            sui_tx_validator_metrics,
1492            admission_queue,
1493            transaction_pool_context,
1494            node_role,
1495        )
1496        .await
1497    }
1498
1499    fn address_prober_enabled(config: &NodeConfig) -> bool {
1500        let prober_enabled = config
1501            .address_prober
1502            .as_ref()
1503            .map(|c| c.enabled())
1504            .unwrap_or(true);
1505        let v3_enabled = config
1506            .p2p_config
1507            .discovery
1508            .as_ref()
1509            .is_some_and(|d| d.use_get_known_peers_v3());
1510        prober_enabled && v3_enabled
1511    }
1512
1513    fn update_address_prober_epoch(
1514        &self,
1515        epoch_store: &AuthorityPerEpochStore,
1516        consensus_manager: &Arc<ConsensusManager>,
1517    ) {
1518        if !epoch_store.is_validator() {
1519            return;
1520        }
1521        if let Some(handle) = &self.address_prober {
1522            handle.update_epoch(
1523                epoch_store.epoch(),
1524                epoch_store.epoch_start_state().get_consensus_committee(),
1525                consensus_manager.clone(),
1526            );
1527        }
1528    }
1529
1530    async fn start_epoch_specific_validator_components(
1531        config: &NodeConfig,
1532        state: Arc<AuthorityState>,
1533        consensus_adapter: Arc<ConsensusAdapter>,
1534        checkpoint_store: Arc<CheckpointStore>,
1535        epoch_store: Arc<AuthorityPerEpochStore>,
1536        state_sync_handle: state_sync::Handle,
1537        randomness_handle: randomness::Handle,
1538        randomness_receiver_handle: Arc<RandomnessRoundReceiverHandle>,
1539        consensus_manager: Arc<ConsensusManager>,
1540        consensus_store_pruner: ConsensusStorePruner,
1541        state_hasher: Weak<GlobalStateHasher>,
1542        backpressure_manager: Arc<BackpressureManager>,
1543        validator_server_handle: Option<ValidatorGrpcServer>,
1544        validator_overload_monitor_handle: Option<JoinHandle<()>>,
1545        checkpoint_metrics: Arc<CheckpointMetrics>,
1546        sui_node_metrics: Arc<SuiNodeMetrics>,
1547        sui_tx_validator_metrics: Arc<SuiTxValidatorMetrics>,
1548        admission_queue: Option<AdmissionQueueContext>,
1549        transaction_pool_context: Option<Arc<TransactionPoolContext>>,
1550        node_role: NodeRole,
1551    ) -> Result<ValidatorComponents> {
1552        let checkpoint_service = Self::build_checkpoint_service(
1553            config,
1554            consensus_adapter.clone(),
1555            checkpoint_store.clone(),
1556            epoch_store.clone(),
1557            state.clone(),
1558            state_sync_handle,
1559            state_hasher,
1560            checkpoint_metrics.clone(),
1561            node_role,
1562        );
1563
1564        // Clear the VSS public key from the previous epoch so any randomness round
1565        // signatures buffer in the channel until the new DKG completes.
1566        randomness_receiver_handle.clear_public_key();
1567
1568        if node_role.runs_consensus() && epoch_store.randomness_state_enabled() {
1569            let authority_key_pair = if node_role.is_validator() {
1570                Some(config.protocol_key_pair())
1571            } else {
1572                None
1573            };
1574            let randomness_manager = RandomnessManager::try_new(
1575                Arc::downgrade(&epoch_store),
1576                Box::new(consensus_adapter.clone()),
1577                randomness_handle,
1578                authority_key_pair,
1579                randomness_receiver_handle.clone(),
1580            )
1581            .await;
1582            if let Some(randomness_manager) = randomness_manager {
1583                epoch_store
1584                    .set_randomness_manager(randomness_manager)
1585                    .await?;
1586            }
1587        }
1588
1589        if node_role.is_validator() {
1590            ExecutionTimeObserver::spawn(
1591                epoch_store.clone(),
1592                Box::new(consensus_adapter.clone()),
1593                config
1594                    .execution_time_observer_config
1595                    .clone()
1596                    .unwrap_or_default(),
1597            );
1598        }
1599
1600        let throughput_calculator = Arc::new(ConsensusThroughputCalculator::new(
1601            None,
1602            state.metrics.clone(),
1603        ));
1604
1605        let consensus_handler_initializer = ConsensusHandlerInitializer::new(
1606            state.clone(),
1607            checkpoint_service.clone(),
1608            epoch_store.clone(),
1609            throughput_calculator,
1610            backpressure_manager,
1611            config.congestion_log.clone(),
1612        );
1613
1614        info!("Starting consensus manager asynchronously");
1615
1616        // Spawn consensus startup asynchronously to avoid blocking other components
1617        tokio::spawn({
1618            let config = config.clone();
1619            let epoch_store = epoch_store.clone();
1620            let sui_tx_validator = SuiTxValidator::new(
1621                state.clone(),
1622                epoch_store.clone(),
1623                checkpoint_service.clone(),
1624                sui_tx_validator_metrics.clone(),
1625            );
1626            let consensus_manager = consensus_manager.clone();
1627            async move {
1628                consensus_manager
1629                    .start(
1630                        &config,
1631                        epoch_store,
1632                        consensus_handler_initializer,
1633                        sui_tx_validator,
1634                        Some(randomness_receiver_handle),
1635                    )
1636                    .await;
1637            }
1638        });
1639        let replay_waiter = consensus_manager.replay_waiter();
1640
1641        info!("Spawning checkpoint service");
1642        let replay_waiter = if std::env::var("DISABLE_REPLAY_WAITER").is_ok() {
1643            None
1644        } else {
1645            Some(replay_waiter)
1646        };
1647        checkpoint_service
1648            .spawn(epoch_store.clone(), replay_waiter)
1649            .await;
1650
1651        if node_role.is_validator() && epoch_store.authenticator_state_enabled() {
1652            Self::start_jwk_updater(
1653                config,
1654                sui_node_metrics,
1655                state.name,
1656                epoch_store.clone(),
1657                consensus_adapter.clone(),
1658            );
1659        }
1660
1661        if let Some(ctx) = &admission_queue {
1662            ctx.rotate_for_epoch(epoch_store);
1663        }
1664
1665        Ok(ValidatorComponents {
1666            validator_server_handle,
1667            validator_overload_monitor_handle,
1668            consensus_manager,
1669            consensus_store_pruner,
1670            consensus_adapter,
1671            checkpoint_metrics,
1672            sui_tx_validator_metrics,
1673            admission_queue,
1674            transaction_pool_context,
1675        })
1676    }
1677
1678    fn build_checkpoint_service(
1679        config: &NodeConfig,
1680        consensus_adapter: Arc<ConsensusAdapter>,
1681        checkpoint_store: Arc<CheckpointStore>,
1682        epoch_store: Arc<AuthorityPerEpochStore>,
1683        state: Arc<AuthorityState>,
1684        state_sync_handle: state_sync::Handle,
1685        state_hasher: Weak<GlobalStateHasher>,
1686        checkpoint_metrics: Arc<CheckpointMetrics>,
1687        node_role: NodeRole,
1688    ) -> Arc<CheckpointService> {
1689        let checkpoint_output: Box<dyn CheckpointOutput> = if node_role.is_validator() {
1690            Box::new(SubmitCheckpointToConsensus::new(
1691                consensus_adapter,
1692                state.secret.clone(),
1693                config.protocol_public_key(),
1694                checkpoint_metrics.clone(),
1695            ))
1696        } else {
1697            Box::new(LogCheckpointOutput::new(checkpoint_metrics.clone()))
1698        };
1699
1700        let certified_checkpoint_output = SendCheckpointToStateSync::new(state_sync_handle);
1701
1702        CheckpointService::build(
1703            state.clone(),
1704            checkpoint_store,
1705            epoch_store,
1706            state.get_transaction_cache_reader().clone(),
1707            state_hasher,
1708            checkpoint_output,
1709            Box::new(certified_checkpoint_output),
1710            checkpoint_metrics,
1711        )
1712    }
1713
1714    fn construct_consensus_adapter(
1715        committee: &Committee,
1716        consensus_config: &ConsensusConfig,
1717        authority: AuthorityName,
1718        ca_metrics: ConsensusAdapterMetrics,
1719        consensus_client: Arc<dyn ConsensusClient>,
1720        checkpoint_store: Arc<CheckpointStore>,
1721        inflight_slot_freed_notify: Arc<tokio::sync::Notify>,
1722    ) -> ConsensusAdapter {
1723        // The consensus adapter allows the authority to send user certificates through consensus.
1724        ConsensusAdapter::new(
1725            consensus_client,
1726            checkpoint_store,
1727            authority,
1728            consensus_config.max_pending_transactions(),
1729            consensus_config.max_pending_transactions() * 2 / committee.num_members(),
1730            ca_metrics,
1731            inflight_slot_freed_notify,
1732        )
1733    }
1734
1735    async fn start_grpc_validator_service(
1736        config: &NodeConfig,
1737        state: Arc<AuthorityState>,
1738        consensus_adapter: Arc<ConsensusAdapter>,
1739        epoch_store: Arc<AuthorityPerEpochStore>,
1740        prometheus_registry: &Registry,
1741        inflight_slot_freed_notify: Arc<tokio::sync::Notify>,
1742        transaction_pool_context: Option<Arc<TransactionPoolContext>>,
1743    ) -> Result<(ValidatorGrpcServer, Option<AdmissionQueueContext>)> {
1744        let overload_config = &config.authority_overload_config;
1745        let admission_queue =
1746            if transaction_pool_context.is_none() && overload_config.admission_queue_enabled {
1747                let manager = Arc::new(AdmissionQueueManager::new(
1748                    consensus_adapter.clone(),
1749                    Arc::new(AdmissionQueueMetrics::new(prometheus_registry)),
1750                    overload_config.admission_queue_capacity_fraction,
1751                    overload_config.admission_queue_failover_timeout,
1752                    inflight_slot_freed_notify,
1753                ));
1754                Some(AdmissionQueueContext::spawn(manager, epoch_store))
1755            } else {
1756                None
1757            };
1758        let user_submission_path = if let Some(context) = transaction_pool_context {
1759            UserSubmissionPath::Pool(context)
1760        } else if let Some(context) = admission_queue.clone() {
1761            UserSubmissionPath::AdmissionQueue(context)
1762        } else {
1763            UserSubmissionPath::Direct
1764        };
1765        let validator_service = ValidatorService::new(
1766            state.clone(),
1767            consensus_adapter,
1768            Arc::new(ValidatorServiceMetrics::new(prometheus_registry)),
1769            config.policy_config.clone().map(|p| p.client_id_source),
1770            user_submission_path,
1771        );
1772
1773        let mut server_conf = mysten_network::config::Config::new();
1774        server_conf.connect_timeout = Some(DEFAULT_GRPC_CONNECT_TIMEOUT);
1775        server_conf.http2_keepalive_interval = Some(DEFAULT_GRPC_CONNECT_TIMEOUT);
1776        server_conf.http2_keepalive_timeout = Some(DEFAULT_GRPC_CONNECT_TIMEOUT);
1777        server_conf.global_concurrency_limit = config.grpc_concurrency_limit;
1778        server_conf.load_shed = config.grpc_load_shed;
1779        let mut server_builder =
1780            ServerBuilder::from_config(&server_conf, GrpcMetrics::new(prometheus_registry));
1781
1782        server_builder = server_builder.add_service(ValidatorServer::new(validator_service));
1783
1784        let tls_config = sui_tls::create_rustls_server_config(
1785            config.network_key_pair().copy().private(),
1786            SUI_TLS_SERVER_NAME.to_string(),
1787        );
1788
1789        let network_address = config.network_address().clone();
1790
1791        let server = ValidatorGrpcServer::new(async move {
1792            let server = server_builder
1793                .bind(&network_address, Some(tls_config))
1794                .await
1795                .unwrap_or_else(|err| panic!("Failed to bind to {network_address}: {err}"));
1796            info!("Listening to traffic on {}", server.local_addr());
1797            server.into_handle()
1798        });
1799        Ok((server, admission_queue))
1800    }
1801
1802    pub fn state(&self) -> Arc<AuthorityState> {
1803        self.state.clone()
1804    }
1805
1806    /// The embedded `sui-rpc-store` index backend, when the node is a
1807    /// fullnode with indexing enabled. Exposes the startup bootstrap
1808    /// decision and per-cohort watermarks for introspection (used by
1809    /// tests to observe restore/resume behavior across restarts without
1810    /// going through the RPC surface).
1811    pub fn embedded_rpc_store(&self) -> Option<&EmbeddedRpcStore> {
1812        self.embedded_rpc_store.as_ref()
1813    }
1814
1815    #[cfg(any(test, msim))]
1816    pub fn connection_monitor_handle_for_testing(
1817        &self,
1818    ) -> &mysten_network::anemo_connection_monitor::ConnectionMonitorHandle {
1819        &self._connection_monitor_handle
1820    }
1821
1822    #[cfg(any(test, msim))]
1823    pub fn address_prober_metrics_for_testing(
1824        &self,
1825    ) -> std::sync::Arc<address_prober::AddressProberMetrics> {
1826        self.address_prober
1827            .as_ref()
1828            .expect("address prober should be running in tests")
1829            .metrics_for_testing()
1830    }
1831
1832    #[cfg(feature = "testing")]
1833    pub fn prometheus_metrics_for_testing(&self) -> Vec<prometheus::proto::MetricFamily> {
1834        self.registry_service.default_registry().gather()
1835    }
1836
1837    pub fn node_role(&self) -> NodeRole {
1838        self.state.load_epoch_store_one_call_per_task().node_role()
1839    }
1840
1841    /// Returns the validator's `ConsensusAdapter` if this node currently has validator
1842    /// components running. The Arc is cloned out and the lock is released immediately
1843    /// so callers never hold the `validator_components` mutex across consensus
1844    /// submission.
1845    pub async fn consensus_adapter(&self) -> Option<Arc<ConsensusAdapter>> {
1846        self.validator_components
1847            .lock()
1848            .await
1849            .as_ref()
1850            .map(|c| c.consensus_adapter.clone())
1851    }
1852
1853    // Only used for testing because of how epoch store is loaded.
1854    pub fn reference_gas_price_for_testing(&self) -> Result<u64, anyhow::Error> {
1855        self.state.reference_gas_price_for_testing()
1856    }
1857
1858    pub fn clone_committee_store(&self) -> Arc<CommitteeStore> {
1859        self.state.committee_store().clone()
1860    }
1861
1862    pub fn clone_checkpoint_store(&self) -> Arc<CheckpointStore> {
1863        self.checkpoint_store.clone()
1864    }
1865
1866    pub fn clone_authority_store(&self) -> Arc<AuthorityStore> {
1867        self.state.authority_store()
1868    }
1869
1870    pub fn clone_consensus_store(
1871        &self,
1872    ) -> Option<Arc<consensus_core::storage::rocksdb_store::RocksDBStore>> {
1873        self.validator_components
1874            .try_lock()
1875            .ok()?
1876            .as_ref()?
1877            .consensus_manager
1878            .consensus_store()
1879    }
1880
1881    /// Clone an AuthorityAggregator currently used in this node, if the node is a fullnode.
1882    /// After reconfig, Transaction Driver builds a new AuthorityAggregator. The caller
1883    /// of this function will mostly likely want to call this again
1884    /// to get a fresh one.
1885    pub fn clone_authority_aggregator(
1886        &self,
1887    ) -> Option<Arc<AuthorityAggregator<NetworkAuthorityClient>>> {
1888        self.transaction_orchestrator
1889            .as_ref()
1890            .map(|to| to.clone_authority_aggregator())
1891    }
1892
1893    pub fn transaction_orchestrator(
1894        &self,
1895    ) -> Option<Arc<TransactionOrchestrator<NetworkAuthorityClient>>> {
1896        self.transaction_orchestrator.clone()
1897    }
1898
1899    /// This function awaits the completion of checkpoint execution of the current epoch,
1900    /// after which it initiates reconfiguration of the entire system.
1901    pub async fn monitor_reconfiguration(
1902        self: Arc<Self>,
1903        mut epoch_store: Arc<AuthorityPerEpochStore>,
1904    ) -> Result<()> {
1905        let checkpoint_executor_metrics =
1906            CheckpointExecutorMetrics::new(&self.registry_service.default_registry());
1907
1908        // Holds the startup-specific deny-config broadcast setting; consumed by the
1909        // first iteration on which deny-config sharing is enabled in the protocol
1910        // config. Subsequent iterations fall back to the epoch-change setting.
1911        let mut broadcast_on_startup: Option<bool> =
1912            Some(self.config.peer_deny_sync_config.broadcast_on_startup);
1913
1914        loop {
1915            let mut hasher_guard = self.global_state_hasher.lock().await;
1916            let hasher = hasher_guard.take().unwrap();
1917            info!(
1918                "Creating checkpoint executor for epoch {}",
1919                epoch_store.epoch()
1920            );
1921            let checkpoint_executor = CheckpointExecutor::new(
1922                epoch_store.clone(),
1923                self.checkpoint_store.clone(),
1924                self.state.clone(),
1925                hasher.clone(),
1926                self.backpressure_manager.clone(),
1927                self.config.checkpoint_executor_config.clone(),
1928                checkpoint_executor_metrics.clone(),
1929                self.subscription_service_checkpoint_sender.clone(),
1930            );
1931
1932            let run_with_range = self.config.run_with_range;
1933
1934            let cur_epoch_store = self.state.load_epoch_store_one_call_per_task();
1935
1936            // Update the current protocol version metric.
1937            self.metrics
1938                .current_protocol_version
1939                .set(cur_epoch_store.protocol_config().version.as_u64() as i64);
1940
1941            // Advertise capabilities to committee, if we are a validator.
1942            // FullNodes that state sync via consensus will also have validator components, by they are not supposed to submit any capabilities.
1943            if let Some(components) = &*self.validator_components.lock().await
1944                && cur_epoch_store.is_validator()
1945            {
1946                // TODO: without this sleep, the consensus message is not delivered reliably.
1947                tokio::time::sleep(Duration::from_millis(1)).await;
1948
1949                let config = cur_epoch_store.protocol_config();
1950                let mut supported_protocol_versions = self
1951                    .config
1952                    .supported_protocol_versions
1953                    .expect("Supported versions should be populated")
1954                    // no need to send digests of versions less than the current version
1955                    .truncate_below(config.version);
1956
1957                while supported_protocol_versions.max > config.version {
1958                    let proposed_protocol_config = ProtocolConfig::get_for_version(
1959                        supported_protocol_versions.max,
1960                        cur_epoch_store.get_chain(),
1961                    );
1962
1963                    if proposed_protocol_config.enable_accumulators()
1964                        && !epoch_store.accumulator_root_exists()
1965                    {
1966                        error!(
1967                            "cannot upgrade to protocol version {:?} because accumulator root does not exist",
1968                            supported_protocol_versions.max
1969                        );
1970                        supported_protocol_versions.max = supported_protocol_versions.max.prev();
1971                    } else {
1972                        break;
1973                    }
1974                }
1975
1976                let binary_config = config.binary_config(None);
1977                let transaction = ConsensusTransaction::new_capability_notification_v2(
1978                    AuthorityCapabilitiesV2::new(
1979                        self.state.name,
1980                        cur_epoch_store.get_chain_identifier().chain(),
1981                        supported_protocol_versions,
1982                        self.state
1983                            .get_available_system_packages(&binary_config)
1984                            .await,
1985                    ),
1986                );
1987                info!(?transaction, "submitting capabilities to consensus");
1988                components.consensus_adapter.submit(
1989                    transaction,
1990                    None,
1991                    &cur_epoch_store,
1992                    None,
1993                    None,
1994                )?;
1995
1996                // Reconcile our shared TransactionDenyConfig vote with the (possibly
1997                // restart-edited) local config.
1998                if cur_epoch_store
1999                    .protocol_config()
2000                    .share_transaction_deny_config_in_consensus()
2001                {
2002                    let sync_cfg = &self.config.peer_deny_sync_config;
2003                    let is_startup = broadcast_on_startup.is_some();
2004                    let should_broadcast = broadcast_on_startup
2005                        .take()
2006                        .unwrap_or(sync_cfg.broadcast_on_epoch_change);
2007                    let manager = self.state.transaction_deny_config_manager();
2008                    let publish = |rules: Option<TransactionDenyRules>| {
2009                        if let Err(e) = manager.submit_broadcast(
2010                            rules,
2011                            &components.consensus_adapter,
2012                            &cur_epoch_store,
2013                        ) {
2014                            warn!("Failed to broadcast transaction deny config: {e:?}");
2015                        }
2016                    };
2017                    let action = deny_config_broadcast_payload(
2018                        manager.local().rules(),
2019                        should_broadcast,
2020                        is_startup,
2021                        manager.may_have_outstanding_broadcast(),
2022                    );
2023                    match action {
2024                        DenyConfigBroadcastAction::Skip => {}
2025                        DenyConfigBroadcastAction::Broadcast(rules) => publish(Some(rules)),
2026                        DenyConfigBroadcastAction::Withdraw => publish(None),
2027                    }
2028                }
2029            }
2030
2031            let stop_condition = checkpoint_executor.run_epoch(run_with_range).await;
2032
2033            if stop_condition == StopReason::RunWithRangeCondition {
2034                SuiNode::shutdown(&self).await;
2035                self.shutdown_channel_tx
2036                    .send(run_with_range)
2037                    .expect("RunWithRangeCondition met but failed to send shutdown message");
2038                return Ok(());
2039            }
2040
2041            // Safe to call because we are in the middle of reconfiguration.
2042            let latest_system_state = self
2043                .state
2044                .get_object_cache_reader()
2045                .get_sui_system_state_object_unsafe()
2046                .expect("Read Sui System State object cannot fail");
2047
2048            #[cfg(msim)]
2049            if !self
2050                .sim_state
2051                .sim_safe_mode_expected
2052                .load(Ordering::Relaxed)
2053            {
2054                debug_assert!(!latest_system_state.safe_mode());
2055            }
2056
2057            #[cfg(not(msim))]
2058            debug_assert!(!latest_system_state.safe_mode());
2059
2060            if let Err(err) = self.end_of_epoch_channel.send(latest_system_state.clone())
2061                && self.state.is_fullnode(&cur_epoch_store)
2062            {
2063                warn!(
2064                    "Failed to send end of epoch notification to subscriber: {:?}",
2065                    err
2066                );
2067            }
2068
2069            cur_epoch_store.record_is_safe_mode_metric(latest_system_state.safe_mode());
2070            let new_epoch_start_state = latest_system_state.into_epoch_start_state();
2071
2072            self.auth_agg.store(Arc::new(
2073                self.auth_agg
2074                    .load()
2075                    .recreate_with_new_epoch_start_state(&new_epoch_start_state),
2076            ));
2077
2078            let next_epoch_committee = new_epoch_start_state.get_sui_committee();
2079            let next_epoch = next_epoch_committee.epoch();
2080            assert_eq!(cur_epoch_store.epoch() + 1, next_epoch);
2081
2082            info!(
2083                next_epoch,
2084                "Finished executing all checkpoints in epoch. About to reconfigure the system."
2085            );
2086
2087            fail_point_async!("reconfig_delay");
2088
2089            cur_epoch_store.record_epoch_reconfig_start_time_metric();
2090
2091            update_peer_addresses(
2092                &self.config,
2093                &self.endpoint_manager,
2094                &new_epoch_start_state,
2095                Some(cur_epoch_store.epoch_start_state()),
2096            );
2097
2098            let mut validator_components_lock_guard = self.validator_components.lock().await;
2099
2100            // The following code handles 4 different cases, depending on whether the node
2101            // was a validator in the previous epoch, and whether the node is a validator
2102            // in the new epoch.
2103            let new_epoch_store = self
2104                .reconfigure_state(
2105                    &self.state,
2106                    &cur_epoch_store,
2107                    next_epoch_committee.clone(),
2108                    new_epoch_start_state,
2109                    hasher.clone(),
2110                )
2111                .await;
2112
2113            let new_role = new_epoch_store.node_role();
2114            if !new_role.is_validator()
2115                && let Some(context) = &self.transaction_pool_context
2116            {
2117                context.set_unavailable(next_epoch);
2118            }
2119
2120            let new_validator_components = if let Some(ValidatorComponents {
2121                validator_server_handle,
2122                validator_overload_monitor_handle,
2123                consensus_manager,
2124                consensus_store_pruner,
2125                consensus_adapter,
2126                checkpoint_metrics,
2127                sui_tx_validator_metrics,
2128                admission_queue,
2129                transaction_pool_context,
2130            }) = validator_components_lock_guard.take()
2131            {
2132                info!("Reconfiguring node (was running consensus).");
2133
2134                fail_point_async!("consensus_transaction_pool_reconfig_before_shutdown");
2135                consensus_manager.shutdown().await;
2136                info!("Consensus has shut down.");
2137
2138                // A node that left the committee must stop serving validator RPCs. The server
2139                // owns per-epoch state (e.g. the admission queue), so leaving it running would
2140                // also keep the previous epoch's store alive for the rest of the process.
2141                let validator_server_handle = match validator_server_handle {
2142                    Some(server) if !new_role.is_validator() => {
2143                        info!("Node is no longer a validator, shutting down validator gRPC server");
2144                        server.shutdown().await;
2145                        None
2146                    }
2147                    other => other,
2148                };
2149
2150                if let Some(handle) = &self.address_prober {
2151                    handle.leave_committee();
2152                }
2153
2154                info!("Epoch store finished reconfiguration.");
2155
2156                // No other components should be holding a strong reference to state hasher
2157                // at this point. Confirm here before we swap in the new hasher.
2158                let global_state_hasher_metrics = Arc::into_inner(hasher)
2159                    .expect("Object state hasher should have no other references at this point")
2160                    .metrics();
2161                let new_hasher = Arc::new(GlobalStateHasher::new(
2162                    self.state.get_global_state_hash_store().clone(),
2163                    global_state_hasher_metrics,
2164                ));
2165                let weak_hasher = Arc::downgrade(&new_hasher);
2166                *hasher_guard = Some(new_hasher);
2167
2168                consensus_store_pruner.prune(next_epoch).await;
2169
2170                if new_role.runs_consensus() {
2171                    info!("Restarting consensus as {new_role}");
2172                    let components = Self::start_epoch_specific_validator_components(
2173                        &self.config,
2174                        self.state.clone(),
2175                        consensus_adapter,
2176                        self.checkpoint_store.clone(),
2177                        new_epoch_store.clone(),
2178                        self.state_sync_handle.clone(),
2179                        self.randomness_handle.clone(),
2180                        self.randomness_receiver_handle.clone(),
2181                        consensus_manager,
2182                        consensus_store_pruner,
2183                        weak_hasher,
2184                        self.backpressure_manager.clone(),
2185                        validator_server_handle,
2186                        validator_overload_monitor_handle,
2187                        checkpoint_metrics,
2188                        self.metrics.clone(),
2189                        sui_tx_validator_metrics,
2190                        admission_queue,
2191                        transaction_pool_context.clone(),
2192                        new_role,
2193                    )
2194                    .await?;
2195                    self.update_address_prober_epoch(
2196                        &new_epoch_store,
2197                        &components.consensus_manager,
2198                    );
2199                    Some(components)
2200                } else {
2201                    info!(
2202                        "This node has new role {new_role} and no longer runs consensus after reconfiguration"
2203                    );
2204                    None
2205                }
2206            } else {
2207                // No other components should be holding a strong reference to state hasher
2208                // at this point. Confirm here before we swap in the new hasher.
2209                let global_state_hasher_metrics = Arc::into_inner(hasher)
2210                    .expect("Object state hasher should have no other references at this point")
2211                    .metrics();
2212                let new_hasher = Arc::new(GlobalStateHasher::new(
2213                    self.state.get_global_state_hash_store().clone(),
2214                    global_state_hasher_metrics,
2215                ));
2216                let weak_hasher = Arc::downgrade(&new_hasher);
2217                *hasher_guard = Some(new_hasher);
2218
2219                if new_role.runs_consensus() {
2220                    info!("Promoting node to {new_role}, starting consensus components");
2221
2222                    let mut components = Self::construct_validator_components(
2223                        self.config.clone(),
2224                        self.state.clone(),
2225                        Arc::new(next_epoch_committee.clone()),
2226                        new_epoch_store.clone(),
2227                        self.checkpoint_store.clone(),
2228                        self.state_sync_handle.clone(),
2229                        self.randomness_handle.clone(),
2230                        weak_hasher,
2231                        self.backpressure_manager.clone(),
2232                        &self.registry_service,
2233                        self.transaction_pool_context.clone(),
2234                        &self.consensus_adapter_metrics,
2235                        self.metrics.clone(),
2236                        self.checkpoint_metrics.clone(),
2237                        new_role,
2238                        self.randomness_receiver_handle.clone(),
2239                    )
2240                    .await?;
2241
2242                    if new_role.is_validator() {
2243                        components.validator_server_handle = Some(
2244                            components
2245                                .validator_server_handle
2246                                .take()
2247                                .unwrap()
2248                                .start()
2249                                .await,
2250                        );
2251
2252                        self.endpoint_manager
2253                            .set_consensus_address_updater(components.consensus_manager.clone());
2254                    }
2255
2256                    self.update_address_prober_epoch(
2257                        &new_epoch_store,
2258                        &components.consensus_manager,
2259                    );
2260                    Some(components)
2261                } else {
2262                    None
2263                }
2264            };
2265            *validator_components_lock_guard = new_validator_components;
2266
2267            if cfg!(msim)
2268                && !matches!(
2269                    self.config
2270                        .authority_store_pruning_config
2271                        .num_epochs_to_retain_for_checkpoints(),
2272                    None | Some(u64::MAX) | Some(0)
2273                )
2274            {
2275                self.state
2276                    .prune_checkpoints_for_eligible_epochs_for_testing(
2277                        self.config.clone(),
2278                        sui_core::authority::authority_store_pruner::AuthorityStorePruningMetrics::new_for_test(),
2279                    )
2280                    .await?;
2281            }
2282
2283            let prev_epoch = epoch_store.epoch();
2284            let prev_epoch_store = Arc::downgrade(&epoch_store);
2285            drop(cur_epoch_store);
2286            epoch_store = new_epoch_store;
2287            spawn_monitored_task!(Self::check_epoch_store_released(
2288                prev_epoch_store,
2289                prev_epoch
2290            ));
2291            info!("Reconfiguration finished");
2292        }
2293    }
2294
2295    /// Verifies that nothing holds on to the previous epoch's `AuthorityPerEpochStore` once
2296    /// reconfiguration is complete. A lingering reference keeps that epoch's DB handles and
2297    /// caches alive for the rest of the process lifetime.
2298    async fn check_epoch_store_released(
2299        prev_epoch_store: Weak<AuthorityPerEpochStore>,
2300        prev_epoch: EpochId,
2301    ) {
2302        let grace_period = epoch_store_release_grace_period();
2303        tokio::time::sleep(grace_period).await;
2304        let strong_count = prev_epoch_store.strong_count();
2305        if strong_count > 0 {
2306            debug_fatal!(
2307                "AuthorityPerEpochStore for epoch {prev_epoch} still has {strong_count} strong \
2308                 references {grace_period:?} after reconfiguration"
2309            );
2310        } else {
2311            info!(prev_epoch, "Previous epoch store released");
2312        }
2313    }
2314
2315    async fn shutdown(&self) {
2316        if let Some(validator_components) = &*self.validator_components.lock().await {
2317            validator_components.consensus_manager.shutdown().await;
2318        }
2319        if let Some(context) = &self.transaction_pool_context {
2320            context.set_unavailable(self.state.load_epoch_store_one_call_per_task().epoch());
2321        }
2322    }
2323
2324    async fn reconfigure_state(
2325        &self,
2326        state: &Arc<AuthorityState>,
2327        cur_epoch_store: &AuthorityPerEpochStore,
2328        next_epoch_committee: Committee,
2329        next_epoch_start_system_state: EpochStartSystemState,
2330        global_state_hasher: Arc<GlobalStateHasher>,
2331    ) -> Arc<AuthorityPerEpochStore> {
2332        let next_epoch = next_epoch_committee.epoch();
2333
2334        let last_checkpoint = self
2335            .checkpoint_store
2336            .get_epoch_last_checkpoint(cur_epoch_store.epoch())
2337            .expect("Error loading last checkpoint for current epoch")
2338            .expect("Could not load last checkpoint for current epoch");
2339
2340        let last_checkpoint_seq = *last_checkpoint.sequence_number();
2341
2342        assert_eq!(
2343            Some(last_checkpoint_seq),
2344            self.checkpoint_store
2345                .get_highest_executed_checkpoint_seq_number()
2346                .expect("Error loading highest executed checkpoint sequence number")
2347        );
2348
2349        let epoch_start_configuration = EpochStartConfiguration::new(
2350            next_epoch_start_system_state,
2351            *last_checkpoint.digest(),
2352            state.get_object_store().as_ref(),
2353            EpochFlag::default_flags_for_new_epoch(&state.config),
2354        )
2355        .expect("EpochStartConfiguration construction cannot fail");
2356
2357        let new_epoch_store = self
2358            .state
2359            .reconfigure(
2360                cur_epoch_store,
2361                self.config.supported_protocol_versions.unwrap(),
2362                next_epoch_committee,
2363                epoch_start_configuration,
2364                global_state_hasher,
2365                &self.config.expensive_safety_check_config,
2366                last_checkpoint_seq,
2367            )
2368            .await
2369            .expect("Reconfigure authority state cannot fail");
2370        info!(next_epoch, "Node State has been reconfigured");
2371        assert_eq!(next_epoch, new_epoch_store.epoch());
2372        self.state.get_reconfig_api().update_epoch_flags_metrics(
2373            cur_epoch_store.epoch_start_config().flags(),
2374            new_epoch_store.epoch_start_config().flags(),
2375        );
2376
2377        new_epoch_store
2378    }
2379
2380    pub fn get_config(&self) -> &NodeConfig {
2381        &self.config
2382    }
2383
2384    pub fn randomness_handle(&self) -> randomness::Handle {
2385        self.randomness_handle.clone()
2386    }
2387
2388    pub fn state_sync_handle(&self) -> state_sync::Handle {
2389        self.state_sync_handle.clone()
2390    }
2391
2392    pub fn endpoint_manager(&self) -> &EndpointManager {
2393        &self.endpoint_manager
2394    }
2395
2396    pub async fn address_prober_report(&self) -> Option<address_prober::ProbeReport> {
2397        match &self.address_prober {
2398            Some(handle) => handle.probe_report().await,
2399            None => None,
2400        }
2401    }
2402
2403    /// Get a short prefix of a digest for metric labels
2404    fn get_digest_prefix(digest: impl std::fmt::Display) -> String {
2405        let digest_str = digest.to_string();
2406        if digest_str.len() >= 8 {
2407            digest_str[0..8].to_string()
2408        } else {
2409            digest_str
2410        }
2411    }
2412
2413    /// Check for previously detected forks and handle them appropriately.
2414    /// For validators with fork recovery config, clear the fork if it matches the recovery config.
2415    /// For all other cases, block node startup if a fork is detected.
2416    async fn check_and_recover_forks(
2417        checkpoint_store: &CheckpointStore,
2418        checkpoint_metrics: &CheckpointMetrics,
2419        fork_recovery: Option<&ForkRecoveryConfig>,
2420        build_version: &str,
2421    ) -> Result<()> {
2422        // Manual recovery from operator-supplied overrides; runs regardless of fork_crash_behavior
2423        // and only acts on the checkpoints / transactions explicitly listed in the config.
2424        if let Some(recovery) = fork_recovery {
2425            Self::try_recover_checkpoint_fork(checkpoint_store, recovery)?;
2426            Self::try_recover_transaction_fork(checkpoint_store, recovery)?;
2427        }
2428
2429        let behavior = fork_recovery
2430            .map(|fr| fr.fork_crash_behavior)
2431            .unwrap_or_default();
2432
2433        match behavior {
2434            ForkCrashBehavior::RecoverOncePerVersion => {
2435                Self::try_recover_forks(checkpoint_store, checkpoint_metrics, build_version)?;
2436            }
2437            ForkCrashBehavior::AwaitForkRecovery | ForkCrashBehavior::ReturnError => {}
2438        }
2439
2440        if let Some(fork_info) = checkpoint_store
2441            .get_checkpoint_fork_detected()
2442            .map_err(|e| {
2443                error!("Failed to check for checkpoint fork: {:?}", e);
2444                e
2445            })?
2446        {
2447            Self::handle_checkpoint_fork(
2448                fork_info.checkpoint_seq,
2449                fork_info.checkpoint_digest,
2450                checkpoint_metrics,
2451                fork_recovery,
2452            )
2453            .await?;
2454        }
2455        if let Some(fork_info) = checkpoint_store
2456            .get_transaction_fork_detected()
2457            .map_err(|e| {
2458                error!("Failed to check for transaction fork: {:?}", e);
2459                e
2460            })?
2461        {
2462            Self::handle_transaction_fork(
2463                fork_info.tx_digest,
2464                fork_info.expected_effects_digest,
2465                fork_info.actual_effects_digest,
2466                checkpoint_metrics,
2467                fork_recovery,
2468            )
2469            .await?;
2470        }
2471
2472        Ok(())
2473    }
2474
2475    /// Manual recovery: for each `seq -> digest` override, if the locally computed checkpoint at
2476    /// `seq` differs, clear locally computed checkpoints from `seq` (and the checkpoint fork marker)
2477    /// so the node rebuilds toward the operator-specified digest.
2478    fn try_recover_checkpoint_fork(
2479        checkpoint_store: &CheckpointStore,
2480        recovery: &ForkRecoveryConfig,
2481    ) -> Result<()> {
2482        if recovery.checkpoint_overrides.is_empty() {
2483            return Ok(());
2484        }
2485
2486        for (seq, expected_digest_str) in &recovery.checkpoint_overrides {
2487            let Ok(expected_digest) = CheckpointDigest::from_str(expected_digest_str) else {
2488                anyhow::bail!(
2489                    "Invalid checkpoint digest override for seq {}: {}",
2490                    seq,
2491                    expected_digest_str
2492                );
2493            };
2494
2495            if let Some(local_summary) = checkpoint_store.get_locally_computed_checkpoint(*seq)? {
2496                let local_digest = sui_types::message_envelope::Message::digest(&local_summary);
2497                if local_digest != expected_digest {
2498                    info!(
2499                        seq,
2500                        local = %Self::get_digest_prefix(local_digest),
2501                        expected = %Self::get_digest_prefix(expected_digest),
2502                        "Fork recovery: clearing locally_computed_checkpoints from {} due to digest mismatch",
2503                        seq
2504                    );
2505                    checkpoint_store
2506                        .clear_locally_computed_checkpoints_from(*seq)
2507                        .context(
2508                            "Failed to clear locally computed checkpoints from override seq",
2509                        )?;
2510                }
2511            }
2512        }
2513
2514        if let Some(fork_info) = checkpoint_store.get_checkpoint_fork_detected()?
2515            && recovery
2516                .checkpoint_overrides
2517                .contains_key(&fork_info.checkpoint_seq)
2518        {
2519            info!(
2520                "Fork recovery enabled: clearing checkpoint fork at seq {} with digest {:?}",
2521                fork_info.checkpoint_seq, fork_info.checkpoint_digest
2522            );
2523            checkpoint_store
2524                .clear_checkpoint_fork_detected()
2525                .expect("Failed to clear checkpoint fork detected marker");
2526        }
2527        Ok(())
2528    }
2529
2530    /// Manual recovery: if the forked transaction is listed in transaction_overrides, clear its fork
2531    /// marker so the node proceeds on restart.
2532    fn try_recover_transaction_fork(
2533        checkpoint_store: &CheckpointStore,
2534        recovery: &ForkRecoveryConfig,
2535    ) -> Result<()> {
2536        if recovery.transaction_overrides.is_empty() {
2537            return Ok(());
2538        }
2539
2540        if let Some(fork_info) = checkpoint_store.get_transaction_fork_detected()?
2541            && recovery
2542                .transaction_overrides
2543                .contains_key(&fork_info.tx_digest.to_string())
2544        {
2545            info!(
2546                "Fork recovery enabled: clearing transaction fork for tx {:?}",
2547                fork_info.tx_digest
2548            );
2549            checkpoint_store
2550                .clear_transaction_fork_detected()
2551                .expect("Failed to clear transaction fork detected marker");
2552        }
2553        Ok(())
2554    }
2555
2556    /// Auto-recovery: clear fork markers (the affected seq/tx is read from the markers) so the
2557    /// node re-derives canonically. A marker is cleared only if both gates pass:
2558    ///
2559    /// - Version gate: the marker was recorded by a different binary version than the one now
2560    ///   running. The binary that forked would deterministically fork again, so clearing under
2561    ///   it would only add a second equivocation; the node hangs until a corrected binary is
2562    ///   deployed.
2563    /// - Certification gate: the marker records the certified checkpoint the node diverged
2564    ///   from. Markers carry it only when detection compared against a certificate already
2565    ///   durably persisted locally, so its presence proves the network certified the canonical
2566    ///   outcome. Recovery is deliberate equivocation — the node may have already signed the
2567    ///   forked result and will sign a different one after re-deriving — which is safe only
2568    ///   under that proof: a quorum certificate is irrevocable (a conflicting certificate would
2569    ///   require f+1 double-signers), so re-signing can no longer influence what finalizes.
2570    ///   Self-divergence markers (the node disagreeing with its own prior result rather than a
2571    ///   certificate) carry no certified reference and never pass; the node halts awaiting
2572    ///   operator intervention.
2573    fn try_recover_forks(
2574        checkpoint_store: &CheckpointStore,
2575        checkpoint_metrics: &CheckpointMetrics,
2576        build_version: &str,
2577    ) -> Result<()> {
2578        if let Some(fork_info) = checkpoint_store.get_checkpoint_fork_detected()? {
2579            if fork_info.binary_version == build_version {
2580                error!(
2581                    checkpoint_seq = fork_info.checkpoint_seq,
2582                    build_version,
2583                    "Fork recovery blocked: this binary version produced the checkpoint fork and \
2584                     would fork again. Halting; deploy a corrected binary to recover."
2585                );
2586                checkpoint_metrics
2587                    .fork_auto_recovery_awaiting_new_binary
2588                    .set(1);
2589            } else if fork_info.certified_checkpoint_digest.is_none() {
2590                // The builder re-derived a previously computed checkpoint differently: the fork
2591                // is against the node's own prior result, not a certified checkpoint, so there
2592                // is no canonical outcome to converge toward.
2593                error!(
2594                    checkpoint_seq = fork_info.checkpoint_seq,
2595                    checkpoint_digest = ?fork_info.checkpoint_digest,
2596                    "Fork recovery blocked: the builder re-derived its own previous checkpoint \
2597                     differently, so there is no certified checkpoint proving the canonical \
2598                     outcome to converge toward. Halting awaiting operator intervention."
2599                );
2600                checkpoint_metrics
2601                    .fork_auto_recovery_blocked_uncertified
2602                    .set(1);
2603            } else {
2604                info!(
2605                    checkpoint_seq = fork_info.checkpoint_seq,
2606                    checkpoint_digest = ?fork_info.checkpoint_digest,
2607                    forked_binary_version = ?fork_info.binary_version,
2608                    build_version,
2609                    "Fork recovery: clearing checkpoint fork and locally computed checkpoints \
2610                     from the forked sequence so the builder rebuilds toward the certified \
2611                     checkpoint"
2612                );
2613                checkpoint_store
2614                    .clear_locally_computed_checkpoints_from(fork_info.checkpoint_seq)
2615                    .context("Failed to clear locally computed checkpoints during fork recovery")?;
2616                checkpoint_store.clear_checkpoint_fork_detected()?;
2617                checkpoint_metrics.checkpoint_fork_auto_recovered.set(1);
2618            }
2619        }
2620
2621        if let Some(fork_info) = checkpoint_store.get_transaction_fork_detected()? {
2622            if fork_info.binary_version == build_version {
2623                error!(
2624                    tx_digest = ?fork_info.tx_digest,
2625                    build_version,
2626                    "Fork recovery blocked: this binary version produced the transaction fork and \
2627                     would fork again. Halting; deploy a corrected binary to recover."
2628                );
2629                checkpoint_metrics
2630                    .fork_auto_recovery_awaiting_new_binary
2631                    .set(1);
2632            } else if fork_info.certified_checkpoint_seq.is_none() {
2633                error!(
2634                    tx_digest = ?fork_info.tx_digest,
2635                    "Fork recovery blocked: the expected effects of the forked transaction did \
2636                     not come from a certified checkpoint (they came from this validator's own \
2637                     previously signed effects), so the network has not provably certified the \
2638                     canonical outcome. Halting awaiting operator intervention."
2639                );
2640                checkpoint_metrics
2641                    .fork_auto_recovery_blocked_uncertified
2642                    .set(1);
2643            } else {
2644                info!(
2645                    tx_digest = ?fork_info.tx_digest,
2646                    expected_effects = ?fork_info.expected_effects_digest,
2647                    actual_effects = ?fork_info.actual_effects_digest,
2648                    certified_checkpoint_seq = ?fork_info.certified_checkpoint_seq,
2649                    forked_binary_version = ?fork_info.binary_version,
2650                    build_version,
2651                    "Fork recovery: clearing transaction fork; re-execution will converge toward \
2652                     the canonical certified effects"
2653                );
2654                checkpoint_store.clear_transaction_fork_detected()?;
2655                checkpoint_metrics.transaction_fork_auto_recovered.set(1);
2656            }
2657        }
2658
2659        Ok(())
2660    }
2661
2662    fn get_current_timestamp() -> u64 {
2663        std::time::SystemTime::now()
2664            .duration_since(std::time::SystemTime::UNIX_EPOCH)
2665            .unwrap()
2666            .as_secs()
2667    }
2668
2669    async fn handle_checkpoint_fork(
2670        checkpoint_seq: u64,
2671        checkpoint_digest: CheckpointDigest,
2672        checkpoint_metrics: &CheckpointMetrics,
2673        fork_recovery: Option<&ForkRecoveryConfig>,
2674    ) -> Result<()> {
2675        checkpoint_metrics
2676            .checkpoint_fork_crash_mode
2677            .with_label_values(&[
2678                &checkpoint_seq.to_string(),
2679                &Self::get_digest_prefix(checkpoint_digest),
2680                &Self::get_current_timestamp().to_string(),
2681            ])
2682            .set(1);
2683
2684        let behavior = fork_recovery
2685            .map(|fr| fr.fork_crash_behavior)
2686            .unwrap_or_default();
2687
2688        match behavior {
2689            ForkCrashBehavior::AwaitForkRecovery | ForkCrashBehavior::RecoverOncePerVersion => {
2690                error!(
2691                    checkpoint_seq = checkpoint_seq,
2692                    checkpoint_digest = ?checkpoint_digest,
2693                    "Checkpoint fork detected! Node startup halted. Sleeping indefinitely."
2694                );
2695                futures::future::pending::<()>().await;
2696                unreachable!("pending() should never return");
2697            }
2698            ForkCrashBehavior::ReturnError => {
2699                error!(
2700                    checkpoint_seq = checkpoint_seq,
2701                    checkpoint_digest = ?checkpoint_digest,
2702                    "Checkpoint fork detected! Returning error."
2703                );
2704                Err(anyhow::anyhow!(
2705                    "Checkpoint fork detected! checkpoint_seq: {}, checkpoint_digest: {:?}",
2706                    checkpoint_seq,
2707                    checkpoint_digest
2708                ))
2709            }
2710        }
2711    }
2712
2713    async fn handle_transaction_fork(
2714        tx_digest: TransactionDigest,
2715        expected_effects_digest: TransactionEffectsDigest,
2716        actual_effects_digest: TransactionEffectsDigest,
2717        checkpoint_metrics: &CheckpointMetrics,
2718        fork_recovery: Option<&ForkRecoveryConfig>,
2719    ) -> Result<()> {
2720        checkpoint_metrics
2721            .transaction_fork_crash_mode
2722            .with_label_values(&[
2723                &Self::get_digest_prefix(tx_digest),
2724                &Self::get_digest_prefix(expected_effects_digest),
2725                &Self::get_digest_prefix(actual_effects_digest),
2726                &Self::get_current_timestamp().to_string(),
2727            ])
2728            .set(1);
2729
2730        let behavior = fork_recovery
2731            .map(|fr| fr.fork_crash_behavior)
2732            .unwrap_or_default();
2733
2734        match behavior {
2735            ForkCrashBehavior::AwaitForkRecovery | ForkCrashBehavior::RecoverOncePerVersion => {
2736                error!(
2737                    tx_digest = ?tx_digest,
2738                    expected_effects_digest = ?expected_effects_digest,
2739                    actual_effects_digest = ?actual_effects_digest,
2740                    "Transaction fork detected! Node startup halted. Sleeping indefinitely."
2741                );
2742                futures::future::pending::<()>().await;
2743                unreachable!("pending() should never return");
2744            }
2745            ForkCrashBehavior::ReturnError => {
2746                error!(
2747                    tx_digest = ?tx_digest,
2748                    expected_effects_digest = ?expected_effects_digest,
2749                    actual_effects_digest = ?actual_effects_digest,
2750                    "Transaction fork detected! Returning error."
2751                );
2752                Err(anyhow::anyhow!(
2753                    "Transaction fork detected! tx_digest: {:?}, expected_effects: {:?}, actual_effects: {:?}",
2754                    tx_digest,
2755                    expected_effects_digest,
2756                    actual_effects_digest
2757                ))
2758            }
2759        }
2760    }
2761}
2762
2763#[cfg(not(msim))]
2764impl SuiNode {
2765    async fn fetch_jwks(
2766        _authority: AuthorityName,
2767        provider: &OIDCProvider,
2768    ) -> SuiResult<Vec<(JwkId, JWK)>> {
2769        use fastcrypto_zkp::bn254::zk_login::fetch_jwks;
2770        use sui_types::error::SuiErrorKind;
2771        let client = reqwest::Client::new();
2772        fetch_jwks(provider, &client, true)
2773            .await
2774            .map_err(|_| SuiErrorKind::JWKRetrievalError.into())
2775    }
2776}
2777
2778#[cfg(msim)]
2779impl SuiNode {
2780    pub fn get_sim_node_id(&self) -> sui_simulator::task::NodeId {
2781        self.sim_state.sim_node.id()
2782    }
2783
2784    pub fn set_safe_mode_expected(&self, new_value: bool) {
2785        info!("Setting safe mode expected to {}", new_value);
2786        self.sim_state
2787            .sim_safe_mode_expected
2788            .store(new_value, Ordering::Relaxed);
2789    }
2790
2791    #[allow(unused_variables)]
2792    async fn fetch_jwks(
2793        authority: AuthorityName,
2794        provider: &OIDCProvider,
2795    ) -> SuiResult<Vec<(JwkId, JWK)>> {
2796        get_jwk_injector()(authority, provider)
2797    }
2798}
2799
2800/// The validator gRPC server. Binding is deferred until `start()` so that the rest of the node
2801/// can finish initializing first; once started, the server runs until `shutdown()` is called.
2802enum ValidatorGrpcServer {
2803    // Mutex is only needed to make the future Send
2804    Unstarted(Mutex<BoxFuture<'static, sui_http::ServerHandle>>),
2805    Started(sui_http::ServerHandle),
2806}
2807
2808impl ValidatorGrpcServer {
2809    pub fn new(bind: impl Future<Output = sui_http::ServerHandle> + Send + 'static) -> Self {
2810        Self::Unstarted(Mutex::new(Box::pin(bind)))
2811    }
2812
2813    pub async fn start(self) -> Self {
2814        match self {
2815            Self::Unstarted(bind) => Self::Started(bind.into_inner().await),
2816            Self::Started(_) => self,
2817        }
2818    }
2819
2820    /// Stops accepting requests and waits for in-flight ones to drain. The serving task is
2821    /// owned by sui_http, so merely dropping this value leaves the server running.
2822    pub async fn shutdown(self) {
2823        if let Self::Started(handle) = self {
2824            handle.trigger_shutdown();
2825            match tokio::time::timeout(
2826                VALIDATOR_GRPC_SERVER_SHUTDOWN_TIMEOUT,
2827                handle.wait_for_shutdown(),
2828            )
2829            .await
2830            {
2831                Ok(()) => info!("Validator gRPC server stopped"),
2832                // Shutdown was triggered, so the server still winds down in the background.
2833                Err(e) => warn!(
2834                    error = ?e,
2835                    "Validator gRPC server did not stop within {VALIDATOR_GRPC_SERVER_SHUTDOWN_TIMEOUT:?}"
2836                ),
2837            }
2838        }
2839    }
2840}
2841
2842/// Updates trusted peer addresses in the p2p network (for nodes configured as validators).
2843/// When `prev_epoch_start_state` is provided, validators that are no longer in the committee
2844/// have their Chain addresses cleared.
2845fn update_peer_addresses(
2846    config: &NodeConfig,
2847    endpoint_manager: &EndpointManager,
2848    epoch_start_state: &EpochStartSystemState,
2849    prev_epoch_start_state: Option<&EpochStartSystemState>,
2850) {
2851    if config.consensus_config().is_none() {
2852        return;
2853    }
2854    let new_peers: HashSet<PeerId> = epoch_start_state
2855        .get_validator_as_p2p_peers(config.protocol_public_key())
2856        .into_iter()
2857        .map(|(peer_id, address)| {
2858            endpoint_manager
2859                .update_endpoint(
2860                    EndpointId::P2p(peer_id),
2861                    AddressSource::Chain,
2862                    vec![address],
2863                )
2864                .expect("Updating peer addresses should not fail");
2865            peer_id
2866        })
2867        .collect();
2868
2869    // Clear Chain addresses for validators that left the committee.
2870    if let Some(prev) = prev_epoch_start_state {
2871        for (peer_id, _) in prev.get_validator_as_p2p_peers(config.protocol_public_key()) {
2872            if !new_peers.contains(&peer_id) {
2873                endpoint_manager
2874                    .update_endpoint(EndpointId::P2p(peer_id), AddressSource::Chain, vec![])
2875                    .expect("Clearing peer addresses should not fail");
2876            }
2877        }
2878    }
2879}
2880
2881/// On-disk directories of index backends that no longer exist: `rpc-index`
2882/// was the `RpcIndexStore` that the embedded `sui-rpc-store` replaced, and
2883/// `indexes` was the `IndexStore` behind the removed JSON-RPC service.
2884const LEGACY_INDEX_STORE_DIRS: [&str; 2] = ["rpc-index", "indexes"];
2885
2886/// Remove the on-disk directories of the legacy index backends.
2887///
2888/// Their data is dead, so remove it on startup so a node upgraded from an
2889/// older version does not leave it lingering and wasting disk. Best-effort: a
2890/// node that never ran a legacy backend has nothing to remove, and a failure
2891/// to remove stale data must not block startup.
2892fn remove_legacy_index_stores(db_path: &Path) {
2893    for dir in LEGACY_INDEX_STORE_DIRS {
2894        let legacy_dir = db_path.join(dir);
2895        match std::fs::remove_dir_all(&legacy_dir) {
2896            Ok(()) => info!("removed legacy {dir} directory {}", legacy_dir.display()),
2897            // The common case: the node never ran the legacy backend, or it
2898            // was already cleaned up on a prior startup.
2899            Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
2900            Err(e) => warn!(
2901                "failed to remove legacy {dir} directory {}: {e:?}",
2902                legacy_dir.display()
2903            ),
2904        }
2905    }
2906}
2907
2908async fn build_http_servers(
2909    state: Arc<AuthorityState>,
2910    store: RocksDbStore,
2911    transaction_orchestrator: &Option<Arc<TransactionOrchestrator<NetworkAuthorityClient>>>,
2912    config: &NodeConfig,
2913    prometheus_registry: &Registry,
2914    server_version: ServerVersion,
2915    node_role: NodeRole,
2916    embedded_rpc_store: Option<&EmbeddedRpcStore>,
2917) -> Result<(
2918    HttpServers,
2919    Option<tokio::sync::broadcast::Sender<Arc<Checkpoint>>>,
2920)> {
2921    // Validators do not expose these APIs
2922    if !node_role.is_fullnode() {
2923        return Ok((HttpServers::default(), None));
2924    }
2925
2926    info!("starting rpc service with config: {:?}", config.rpc);
2927
2928    // When the embedded rpc-store is active, gate checkpoint delivery on the
2929    // index so a client that waits for a checkpoint can immediately read its
2930    // indexed state (matching the legacy synchronously-committed index).
2931    let indexed_checkpoint = embedded_rpc_store.map(|embedded| embedded.indexed_checkpoint_fn());
2932    let subscription_watermark_interval = config
2933        .rpc
2934        .as_ref()
2935        .and_then(|rpc| rpc.subscription_watermark_interval);
2936    let subscription_max_subscribers = config
2937        .rpc
2938        .as_ref()
2939        .and_then(|rpc| rpc.subscription_max_subscribers);
2940    let subscription_shards = config.rpc.as_ref().and_then(|rpc| rpc.subscription_shards);
2941    let (subscription_service_checkpoint_sender, subscription_service_handle) =
2942        SubscriptionService::build(
2943            prometheus_registry,
2944            indexed_checkpoint,
2945            subscription_watermark_interval,
2946            subscription_max_subscribers,
2947            subscription_shards,
2948        );
2949    let rpc_router = {
2950        // Serve the index read paths from the embedded rpc-store when it
2951        // is enabled. Raw chain data comes from the perpetual / checkpoint
2952        // stores either way.
2953        let reader: Arc<dyn RpcStateReader> = match embedded_rpc_store {
2954            Some(embedded) => Arc::new(RpcStoreReadStore::new(
2955                state.clone(),
2956                store,
2957                embedded.reader(),
2958            )),
2959            None => Arc::new(RestReadStore::new(state.clone(), store)),
2960        };
2961        let mut rpc_service = sui_rpc_api::RpcService::new(reader);
2962        rpc_service.with_server_version(server_version);
2963
2964        if let Some(config) = config.rpc.clone() {
2965            config.validate()?;
2966            rpc_service.with_config(config);
2967        }
2968
2969        rpc_service.with_metrics(prometheus_registry);
2970        rpc_service.with_subscription_service(subscription_service_handle);
2971
2972        if let Some(transaction_orchestrator) = transaction_orchestrator {
2973            rpc_service.with_executor(transaction_orchestrator.clone());
2974            // The driver knows which validators this node has been submitting to successfully,
2975            // which is what simulate names as a transaction's allowed proposers. Without a
2976            // selector, simulate leaves transactions unrestricted.
2977            if config.enable_simulate_allowed_proposers {
2978                rpc_service
2979                    .with_proposer_selector(transaction_orchestrator.transaction_driver().clone());
2980            }
2981        }
2982
2983        rpc_service.into_router().await
2984    };
2985
2986    let layers = ServiceBuilder::new()
2987        .map_request(|mut request: axum::http::Request<_>| {
2988            if let Some(connect_info) = request.extensions().get::<sui_http::ConnectInfo>() {
2989                let axum_connect_info = axum::extract::ConnectInfo(connect_info.remote_addr);
2990                request.extensions_mut().insert(axum_connect_info);
2991            }
2992            request
2993        })
2994        .layer(axum::middleware::from_fn(server_timing_middleware))
2995        // Setup a permissive CORS policy
2996        .layer(
2997            tower_http::cors::CorsLayer::new()
2998                .allow_methods([http::Method::GET, http::Method::POST])
2999                .allow_origin(tower_http::cors::Any)
3000                .allow_headers(tower_http::cors::Any)
3001                .expose_headers(tower_http::cors::Any),
3002        );
3003
3004    let router = rpc_router.layer(layers);
3005
3006    // On top of sui-http's hardened defaults (bounded concurrent streams;
3007    // transport keepalives stay disabled by default), bound connection
3008    // lifetime: GOAWAY at the configured age and force-close after the grace
3009    // period. The hard close is the only server-side mechanism that reclaims
3010    // streams wedged behind HTTP/2 flow-control windows that a stalled peer
3011    // never reopens, and connection age also bounds how long a vanished peer
3012    // can pin connection state, which keepalives would otherwise detect.
3013    let server_config = {
3014        let rpc_config = config.rpc().cloned().unwrap_or_default();
3015        let mut server_config = sui_http::Config::default()
3016            .max_connection_age_grace(rpc_config.max_connection_age_grace());
3017        if let Some(age) = rpc_config.max_connection_age() {
3018            server_config = server_config.max_connection_age(age);
3019        }
3020        server_config
3021    };
3022
3023    let https = if let Some((tls_config, https_address)) = config
3024        .rpc()
3025        .and_then(|config| config.tls_config().map(|tls| (tls, config.https_address())))
3026    {
3027        let tls_server_config = https_rustls_config(tls_config.cert(), tls_config.key())?;
3028        let https = sui_http::Builder::new()
3029            .config(server_config.clone())
3030            .tls_config(tls_server_config)
3031            .serve(https_address, router.clone())
3032            .map_err(|e| anyhow::anyhow!(e))?;
3033
3034        info!(
3035            https_address =? https.local_addr(),
3036            "HTTPS rpc server listening on {}",
3037            https.local_addr()
3038        );
3039
3040        Some(https)
3041    } else {
3042        None
3043    };
3044
3045    let http = sui_http::Builder::new()
3046        .config(server_config)
3047        .serve(&config.json_rpc_address, router)
3048        .map_err(|e| anyhow::anyhow!(e))?;
3049
3050    info!(
3051        http_address =? http.local_addr(),
3052        "HTTP rpc server listening on {}",
3053        http.local_addr()
3054    );
3055
3056    Ok((
3057        HttpServers {
3058            http: Some(http),
3059            https,
3060        },
3061        Some(subscription_service_checkpoint_sender),
3062    ))
3063}
3064
3065/// Action to take with our deny-config network vote at a startup or epoch-change
3066/// transition, after reconciling against the current local config.
3067///
3068/// An oversized local config can't go on the wire (the consensus validator would reject
3069/// it), so it is treated as unshareable: on startup a stale prior vote is withdrawn
3070/// rather than left to diverge from the local config.
3071#[derive(Debug, PartialEq, Eq)]
3072enum DenyConfigBroadcastAction {
3073    /// Leave any existing on-network vote untouched.
3074    Skip,
3075    /// Publish these rules as our current vote.
3076    Broadcast(TransactionDenyRules),
3077    /// Publish a withdrawal of any prior vote.
3078    Withdraw,
3079}
3080
3081fn deny_config_broadcast_payload(
3082    local_rules: &TransactionDenyRules,
3083    broadcast: bool,
3084    is_startup: bool,
3085    may_have_prior_broadcast: bool,
3086) -> DenyConfigBroadcastAction {
3087    let shareable = local_rules.is_empty() || local_rules.check_share_limits().is_ok();
3088    if broadcast && shareable {
3089        // An empty local config broadcasts a withdrawal.
3090        if local_rules.is_empty() {
3091            DenyConfigBroadcastAction::Withdraw
3092        } else {
3093            DenyConfigBroadcastAction::Broadcast(local_rules.clone())
3094        }
3095    } else if is_startup && may_have_prior_broadcast {
3096        DenyConfigBroadcastAction::Withdraw
3097    } else {
3098        DenyConfigBroadcastAction::Skip
3099    }
3100}
3101
3102/// Builds the HTTPS RPC server's rustls config from PEM files, pinning the
3103/// ring crypto provider.
3104///
3105/// `sui_http::Builder::tls_single_cert` resolves the provider from rustls
3106/// crate features and panics at runtime when more than one provider feature is
3107/// enabled in the final binary (e.g. `aws-lc-rs` is pulled in through
3108/// `aws-config` in the `sui` CLI), so the provider is pinned explicitly here
3109/// instead.
3110fn https_rustls_config(cert: &str, key: &str) -> Result<sui_http::rustls::ServerConfig> {
3111    use sui_http::rustls;
3112    use sui_http::rustls::pki_types::pem::PemObject;
3113
3114    let certs = rustls::pki_types::CertificateDer::pem_file_iter(cert)
3115        .with_context(|| format!("failed to read TLS certificate chain from {cert}"))?
3116        .collect::<Result<Vec<_>, _>>()
3117        .with_context(|| format!("failed to parse TLS certificate chain from {cert}"))?;
3118    let private_key = rustls::pki_types::PrivateKeyDer::from_pem_file(key)
3119        .with_context(|| format!("failed to read TLS private key from {key}"))?;
3120    let config = rustls::ServerConfig::builder_with_provider(Arc::new(
3121        rustls::crypto::ring::default_provider(),
3122    ))
3123    .with_protocol_versions(rustls::DEFAULT_VERSIONS)?
3124    .with_no_client_auth()
3125    .with_single_cert(certs, private_key)?;
3126    Ok(config)
3127}
3128
3129#[derive(Default)]
3130struct HttpServers {
3131    #[allow(unused)]
3132    http: Option<sui_http::ServerHandle>,
3133    #[allow(unused)]
3134    https: Option<sui_http::ServerHandle>,
3135}
3136
3137#[cfg(test)]
3138mod tests {
3139    use super::*;
3140    use prometheus::Registry;
3141    use std::collections::BTreeMap;
3142    use sui_config::node::{ForkCrashBehavior, ForkRecoveryConfig};
3143    use sui_core::checkpoints::{CheckpointMetrics, CheckpointStore};
3144    use sui_types::digests::{CheckpointDigest, TransactionDigest, TransactionEffectsDigest};
3145
3146    #[test]
3147    fn deny_config_broadcast_payload_decisions() {
3148        let empty = TransactionDenyRules::default();
3149        let populated = TransactionDenyRules {
3150            package_publish_disabled: true,
3151            ..Default::default()
3152        };
3153        let oversized = TransactionDenyRules {
3154            zklogin_disabled_providers: std::iter::once(
3155                "x".repeat(TransactionDenyRules::MAX_ZKLOGIN_PROVIDER_LENGTH + 1),
3156            )
3157            .collect(),
3158            ..Default::default()
3159        };
3160
3161        // broadcast=true: (re-)broadcast the current local config.
3162        assert_eq!(
3163            deny_config_broadcast_payload(&populated, true, true, false),
3164            DenyConfigBroadcastAction::Broadcast(populated.clone()),
3165        );
3166        // broadcast=true with an empty local config: broadcast a withdrawal.
3167        assert_eq!(
3168            deny_config_broadcast_payload(&empty, true, true, false),
3169            DenyConfigBroadcastAction::Withdraw,
3170        );
3171        // Not broadcasting, startup, prior vote -> withdraw it.
3172        assert_eq!(
3173            deny_config_broadcast_payload(&populated, false, true, true),
3174            DenyConfigBroadcastAction::Withdraw,
3175        );
3176        // Not broadcasting, startup, no prior vote -> nothing.
3177        assert_eq!(
3178            deny_config_broadcast_payload(&populated, false, true, false),
3179            DenyConfigBroadcastAction::Skip,
3180        );
3181        // Epoch change (not startup) without broadcast -> nothing, even with a prior vote.
3182        assert_eq!(
3183            deny_config_broadcast_payload(&populated, false, false, true),
3184            DenyConfigBroadcastAction::Skip,
3185        );
3186        // An oversized local config is unshareable: broadcast=true falls back to
3187        // withdrawing a prior vote instead of attempting an (oversized) broadcast.
3188        assert_eq!(
3189            deny_config_broadcast_payload(&oversized, true, true, true),
3190            DenyConfigBroadcastAction::Withdraw,
3191        );
3192        assert_eq!(
3193            deny_config_broadcast_payload(&oversized, true, true, false),
3194            DenyConfigBroadcastAction::Skip,
3195        );
3196    }
3197
3198    // Present legacy `rpc-index` and `indexes` directories are removed, while
3199    // their siblings (such as the perpetual `store`) are left untouched, and a
3200    // missing directory is a no-op.
3201    #[test]
3202    fn removes_only_the_legacy_index_directories() {
3203        let db = tempfile::tempdir().unwrap();
3204        let rpc_index = db.path().join("rpc-index");
3205        let indexes = db.path().join("indexes");
3206        let sibling = db.path().join("store");
3207        for dir in [&rpc_index, &indexes, &sibling] {
3208            std::fs::create_dir(dir).unwrap();
3209            std::fs::write(dir.join("CURRENT"), b"stale").unwrap();
3210        }
3211
3212        remove_legacy_index_stores(db.path());
3213        assert!(
3214            !rpc_index.exists(),
3215            "legacy rpc-index directory should be gone"
3216        );
3217        assert!(!indexes.exists(), "legacy indexes directory should be gone");
3218        assert!(sibling.exists(), "sibling stores must be left untouched");
3219
3220        // Idempotent: a second run (nothing to remove) does not error or touch
3221        // the siblings.
3222        remove_legacy_index_stores(db.path());
3223        assert!(!rpc_index.exists());
3224        assert!(!indexes.exists());
3225        assert!(sibling.exists());
3226    }
3227
3228    // Halt / ReturnError never clear markers; ReturnError surfaces the fork as a startup error.
3229    #[tokio::test]
3230    async fn test_return_error_does_not_recover() {
3231        let checkpoint_store = CheckpointStore::new_for_tests();
3232        let checkpoint_metrics = CheckpointMetrics::new(&Registry::new());
3233        let cfg = ForkRecoveryConfig {
3234            transaction_overrides: Default::default(),
3235            checkpoint_overrides: Default::default(),
3236            fork_crash_behavior: ForkCrashBehavior::ReturnError,
3237        };
3238
3239        // Checkpoint fork.
3240        checkpoint_store
3241            .record_checkpoint_fork_detected(
3242                42,
3243                CheckpointDigest::random(),
3244                Some(CheckpointDigest::random()),
3245            )
3246            .unwrap();
3247        let r = SuiNode::check_and_recover_forks(
3248            &checkpoint_store,
3249            &checkpoint_metrics,
3250            Some(&cfg),
3251            "v1",
3252        )
3253        .await;
3254        assert!(
3255            r.unwrap_err()
3256                .to_string()
3257                .contains("Checkpoint fork detected")
3258        );
3259        assert!(
3260            checkpoint_store
3261                .get_checkpoint_fork_detected()
3262                .unwrap()
3263                .is_some()
3264        );
3265        checkpoint_store.clear_checkpoint_fork_detected().unwrap();
3266
3267        // Transaction fork.
3268        checkpoint_store
3269            .record_transaction_fork_detected(
3270                TransactionDigest::random(),
3271                TransactionEffectsDigest::random(),
3272                TransactionEffectsDigest::random(),
3273                Some(1),
3274            )
3275            .unwrap();
3276        let r = SuiNode::check_and_recover_forks(
3277            &checkpoint_store,
3278            &checkpoint_metrics,
3279            Some(&cfg),
3280            "v1",
3281        )
3282        .await;
3283        assert!(
3284            r.unwrap_err()
3285                .to_string()
3286                .contains("Transaction fork detected")
3287        );
3288    }
3289
3290    // A fork marker carrying the currently running binary version is never cleared — the binary
3291    // that forked would deterministically fork again — so the node hangs until a corrected
3292    // binary (different version) runs recovery.
3293    #[tokio::test]
3294    async fn test_same_binary_version_does_not_recover() {
3295        let checkpoint_store = CheckpointStore::new_for_tests();
3296        let checkpoint_metrics = CheckpointMetrics::new(&Registry::new());
3297        let seq = 7;
3298        // The fork is recorded by binary "v1", against a certified checkpoint.
3299        checkpoint_store.set_binary_version("v1");
3300        checkpoint_store
3301            .record_checkpoint_fork_detected(
3302                seq,
3303                CheckpointDigest::random(),
3304                Some(CheckpointDigest::random()),
3305            )
3306            .unwrap();
3307
3308        // Restarting the same binary: recovery refused despite certification.
3309        SuiNode::try_recover_forks(&checkpoint_store, &checkpoint_metrics, "v1").unwrap();
3310        assert!(
3311            checkpoint_store
3312                .get_checkpoint_fork_detected()
3313                .unwrap()
3314                .is_some()
3315        );
3316        assert_eq!(
3317            checkpoint_metrics
3318                .fork_auto_recovery_awaiting_new_binary
3319                .get(),
3320            1
3321        );
3322        assert_eq!(checkpoint_metrics.checkpoint_fork_auto_recovered.get(), 0);
3323
3324        // Corrected binary (new version): recovers.
3325        SuiNode::try_recover_forks(&checkpoint_store, &checkpoint_metrics, "v2").unwrap();
3326        assert!(
3327            checkpoint_store
3328                .get_checkpoint_fork_detected()
3329                .unwrap()
3330                .is_none()
3331        );
3332        assert_eq!(checkpoint_metrics.checkpoint_fork_auto_recovered.get(), 1);
3333    }
3334
3335    // The default behavior (RecoverOncePerVersion) recovers with no fork-recovery config present.
3336    #[tokio::test]
3337    async fn test_default_recovers() {
3338        let checkpoint_store = CheckpointStore::new_for_tests();
3339        let checkpoint_metrics = CheckpointMetrics::new(&Registry::new());
3340
3341        let tx_digest = TransactionDigest::random();
3342        // The fork was recorded by binary "v1"; its expected effects came from certified
3343        // checkpoint 3, so recovery under "v2" is permitted.
3344        checkpoint_store.set_binary_version("v1");
3345        checkpoint_store
3346            .record_transaction_fork_detected(
3347                tx_digest,
3348                TransactionEffectsDigest::random(),
3349                TransactionEffectsDigest::random(),
3350                Some(3),
3351            )
3352            .unwrap();
3353
3354        let r =
3355            SuiNode::check_and_recover_forks(&checkpoint_store, &checkpoint_metrics, None, "v2")
3356                .await;
3357        assert!(r.is_ok());
3358        assert!(
3359            checkpoint_store
3360                .get_transaction_fork_detected()
3361                .unwrap()
3362                .is_none()
3363        );
3364        assert_eq!(checkpoint_metrics.transaction_fork_auto_recovered.get(), 1);
3365    }
3366
3367    // checkpoint_overrides clears only the checkpoint fork marker (when the forked seq is listed); it
3368    // is decoupled from the transaction fork marker, which is cleared by transaction_overrides.
3369    #[tokio::test]
3370    async fn test_checkpoint_overrides_clear_checkpoint_marker_only() {
3371        let checkpoint_store = CheckpointStore::new_for_tests();
3372        let seq = 9;
3373
3374        checkpoint_store
3375            .record_checkpoint_fork_detected(
3376                seq,
3377                CheckpointDigest::random(),
3378                Some(CheckpointDigest::random()),
3379            )
3380            .unwrap();
3381        checkpoint_store
3382            .record_transaction_fork_detected(
3383                TransactionDigest::random(),
3384                TransactionEffectsDigest::random(),
3385                TransactionEffectsDigest::random(),
3386                None,
3387            )
3388            .unwrap();
3389
3390        // No overrides: both markers are left intact.
3391        SuiNode::try_recover_checkpoint_fork(&checkpoint_store, &ForkRecoveryConfig::default())
3392            .unwrap();
3393        assert!(
3394            checkpoint_store
3395                .get_checkpoint_fork_detected()
3396                .unwrap()
3397                .is_some()
3398        );
3399        assert!(
3400            checkpoint_store
3401                .get_transaction_fork_detected()
3402                .unwrap()
3403                .is_some()
3404        );
3405
3406        // Override for the forked seq: clears the checkpoint marker but leaves the transaction marker.
3407        let mut checkpoint_overrides = BTreeMap::new();
3408        checkpoint_overrides.insert(seq, CheckpointDigest::random().to_string());
3409        let cfg = ForkRecoveryConfig {
3410            transaction_overrides: Default::default(),
3411            checkpoint_overrides,
3412            fork_crash_behavior: ForkCrashBehavior::AwaitForkRecovery,
3413        };
3414        SuiNode::try_recover_checkpoint_fork(&checkpoint_store, &cfg).unwrap();
3415        assert!(
3416            checkpoint_store
3417                .get_checkpoint_fork_detected()
3418                .unwrap()
3419                .is_none()
3420        );
3421        assert!(
3422            checkpoint_store
3423                .get_transaction_fork_detected()
3424                .unwrap()
3425                .is_some(),
3426            "checkpoint_overrides must not touch the transaction fork marker"
3427        );
3428    }
3429
3430    // transaction_overrides clears the transaction fork marker when the forked tx is listed.
3431    #[tokio::test]
3432    async fn test_transaction_overrides_clear_transaction_marker() {
3433        let checkpoint_store = CheckpointStore::new_for_tests();
3434        let tx_digest = TransactionDigest::random();
3435        checkpoint_store
3436            .record_transaction_fork_detected(
3437                tx_digest,
3438                TransactionEffectsDigest::random(),
3439                TransactionEffectsDigest::random(),
3440                None,
3441            )
3442            .unwrap();
3443
3444        // Unrelated override: marker stays.
3445        let mut transaction_overrides = BTreeMap::new();
3446        transaction_overrides.insert(TransactionDigest::random().to_string(), String::new());
3447        let cfg = ForkRecoveryConfig {
3448            transaction_overrides,
3449            checkpoint_overrides: Default::default(),
3450            fork_crash_behavior: ForkCrashBehavior::AwaitForkRecovery,
3451        };
3452        SuiNode::try_recover_transaction_fork(&checkpoint_store, &cfg).unwrap();
3453        assert!(
3454            checkpoint_store
3455                .get_transaction_fork_detected()
3456                .unwrap()
3457                .is_some()
3458        );
3459
3460        // Override for the forked tx: marker cleared.
3461        let mut transaction_overrides = BTreeMap::new();
3462        transaction_overrides.insert(tx_digest.to_string(), String::new());
3463        let cfg = ForkRecoveryConfig {
3464            transaction_overrides,
3465            checkpoint_overrides: Default::default(),
3466            fork_crash_behavior: ForkCrashBehavior::AwaitForkRecovery,
3467        };
3468        SuiNode::try_recover_transaction_fork(&checkpoint_store, &cfg).unwrap();
3469        assert!(
3470            checkpoint_store
3471                .get_transaction_fork_detected()
3472                .unwrap()
3473                .is_none()
3474        );
3475    }
3476
3477    // Under RecoverOncePerVersion, a checkpoint override clears the fork via the manual path
3478    // even when the auto path would refuse (here: the fork was recorded by the currently running
3479    // binary version and the sequence is not certified).
3480    #[tokio::test]
3481    async fn test_override_clears_fork_auto_recovery_refuses() {
3482        let checkpoint_store = CheckpointStore::new_for_tests();
3483        let checkpoint_metrics = CheckpointMetrics::new(&Registry::new());
3484        let seq = 5;
3485        checkpoint_store.set_binary_version("v1");
3486        checkpoint_store
3487            .record_checkpoint_fork_detected(
3488                seq,
3489                CheckpointDigest::random(),
3490                Some(CheckpointDigest::random()),
3491            )
3492            .unwrap();
3493
3494        let mut checkpoint_overrides = BTreeMap::new();
3495        checkpoint_overrides.insert(seq, CheckpointDigest::random().to_string());
3496        let cfg = ForkRecoveryConfig {
3497            transaction_overrides: Default::default(),
3498            checkpoint_overrides,
3499            fork_crash_behavior: ForkCrashBehavior::RecoverOncePerVersion,
3500        };
3501
3502        SuiNode::check_and_recover_forks(&checkpoint_store, &checkpoint_metrics, Some(&cfg), "v1")
3503            .await
3504            .unwrap();
3505
3506        assert!(
3507            checkpoint_store
3508                .get_checkpoint_fork_detected()
3509                .unwrap()
3510                .is_none()
3511        );
3512    }
3513
3514    // A self-divergence checkpoint fork (the builder re-derived its own previous checkpoint
3515    // differently; no certified digest in the marker) is never auto-recovered, even under a new
3516    // binary version, because neither result is proven canonical. It requires operator
3517    // overrides.
3518    #[tokio::test]
3519    async fn test_self_divergence_checkpoint_fork_blocks_recovery() {
3520        let checkpoint_store = CheckpointStore::new_for_tests();
3521        let checkpoint_metrics = CheckpointMetrics::new(&Registry::new());
3522        let seq = 17;
3523        checkpoint_store.set_binary_version("v1");
3524        checkpoint_store
3525            .record_checkpoint_fork_detected(seq, CheckpointDigest::random(), None)
3526            .unwrap();
3527
3528        SuiNode::try_recover_forks(&checkpoint_store, &checkpoint_metrics, "v2").unwrap();
3529
3530        assert!(
3531            checkpoint_store
3532                .get_checkpoint_fork_detected()
3533                .unwrap()
3534                .is_some()
3535        );
3536        assert_eq!(
3537            checkpoint_metrics
3538                .fork_auto_recovery_blocked_uncertified
3539                .get(),
3540            1
3541        );
3542        assert_eq!(checkpoint_metrics.checkpoint_fork_auto_recovered.get(), 0);
3543    }
3544
3545    // A transaction fork whose expected effects did not come from a certified checkpoint (i.e.
3546    // they came from this validator's own previously signed effects) is never auto-recovered,
3547    // even on a new binary version.
3548    #[tokio::test]
3549    async fn test_uncertified_transaction_fork_blocks_recovery() {
3550        let checkpoint_store = CheckpointStore::new_for_tests();
3551        let checkpoint_metrics = CheckpointMetrics::new(&Registry::new());
3552        checkpoint_store.set_binary_version("v1");
3553        checkpoint_store
3554            .record_transaction_fork_detected(
3555                TransactionDigest::random(),
3556                TransactionEffectsDigest::random(),
3557                TransactionEffectsDigest::random(),
3558                None,
3559            )
3560            .unwrap();
3561
3562        SuiNode::try_recover_forks(&checkpoint_store, &checkpoint_metrics, "v2").unwrap();
3563        SuiNode::try_recover_forks(&checkpoint_store, &checkpoint_metrics, "v3").unwrap();
3564
3565        assert!(
3566            checkpoint_store
3567                .get_transaction_fork_detected()
3568                .unwrap()
3569                .is_some()
3570        );
3571        assert_eq!(
3572            checkpoint_metrics
3573                .fork_auto_recovery_blocked_uncertified
3574                .get(),
3575            1
3576        );
3577    }
3578}