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