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