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