1#[cfg(msim)]
6use std::sync::atomic::Ordering;
7use std::{
8 collections::HashMap,
9 fmt,
10 future::Future,
11 num::NonZeroUsize,
12 sync::{Arc, Weak},
13 time::Duration,
14};
15
16use anemo::Network;
17use anemo_tower::{
18 callback::CallbackLayer,
19 trace::{DefaultMakeSpan, DefaultOnFailure, TraceLayer},
20};
21use anyhow::{Result, anyhow};
22use arc_swap::ArcSwap;
23use futures::future::BoxFuture;
24pub use handle::IotaNodeHandle;
25use iota_common::{debug_fatal, fatal};
26use iota_config::{
27 ConsensusConfig, NodeConfig, node::RunWithRange, node_config_metrics::NodeConfigMetrics,
28};
29use iota_core::{
30 authority::{
31 AuthorityState, AuthorityStore, ExecutionEnv, RandomnessRoundReceiver,
32 authority_per_epoch_store::AuthorityPerEpochStore,
33 authority_store_tables::{AuthorityPerpetualTables, AuthorityPerpetualTablesOptions},
34 backpressure::BackpressureManager,
35 epoch_start_configuration::{EpochFlag, EpochStartConfigTrait, EpochStartConfiguration},
36 shared_object_version_manager::Schedulable,
37 },
38 authority_aggregator::{
39 AggregatorSendCapabilityNotificationError, AuthAggMetrics, AuthorityAggregator,
40 },
41 authority_client::NetworkAuthorityClient,
42 authority_server::{
43 ValidatorService, ValidatorServiceMetrics, soft_lock::PreConsensusSoftLocks,
44 },
45 checkpoint_progress_tracker::CheckpointProgressTracker,
46 checkpoints::{
47 CheckpointMetrics, CheckpointService, CheckpointStore, FullCheckpointContentsCache,
48 FullCheckpointContentsCacheMetrics, SendCheckpointToStateSync, SubmitCheckpointToConsensus,
49 checkpoint_executor::{CheckpointExecutor, StopReason, metrics::CheckpointExecutorMetrics},
50 },
51 connection_monitor::ConnectionMonitor,
52 consensus_adapter::{
53 CheckConnection, ConnectionMonitorStatus, ConsensusAdapter, ConsensusAdapterMetrics,
54 ConsensusClient,
55 },
56 consensus_handler::ConsensusHandlerInitializer,
57 consensus_manager::{ConsensusManager, ConsensusManagerTrait, UpdatableConsensusClient},
58 consensus_validator::{IotaTxValidator, IotaTxValidatorMetrics},
59 epoch::{
60 committee_store::CommitteeStore, consensus_store_pruner::ConsensusStorePruner,
61 epoch_metrics::EpochMetrics, randomness::RandomnessManager,
62 reconfiguration::ReconfigurationInitiator,
63 },
64 execution_cache::build_execution_cache,
65 execution_scheduler::ExecutionSchedulerAPI,
66 global_state_hasher::{GlobalStateHashMetrics, GlobalStateHasher},
67 grpc_indexes::{GRPC_INDEXES_DIR, GrpcIndexesStore},
68 jsonrpc_index::IndexStore,
69 module_cache_metrics::ResolverMetrics,
70 overload_monitor::{consensus_queue_overload_monitor, overload_monitor},
71 safe_client::SafeClientMetricsBase,
72 signature_verifier::SignatureVerifierMetrics,
73 storage::{GrpcReadStore, RocksDbStore},
74 transaction_orchestrator::TransactionOrchestrator,
75 validator_tx_finalizer::ValidatorTxFinalizer,
76};
77use iota_genesis_common::MigrationTxDataExt;
78use iota_grpc_server::{GrpcReader, GrpcServerHandle, start_grpc_server};
79use iota_json_rpc::{
80 JsonRpcServerBuilder, coin_api::CoinReadApi, governance_api::GovernanceReadApi,
81 indexer_api::IndexerApi, move_utils::MoveUtils, read_api::ReadApi,
82 transaction_builder_api::TransactionBuilderApi,
83 transaction_execution_api::TransactionExecutionApi,
84};
85use iota_json_rpc_api::JsonRpcMetrics;
86use iota_macros::{fail_point, fail_point_async, replay_log};
87use iota_metrics::{
88 RegistryID, RegistryService,
89 metrics_network::{MetricsMakeCallbackHandler, NetworkConnectionMetrics, NetworkMetrics},
90 server_timing_middleware, spawn_monitored_task,
91};
92use iota_names::config::IotaNamesConfig;
93use iota_network::{
94 api::{ValidatorPeerServer, ValidatorServer, ValidatorV2Server},
95 discovery,
96 discovery::TrustedPeerChangeEvent,
97 randomness, state_sync,
98};
99use iota_network_stack::server::{IOTA_TLS_SERVER_NAME, ServerBuilder};
100use iota_node_transaction_builder::NodeTransactionBuilderLedgerClient;
101use iota_protocol_config::{ProtocolConfig, ProtocolVersion};
102use iota_sdk_types::{
103 RandomnessRound,
104 crypto::{Intent, IntentMessage, IntentScope},
105};
106use iota_snapshot::uploader::StateSnapshotUploader;
107use iota_storage::{
108 http_key_value_store::HttpKVStore,
109 key_value_store::{FallbackTransactionKVStore, TransactionKeyValueStore},
110 key_value_store_metrics::KeyValueStoreMetrics,
111};
112use iota_types::{
113 base_types::{AuthorityName, ConciseableName, EpochId},
114 committee::Committee,
115 crypto::{AuthoritySignature, IotaAuthoritySignature, KeypairTraits},
116 digests::ChainIdentifier,
117 error::{IotaError, IotaResult},
118 executable_transaction::VerifiedExecutableTransaction,
119 full_checkpoint_content::CheckpointData,
120 iota_system_state::{
121 IotaSystemState, IotaSystemStateTrait,
122 epoch_start_iota_system_state::{EpochStartSystemState, EpochStartSystemStateTrait},
123 },
124 messages_checkpoint::CheckpointSummaryExt,
125 messages_consensus::{
126 AuthorityCapabilitiesV1, ConsensusTransaction, ConsensusTransactionKind,
127 SignedAuthorityCapabilitiesV1, TransactionDenyRuleProposal,
128 },
129 messages_grpc::HandleCapabilityNotificationRequestV1,
130 quorum_driver_types::QuorumDriverEffectsQueueResult,
131 supported_protocol_versions::SupportedProtocolVersions,
132 transaction::{SenderSignedTransactionAPI, TransactionEnvelope, VerifiedCertificate},
133};
134use prometheus_filtered::Registry;
135#[cfg(msim)]
136use simulator::*;
137use tap::tap::TapFallible;
138use tokio::{
139 sync::{Mutex, broadcast, mpsc, watch},
140 task::{JoinHandle, JoinSet},
141};
142use tokio_util::sync::CancellationToken;
143use tower::ServiceBuilder;
144use tracing::{Instrument, debug, error, error_span, info, trace_span, warn};
145use typed_store::{
146 DBMetrics,
147 rocks::{check_and_mark_db_corruption, default_db_options, unmark_db_corruption},
148};
149
150use crate::metrics::{GrpcMetrics, IotaNodeMetrics};
151
152pub mod admin;
153mod handle;
154pub mod metrics;
155
156pub struct ValidatorComponents {
157 validator_server_handle: SpawnOnce,
158 validator_overload_monitor_handle: Option<JoinHandle<()>>,
159 consensus_queue_overload_monitor_handle: Option<JoinHandle<()>>,
164 soft_lock_sweep_handle: JoinHandle<()>,
167 overload_notifier_handle: Option<JoinHandle<()>>,
168 consensus_manager: Arc<ConsensusManager>,
169 consensus_store_pruner: ConsensusStorePruner,
170 consensus_adapter: Arc<ConsensusAdapter>,
171 soft_locks: Arc<PreConsensusSoftLocks>,
172 checkpoint_service_tasks: JoinSet<()>,
174 checkpoint_metrics: Arc<CheckpointMetrics>,
175 iota_tx_validator_metrics: Arc<IotaTxValidatorMetrics>,
176 validator_registry_id: RegistryID,
177}
178
179#[cfg(msim)]
180mod simulator {
181 use std::sync::atomic::AtomicBool;
182
183 pub(super) struct SimState {
184 pub sim_node: iota_simulator::runtime::NodeHandle,
185 pub sim_safe_mode_expected: AtomicBool,
186 _leak_detector: iota_simulator::NodeLeakDetector,
187 }
188
189 impl Default for SimState {
190 fn default() -> Self {
191 Self {
192 sim_node: iota_simulator::runtime::NodeHandle::current(),
193 sim_safe_mode_expected: AtomicBool::new(false),
194 _leak_detector: iota_simulator::NodeLeakDetector::new(),
195 }
196 }
197 }
198}
199
200#[derive(Clone)]
201pub struct ServerVersion {
202 pub bin: &'static str,
203 pub version: &'static str,
204}
205
206impl ServerVersion {
207 pub fn new(bin: &'static str, version: &'static str) -> Self {
208 Self { bin, version }
209 }
210}
211
212impl std::fmt::Display for ServerVersion {
213 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
214 f.write_str(self.bin)?;
215 f.write_str("/")?;
216 f.write_str(self.version)
217 }
218}
219
220pub struct IotaNode {
221 config: NodeConfig,
222 validator_components: Mutex<Option<ValidatorComponents>>,
223 _http_server: Option<iota_http::ServerHandle>,
225 state: Arc<AuthorityState>,
226 transaction_orchestrator: Option<Arc<TransactionOrchestrator<NetworkAuthorityClient>>>,
227 registry_service: RegistryService,
228 metrics: Arc<IotaNodeMetrics>,
229
230 _discovery: discovery::Handle,
231 state_sync_handle: state_sync::Handle,
232 randomness_handle: randomness::Handle,
233 checkpoint_store: Arc<CheckpointStore>,
234 state_sync_store: RocksDbStore,
235 global_state_hasher: Mutex<Option<Arc<GlobalStateHasher>>>,
236 connection_monitor_status: Arc<ConnectionMonitorStatus>,
237
238 end_of_epoch_channel: broadcast::Sender<IotaSystemState>,
240
241 trusted_peer_change_tx: watch::Sender<TrustedPeerChangeEvent>,
244
245 backpressure_manager: Arc<BackpressureManager>,
246
247 checkpoint_progress_tracker: Arc<CheckpointProgressTracker>,
248
249 #[cfg(msim)]
250 sim_state: SimState,
251
252 _state_snapshot_uploader_handle: Option<broadcast::Sender<()>>,
253 shutdown_channel_tx: broadcast::Sender<Option<RunWithRange>>,
255
256 grpc_server_handle: Mutex<Option<GrpcServerHandle>>,
258
259 auth_agg: Arc<ArcSwap<AuthorityAggregator<NetworkAuthorityClient>>>,
265
266 serving_rt_handle: tokio::runtime::Handle,
269}
270
271impl fmt::Debug for IotaNode {
272 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
273 f.debug_struct("IotaNode")
274 .field("name", &self.state.name.concise())
275 .finish()
276 }
277}
278
279impl IotaNode {
280 pub async fn start(
291 config: NodeConfig,
292 registry_service: RegistryService,
293 ) -> Result<Arc<IotaNode>> {
294 Self::start_async(
295 config,
296 registry_service,
297 ServerVersion::new("iota-node", "unknown"),
298 tokio::runtime::Handle::current(),
299 )
300 .await
301 }
302
303 fn start_overload_notifier(
308 config: &NodeConfig,
309 state: Arc<AuthorityState>,
310 epoch_store: Arc<AuthorityPerEpochStore>,
311 consensus_adapter: Arc<ConsensusAdapter>,
312 ) -> Option<JoinHandle<()>> {
313 if !epoch_store.protocol_config().enable_pcool_flow() {
314 return None;
315 }
316
317 let poll_interval = config.authority_overload_config.overload_monitor_interval;
318 let authority_name = state.name;
319
320 Some(spawn_monitored_task!(async move {
321 let mut last_notified_percentage: u32 = epoch_store
323 .load_overload_notification(&authority_name)
324 .unwrap_or(0) as u32;
325 loop {
326 tokio::time::sleep(poll_interval).await;
327 let current = state
328 .overload_info
329 .local_load_shedding_percentage
330 .load(std::sync::atomic::Ordering::Relaxed);
331 if current != last_notified_percentage {
332 last_notified_percentage = current;
333 let transaction = ConsensusTransaction::new_overload_notification_v1(
334 authority_name,
335 current as u8,
336 );
337 if let Err(e) = consensus_adapter.submit(transaction, None, &epoch_store) {
338 tracing::warn!(
339 "Failed to submit overload notification to consensus: {:?}",
340 e
341 );
342 }
343 }
344 }
345 }))
346 }
347
348 pub async fn start_async(
349 mut config: NodeConfig,
350 registry_service: RegistryService,
351 server_version: ServerVersion,
352 serving_rt_handle: tokio::runtime::Handle,
353 ) -> Result<Arc<IotaNode>> {
354 config.validate()?;
355 NodeConfigMetrics::new(®istry_service.default_registry()).record_metrics(&config);
356 if config.supported_protocol_versions.is_none() {
357 info!(
358 "populating config.supported_protocol_versions with default {:?}",
359 SupportedProtocolVersions::SYSTEM_DEFAULT
360 );
361 config.supported_protocol_versions = Some(SupportedProtocolVersions::SYSTEM_DEFAULT);
362 }
363
364 let run_with_range = config.run_with_range;
365 let is_validator = config.is_validator();
366 let is_full_node = !is_validator;
367 let prometheus_registry = registry_service.default_registry();
368
369 info!(node =? config.authority_public_key(),
370 "Initializing iota-node listening on {}", config.network_address
371 );
372
373 let genesis = config.genesis()?.clone();
374
375 let chain_identifier = ChainIdentifier::from(*genesis.checkpoint().digest());
376 info!("IOTA chain identifier: {chain_identifier}");
377
378 let db_corrupted_path = &config.db_path().join("status");
380 if let Err(err) = check_and_mark_db_corruption(db_corrupted_path) {
381 panic!("Failed to check database corruption: {err}");
382 }
383
384 DBMetrics::init(&prometheus_registry);
386
387 iota_metrics::init_metrics(&prometheus_registry);
389 #[cfg(not(msim))]
392 iota_metrics::thread_stall_monitor::start_thread_stall_monitor();
393
394 #[cfg(not(msim))]
398 {
399 let runtime_monitor_metrics =
400 iota_metrics::runtime_metrics::RuntimeMonitorMetrics::new(&prometheus_registry);
401 iota_metrics::runtime_metrics::start_runtime_monitor(
402 "iota_node",
403 &tokio::runtime::Handle::current(),
404 runtime_monitor_metrics.clone(),
405 );
406 iota_metrics::runtime_metrics::start_runtime_monitor(
407 "serving",
408 &serving_rt_handle,
409 runtime_monitor_metrics,
410 );
411 }
412
413 prometheus_registry
415 .register(iota_metrics::uptime_metric(
416 if is_validator {
417 "validator"
418 } else {
419 "fullnode"
420 },
421 server_version.version,
422 &chain_identifier.to_string(),
423 ))
424 .expect("Failed registering uptime metric");
425
426 let migration_tx_data = if genesis.contains_migrations() {
429 Some(config.load_migration_tx_data()?)
432 } else {
433 None
434 };
435
436 let secret = Arc::pin(config.authority_key_pair().copy());
437 let genesis_committee = genesis.committee()?;
438 let committee_store = Arc::new(CommitteeStore::new(
439 config.db_path().join("epochs"),
440 &genesis_committee,
441 None,
442 ));
443
444 let enable_write_stall = config.enable_db_write_stall.unwrap_or(is_validator);
446 let perpetual_tables_options = AuthorityPerpetualTablesOptions { enable_write_stall };
447 let perpetual_tables = Arc::new(AuthorityPerpetualTables::open(
448 &config.db_path().join("store"),
449 Some(perpetual_tables_options),
450 ));
451 let is_genesis = perpetual_tables
452 .database_is_empty()
453 .expect("Database read should not fail at init.");
454 let checkpoint_store = CheckpointStore::new_with_contents_cache(
455 &config.db_path().join("checkpoints"),
456 FullCheckpointContentsCache::new(
457 config
458 .full_checkpoint_contents_cache_size_mb
459 .saturating_mul(1024 * 1024),
460 FullCheckpointContentsCacheMetrics::new(&prometheus_registry),
461 ),
462 );
463 let backpressure_manager =
464 BackpressureManager::new_from_checkpoint_store(&checkpoint_store);
465
466 let perpetual_tables_for_progress = perpetual_tables.clone();
467 let store = AuthorityStore::open(
468 perpetual_tables,
469 &genesis,
470 &config,
471 &prometheus_registry,
472 migration_tx_data.as_ref(),
473 )
474 .await?;
475
476 let cur_epoch = store.get_recovery_epoch_at_restart()?;
477 let committee = committee_store
478 .get_committee(&cur_epoch)?
479 .expect("Committee of the current epoch must exist");
480 let epoch_start_configuration = store
481 .get_epoch_start_configuration()?
482 .expect("EpochStartConfiguration of the current epoch must exist");
483 let cache_metrics = Arc::new(ResolverMetrics::new(&prometheus_registry));
484 let signature_verifier_metrics = SignatureVerifierMetrics::new(&prometheus_registry);
485
486 let cache_traits = build_execution_cache(
487 &config.execution_cache_config,
488 &prometheus_registry,
489 &store,
490 backpressure_manager.clone(),
491 );
492
493 let auth_agg = {
494 let safe_client_metrics_base = SafeClientMetricsBase::new(&prometheus_registry);
495 let auth_agg_metrics = Arc::new(AuthAggMetrics::new(&prometheus_registry));
496 Arc::new(ArcSwap::new(Arc::new(
497 AuthorityAggregator::new_from_epoch_start_state(
498 epoch_start_configuration.epoch_start_state(),
499 &committee_store,
500 safe_client_metrics_base,
501 auth_agg_metrics,
502 ),
503 )))
504 };
505
506 let chain = match config.chain_override_for_testing {
507 Some(chain) => chain,
508 None => chain_identifier.chain(),
509 };
510
511 let epoch_options = default_db_options().optimize_db_for_write_throughput(4);
512 let epoch_store = AuthorityPerEpochStore::new(
513 config.authority_public_key(),
514 committee.clone(),
515 &config.db_path().join("store"),
516 Some(epoch_options.options),
517 EpochMetrics::new(®istry_service.default_registry()),
518 epoch_start_configuration,
519 cache_traits.backing_package_store.clone(),
520 cache_metrics,
521 signature_verifier_metrics,
522 &config.expensive_safety_check_config,
523 (chain_identifier, chain),
524 checkpoint_store
525 .get_highest_executed_checkpoint_seq_number()
526 .expect("checkpoint store read cannot fail")
527 .unwrap_or(0),
528 )?;
529
530 info!("created epoch store");
531
532 replay_log!(
533 "Beginning replay run. Epoch: {:?}, Protocol config: {:?}",
534 epoch_store.epoch(),
535 epoch_store.protocol_config()
536 );
537
538 if is_genesis {
540 info!("checking IOTA conservation at genesis");
541 cache_traits
546 .reconfig_api
547 .try_expensive_check_iota_conservation(&epoch_store, None)
548 .expect("IOTA conservation check cannot fail at genesis");
549 }
550
551 let effective_buffer_stake = epoch_store.get_effective_buffer_stake_bps();
552 let default_buffer_stake = epoch_store
553 .protocol_config()
554 .buffer_stake_for_protocol_upgrade_bps();
555 if effective_buffer_stake != default_buffer_stake {
556 warn!(
557 ?effective_buffer_stake,
558 ?default_buffer_stake,
559 "buffer_stake_for_protocol_upgrade_bps is currently overridden"
560 );
561 }
562
563 checkpoint_store.insert_genesis_checkpoint(
564 genesis.checkpoint(),
565 genesis.checkpoint_contents().clone(),
566 &epoch_store,
567 );
568
569 unmark_db_corruption(db_corrupted_path)?;
571
572 info!("creating state sync store");
573 let state_sync_store = RocksDbStore::new(
574 cache_traits.clone(),
575 committee_store.clone(),
576 checkpoint_store.clone(),
577 );
578
579 let index_store = if is_full_node && config.enable_index_processing {
580 info!("creating index store");
581 Some(Arc::new(IndexStore::new(
582 config.db_path().join("indexes"),
583 &prometheus_registry,
584 epoch_store
585 .protocol_config()
586 .max_move_identifier_len_as_option(),
587 )))
588 } else {
589 None
590 };
591
592 let grpc_indexes_store = if is_full_node && config.enable_grpc_api {
593 Some(Arc::new(
594 GrpcIndexesStore::new(
595 config.db_path().join(GRPC_INDEXES_DIR),
596 Arc::clone(&store),
597 &checkpoint_store,
598 )
599 .await,
600 ))
601 } else {
602 None
603 };
604
605 checkpoint_store
607 .seed_epoch_info(&store, chain_identifier)
608 .expect("failed to seed the epoch_info chain");
609
610 info!("creating archive reader");
611 let (trusted_peer_change_tx, trusted_peer_change_rx) = watch::channel(Default::default());
616 let (randomness_tx, randomness_rx) = mpsc::channel(
617 config
618 .p2p_config
619 .randomness
620 .clone()
621 .unwrap_or_default()
622 .mailbox_capacity(),
623 );
624 let (p2p_network, discovery_handle, state_sync_handle, randomness_handle) =
625 Self::create_p2p_network(
626 &config,
627 state_sync_store.clone(),
628 chain_identifier,
629 trusted_peer_change_rx,
630 randomness_tx,
631 &prometheus_registry,
632 )?;
633
634 send_trusted_peer_change(
637 &config,
638 &trusted_peer_change_tx,
639 epoch_store.epoch_start_state(),
640 );
641
642 info!("start snapshot upload");
643 let state_snapshot_handle =
645 Self::start_state_snapshot(&config, &prometheus_registry, checkpoint_store.clone())?;
646
647 let checkpoint_progress_tracker = Arc::new(CheckpointProgressTracker::new());
648
649 let mut genesis_objects = genesis.objects().to_vec();
650 if let Some(migration_tx_data) = migration_tx_data.as_ref() {
651 genesis_objects.extend(migration_tx_data.get_objects());
652 }
653
654 let authority_name = config.authority_public_key();
655 let validator_tx_finalizer =
656 config
657 .enable_validator_tx_finalizer
658 .then_some(Arc::new(ValidatorTxFinalizer::new(
659 auth_agg.clone(),
660 authority_name,
661 &prometheus_registry,
662 )));
663
664 info!("create authority state");
665 let state = AuthorityState::new(
666 authority_name,
667 secret,
668 config.supported_protocol_versions.unwrap(),
669 store.clone(),
670 cache_traits.clone(),
671 epoch_store.clone(),
672 committee_store.clone(),
673 index_store.clone(),
674 grpc_indexes_store,
675 checkpoint_store.clone(),
676 &prometheus_registry,
677 &genesis_objects,
678 config.clone(),
679 validator_tx_finalizer,
680 chain_identifier,
681 Some(checkpoint_progress_tracker.clone()),
682 config.policy_config.clone(),
683 config.firewall_config.clone(),
684 )
685 .await;
686
687 if epoch_store.epoch() == 0 {
689 let genesis_tx = &genesis.transaction();
690 let span = error_span!("genesis_txn", tx_digest = ?genesis_tx.digest());
691 Self::execute_transaction_immediately_at_zero_epoch(
693 &state,
694 &epoch_store,
695 genesis_tx,
696 span,
697 )
698 .await;
699
700 if let Some(migration_tx_data) = migration_tx_data {
702 for (tx_digest, (tx, _, _)) in migration_tx_data.txs_data() {
703 let span = error_span!("migration_txn", tx_digest = ?tx_digest);
704 Self::execute_transaction_immediately_at_zero_epoch(
705 &state,
706 &epoch_store,
707 tx,
708 span,
709 )
710 .await;
711 }
712 }
713 }
714
715 RandomnessRoundReceiver::spawn(state.clone(), randomness_rx);
718
719 if config
720 .expensive_safety_check_config
721 .enable_secondary_index_checks()
722 {
723 if let Some(indexes) = state.indexes.clone() {
724 iota_core::verify_indexes::verify_indexes(
725 state.get_global_state_hash_store().as_ref(),
726 indexes,
727 )
728 .expect("secondary indexes are inconsistent");
729 }
730 }
731
732 let (end_of_epoch_channel, end_of_epoch_receiver) =
733 broadcast::channel(config.end_of_epoch_broadcast_channel_capacity);
734
735 let transaction_orchestrator = if is_full_node && run_with_range.is_none() {
736 Some(Arc::new(TransactionOrchestrator::new_with_auth_aggregator(
737 auth_agg.load_full(),
738 state.clone(),
739 end_of_epoch_receiver,
740 &config.db_path(),
741 &prometheus_registry,
742 Some(&config),
743 )))
744 } else {
745 None
746 };
747
748 let http_server = serving_rt_handle
752 .spawn({
753 let state = state.clone();
754 let transaction_orchestrator = transaction_orchestrator.clone();
755 let config = config.clone();
756 let prometheus_registry = prometheus_registry.clone();
757 async move {
758 build_http_server(
759 state,
760 &transaction_orchestrator,
761 &config,
762 &prometheus_registry,
763 )
764 .await
765 }
766 })
767 .await
768 .expect("Failed to join JSON-RPC server startup task")?;
769
770 let global_state_hasher = Arc::new(GlobalStateHasher::new(
771 cache_traits.global_state_hash_store.clone(),
772 GlobalStateHashMetrics::new(&prometheus_registry),
773 ));
774
775 let authority_names_to_peer_ids = epoch_store
776 .epoch_start_state()
777 .get_authority_names_to_peer_ids();
778
779 let network_connection_metrics =
780 NetworkConnectionMetrics::new("iota", ®istry_service.default_registry());
781
782 let authority_names_to_peer_ids = ArcSwap::from_pointee(authority_names_to_peer_ids);
783
784 let (_connection_monitor_handle, connection_statuses) = ConnectionMonitor::spawn(
785 p2p_network.downgrade(),
786 network_connection_metrics,
787 HashMap::new(),
788 None,
789 );
790
791 let connection_monitor_status = ConnectionMonitorStatus {
792 connection_statuses,
793 authority_names_to_peer_ids,
794 };
795
796 let connection_monitor_status = Arc::new(connection_monitor_status);
797 let iota_node_metrics =
798 Arc::new(IotaNodeMetrics::new(®istry_service.default_registry()));
799
800 iota_node_metrics
801 .binary_max_protocol_version
802 .set(ProtocolVersion::MAX.as_u64() as i64);
803 iota_node_metrics
804 .configured_max_protocol_version
805 .set(config.supported_protocol_versions.unwrap().max.as_u64() as i64);
806
807 let executor: Option<Arc<dyn iota_types::transaction_executor::TransactionExecutor>> =
811 transaction_orchestrator
812 .clone()
813 .map(|o| o as Arc<dyn iota_types::transaction_executor::TransactionExecutor>);
814
815 let grpc_server_handle = serving_rt_handle
818 .spawn({
819 let config = config.clone();
820 let state = state.clone();
821 let state_sync_store = state_sync_store.clone();
822 let prometheus_registry = prometheus_registry.clone();
823 async move {
824 build_grpc_server(
825 &config,
826 state,
827 state_sync_store,
828 executor,
829 &prometheus_registry,
830 server_version,
831 )
832 .await
833 }
834 })
835 .await
836 .expect("Failed to join gRPC server startup task")?;
837
838 let validator_components = if state.is_committee_validator(&epoch_store) {
839 let (components, _) = futures::join!(
840 Self::construct_validator_components(
841 config.clone(),
842 state.clone(),
843 committee,
844 epoch_store.clone(),
845 checkpoint_store.clone(),
846 state_sync_handle.clone(),
847 randomness_handle.clone(),
848 Arc::downgrade(&global_state_hasher),
849 backpressure_manager.clone(),
850 connection_monitor_status.clone(),
851 ®istry_service,
852 serving_rt_handle.clone(),
853 ),
854 Self::reexecute_pending_consensus_certs(&epoch_store, &state,)
855 );
856 let mut components = components?;
857
858 components.consensus_adapter.submit_recovered(&epoch_store);
859
860 components.validator_server_handle = components.validator_server_handle.start().await;
862
863 Some(components)
864 } else {
865 None
866 };
867
868 let (shutdown_channel, _) = broadcast::channel::<Option<RunWithRange>>(1);
870
871 let node = Self {
872 config,
873 validator_components: Mutex::new(validator_components),
874 _http_server: http_server,
875 state,
876 transaction_orchestrator,
877 registry_service,
878 metrics: iota_node_metrics,
879
880 _discovery: discovery_handle,
881 state_sync_handle,
882 randomness_handle,
883 checkpoint_store,
884 state_sync_store,
885 global_state_hasher: Mutex::new(Some(global_state_hasher)),
886 end_of_epoch_channel,
887 connection_monitor_status,
888 trusted_peer_change_tx,
889 backpressure_manager,
890 checkpoint_progress_tracker: checkpoint_progress_tracker.clone(),
891
892 #[cfg(msim)]
893 sim_state: Default::default(),
894
895 _state_snapshot_uploader_handle: state_snapshot_handle,
896 shutdown_channel_tx: shutdown_channel,
897
898 grpc_server_handle: Mutex::new(grpc_server_handle),
899
900 auth_agg,
901
902 serving_rt_handle,
903 };
904
905 info!("IotaNode started!");
906 let node = Arc::new(node);
907 let node_copy = node.clone();
908 spawn_monitored_task!(async move {
909 let result = Self::monitor_reconfiguration(node_copy, epoch_store).await;
910 if let Err(error) = result {
911 warn!("Reconfiguration finished with error {:?}", error);
912 }
913 });
914
915 node.checkpoint_progress_tracker
916 .spawn_logging_task(node.checkpoint_store.clone(), perpetual_tables_for_progress);
917
918 Ok(node)
919 }
920
921 pub fn subscribe_to_epoch_change(&self) -> broadcast::Receiver<IotaSystemState> {
922 self.end_of_epoch_channel.subscribe()
923 }
924
925 pub fn subscribe_to_shutdown_channel(&self) -> broadcast::Receiver<Option<RunWithRange>> {
926 self.shutdown_channel_tx.subscribe()
927 }
928
929 pub fn current_epoch_for_testing(&self) -> EpochId {
930 self.state.current_epoch_for_testing()
931 }
932
933 pub async fn close_epoch(&self, epoch_store: &Arc<AuthorityPerEpochStore>) -> IotaResult {
935 info!("close_epoch (current epoch = {})", epoch_store.epoch());
936 self.validator_components
937 .lock()
938 .await
939 .as_ref()
940 .ok_or_else(|| IotaError::from("Node is not a validator"))?
941 .consensus_adapter
942 .close_epoch(epoch_store);
943 Ok(())
944 }
945
946 pub fn clear_override_protocol_upgrade_buffer_stake(&self, epoch: EpochId) -> IotaResult {
947 self.state
948 .clear_override_protocol_upgrade_buffer_stake(epoch)
949 }
950
951 pub fn set_override_protocol_upgrade_buffer_stake(
952 &self,
953 epoch: EpochId,
954 buffer_stake_bps: u64,
955 ) -> IotaResult {
956 self.state
957 .set_override_protocol_upgrade_buffer_stake(epoch, buffer_stake_bps)
958 }
959
960 pub async fn close_epoch_for_testing(&self) -> IotaResult {
963 let epoch_store = self.state.epoch_store_for_testing();
964 self.close_epoch(&epoch_store).await
965 }
966
967 fn start_state_snapshot(
971 config: &NodeConfig,
972 prometheus_registry: &Registry,
973 checkpoint_store: Arc<CheckpointStore>,
974 ) -> Result<Option<tokio::sync::broadcast::Sender<()>>> {
975 if let Some(remote_store_config) = &config.state_snapshot_write_config.object_store_config {
976 debug_assert!(
977 !config.is_validator(),
978 "`NodeConfig::validate` rejects snapshot upload on a validator"
979 );
980 let snapshot_uploader = StateSnapshotUploader::new(
981 &config.db_checkpoint_path(),
982 &config.snapshot_path(),
983 remote_store_config.clone(),
984 config.state_snapshot_write_config.concurrency,
985 60,
986 prometheus_registry,
987 checkpoint_store,
988 )?;
989 Ok(Some(snapshot_uploader.start()))
990 } else {
991 Ok(None)
992 }
993 }
994
995 fn create_p2p_network(
996 config: &NodeConfig,
997 state_sync_store: RocksDbStore,
998 chain_identifier: ChainIdentifier,
999 trusted_peer_change_rx: watch::Receiver<TrustedPeerChangeEvent>,
1000 randomness_tx: mpsc::Sender<(EpochId, RandomnessRound, Vec<u8>)>,
1001 prometheus_registry: &Registry,
1002 ) -> Result<(
1003 Network,
1004 discovery::Handle,
1005 state_sync::Handle,
1006 randomness::Handle,
1007 )> {
1008 let (state_sync, state_sync_server) = state_sync::Builder::new()
1009 .config(config.p2p_config.state_sync.clone().unwrap_or_default())
1010 .store(state_sync_store)
1011 .checkpoint_archive_config(config.checkpoint_archive_config().cloned())
1012 .with_metrics(prometheus_registry)
1013 .build();
1014
1015 let (discovery, discovery_server) = discovery::Builder::new(trusted_peer_change_rx)
1016 .config(config.p2p_config.clone())
1017 .build();
1018
1019 let (randomness, randomness_router) =
1020 randomness::Builder::new(config.authority_public_key(), randomness_tx)
1021 .config(config.p2p_config.randomness.clone().unwrap_or_default())
1022 .with_metrics(prometheus_registry)
1023 .build();
1024
1025 let p2p_network = {
1026 let routes = anemo::Router::new()
1027 .add_rpc_service(discovery_server)
1028 .add_rpc_service(state_sync_server);
1029 let routes = routes.merge(randomness_router);
1030
1031 let inbound_network_metrics =
1032 NetworkMetrics::new("iota", "inbound", prometheus_registry);
1033 let outbound_network_metrics =
1034 NetworkMetrics::new("iota", "outbound", prometheus_registry);
1035
1036 let service = ServiceBuilder::new()
1037 .layer(
1038 TraceLayer::new_for_server_errors()
1039 .make_span_with(DefaultMakeSpan::new().level(tracing::Level::INFO))
1040 .on_failure(DefaultOnFailure::new().level(tracing::Level::WARN)),
1041 )
1042 .layer(CallbackLayer::new(MetricsMakeCallbackHandler::new(
1043 Arc::new(inbound_network_metrics),
1044 config.p2p_config.excessive_message_size(),
1045 )))
1046 .service(routes);
1047
1048 let outbound_layer = ServiceBuilder::new()
1049 .layer(
1050 TraceLayer::new_for_client_and_server_errors()
1051 .make_span_with(DefaultMakeSpan::new().level(tracing::Level::INFO))
1052 .on_failure(DefaultOnFailure::new().level(tracing::Level::DEBUG)),
1053 )
1054 .layer(CallbackLayer::new(MetricsMakeCallbackHandler::new(
1055 Arc::new(outbound_network_metrics),
1056 config.p2p_config.excessive_message_size(),
1057 )))
1058 .into_inner();
1059
1060 let mut anemo_config = config.p2p_config.anemo_config.clone().unwrap_or_default();
1061 anemo_config.max_request_frame_size = Some(1 << 20);
1064 anemo_config.max_response_frame_size = Some(128 << 20);
1067
1068 let mut quic_config = anemo_config.quic.unwrap_or_default();
1071 if quic_config.socket_send_buffer_size.is_none() {
1072 quic_config.socket_send_buffer_size = Some(20 << 20);
1073 }
1074 if quic_config.socket_receive_buffer_size.is_none() {
1075 quic_config.socket_receive_buffer_size = Some(20 << 20);
1076 }
1077 quic_config.allow_failed_socket_buffer_size_setting = true;
1078
1079 if quic_config.max_concurrent_bidi_streams.is_none() {
1082 quic_config.max_concurrent_bidi_streams = Some(500);
1083 }
1084 if quic_config.max_concurrent_uni_streams.is_none() {
1085 quic_config.max_concurrent_uni_streams = Some(500);
1086 }
1087 if quic_config.stream_receive_window.is_none() {
1088 quic_config.stream_receive_window = Some(100 << 20);
1089 }
1090 if quic_config.receive_window.is_none() {
1091 quic_config.receive_window = Some(200 << 20);
1092 }
1093 if quic_config.send_window.is_none() {
1094 quic_config.send_window = Some(200 << 20);
1095 }
1096 if quic_config.crypto_buffer_size.is_none() {
1097 quic_config.crypto_buffer_size = Some(1 << 20);
1098 }
1099 if quic_config.max_idle_timeout_ms.is_none() {
1100 quic_config.max_idle_timeout_ms = Some(10_000);
1101 }
1102 if quic_config.keep_alive_interval_ms.is_none() {
1103 quic_config.keep_alive_interval_ms = Some(5_000);
1104 }
1105 anemo_config.quic = Some(quic_config);
1106
1107 let server_name = format!("iota-{chain_identifier}");
1108 let network = Network::bind(config.p2p_config.listen_address)
1109 .server_name(&server_name)
1110 .private_key(config.network_key_pair().copy().private().0.to_bytes())
1111 .config(anemo_config)
1112 .outbound_request_layer(outbound_layer)
1113 .start(service)?;
1114 info!(
1115 server_name = server_name,
1116 "P2p network started on {}",
1117 network.local_addr()
1118 );
1119
1120 network
1121 };
1122
1123 let discovery_handle =
1124 discovery.start(p2p_network.clone(), config.network_key_pair().copy());
1125 let state_sync_handle = state_sync.start(p2p_network.clone());
1126 let randomness_handle = randomness.start(p2p_network.clone());
1127
1128 Ok((
1129 p2p_network,
1130 discovery_handle,
1131 state_sync_handle,
1132 randomness_handle,
1133 ))
1134 }
1135
1136 async fn construct_validator_components(
1139 config: NodeConfig,
1140 state: Arc<AuthorityState>,
1141 committee: Arc<Committee>,
1142 epoch_store: Arc<AuthorityPerEpochStore>,
1143 checkpoint_store: Arc<CheckpointStore>,
1144 state_sync_handle: state_sync::Handle,
1145 randomness_handle: randomness::Handle,
1146 global_state_hasher: Weak<GlobalStateHasher>,
1147 backpressure_manager: Arc<BackpressureManager>,
1148 connection_monitor_status: Arc<ConnectionMonitorStatus>,
1149 registry_service: &RegistryService,
1150 serving_rt_handle: tokio::runtime::Handle,
1151 ) -> Result<ValidatorComponents> {
1152 let mut config_clone = config.clone();
1153 let consensus_config = config_clone
1154 .consensus_config
1155 .as_mut()
1156 .ok_or_else(|| anyhow!("Validator is missing consensus config"))?;
1157 let validator_registry = registry_service.new_registry_custom(None, None)?;
1158 let validator_registry_id = registry_service.add(validator_registry.clone());
1159
1160 let client = Arc::new(UpdatableConsensusClient::new());
1161 let consensus_adapter = Arc::new(Self::construct_consensus_adapter(
1162 &committee,
1163 consensus_config,
1164 state.name,
1165 connection_monitor_status.clone(),
1166 &validator_registry,
1167 client.clone(),
1168 checkpoint_store.clone(),
1169 ));
1170 let consensus_manager = Arc::new(ConsensusManager::new(
1171 &config,
1172 consensus_config,
1173 registry_service,
1174 &validator_registry,
1175 client,
1176 ));
1177
1178 let consensus_store_pruner = ConsensusStorePruner::new(
1181 consensus_manager.get_storage_base_path(),
1182 consensus_config.db_retention_epochs(),
1183 consensus_config.db_pruner_period(),
1184 &validator_registry,
1185 );
1186
1187 let soft_locks = Arc::new(if config.enable_soft_locking {
1188 PreConsensusSoftLocks::new()
1189 } else {
1190 info!("pre-consensus soft-locking disabled via node config");
1191 PreConsensusSoftLocks::disabled()
1192 });
1193
1194 let checkpoint_metrics = CheckpointMetrics::new(&validator_registry);
1195 let iota_tx_validator_metrics = IotaTxValidatorMetrics::new(&validator_registry);
1196 let validator_service_metrics = Arc::new(ValidatorServiceMetrics::new(&validator_registry));
1197
1198 let soft_lock_sweep_handle = PreConsensusSoftLocks::spawn_sweep(
1209 Arc::downgrade(&soft_locks),
1210 validator_service_metrics.clone(),
1211 );
1212
1213 let validator_server_handle = Self::start_grpc_validator_service(
1214 &config,
1215 state.clone(),
1216 consensus_adapter.clone(),
1217 &validator_registry,
1218 soft_locks.clone(),
1219 validator_service_metrics.clone(),
1220 serving_rt_handle,
1221 )
1222 .await?;
1223
1224 let validator_overload_monitor_handle = if config
1227 .authority_overload_config
1228 .max_load_shedding_percentage
1229 > 0
1230 {
1231 let authority_state = Arc::downgrade(&state);
1232 let overload_config = config.authority_overload_config.clone();
1233 fail_point!("starting_overload_monitor");
1234 Some(spawn_monitored_task!(overload_monitor(
1235 authority_state,
1236 overload_config,
1237 )))
1238 } else {
1239 None
1240 };
1241
1242 let consensus_queue_overload_monitor_handle =
1249 if epoch_store.protocol_config().enable_pcool_flow() {
1250 let consensus_queue_monitor_authority_state = Arc::downgrade(&state);
1251 let consensus_queue_monitor_consensus_adapter = Arc::downgrade(&consensus_adapter);
1252 let consensus_queue_monitor_interval =
1253 config.authority_overload_config.overload_monitor_interval;
1254 Some(spawn_monitored_task!(consensus_queue_overload_monitor(
1255 consensus_queue_monitor_authority_state,
1256 consensus_queue_monitor_consensus_adapter,
1257 consensus_queue_monitor_interval,
1258 )))
1259 } else {
1260 None
1261 };
1262
1263 Self::start_epoch_specific_validator_components(
1264 &config,
1265 state.clone(),
1266 consensus_adapter,
1267 checkpoint_store,
1268 epoch_store,
1269 state_sync_handle,
1270 randomness_handle,
1271 consensus_manager,
1272 consensus_store_pruner,
1273 global_state_hasher,
1274 backpressure_manager,
1275 soft_locks,
1276 validator_server_handle,
1277 validator_overload_monitor_handle,
1278 consensus_queue_overload_monitor_handle,
1279 soft_lock_sweep_handle,
1280 checkpoint_metrics,
1281 iota_tx_validator_metrics,
1282 validator_registry_id,
1283 )
1284 .await
1285 }
1286
1287 async fn start_epoch_specific_validator_components(
1290 config: &NodeConfig,
1291 state: Arc<AuthorityState>,
1292 consensus_adapter: Arc<ConsensusAdapter>,
1293 checkpoint_store: Arc<CheckpointStore>,
1294 epoch_store: Arc<AuthorityPerEpochStore>,
1295 state_sync_handle: state_sync::Handle,
1296 randomness_handle: randomness::Handle,
1297 consensus_manager: Arc<ConsensusManager>,
1298 consensus_store_pruner: ConsensusStorePruner,
1299 global_state_hasher: Weak<GlobalStateHasher>,
1300 backpressure_manager: Arc<BackpressureManager>,
1301 soft_locks: Arc<PreConsensusSoftLocks>,
1302 validator_server_handle: SpawnOnce,
1303 validator_overload_monitor_handle: Option<JoinHandle<()>>,
1304 consensus_queue_overload_monitor_handle: Option<JoinHandle<()>>,
1305 soft_lock_sweep_handle: JoinHandle<()>,
1306 checkpoint_metrics: Arc<CheckpointMetrics>,
1307 iota_tx_validator_metrics: Arc<IotaTxValidatorMetrics>,
1308 validator_registry_id: RegistryID,
1309 ) -> Result<ValidatorComponents> {
1310 let checkpoint_service = Self::build_checkpoint_service(
1311 config,
1312 consensus_adapter.clone(),
1313 checkpoint_store.clone(),
1314 epoch_store.clone(),
1315 state.clone(),
1316 state_sync_handle,
1317 global_state_hasher,
1318 checkpoint_metrics.clone(),
1319 );
1320
1321 let low_scoring_authorities = Arc::new(ArcSwap::new(Arc::new(HashMap::new())));
1326
1327 consensus_adapter.swap_low_scoring_authorities(low_scoring_authorities.clone());
1328
1329 soft_locks.clear();
1334 epoch_store.set_soft_locks(soft_locks.clone());
1335
1336 let randomness_manager = match RandomnessManager::try_new(
1341 Arc::downgrade(&epoch_store),
1342 Box::new(consensus_adapter.clone()),
1343 randomness_handle,
1344 config.authority_key_pair(),
1345 )
1346 .await
1347 {
1348 Ok(randomness_manager) => randomness_manager,
1349 Err(err) => {
1350 fatal!(
1351 "validator cannot start epoch {} without a randomness manager: {err}",
1352 epoch_store.epoch()
1353 );
1354 }
1355 };
1356 epoch_store
1357 .set_randomness_manager(randomness_manager)
1358 .await?;
1359
1360 let consensus_handler_initializer = ConsensusHandlerInitializer::new(
1361 state.clone(),
1362 checkpoint_service.clone(),
1363 epoch_store.clone(),
1364 low_scoring_authorities,
1365 backpressure_manager,
1366 );
1367
1368 info!("Starting consensus manager asynchronously");
1369
1370 tokio::spawn({
1372 let config = config.clone();
1373 let epoch_store = epoch_store.clone();
1374 let iota_tx_validator = IotaTxValidator::new(
1375 epoch_store.clone(),
1376 checkpoint_service.clone(),
1377 iota_tx_validator_metrics.clone(),
1378 );
1379 let consensus_manager = consensus_manager.clone();
1380 async move {
1381 consensus_manager
1382 .start(
1383 &config,
1384 epoch_store,
1385 consensus_handler_initializer,
1386 iota_tx_validator,
1387 )
1388 .await;
1389 }
1390 });
1391 let replay_waiter = consensus_manager.replay_waiter();
1392
1393 info!("Spawning checkpoint service");
1394 let replay_waiter = if std::env::var("DISABLE_REPLAY_WAITER").is_ok() {
1395 None
1396 } else {
1397 Some(replay_waiter)
1398 };
1399 let checkpoint_service_tasks = checkpoint_service.spawn(replay_waiter).await;
1400
1401 let overload_notifier_handle = Self::start_overload_notifier(
1402 config,
1403 state.clone(),
1404 epoch_store.clone(),
1405 consensus_adapter.clone(),
1406 );
1407
1408 Ok(ValidatorComponents {
1409 validator_server_handle,
1410 validator_overload_monitor_handle,
1411 consensus_queue_overload_monitor_handle,
1412 soft_lock_sweep_handle,
1413 overload_notifier_handle,
1414 consensus_manager,
1415 consensus_store_pruner,
1416 consensus_adapter,
1417 soft_locks,
1418 checkpoint_service_tasks,
1419 checkpoint_metrics,
1420 iota_tx_validator_metrics,
1421 validator_registry_id,
1422 })
1423 }
1424
1425 fn build_checkpoint_service(
1432 config: &NodeConfig,
1433 consensus_adapter: Arc<ConsensusAdapter>,
1434 checkpoint_store: Arc<CheckpointStore>,
1435 epoch_store: Arc<AuthorityPerEpochStore>,
1436 state: Arc<AuthorityState>,
1437 state_sync_handle: state_sync::Handle,
1438 global_state_hasher: Weak<GlobalStateHasher>,
1439 checkpoint_metrics: Arc<CheckpointMetrics>,
1440 ) -> Arc<CheckpointService> {
1441 let epoch_start_timestamp_ms = epoch_store.epoch_start_state().epoch_start_timestamp_ms();
1442 let epoch_duration_ms = epoch_store.epoch_start_state().epoch_duration_ms();
1443
1444 debug!(
1445 "Starting checkpoint service with epoch start timestamp {}
1446 and epoch duration {}",
1447 epoch_start_timestamp_ms, epoch_duration_ms
1448 );
1449
1450 let checkpoint_output = Box::new(SubmitCheckpointToConsensus {
1451 sender: consensus_adapter,
1452 signer: state.secret.clone(),
1453 authority: config.authority_public_key(),
1454 next_reconfiguration_timestamp_ms: epoch_start_timestamp_ms
1455 .checked_add(epoch_duration_ms)
1456 .expect("Overflow calculating next_reconfiguration_timestamp_ms"),
1457 metrics: checkpoint_metrics.clone(),
1458 });
1459
1460 let certified_checkpoint_output = SendCheckpointToStateSync::new(state_sync_handle);
1461 let max_tx_per_checkpoint = max_tx_per_checkpoint(epoch_store.protocol_config());
1462 let max_checkpoint_size_bytes =
1463 epoch_store.protocol_config().max_checkpoint_size_bytes() as usize;
1464
1465 CheckpointService::build(
1466 state.clone(),
1467 checkpoint_store,
1468 epoch_store,
1469 state.get_transaction_cache_reader().clone(),
1470 global_state_hasher,
1471 checkpoint_output,
1472 Box::new(certified_checkpoint_output),
1473 checkpoint_metrics,
1474 max_tx_per_checkpoint,
1475 max_checkpoint_size_bytes,
1476 )
1477 }
1478
1479 fn construct_consensus_adapter(
1480 committee: &Committee,
1481 consensus_config: &ConsensusConfig,
1482 authority: AuthorityName,
1483 connection_monitor_status: Arc<ConnectionMonitorStatus>,
1484 prometheus_registry: &Registry,
1485 consensus_client: Arc<dyn ConsensusClient>,
1486 checkpoint_store: Arc<CheckpointStore>,
1487 ) -> ConsensusAdapter {
1488 let ca_metrics = ConsensusAdapterMetrics::new(prometheus_registry);
1489 ConsensusAdapter::new(
1493 consensus_client,
1494 checkpoint_store,
1495 authority,
1496 connection_monitor_status,
1497 consensus_config.max_pending_transactions(),
1498 consensus_config.max_pending_transactions() * 2 / committee.num_members(),
1499 consensus_config.max_submit_position,
1500 consensus_config.submit_delay_step_override(),
1501 ca_metrics,
1502 consensus_config.graduated_load_shedding_soft_limit_pct(),
1503 )
1504 }
1505
1506 async fn start_grpc_validator_service(
1507 config: &NodeConfig,
1508 state: Arc<AuthorityState>,
1509 consensus_adapter: Arc<ConsensusAdapter>,
1510 prometheus_registry: &Registry,
1511 soft_locks: Arc<PreConsensusSoftLocks>,
1512 validator_service_metrics: Arc<ValidatorServiceMetrics>,
1513 serving_rt_handle: tokio::runtime::Handle,
1514 ) -> Result<SpawnOnce> {
1515 let validator_service = ValidatorService::new(
1516 state,
1517 consensus_adapter,
1518 validator_service_metrics,
1519 config.policy_config.clone().map(|p| p.client_id_source),
1520 soft_locks,
1521 );
1522
1523 let concurrency_limit = config.grpc_concurrency_limit_per_core.saturating_mul(
1530 NonZeroUsize::new(iota_core::runtime::available_cpu_cores())
1531 .unwrap_or(NonZeroUsize::MIN),
1532 );
1533 let load_shed = config.grpc_load_shed.unwrap_or_default();
1534
1535 const VALIDATOR_GRPC_KEEPALIVE: Duration = Duration::from_secs(60);
1540 const VALIDATOR_GRPC_MAX_CONCURRENT_STREAMS: u32 = 1000;
1545
1546 let mut server_conf = iota_network_stack::config::Config::new();
1547 server_conf.http2_keepalive_interval = Some(VALIDATOR_GRPC_KEEPALIVE);
1548 server_conf.http2_keepalive_timeout = Some(VALIDATOR_GRPC_KEEPALIVE);
1549 server_conf.http2_max_concurrent_streams = Some(VALIDATOR_GRPC_MAX_CONCURRENT_STREAMS);
1550 let server_builder =
1551 ServerBuilder::from_config(&server_conf, GrpcMetrics::new(prometheus_registry))
1552 .add_service_with_concurrency_limit(
1553 ValidatorServer::new(validator_service.clone()),
1554 concurrency_limit,
1555 load_shed,
1556 )
1557 .add_service_with_concurrency_limit(
1558 ValidatorV2Server::new(validator_service.clone()),
1559 concurrency_limit,
1560 load_shed,
1561 )
1562 .add_service_with_concurrency_limit(
1563 ValidatorPeerServer::new(validator_service),
1564 concurrency_limit,
1565 load_shed,
1566 );
1567
1568 let tls_config = iota_tls::create_rustls_server_config(
1569 config.network_key_pair().copy().private(),
1570 IOTA_TLS_SERVER_NAME.to_string(),
1571 );
1572
1573 let network_address = config.network_address().clone();
1574
1575 let bind_future = async move {
1576 let server = server_builder
1577 .bind(&network_address, Some(tls_config))
1578 .await
1579 .map_err(|err| anyhow!("Failed to bind to {network_address}: {err}"))?;
1580
1581 let local_addr = server.local_addr();
1582 info!("Listening to traffic on {local_addr}");
1583
1584 Ok(server)
1585 };
1586
1587 Ok(SpawnOnce::new(bind_future, serving_rt_handle))
1588 }
1589
1590 async fn reexecute_pending_consensus_certs(
1608 epoch_store: &Arc<AuthorityPerEpochStore>,
1609 state: &Arc<AuthorityState>,
1610 ) {
1611 let mut pending_consensus_certificates = Vec::new();
1612 let mut additional_certs = Vec::new();
1613
1614 for tx in epoch_store.get_all_pending_consensus_transactions() {
1615 match tx.kind {
1616 ConsensusTransactionKind::CertifiedTransaction(tx)
1626 if !tx.contains_shared_object() =>
1627 {
1628 let tx = *tx;
1629 let tx = VerifiedExecutableTransaction::new_from_certificate(
1632 VerifiedCertificate::new_unchecked(tx),
1633 );
1634 if let Some(fx_digest) = epoch_store
1637 .get_signed_effects_digest(tx.digest())
1638 .expect("db error")
1639 {
1640 pending_consensus_certificates.push((
1641 Schedulable::Transaction(tx),
1642 ExecutionEnv::new().with_expected_effects_digest(fx_digest),
1643 ));
1644 } else {
1645 additional_certs.push((Schedulable::Transaction(tx), ExecutionEnv::new()));
1646 }
1647 }
1648 _ => (),
1649 }
1650 }
1651
1652 let digests = pending_consensus_certificates
1653 .iter()
1654 .map(|(tx, _)| *tx.key().unwrap_digest())
1657 .collect::<Vec<_>>();
1658
1659 info!(
1660 "reexecuting {} pending consensus certificates: {:?}",
1661 digests.len(),
1662 digests
1663 );
1664
1665 state
1666 .execution_scheduler()
1667 .enqueue(pending_consensus_certificates, epoch_store);
1668 state
1669 .execution_scheduler()
1670 .enqueue(additional_certs, epoch_store);
1671
1672 let timeout = if cfg!(msim) { 120 } else { 60 };
1678 if tokio::time::timeout(
1679 std::time::Duration::from_secs(timeout),
1680 state
1681 .get_transaction_cache_reader()
1682 .try_notify_read_executed_effects_digests(
1683 "IotaNode::notify_read_executed_effects_digests",
1684 &digests,
1685 ),
1686 )
1687 .await
1688 .is_err()
1689 {
1690 if let Ok(executed_effects_digests) = state
1692 .get_transaction_cache_reader()
1693 .try_multi_get_executed_effects_digests(&digests)
1694 {
1695 let pending_digests = digests
1696 .iter()
1697 .zip(executed_effects_digests.iter())
1698 .filter_map(|(digest, executed_effects_digest)| {
1699 if executed_effects_digest.is_none() {
1700 Some(digest)
1701 } else {
1702 None
1703 }
1704 })
1705 .collect::<Vec<_>>();
1706 debug_fatal!(
1707 "Timed out waiting for effects digests to be executed: {:?}",
1708 pending_digests
1709 );
1710 } else {
1711 debug_fatal!(
1712 "Timed out waiting for effects digests to be executed, digests not found"
1713 );
1714 }
1715 }
1716 }
1717
1718 pub fn state(&self) -> Arc<AuthorityState> {
1719 self.state.clone()
1720 }
1721
1722 pub fn reference_gas_price_for_testing(&self) -> Result<u64, anyhow::Error> {
1724 self.state.reference_gas_price_for_testing()
1725 }
1726
1727 pub fn clone_committee_store(&self) -> Arc<CommitteeStore> {
1728 self.state.committee_store().clone()
1729 }
1730
1731 pub fn clone_authority_aggregator(
1739 &self,
1740 ) -> Option<Arc<AuthorityAggregator<NetworkAuthorityClient>>> {
1741 self.transaction_orchestrator
1742 .as_ref()
1743 .map(|to| to.clone_authority_aggregator())
1744 }
1745
1746 pub fn transaction_orchestrator(
1747 &self,
1748 ) -> Option<Arc<TransactionOrchestrator<NetworkAuthorityClient>>> {
1749 self.transaction_orchestrator.clone()
1750 }
1751
1752 pub fn transaction_builder_ledger_client(&self) -> NodeTransactionBuilderLedgerClient {
1755 let reader = Arc::new(GrpcReadStore::new(
1756 self.state.clone(),
1757 self.state_sync_store.clone(),
1758 ));
1759 NodeTransactionBuilderLedgerClient::new(reader)
1760 }
1761
1762 pub fn subscribe_to_transaction_orchestrator_effects(
1765 &self,
1766 ) -> Result<tokio::sync::broadcast::Receiver<QuorumDriverEffectsQueueResult>> {
1767 self.transaction_orchestrator
1768 .as_ref()
1769 .ok_or_else(|| {
1770 anyhow::anyhow!("Transaction Orchestrator is not enabled in this node.")
1771 })?
1772 .subscribe_to_effects_queue()
1773 .ok_or_else(|| {
1774 anyhow::anyhow!(
1775 "Effects queue is not available: the quorum driver is not the currently \
1776 served flow on this node."
1777 )
1778 })
1779 }
1780
1781 pub async fn monitor_reconfiguration(
1787 self: Arc<Self>,
1788 mut epoch_store: Arc<AuthorityPerEpochStore>,
1789 ) -> Result<()> {
1790 let checkpoint_executor_metrics =
1791 CheckpointExecutorMetrics::new(&self.registry_service.default_registry());
1792
1793 loop {
1794 let mut hasher_guard = self.global_state_hasher.lock().await;
1795 let hasher = hasher_guard.take().unwrap();
1796 info!(
1797 "Creating checkpoint executor for epoch {}",
1798 epoch_store.epoch()
1799 );
1800
1801 let data_sender = if let Ok(guard) = self.grpc_server_handle.try_lock() {
1803 guard.as_ref().map(|handle| {
1804 let tx = handle.checkpoint_data_broadcaster().clone();
1805 Box::new(move |data: &CheckpointData| {
1806 tx.send_traced(data);
1807 }) as Box<dyn Fn(&CheckpointData) + Send + Sync>
1808 })
1809 } else {
1810 None
1811 };
1812
1813 let checkpoint_executor = CheckpointExecutor::new(
1814 epoch_store.clone(),
1815 self.checkpoint_store.clone(),
1816 self.state.clone(),
1817 hasher.clone(),
1818 self.backpressure_manager.clone(),
1819 self.config.checkpoint_executor_config.clone(),
1820 checkpoint_executor_metrics.clone(),
1821 data_sender,
1822 Some(self.checkpoint_progress_tracker.clone()),
1823 );
1824
1825 let run_with_range = self.config.run_with_range;
1826
1827 let cur_epoch_store = self.state.load_epoch_store_one_call_per_task();
1828
1829 self.metrics
1831 .current_protocol_version
1832 .set(cur_epoch_store.protocol_config().version.as_u64() as i64);
1833
1834 if let Some(components) = &*self.validator_components.lock().await {
1836 tokio::time::sleep(Duration::from_millis(1)).await;
1838
1839 let config = cur_epoch_store.protocol_config();
1840 let transaction = ConsensusTransaction::new_capability_notification_v1(
1841 AuthorityCapabilitiesV1::new(
1842 self.state.name,
1843 cur_epoch_store.get_chain(),
1844 self.config
1845 .supported_protocol_versions
1846 .expect("Supported versions should be populated")
1847 .truncate_below(config.version),
1849 self.state.get_available_system_packages(config).await,
1850 ),
1851 );
1852 info!(?transaction, "submitting capabilities to consensus");
1853 components
1854 .consensus_adapter
1855 .submit(transaction, None, &cur_epoch_store)?;
1856
1857 if config.deny_rule_governance() {
1863 let proposed_rules = self.config.transaction_deny_config.to_deny_rule_set();
1864 let recorded = cur_epoch_store.recorded_deny_rule_proposal(&self.state.name);
1865 let transaction = ConsensusTransaction::new_transaction_deny_rule_proposal(
1866 TransactionDenyRuleProposal::new(
1867 self.state.name,
1868 proposed_rules,
1869 recorded.map(|p| p.generation),
1870 ),
1871 );
1872 info!(
1873 tracking_id = ?transaction.get_tracking_id(),
1874 "submitting deny rule proposal to consensus"
1875 );
1876 components
1877 .consensus_adapter
1878 .submit(transaction, None, &cur_epoch_store)?;
1879 }
1880 } else if self.state.is_active_validator(&cur_epoch_store)
1881 && cur_epoch_store
1882 .protocol_config()
1883 .track_non_committee_eligible_validators()
1884 {
1885 let epoch_store = cur_epoch_store.clone();
1889 let node_clone = self.clone();
1890 spawn_monitored_task!(epoch_store.clone().within_alive_epoch(async move {
1891 node_clone
1892 .send_signed_capability_notification_to_committee_with_retry(&epoch_store)
1893 .instrument(trace_span!(
1894 "send_signed_capability_notification_to_committee_with_retry"
1895 ))
1896 .await;
1897 }));
1898 }
1899
1900 let stop_condition = checkpoint_executor.run_epoch(run_with_range).await;
1901
1902 if stop_condition == StopReason::RunWithRangeCondition {
1903 IotaNode::shutdown(&self).await;
1904 self.shutdown_channel_tx
1905 .send(run_with_range)
1906 .expect("RunWithRangeCondition met but failed to send shutdown message");
1907 return Ok(());
1908 }
1909
1910 let latest_system_state = self
1912 .state
1913 .get_object_cache_reader()
1914 .try_get_iota_system_state_object_unsafe()
1915 .expect("Read IOTA System State object cannot fail");
1916
1917 #[cfg(msim)]
1918 if !self
1919 .sim_state
1920 .sim_safe_mode_expected
1921 .load(Ordering::Relaxed)
1922 {
1923 debug_assert!(!latest_system_state.safe_mode());
1924 }
1925
1926 #[cfg(not(msim))]
1927 debug_assert!(!latest_system_state.safe_mode());
1928
1929 if let Err(err) = self.end_of_epoch_channel.send(latest_system_state.clone()) {
1930 if self.state.is_fullnode(&cur_epoch_store) {
1931 warn!(
1932 "Failed to send end of epoch notification to subscriber: {:?}",
1933 err
1934 );
1935 }
1936 }
1937
1938 cur_epoch_store.record_is_safe_mode_metric(latest_system_state.safe_mode());
1939 let new_epoch_start_state = latest_system_state.into_epoch_start_state();
1940
1941 self.auth_agg.store(Arc::new(
1942 self.auth_agg
1943 .load()
1944 .recreate_with_new_epoch_start_state(&new_epoch_start_state),
1945 ));
1946
1947 let next_epoch_committee = new_epoch_start_state.get_iota_committee();
1948 let next_epoch = next_epoch_committee.epoch();
1949 assert_eq!(cur_epoch_store.epoch() + 1, next_epoch);
1950
1951 info!(
1952 next_epoch,
1953 "Finished executing all checkpoints in epoch. About to reconfigure the system."
1954 );
1955
1956 fail_point_async!("reconfig_delay");
1957
1958 let authority_names_to_peer_ids =
1963 new_epoch_start_state.get_authority_names_to_peer_ids();
1964 self.connection_monitor_status
1965 .update_mapping_for_epoch(authority_names_to_peer_ids);
1966
1967 cur_epoch_store.record_epoch_reconfig_start_time_metric();
1968
1969 send_trusted_peer_change(
1970 &self.config,
1971 &self.trusted_peer_change_tx,
1972 &new_epoch_start_state,
1973 );
1974
1975 let mut validator_components_lock_guard = self.validator_components.lock().await;
1976
1977 let new_epoch_store = self
1981 .reconfigure_state(
1982 &self.state,
1983 &cur_epoch_store,
1984 next_epoch_committee.clone(),
1985 new_epoch_start_state,
1986 hasher.clone(),
1987 )
1988 .await?;
1989
1990 let new_validator_components = if let Some(ValidatorComponents {
1991 validator_server_handle,
1992 validator_overload_monitor_handle,
1993 consensus_queue_overload_monitor_handle,
1994 soft_lock_sweep_handle,
1995 overload_notifier_handle,
1996 consensus_manager,
1997 consensus_store_pruner,
1998 consensus_adapter,
1999 soft_locks,
2000 mut checkpoint_service_tasks,
2001 checkpoint_metrics,
2002 iota_tx_validator_metrics,
2003 validator_registry_id,
2004 }) = validator_components_lock_guard.take()
2005 {
2006 info!("Reconfiguring the validator.");
2007 if let Some(handle) = overload_notifier_handle {
2010 handle.abort();
2011 }
2012 checkpoint_service_tasks.abort_all();
2017 while let Some(result) = checkpoint_service_tasks.join_next().await {
2018 if let Err(err) = result {
2019 if err.is_panic() {
2020 std::panic::resume_unwind(err.into_panic());
2021 }
2022 warn!("Error in checkpoint service task: {:?}", err);
2023 }
2024 }
2025 info!("Checkpoint service has shut down.");
2026
2027 consensus_manager.shutdown().await;
2028 info!("Consensus has shut down.");
2029
2030 info!("Epoch store finished reconfiguration.");
2031
2032 let global_state_hasher_metrics = Arc::into_inner(hasher)
2035 .expect("Object state hasher should have no other references at this point")
2036 .metrics();
2037 let new_hasher = Arc::new(GlobalStateHasher::new(
2038 self.state.get_global_state_hash_store().clone(),
2039 global_state_hasher_metrics,
2040 ));
2041 let weak_hasher = Arc::downgrade(&new_hasher);
2042 *hasher_guard = Some(new_hasher);
2043
2044 consensus_store_pruner.prune(next_epoch).await;
2045
2046 if self.state.is_committee_validator(&new_epoch_store) {
2047 Some(
2049 Self::start_epoch_specific_validator_components(
2050 &self.config,
2051 self.state.clone(),
2052 consensus_adapter,
2053 self.checkpoint_store.clone(),
2054 new_epoch_store.clone(),
2055 self.state_sync_handle.clone(),
2056 self.randomness_handle.clone(),
2057 consensus_manager,
2058 consensus_store_pruner,
2059 weak_hasher,
2060 self.backpressure_manager.clone(),
2061 soft_locks,
2062 validator_server_handle,
2063 validator_overload_monitor_handle,
2064 consensus_queue_overload_monitor_handle,
2065 soft_lock_sweep_handle,
2066 checkpoint_metrics,
2067 iota_tx_validator_metrics,
2068 validator_registry_id,
2069 )
2070 .await?,
2071 )
2072 } else {
2073 info!("This node is no longer a validator after reconfiguration");
2074 if self.registry_service.remove(validator_registry_id) {
2075 debug!("Removed validator metrics registry");
2076 } else {
2077 warn!("Failed to remove validator metrics registry");
2078 }
2079 validator_server_handle.shutdown();
2080 debug!("Validator grpc server shutdown triggered");
2081
2082 None
2083 }
2084 } else {
2085 let global_state_hasher_metrics = Arc::into_inner(hasher)
2088 .expect("Object state hasher should have no other references at this point")
2089 .metrics();
2090 let new_hasher = Arc::new(GlobalStateHasher::new(
2091 self.state.get_global_state_hash_store().clone(),
2092 global_state_hasher_metrics,
2093 ));
2094 let weak_hasher = Arc::downgrade(&new_hasher);
2095 *hasher_guard = Some(new_hasher);
2096
2097 if self.state.is_committee_validator(&new_epoch_store) {
2098 info!("Promoting the node from fullnode to validator, starting grpc server");
2099
2100 let mut components = Self::construct_validator_components(
2101 self.config.clone(),
2102 self.state.clone(),
2103 Arc::new(next_epoch_committee.clone()),
2104 new_epoch_store.clone(),
2105 self.checkpoint_store.clone(),
2106 self.state_sync_handle.clone(),
2107 self.randomness_handle.clone(),
2108 weak_hasher,
2109 self.backpressure_manager.clone(),
2110 self.connection_monitor_status.clone(),
2111 &self.registry_service,
2112 self.serving_rt_handle.clone(),
2113 )
2114 .await?;
2115
2116 components.validator_server_handle =
2117 components.validator_server_handle.start().await;
2118
2119 Some(components)
2120 } else {
2121 None
2122 }
2123 };
2124 *validator_components_lock_guard = new_validator_components;
2125
2126 cur_epoch_store.release_db_handles();
2129
2130 drop(cur_epoch_store);
2134
2135 self.state.epoch_db_pruner().prune_old_epoch_dbs().await;
2139
2140 if cfg!(msim)
2141 && !matches!(
2142 self.config
2143 .authority_store_pruning_config
2144 .num_epochs_to_retain_for_checkpoints(),
2145 None | Some(u64::MAX) | Some(0)
2146 )
2147 {
2148 self.state
2149 .prune_checkpoints_for_eligible_epochs_for_testing(
2150 self.config.clone(),
2151 iota_core::authority::authority_store_pruner::AuthorityStorePruningMetrics::new_for_test(),
2152 )
2153 .await?;
2154 }
2155
2156 epoch_store = new_epoch_store;
2157 info!("Reconfiguration finished");
2158 }
2159 }
2160
2161 async fn shutdown(&self) {
2162 if let Some(validator_components) = &*self.validator_components.lock().await {
2163 validator_components.consensus_manager.shutdown().await;
2164 }
2165
2166 if let Some(grpc_handle) = self.grpc_server_handle.lock().await.take() {
2168 info!("Shutting down gRPC server");
2169 if let Err(e) = grpc_handle.shutdown().await {
2170 warn!("Failed to gracefully shutdown gRPC server: {e}");
2171 }
2172 }
2173 }
2174
2175 async fn reconfigure_state(
2178 &self,
2179 state: &Arc<AuthorityState>,
2180 cur_epoch_store: &AuthorityPerEpochStore,
2181 next_epoch_committee: Committee,
2182 next_epoch_start_system_state: EpochStartSystemState,
2183 global_state_hasher: Arc<GlobalStateHasher>,
2184 ) -> IotaResult<Arc<AuthorityPerEpochStore>> {
2185 let next_epoch = next_epoch_committee.epoch();
2186
2187 let last_checkpoint = self
2188 .checkpoint_store
2189 .get_epoch_last_checkpoint(cur_epoch_store.epoch())
2190 .expect("Error loading last checkpoint for current epoch")
2191 .expect("Could not load last checkpoint for current epoch");
2192 let epoch_supply_change = last_checkpoint
2193 .end_of_epoch_data
2194 .as_ref()
2195 .ok_or_else(|| {
2196 IotaError::from("last checkpoint in epoch should contain end of epoch data")
2197 })?
2198 .epoch_supply_change;
2199
2200 let last_checkpoint_seq = last_checkpoint.sequence_number();
2201
2202 assert_eq!(
2203 Some(last_checkpoint_seq),
2204 self.checkpoint_store
2205 .get_highest_executed_checkpoint_seq_number()
2206 .expect("Error loading highest executed checkpoint sequence number")
2207 );
2208
2209 let epoch_start_configuration = EpochStartConfiguration::new(
2210 next_epoch_start_system_state,
2211 *last_checkpoint.digest(),
2212 state.get_object_store().as_ref(),
2213 EpochFlag::default_flags_for_new_epoch(&state.config),
2214 )
2215 .expect("EpochStartConfiguration construction cannot fail");
2216
2217 let new_epoch_store = self
2218 .state
2219 .reconfigure(
2220 cur_epoch_store,
2221 self.config.supported_protocol_versions.unwrap(),
2222 next_epoch_committee,
2223 epoch_start_configuration,
2224 global_state_hasher,
2225 &self.config.expensive_safety_check_config,
2226 epoch_supply_change,
2227 last_checkpoint_seq,
2228 )
2229 .await
2230 .expect("Reconfigure authority state cannot fail");
2231 info!(next_epoch, "Node State has been reconfigured");
2232 assert_eq!(next_epoch, new_epoch_store.epoch());
2233 self.state.get_reconfig_api().update_epoch_flags_metrics(
2234 cur_epoch_store.epoch_start_config().flags(),
2235 new_epoch_store.epoch_start_config().flags(),
2236 );
2237
2238 Ok(new_epoch_store)
2239 }
2240
2241 pub fn get_config(&self) -> &NodeConfig {
2242 &self.config
2243 }
2244
2245 async fn execute_transaction_immediately_at_zero_epoch(
2246 state: &Arc<AuthorityState>,
2247 epoch_store: &Arc<AuthorityPerEpochStore>,
2248 tx: &TransactionEnvelope,
2249 span: tracing::Span,
2250 ) {
2251 let _guard = span.enter();
2252 let transaction =
2253 iota_types::executable_transaction::VerifiedExecutableTransaction::new_unchecked(
2254 iota_types::executable_transaction::ExecutableTransaction::new_from_data_and_sig(
2255 tx.data().clone(),
2256 iota_types::executable_transaction::CertificateProof::Checkpoint(0, 0),
2257 ),
2258 );
2259 state
2260 .try_execute_immediately(&transaction, ExecutionEnv::new(), epoch_store)
2261 .unwrap();
2262 }
2263
2264 pub fn randomness_handle(&self) -> randomness::Handle {
2265 self.randomness_handle.clone()
2266 }
2267
2268 pub(crate) fn registry_service(&self) -> &RegistryService {
2271 &self.registry_service
2272 }
2273
2274 async fn send_signed_capability_notification_to_committee_with_retry(
2280 &self,
2281 epoch_store: &Arc<AuthorityPerEpochStore>,
2282 ) {
2283 const INITIAL_RETRY_INTERVAL_SECS: u64 = 5;
2284 const RETRY_INTERVAL_INCREMENT_SECS: u64 = 5;
2285 const MAX_RETRY_INTERVAL_SECS: u64 = 300; let config = epoch_store.protocol_config();
2289
2290 let capabilities = AuthorityCapabilitiesV1::new(
2292 self.state.name,
2293 epoch_store.get_chain(),
2294 self.config
2295 .supported_protocol_versions
2296 .expect("Supported versions should be populated")
2297 .truncate_below(config.version),
2298 self.state.get_available_system_packages(config).await,
2299 );
2300
2301 let signature = AuthoritySignature::new_secure(
2303 &IntentMessage::new(
2304 Intent::iota_app(IntentScope::AuthorityCapabilities),
2305 &capabilities,
2306 ),
2307 &epoch_store.epoch(),
2308 self.config.authority_key_pair(),
2309 );
2310
2311 let request = HandleCapabilityNotificationRequestV1 {
2312 message: SignedAuthorityCapabilitiesV1::new_from_data_and_sig(capabilities, signature),
2313 };
2314
2315 let mut retry_interval = Duration::from_secs(INITIAL_RETRY_INTERVAL_SECS);
2316
2317 loop {
2318 let auth_agg = self.auth_agg.load();
2319 match auth_agg
2320 .send_capability_notification_to_quorum(request.clone())
2321 .await
2322 {
2323 Ok(_) => {
2324 info!("Successfully sent capability notification to committee");
2325 break;
2326 }
2327 Err(err) => {
2328 match &err {
2329 AggregatorSendCapabilityNotificationError::RetryableNotification {
2330 errors,
2331 } => {
2332 warn!(
2333 "Failed to send capability notification to committee (retryable error), will retry in {:?}: {:?}",
2334 retry_interval, errors
2335 );
2336 }
2337 AggregatorSendCapabilityNotificationError::NonRetryableNotification {
2338 errors,
2339 } => {
2340 error!(
2341 "Failed to send capability notification to committee (non-retryable error): {:?}",
2342 errors
2343 );
2344 break;
2345 }
2346 };
2347
2348 tokio::time::sleep(retry_interval).await;
2350
2351 retry_interval = std::cmp::min(
2353 retry_interval + Duration::from_secs(RETRY_INTERVAL_INCREMENT_SECS),
2354 Duration::from_secs(MAX_RETRY_INTERVAL_SECS),
2355 );
2356 }
2357 }
2358 }
2359 }
2360}
2361
2362#[cfg(msim)]
2363impl IotaNode {
2364 pub fn get_sim_node_id(&self) -> iota_simulator::task::NodeId {
2365 self.sim_state.sim_node.id()
2366 }
2367
2368 pub fn set_safe_mode_expected(&self, new_value: bool) {
2369 info!("Setting safe mode expected to {}", new_value);
2370 self.sim_state
2371 .sim_safe_mode_expected
2372 .store(new_value, Ordering::Relaxed);
2373 }
2374}
2375
2376enum SpawnOnce {
2377 Unstarted(
2379 Mutex<BoxFuture<'static, Result<iota_network_stack::server::Server>>>,
2380 tokio::runtime::Handle,
2381 ),
2382 #[allow(unused)]
2383 Started(iota_http::ServerHandle),
2384}
2385
2386impl SpawnOnce {
2387 pub fn new(
2388 future: impl Future<Output = Result<iota_network_stack::server::Server>> + Send + 'static,
2389 serving_rt_handle: tokio::runtime::Handle,
2390 ) -> Self {
2391 Self::Unstarted(Mutex::new(Box::pin(future)), serving_rt_handle)
2392 }
2393
2394 pub async fn start(self) -> Self {
2395 match self {
2396 Self::Unstarted(future, serving_rt_handle) => {
2397 let (handle_tx, handle_rx) = tokio::sync::oneshot::channel();
2401 serving_rt_handle.spawn(async move {
2402 let server = future.into_inner().await.unwrap_or_else(|err| {
2403 panic!("Failed to start validator gRPC server: {err}")
2404 });
2405 if handle_tx.send(server.handle().clone()).is_err() {
2406 return;
2407 }
2408 match server.serve().await {
2409 Ok(()) => info!("Server stopped"),
2410 Err(err) => info!("Server stopped: {err}"),
2411 }
2412 });
2413 let handle = handle_rx
2414 .await
2415 .expect("validator gRPC server exited before returning its handle");
2416 Self::Started(handle)
2417 }
2418 Self::Started(_) => self,
2419 }
2420 }
2421
2422 pub fn shutdown(self) {
2423 if let SpawnOnce::Started(handle) = self {
2424 handle.trigger_shutdown();
2425 }
2426 }
2427}
2428
2429fn send_trusted_peer_change(
2432 config: &NodeConfig,
2433 sender: &watch::Sender<TrustedPeerChangeEvent>,
2434 new_epoch_start_state: &EpochStartSystemState,
2435) {
2436 let new_committee =
2437 new_epoch_start_state.get_validator_as_p2p_peers(config.authority_public_key());
2438
2439 sender.send_modify(|event| {
2440 core::mem::swap(&mut event.new_committee, &mut event.old_committee);
2441 event.new_committee = new_committee;
2442 })
2443}
2444
2445fn build_kv_store(
2446 state: &Arc<AuthorityState>,
2447 config: &NodeConfig,
2448 registry: &Registry,
2449) -> Result<Arc<TransactionKeyValueStore>> {
2450 let metrics = KeyValueStoreMetrics::new(registry);
2451 let db_store = TransactionKeyValueStore::new("rocksdb", metrics.clone(), state.clone());
2452
2453 let base_url = &config.transaction_kv_store_read_config.base_url;
2454
2455 if base_url.is_empty() {
2456 info!("no http kv store url provided, using local db only");
2457 return Ok(Arc::new(db_store));
2458 }
2459
2460 base_url.parse::<url::Url>().tap_err(|e| {
2461 error!(
2462 "failed to parse config.transaction_kv_store_config.base_url ({:?}) as url: {}",
2463 base_url, e
2464 )
2465 })?;
2466
2467 let http_store = HttpKVStore::new_kv(
2468 base_url,
2469 config.transaction_kv_store_read_config.cache_size,
2470 metrics.clone(),
2471 )?;
2472 info!("using local key-value store with fallback to http key-value store");
2473 Ok(Arc::new(FallbackTransactionKVStore::new_kv(
2474 db_store,
2475 http_store,
2476 metrics,
2477 "json_rpc_fallback",
2478 )))
2479}
2480
2481async fn build_grpc_server(
2491 config: &NodeConfig,
2492 state: Arc<AuthorityState>,
2493 state_sync_store: RocksDbStore,
2494 executor: Option<Arc<dyn iota_types::transaction_executor::TransactionExecutor>>,
2495 prometheus_registry: &Registry,
2496 server_version: ServerVersion,
2497) -> Result<Option<GrpcServerHandle>> {
2498 if config.is_validator() || !config.enable_grpc_api {
2502 return Ok(None);
2503 }
2504
2505 let grpc_config = config
2506 .grpc_api_config
2507 .as_ref()
2508 .expect("`NodeConfig::validate` rejects an enabled gRPC API without a config");
2509
2510 let chain_id = state.get_chain_identifier();
2512
2513 let grpc_read_store = Arc::new(GrpcReadStore::new(state.clone(), state_sync_store));
2514
2515 let shutdown_token = CancellationToken::new();
2517
2518 let grpc_reader = Arc::new(GrpcReader::new(
2520 grpc_read_store,
2521 Some(server_version.to_string()),
2522 ));
2523
2524 let grpc_server_metrics = iota_grpc_server::GrpcServerMetrics::new(prometheus_registry);
2526 let client_id_source = config
2527 .policy_config
2528 .as_ref()
2529 .map(|p| p.client_id_source.clone());
2530
2531 let handle = start_grpc_server(
2532 grpc_reader,
2533 executor,
2534 grpc_config.clone(),
2535 shutdown_token,
2536 chain_id,
2537 Some(grpc_server_metrics),
2538 state.traffic_controller.clone(),
2539 client_id_source,
2540 )
2541 .await?;
2542
2543 Ok(Some(handle))
2544}
2545
2546pub async fn build_http_server(
2561 state: Arc<AuthorityState>,
2562 transaction_orchestrator: &Option<Arc<TransactionOrchestrator<NetworkAuthorityClient>>>,
2563 config: &NodeConfig,
2564 prometheus_registry: &Registry,
2565) -> Result<Option<iota_http::ServerHandle>> {
2566 if config.is_validator() {
2568 return Ok(None);
2569 }
2570
2571 let mut router = axum::Router::new();
2572
2573 let json_rpc_router = {
2574 let traffic_controller = state.traffic_controller.clone();
2575 let mut server = JsonRpcServerBuilder::new(
2576 env!("CARGO_PKG_VERSION"),
2577 prometheus_registry,
2578 traffic_controller,
2579 config.policy_config.clone(),
2580 );
2581
2582 let kv_store = build_kv_store(&state, config, prometheus_registry)?;
2583
2584 let metrics = Arc::new(JsonRpcMetrics::new(prometheus_registry));
2585 server.register_module(ReadApi::new(
2586 state.clone(),
2587 kv_store.clone(),
2588 metrics.clone(),
2589 ))?;
2590 server.register_module(CoinReadApi::new(
2591 state.clone(),
2592 kv_store.clone(),
2593 metrics.clone(),
2594 )?)?;
2595
2596 if config.run_with_range.is_none() {
2599 server.register_module(TransactionBuilderApi::new(state.clone()))?;
2600 }
2601 server.register_module(GovernanceReadApi::new(state.clone(), metrics.clone()))?;
2602
2603 if let Some(transaction_orchestrator) = transaction_orchestrator {
2604 server.register_module(TransactionExecutionApi::new(
2605 state.clone(),
2606 transaction_orchestrator.clone(),
2607 metrics.clone(),
2608 ))?;
2609 }
2610
2611 let iota_names_config = config
2612 .iota_names_config
2613 .clone()
2614 .unwrap_or_else(|| IotaNamesConfig::from_chain(&state.get_chain_identifier().chain()));
2615
2616 server.register_module(IndexerApi::new(
2617 state.clone(),
2618 ReadApi::new(state.clone(), kv_store.clone(), metrics.clone()),
2619 kv_store,
2620 metrics,
2621 iota_names_config,
2622 config.indexer_max_subscriptions,
2623 ))?;
2624 server.register_module(MoveUtils::new(state.clone()))?;
2625
2626 let server_type = config.jsonrpc_server_type();
2627
2628 server.to_router(server_type).await?
2629 };
2630
2631 router = router.merge(json_rpc_router);
2632
2633 router = router
2634 .route("/health", axum::routing::get(health_check_handler))
2635 .route_layer(axum::Extension(state));
2636
2637 let layers = ServiceBuilder::new()
2638 .map_request(|mut request: axum::http::Request<_>| {
2639 if let Some(connect_info) = request.extensions().get::<iota_http::ConnectInfo>() {
2640 let axum_connect_info = axum::extract::ConnectInfo(connect_info.remote_addr);
2641 request.extensions_mut().insert(axum_connect_info);
2642 }
2643 request
2644 })
2645 .layer(axum::middleware::from_fn(server_timing_middleware));
2646
2647 router = router.layer(layers);
2648
2649 let handle = iota_http::Builder::new()
2650 .serve(&config.json_rpc_address, router)
2651 .map_err(|e| anyhow::anyhow!("{e}"))?;
2652 info!(local_addr =? handle.local_addr(), "IOTA JSON-RPC server listening on {}", handle.local_addr());
2653
2654 Ok(Some(handle))
2655}
2656
2657#[derive(Debug, serde::Serialize, serde::Deserialize)]
2658pub struct Threshold {
2659 pub threshold_seconds: Option<u32>,
2660}
2661
2662async fn health_check_handler(
2663 axum::extract::Query(Threshold { threshold_seconds }): axum::extract::Query<Threshold>,
2664 axum::Extension(state): axum::Extension<Arc<AuthorityState>>,
2665) -> impl axum::response::IntoResponse {
2666 if let Some(threshold_seconds) = threshold_seconds {
2667 let summary = match state
2669 .get_checkpoint_store()
2670 .get_highest_executed_checkpoint()
2671 {
2672 Ok(Some(summary)) => summary,
2673 Ok(None) => {
2674 warn!("Highest executed checkpoint not found");
2675 return (axum::http::StatusCode::SERVICE_UNAVAILABLE, "down");
2676 }
2677 Err(err) => {
2678 warn!("Failed to retrieve highest executed checkpoint: {:?}", err);
2679 return (axum::http::StatusCode::SERVICE_UNAVAILABLE, "down");
2680 }
2681 };
2682
2683 let latest_chain_time = summary.timestamp();
2685 let threshold =
2686 std::time::SystemTime::now() - Duration::from_secs(threshold_seconds as u64);
2687
2688 if latest_chain_time < threshold {
2690 warn!(
2691 ?latest_chain_time,
2692 ?threshold,
2693 "failing health check due to checkpoint lag"
2694 );
2695 return (axum::http::StatusCode::SERVICE_UNAVAILABLE, "down");
2696 }
2697 }
2698 (axum::http::StatusCode::OK, "up")
2700}
2701
2702#[cfg(not(test))]
2703fn max_tx_per_checkpoint(protocol_config: &ProtocolConfig) -> usize {
2704 protocol_config.max_transactions_per_checkpoint() as usize
2705}
2706
2707#[cfg(test)]
2708fn max_tx_per_checkpoint(_: &ProtocolConfig) -> usize {
2709 2
2710}
2711
2712#[cfg(all(test, not(msim)))]
2715mod runtime_split_tests {
2716 use std::{
2717 sync::mpsc,
2718 time::{Duration, Instant},
2719 };
2720
2721 use anyhow::anyhow;
2722
2723 use super::SpawnOnce;
2724
2725 fn runtime(name: &'static str) -> tokio::runtime::Runtime {
2728 tokio::runtime::Builder::new_multi_thread()
2729 .worker_threads(1)
2730 .thread_name(name)
2731 .enable_all()
2732 .build()
2733 .unwrap()
2734 }
2735
2736 fn current_pool() -> String {
2738 std::thread::current()
2739 .name()
2740 .unwrap_or("<unnamed>")
2741 .to_string()
2742 }
2743
2744 #[test]
2755 fn spawn_once_binds_on_serving_not_the_core_runtime() {
2756 let node = runtime("node-core");
2757 let serving = runtime("serving");
2758 let (tx, rx) = mpsc::channel::<(&'static str, String)>();
2759
2760 let bind_tx = tx.clone();
2763 let bind_future = async move {
2764 let _ = bind_tx.send(("bind", current_pool()));
2765 let addr = "/ip4/127.0.0.1/tcp/0/http".parse().unwrap();
2766 let server = iota_network_stack::config::Config::new()
2767 .server_builder()
2768 .bind(&addr, None)
2769 .await
2770 .map_err(|e| anyhow!("bind failed: {e}"))?;
2771 Ok(server)
2772 };
2773 let once = SpawnOnce::new(bind_future, serving.handle().clone());
2774
2775 let caller_tx = tx.clone();
2777 node.spawn(async move {
2778 let _ = caller_tx.send(("caller", current_pool()));
2779 let _ = once.start().await;
2780 });
2781
2782 let (mut caller_pool, mut bind_pool) = (None, None);
2784 let deadline = Instant::now() + Duration::from_secs(10);
2785 while (caller_pool.is_none() || bind_pool.is_none()) && Instant::now() < deadline {
2786 if let Ok((which, pool)) = rx.recv_timeout(Duration::from_millis(200)) {
2787 match which {
2788 "caller" => caller_pool = Some(pool),
2789 "bind" => bind_pool = Some(pool),
2790 _ => {}
2791 }
2792 }
2793 }
2794 let caller_pool = caller_pool.expect("start() never ran");
2795 let bind_pool = bind_pool.expect("bind future never ran");
2796
2797 assert!(
2798 caller_pool.starts_with("node-core"),
2799 "start() should run on the node-core runtime, ran on: {caller_pool}"
2800 );
2801 assert!(
2802 bind_pool.starts_with("serving"),
2803 "server bind must run on the serving runtime, ran on: {bind_pool}"
2804 );
2805 assert_ne!(
2806 caller_pool, bind_pool,
2807 "the fix must move the bind off the caller (node-core) runtime onto serving"
2808 );
2809 }
2810}
2811
2812#[cfg(test)]
2813mod config_tests {
2814 use iota_config::NodeConfig;
2815 use iota_metrics::RegistryService;
2816 use prometheus_filtered::Registry;
2817
2818 use super::IotaNode;
2819
2820 #[tokio::test]
2824 async fn start_rejects_a_config_no_node_could_start_with() {
2825 let mut config: NodeConfig = serde_yaml::from_str(
2826 r#"
2827db-path: /nonexistent/db
2828network-address: /dns/localhost/tcp/8080/http
2829metrics-address: "0.0.0.0:9184"
2830json-rpc-address: "0.0.0.0:9000"
2831genesis:
2832 genesis-file-location: /nonexistent/genesis.blob
2833"#,
2834 )
2835 .unwrap();
2836 config.enable_grpc_api = true;
2837 config.grpc_api_config = None;
2838
2839 let err = IotaNode::start(config, RegistryService::new(Registry::new()))
2840 .await
2841 .unwrap_err();
2842
2843 let err = format!("{err:#}");
2844 assert!(err.contains("`grpc-api-config` is `null`"), "{err}");
2845 }
2846}