Skip to main content

sui_node/
lib.rs

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