Skip to main content

iota_node/
lib.rs

1// Copyright (c) Mysten Labs, Inc.
2// Modifications Copyright (c) 2024 IOTA Stiftung
3// SPDX-License-Identifier: Apache-2.0
4
5#[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    /// Handle for the consensus queue overload monitor task, present only
160    /// when the certificate-less (P-COOL) flow is enabled. The
161    /// task self-terminates via `Weak` references; this handle exists purely
162    /// for ownership clarity.
163    consensus_queue_overload_monitor_handle: Option<JoinHandle<()>>,
164    /// Handle for the soft-lock expiry sweep task. The task self-terminates
165    /// via a `Weak` reference; this handle exists purely for ownership clarity.
166    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    // Keeping the handle to the checkpoint service tasks to shut them down during reconfiguration.
173    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    /// The http server responsible for serving JSON-RPC
224    _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    /// Broadcast channel to send the starting system state for the next epoch.
239    end_of_epoch_channel: broadcast::Sender<IotaSystemState>,
240
241    /// Broadcast channel to notify [`DiscoveryEventLoop`] for new validator
242    /// peers.
243    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    // Channel to allow signaling upstream to shutdown iota-node
254    shutdown_channel_tx: broadcast::Sender<Option<RunWithRange>>,
255
256    /// Handle to the gRPC server for gRPC streaming and graceful shutdown
257    grpc_server_handle: Mutex<Option<GrpcServerHandle>>,
258
259    /// AuthorityAggregator of the network, created at start and beginning of
260    /// each epoch. Use ArcSwap so that we could mutate it without taking
261    /// mut reference.
262    // TODO: Eventually we can make this auth aggregator a shared reference so that this
263    // update will automatically propagate to other uses.
264    auth_agg: Arc<ArcSwap<AuthorityAggregator<NetworkAuthorityClient>>>,
265
266    /// Runtime that hosts the client-facing servers and their per-request
267    /// handlers, isolating external request load from the node core.
268    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    /// Starts a node that hosts the client-facing servers on the caller's
281    /// runtime, alongside everything else.
282    ///
283    /// This is intentional: this entry point serves the in-process nodes of
284    /// iota-swarm, where a separate serving runtime is either impossible
285    /// (simtests must keep every task on the simulator's deterministic
286    /// scheduler) or not worth the threads (thread-mode swarm runs many nodes
287    /// per process). Only the `iota-node` binary isolates client-facing
288    /// request handling on the dedicated serving runtime of `IotaRuntimes`,
289    /// via [`IotaNode::start_async`].
290    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    /// Starts a background task that polls the authority's load shedding
304    /// percentage and broadcasts changes to other validators via consensus.
305    /// Returns the task handle if the feature flag is enabled, or `None`
306    /// otherwise.
307    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            // Seed from the percentage this authority last broadcasted
322            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(&registry_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        // Check and set the db_corrupted flag
379        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        // Initialize metrics to track db usage before creating any stores
385        DBMetrics::init(&prometheus_registry);
386
387        // Initialize IOTA metrics.
388        iota_metrics::init_metrics(&prometheus_registry);
389        // Unsupported (because of the use of static variable) and unnecessary in
390        // simtests.
391        #[cfg(not(msim))]
392        iota_metrics::thread_stall_monitor::start_thread_stall_monitor();
393
394        // Monitor the node-core and serving runtimes so that worker-thread
395        // starvation between them is observable. Gated out of simtests, where
396        // tokio runs under the deterministic simulator.
397        #[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        // Register uptime metric
414        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        // If genesis come with some migration data then load them into memory from the
427        // file path specified in config.
428        let migration_tx_data = if genesis.contains_migrations() {
429            // Here the load already verifies that the content of the migration blob is
430            // valid in respect to the content found in genesis
431            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        // By default, only enable write stall on validators for perpetual db.
445        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(&registry_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        // the database is empty at genesis time
539        if is_genesis {
540            info!("checking IOTA conservation at genesis");
541            // When we are opening the db table, the only time when it's safe to
542            // check IOTA conservation is at genesis. Otherwise we may be in the middle of
543            // an epoch and the IOTA conservation check will fail. This also initialize
544            // the expected_network_iota_amount table.
545            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        // Database has everything from genesis, set corrupted key to 0
570        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        // Seed the open epoch's `epoch_info` row before services start.
606        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        // Create network
612        // TODO only configure validators as seed/preferred peers for validators and not
613        // for fullnodes once we've had a chance to re-work fullnode
614        // configuration generation.
615        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        // We must explicitly send this instead of relying on the initial value to
635        // trigger watch value change, so that state-sync is able to process it.
636        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        // Start uploading state snapshot to remote store
644        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        // ensure genesis and migration txs were executed
688        if epoch_store.epoch() == 0 {
689            let genesis_tx = &genesis.transaction();
690            let span = error_span!("genesis_txn", tx_digest = ?genesis_tx.digest());
691            // Execute genesis transaction
692            Self::execute_transaction_immediately_at_zero_epoch(
693                &state,
694                &epoch_store,
695                genesis_tx,
696                span,
697            )
698            .await;
699
700            // Execute migration transactions if present
701            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        // Start the loop that receives new randomness and generates transactions for
716        // it.
717        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        // Run the JSON-RPC server (and its per-request handlers) on the serving
749        // runtime. `iota_http::Builder::serve` spawns the accept loop via
750        // `Handle::current()`, so the builder must execute on the serving runtime.
751        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", &registry_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(&registry_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        // Convert transaction orchestrator to executor trait object for gRPC server
808        // Note that the transaction_orchestrator (so as executor) will be None if it is
809        // a validator node or run_with_range is set
810        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        // Run the gRPC read API server (and its per-request handlers) on the
816        // serving runtime, for the same reason as the JSON-RPC server above.
817        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                    &registry_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            // Start the gRPC server
861            components.validator_server_handle = components.validator_server_handle.start().await;
862
863            Some(components)
864        } else {
865            None
866        };
867
868        // setup shutdown channel
869        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    // Init reconfig process by starting to reject user certs
934    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    // Testing-only API to start epoch close process.
961    // For production code, please use the non-testing version.
962    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    /// Creates a StateSnapshotUploader and starts it if
968    /// `state-snapshot-write-config.object-store-config` is set. Snapshot
969    /// publication is a fullnode-only role.
970    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            // Inbound requests on this network are small (signatures, queries, summaries).
1062            // Cap request frames at 1 MiB.
1063            anemo_config.max_request_frame_size = Some(1 << 20);
1064            // Responses can be larger (checkpoint contents).
1065            // Cap response frames at 128 MiB.
1066            anemo_config.max_response_frame_size = Some(128 << 20);
1067
1068            // Set a higher default value for socket send/receive buffers if not already
1069            // configured.
1070            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            // Set high-performance defaults for quinn transport.
1080            // With 200MiB buffer size and ~500ms RTT, max throughput ~400MiB/s.
1081            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    /// Asynchronously constructs and initializes the components necessary for
1137    /// the validator node.
1138    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        // This only gets started up once, not on every epoch. (Make call to remove
1179        // every epoch.)
1180        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        // Spawn the soft-lock sweep once for the lifetime of this validator
1199        // instance. The task holds only a `Weak<PreConsensusSoftLocks>` so it
1200        // stops itself automatically: each iteration it tries to upgrade the
1201        // weak reference, and when all strong `Arc` owners have been dropped
1202        // (i.e. `ValidatorComponents` is destructured and the old epoch store
1203        // is released after an epoch transition that removes us from the
1204        // committee) the upgrade returns `None` and the loop exits. No explicit
1205        // `abort()` is needed. The same `Arc<PreConsensusSoftLocks>` is reused
1206        // across epoch transitions (see `start_epoch_specific_validator_components`),
1207        // so the task keeps running uninterrupted while the node remains a validator.
1208        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        // Starts an overload monitor that monitors the execution of the authority.
1225        // Don't start the overload monitor when max_load_shedding_percentage is 0.
1226        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        // Starts a monitor that periodically refreshes the
1243        // `consensus_queue_load_shedding_percentage` metric. Without this, the
1244        // metric goes stale once gRPC traffic stops (the only other update
1245        // path is `AuthorityState::check_consensus_queue_graduated_limits`, called on
1246        // each inbound tx). Used in the certificate-less (P-COOL)
1247        // mode.
1248        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    /// Initializes and starts components specific to the current
1288    /// epoch for the validator node.
1289    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        // create a new map that gets injected into both the consensus handler and the
1322        // consensus adapter the consensus handler will write values forwarded
1323        // from consensus, and the consensus adapter will read the values to
1324        // make decisions about which validator submits a transaction to consensus
1325        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        // Wire pre-consensus soft locks to the epoch store so that
1330        // post-consensus processing can release locks once permanent locks are
1331        // quarantined. Clear stale locks from the previous epoch and spawn a
1332        // background sweep task.
1333        soft_locks.clear();
1334        epoch_store.set_soft_locks(soft_locks.clone());
1335
1336        // A validator cannot participate without a randomness manager. Aborting here
1337        // fails loudly at the real cause (startup or reconfiguration) instead
1338        // of leaving the node to panic later on the first consensus commit that
1339        // unconditionally expects one.
1340        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        // Spawn consensus startup asynchronously to avoid blocking other components
1371        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    /// Starts the checkpoint service for the validator node, initializing
1426    /// necessary components and settings.
1427    /// The function ensures proper initialization of the checkpoint service,
1428    /// preparing it to handle checkpoint creation and submission to consensus,
1429    /// while also setting up the necessary monitoring and synchronization
1430    /// mechanisms.
1431    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        // The consensus adapter allows the authority to send user certificates through
1490        // consensus.
1491
1492        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        // Each service gets its own concurrency limit so that a flood of client
1524        // transaction submissions (Validator / ValidatorV2) cannot crowd the
1525        // validator-peer RPCs sharing this listener out of admission slots.
1526        // The config value is per core, so the same config scales with the
1527        // hardware; the effective limit is computed on the machine the server
1528        // actually runs on.
1529        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        // HTTP/2 keepalive is the only mechanism that closes a connection whose
1536        // peer has gone away without closing it: the server pings after this
1537        // long without inbound frames and drops the connection when the ping
1538        // goes unanswered for as long again.
1539        const VALIDATOR_GRPC_KEEPALIVE: Duration = Duration::from_secs(60);
1540        // Bounds the streams one connection may hold open, so a single peer
1541        // cannot fill a service's admission slots on its own. A fullnode sends
1542        // every request to this validator over one connection, so the cap must
1543        // stay well above a busy fullnode's peak concurrency.
1544        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    /// Re-executes pending consensus certificates, which may not have been
1591    /// committed to disk before the node restarted. This is necessary for
1592    /// the following reasons:
1593    ///
1594    /// 1. For any transaction for which we returned signed effects to a client,
1595    ///    we must ensure that we have re-executed the transaction before we
1596    ///    begin accepting grpc requests. Otherwise we would appear to have
1597    ///    forgotten about the transaction.
1598    /// 2. While this is running, we are concurrently waiting for all previously
1599    ///    built checkpoints to be rebuilt. Since there may be dependencies in
1600    ///    either direction (from checkpointed consensus transactions to pending
1601    ///    consensus transactions, or vice versa), we must re-execute pending
1602    ///    consensus transactions to ensure that both processes can complete.
1603    /// 3. Also note that for any pending consensus transactions for which we
1604    ///    wrote a signed effects digest to disk, we must re-execute using that
1605    ///    digest as the expected effects digest, to ensure that we cannot
1606    ///    arrive at different effects than what we previously signed.
1607    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                // TODO: what to do with UserTransactionV1 here? It seems like this only applies to
1617                //  optimistically executed owned-object transactions that possibly didn't go
1618                //  through  consensus before the node restarted. UserTransactionsV1
1619                //  always needs to go  through consensus, so it will be replayed
1620                //  there, just like shared object  transactions.
1621                //
1622                // Shared object txns
1623                // cannot be re-executed at this  point, because we must wait for
1624                // consensus replay to assign shared  object versions.
1625                ConsensusTransactionKind::CertifiedTransaction(tx)
1626                    if !tx.contains_shared_object() =>
1627                {
1628                    let tx = *tx;
1629                    // new_unchecked is safe because we never submit a transaction to consensus
1630                    // without verifying it
1631                    let tx = VerifiedExecutableTransaction::new_from_certificate(
1632                        VerifiedCertificate::new_unchecked(tx),
1633                    );
1634                    // we only need to re-execute if we previously signed the effects (which
1635                    // indicates we returned the effects to a client).
1636                    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            // unwrap_digest okay because only user certs are in
1655            // pending_consensus_certificates
1656            .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        // If this times out, the validator will still almost certainly start up fine.
1673        // But, it is possible that it may temporarily "forget" about
1674        // transactions that it had previously executed. This could confuse
1675        // clients in some circumstances. However, the transactions are still in
1676        // pending_consensus_certificates, so we cannot lose any finality guarantees.
1677        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            // Log all the digests that were not executed to help debugging.
1691            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    // Only used for testing because of how epoch store is loaded.
1723    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_store(&self) -> Arc<AuthorityStore> {
1732    // self.state.db()
1733    // }
1734
1735    /// Clone an AuthorityAggregator from the transaction orchestrator, if
1736    /// this is a fullnode. The snapshot goes stale after an epoch change;
1737    /// call again for a fresh one.
1738    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    /// Read-only client for the SDK's `TransactionBuilder` backed by this
1753    /// node's local state instead of a remote endpoint.
1754    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    /// Subscribe to the quorum driver's effects stream; errors while the
1763    /// quorum driver is not the currently served flow on this node.
1764    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    /// This function awaits the completion of checkpoint execution of the
1782    /// current epoch, after which it initiates reconfiguration of the
1783    /// entire system. This function also handles role changes for the node when
1784    /// epoch changes and, if the node is a validator, advertises capabilities
1785    /// and submits the local deny rule proposal to consensus.
1786    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            // Create closures that handle gRPC type conversion
1802            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            // Update the current protocol version metric.
1830            self.metrics
1831                .current_protocol_version
1832                .set(cur_epoch_store.protocol_config().version.as_u64() as i64);
1833
1834            // Advertise capabilities to committee, if we are a validator.
1835            if let Some(components) = &*self.validator_components.lock().await {
1836                // TODO: without this sleep, the consensus message is not delivered reliably.
1837                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                            // no need to send digests of versions less than the current version
1848                            .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                // Announce the local deny rules. Recorded proposals are
1858                // epoch-scoped, so this re-announces on every epoch change.
1859                // The empty set is announced too: every committee member
1860                // attests its configuration each epoch, so silence means
1861                // offline rather than "no rules".
1862                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                // Send signed capabilities to committee validators if we are a non-committee
1886                // validator in a separate task to not block the caller. Sending is done only if
1887                // the feature flag supporting it is enabled.
1888                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            // Safe to call because we are in the middle of reconfiguration.
1911            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            // We save the connection monitor status map regardless of validator / fullnode
1959            // status so that we don't need to restart the connection monitor
1960            // every epoch. Update the mappings that will be used by the
1961            // consensus adapter if it exists or is about to be created.
1962            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            // The following code handles 4 different cases, depending on whether the node
1978            // was a validator in the previous epoch, and whether the node is a validator
1979            // in the new epoch.
1980            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                // Cancel the old overload notifier task so a new one can be
2008                // started for the next epoch.
2009                if let Some(handle) = overload_notifier_handle {
2010                    handle.abort();
2011                }
2012                // Cancel the old checkpoint service tasks.
2013                // Waiting for checkpoint builder to finish gracefully is not possible, because
2014                // it may wait on transactions while consensus on peers have
2015                // already shut down.
2016                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                // No other components should be holding a strong reference to state hasher
2033                // at this point. Confirm here before we swap in the new hasher.
2034                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                    // Only restart consensus if this node is still a validator in the new epoch.
2048                    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                // No other components should be holding a strong reference to state hasher
2086                // at this point. Confirm here before we swap in the new hasher.
2087                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            // Force releasing current epoch store DB handle, because the
2127            // Arc<AuthorityPerEpochStore> may linger.
2128            cur_epoch_store.release_db_handles();
2129
2130            // Drop the old epoch store to free its in-memory structures
2131            // (ConsensusOutputCache, ConsensusQuarantine, DashMaps, etc.).
2132            // The DB tables were already released above.
2133            drop(cur_epoch_store);
2134
2135            // Prune old epoch databases after each epoch transition to prevent
2136            // accumulation of RocksDB instances during fast catch-up sync
2137            // (e.g. syncing from genesis).
2138            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        // Shutdown the gRPC server if it's running
2167        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    /// Asynchronously reconfigures the state of the authority node for the next
2176    /// epoch.
2177    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    /// Returns the registry service holding the node's Prometheus registries
2269    /// and their shared exposure filter.
2270    pub(crate) fn registry_service(&self) -> &RegistryService {
2271        &self.registry_service
2272    }
2273
2274    /// Sends signed capability notification to committee validators for
2275    /// non-committee validators. This method implements retry logic to handle
2276    /// failed attempts to send the notification. It will retry sending the
2277    /// notification with an increasing interval until it receives a successful
2278    /// response from a f+1 committee members or 2f+1 non-retryable errors.
2279    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; // 5 minutes
2286
2287        // Create the capability notification once
2288        let config = epoch_store.protocol_config();
2289
2290        // Create the capability notification
2291        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        // Sign the capabilities using the authority key pair from config
2302        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                    // Wait before retrying
2349                    tokio::time::sleep(retry_interval).await;
2350
2351                    // Increase retry interval for the next attempt, capped at max
2352                    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    // Mutex is only needed to make SpawnOnce Sync
2378    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                // bind() and serve() must execute on the serving runtime:
2398                // iota_http::Builder::serve captures Handle::current() there for
2399                // the accept loop and every request handler.
2400                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
2429/// Notify [`DiscoveryEventLoop`] that a new list of trusted peers are now
2430/// available.
2431fn 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
2481/// Builds and starts the gRPC server for the IOTA node based on the node's
2482/// configuration.
2483///
2484/// Returns `None` on a validator and when the gRPC API is disabled.
2485///
2486/// # Panics
2487///
2488/// Panics if the gRPC API is enabled without a `grpc-api-config`.
2489/// [`NodeConfig::validate`] rejects such a config at startup.
2490async 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    // Validators do not expose gRPC APIs. This return must agree with the
2499    // gRPC rule in `NodeConfig::validate`, which is what makes the `expect`
2500    // below unreachable.
2501    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    // Get chain identifier from state directly
2511    let chain_id = state.get_chain_identifier();
2512
2513    let grpc_read_store = Arc::new(GrpcReadStore::new(state.clone(), state_sync_store));
2514
2515    // Create cancellation token for proper shutdown hierarchy
2516    let shutdown_token = CancellationToken::new();
2517
2518    // Create GrpcReader
2519    let grpc_reader = Arc::new(GrpcReader::new(
2520        grpc_read_store,
2521        Some(server_version.to_string()),
2522    ));
2523
2524    // Create gRPC server metrics
2525    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
2546/// Builds and starts the HTTP server for the IOTA node, exposing the JSON-RPC
2547/// API based on the node's configuration.
2548///
2549/// This function performs the following tasks:
2550/// 1. Checks if the node is a validator by inspecting the consensus
2551///    configuration; if so, it returns early as validators do not expose these
2552///    APIs.
2553/// 2. Creates an Axum router to handle HTTP requests.
2554/// 3. Initializes the JSON-RPC server and registers various RPC modules based
2555///    on the node's state and configuration, including CoinApi,
2556///    TransactionBuilderApi, GovernanceApi, TransactionExecutionApi, and
2557///    IndexerApi.
2558/// 4. Binds the server to the specified JSON-RPC address and starts listening
2559///    for incoming connections.
2560pub 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    // Validators do not expose these APIs
2567    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 run_with_range is enabled we want to prevent any transactions
2597        // run_with_range = None is normal operating conditions
2598        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        // Attempt to get the latest checkpoint
2668        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        // Calculate the threshold time based on the provided threshold_seconds
2684        let latest_chain_time = summary.timestamp();
2685        let threshold =
2686            std::time::SystemTime::now() - Duration::from_secs(threshold_seconds as u64);
2687
2688        // Check if the latest checkpoint is within the threshold
2689        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    // if health endpoint is responding and no threshold is given, respond success
2699    (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// Not msim: this test asserts routing across real OS worker-thread pools by
2713// name, which the deterministic simulator collapses onto a single thread.
2714#[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    /// A single-worker-thread runtime whose worker thread carries `name`, so a
2726    /// task can tell which runtime it is running on via `current_pool()`.
2727    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    /// Name of the runtime whose worker thread is executing this code.
2737    fn current_pool() -> String {
2738        std::thread::current()
2739            .name()
2740            .unwrap_or("<unnamed>")
2741            .to_string()
2742    }
2743
2744    /// Regression guard for the runtime split. `SpawnOnce::start()` is invoked
2745    /// on the node-core runtime (as `IotaNode::start_async` does), but it
2746    /// must run the server *bind* on the serving runtime:
2747    /// `iota_http::Builder::serve` captures `Handle::current()` there for
2748    /// the accept loop and every request handler.
2749    ///
2750    /// The test records the runtime `start()` runs on and the runtime the bind
2751    /// runs on, and asserts the former is node-core and the latter is
2752    /// serving (so they differ). With the pre-fix inline bind the bind ran
2753    /// on the caller (node-core) runtime, and this test fails.
2754    #[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        // The "bind" future records where it runs, then binds a minimal real
2761        // server (health service only) on an ephemeral port.
2762        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        // Drive start() ON the node-core runtime and record the runtime it runs on.
2776        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        // Collect both readings.
2783        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    /// `start_async` validates the config before it does anything else. That
2821    /// keeps the `expect` in `build_grpc_server` and the `debug_assert` in
2822    /// `start_state_snapshot` unreachable.
2823    #[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}