Skip to main content

sui_node/
lib.rs

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