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