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