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