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