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