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