Skip to main content

iota_core/
authority.rs

1// Copyright (c) 2021, Facebook, Inc. and its affiliates
2// Copyright (c) Mysten Labs, Inc.
3// Modifications Copyright (c) 2024 IOTA Stiftung
4// SPDX-License-Identifier: Apache-2.0
5
6use std::{
7    collections::{BTreeMap, HashMap, HashSet},
8    fs::File,
9    io::Write,
10    path::{Path, PathBuf},
11    pin::Pin,
12    sync::{Arc, atomic::Ordering},
13    time::{Duration, SystemTime, UNIX_EPOCH},
14    vec,
15};
16
17use arc_swap::{ArcSwap, Guard};
18use async_trait::async_trait;
19use authority_per_epoch_store::TxLockGuard;
20pub use authority_store::{AuthorityStore, ResolverWrapper, UpdateType};
21use fastcrypto::{
22    encoding::{Base58, Encoding},
23    hash::MultisetHash,
24};
25use iota_common::{debug_fatal, fatal};
26use iota_config::{
27    NodeConfig,
28    genesis::Genesis,
29    node::{AuthorityOverloadConfig, ExpensiveSafetyCheckConfig, StateDebugDumpConfig},
30};
31use iota_framework::{BuiltInFramework, SystemPackage as FrameworkSystemPackage};
32use iota_json_rpc_types::{
33    EventFilter, IotaEvent, IotaMoveValue, IotaObjectDataFilter, IotaTransactionBlockEffects,
34    IotaTransactionBlockEvents, TransactionFilter,
35};
36use iota_macros::{fail_point, fail_point_async, fail_point_if};
37use iota_metrics::{
38    TX_TYPE_SHARED_OBJ_TX, TX_TYPE_SINGLE_WRITER_TX, monitored_scope, spawn_monitored_task,
39};
40use iota_sdk_types::{
41    Address, CheckpointCommitment, CheckpointContents, CheckpointContentsDigest, CheckpointDigest,
42    CheckpointSummary, Digest, EndOfEpochTransactionKind, ExecutionStatus, GasCostSummary,
43    InputSharedObject, MoveAuthenticator, ObjectDigest, ObjectId, ObjectReference, Owner,
44    RandomnessRound, SenderSignedTransaction, StructTag, SystemPackage, Transaction,
45    TransactionDigest, TransactionEffects, TransactionEffectsDigest, TransactionEvents,
46    TransactionKind, TypeTag, Version, WriteKind,
47    crypto::{Intent, IntentScope},
48};
49use iota_storage::{
50    key_value_store::{
51        KVStoreTransactionData, TransactionKeyValueStore, TransactionKeyValueStoreTrait,
52    },
53    key_value_store_metrics::KeyValueStoreMetrics,
54};
55use iota_traffic_controller::{TrafficController, metrics::TrafficControllerMetrics};
56use iota_transaction_checks::VerifierLimitsSource;
57#[cfg(msim)]
58use iota_types::committee::CommitteeTrait;
59use iota_types::{
60    account_abstraction::authenticator_function::{
61        AuthenticatorFunctionRef, AuthenticatorFunctionRefForExecution,
62        MoveAuthenticatorForExecution, MoveAuthenticatorsForExecution,
63        authenticator_function_ref_v1_from_dynamic_field_object,
64        derive_authenticator_function_ref_v1_dynamic_field_id, extract_auth_fun_refs,
65    },
66    auth_context::AuthContextData,
67    base_types::{AuthorityName, ConciseableName, ObjectInfo, ObjectType, VersionNumber},
68    committee::{Committee, EpochId, ProtocolVersion},
69    crypto::{AggregateAuthorityPublicKey, AuthoritySignInfo, AuthoritySignature, Signer},
70    deny_list_v1::check_coin_deny_list_v1,
71    deny_rule_governance::DenyRuleConfig,
72    digests::ChainIdentifier,
73    dynamic_field::{DynamicFieldInfo, DynamicFieldName, visitor as DFV},
74    effects::{
75        SignedTransactionEffects, TransactionEffectsAPI, TransactionEffectsExt,
76        VerifiedSignedTransactionEffects,
77    },
78    error::{ExecutionError, ExecutionErrorKind, IotaError, IotaResult, UserInputError},
79    event::{EventID, SystemEpochInfoEvent},
80    executable_transaction::VerifiedExecutableTransaction,
81    execution_config_utils::to_binary_config,
82    fp_ensure,
83    gas::IotaGasStatus,
84    gas_coin::mock_simulation_gas_coin,
85    inner_temporary_store::{
86        InnerTemporaryStore, ObjectMap, PackageStoreWithFallback, TxCoins, WrittenObjects,
87    },
88    iota_sdk_types_conversions::type_tag_core_to_sdk,
89    iota_system_state::{
90        IotaSystemState, IotaSystemStateTrait,
91        epoch_start_iota_system_state::EpochStartSystemStateTrait, get_iota_system_state,
92    },
93    layout_resolver::{LayoutResolver, into_struct_layout},
94    message_envelope::Message,
95    messages_checkpoint::{
96        CertifiedCheckpointSummary, CheckpointContentsExt, CheckpointRequest, CheckpointResponse,
97        CheckpointSequenceNumber, CheckpointSummaryResponse, CheckpointTimestamp,
98        ECMHLiveObjectSetDigest, VerifiedCheckpoint,
99    },
100    messages_consensus::AuthorityCapabilitiesV1,
101    messages_grpc::{
102        HandleTransactionResponse, LayoutGenerationOption, ObjectInfoRequest,
103        ObjectInfoRequestKind, ObjectInfoResponse, TransactionInfoRequest, TransactionInfoResponse,
104        TransactionStatus,
105    },
106    metrics::{BytecodeVerifierMetrics, LimitsMetrics},
107    move_authenticator::MoveAuthenticatorExt,
108    object::{Object, ObjectRead, PastObjectRead, bounded_visitor::BoundedVisitor},
109    storage::{BackingPackageStore, BackingStore, ObjectKey, ObjectOrTombstone, ObjectStore},
110    supported_protocol_versions::{
111        ProtocolConfig, SupportedProtocolVersions, SupportedProtocolVersionsWithHashes,
112    },
113    traffic_control::{PolicyConfig, RemoteFirewallConfig, TrafficControlReconfigParams},
114    transaction::*,
115    transaction_executor::{SimulateTransactionResult, VmChecks},
116};
117use itertools::Itertools;
118use move_binary_format::{CompiledModule, binary_config::BinaryConfig};
119use move_core_types::{
120    account_address::AccountAddress, annotated_value::MoveStructLayout, language_storage::ModuleId,
121};
122use parking_lot::Mutex;
123use prometheus_filtered::{
124    Histogram, HistogramVec, IntCounter, IntCounterVec, IntGauge, IntGaugeVec, MetricLevel,
125    Registry, register_histogram_vec_with_registry, register_histogram_with_registry,
126    register_int_counter_vec_with_registry, register_int_counter_with_registry,
127    register_int_gauge_vec_with_registry, register_int_gauge_with_registry,
128};
129use serde::{Deserialize, Serialize, de::DeserializeOwned};
130use tap::TapFallible;
131use tokio::{
132    sync::{RwLock, mpsc, mpsc::unbounded_channel, oneshot},
133    task::JoinHandle,
134};
135use tracing::{debug, error, info, instrument, trace, warn};
136use typed_store::TypedStoreError;
137
138use self::{
139    authority_store::ExecutionLockWriteGuard, authority_store_pruner::AuthorityStorePruningMetrics,
140};
141#[cfg(msim)]
142pub use crate::checkpoints::checkpoint_executor::utils::{
143    CheckpointTimeoutConfig, init_checkpoint_timeout_config,
144};
145use crate::{
146    authority::{
147        authority_per_epoch_store::{AuthorityPerEpochStore, TxGuard},
148        authority_per_epoch_store_pruner::AuthorityPerEpochStorePruner,
149        authority_store::{ExecutionLockReadGuard, ObjectLockStatus},
150        authority_store_pruner::{AuthorityStorePruner, EPOCH_DURATION_MS_FOR_TESTING},
151        epoch_start_configuration::{EpochStartConfigTrait, EpochStartConfiguration},
152        shared_object_version_manager::{AssignedVersions, Schedulable},
153    },
154    authority_client::NetworkAuthorityClient,
155    checkpoint_progress_tracker::CheckpointProgressTracker,
156    checkpoints::{CheckpointBuilderError, CheckpointBuilderResult, CheckpointStore},
157    congestion_tracker::CongestionTracker,
158    consensus_adapter::ConsensusAdapter,
159    epoch::committee_store::CommitteeStore,
160    epoch_end_db_snapshot::{EpochEndDbSnapshotHandle, HandOver},
161    execution_cache::{
162        CheckpointCache, ExecutionCacheCommit, ExecutionCacheReconfigAPI,
163        ExecutionCacheTraitPointers, ExecutionCacheWrite, ObjectCacheRead, StateSyncAPI,
164        TransactionCacheRead,
165    },
166    execution_driver::execution_process,
167    execution_scheduler::{ExecutionSchedulerAPI, ExecutionSchedulerWrapper},
168    global_state_hasher::{GlobalStateHashStore, GlobalStateHasher},
169    grpc_indexes::GrpcIndexesStore,
170    jsonrpc_index::{CoinInfo, IndexStore, ObjectIndexChanges},
171    metrics::{LatencyObserver, RateTracker},
172    module_cache_metrics::ResolverMetrics,
173    overload_monitor::{
174        AuthorityOverloadInfo, compute_graduated_load_shedding_percentage,
175        overload_monitor_accept_tx,
176    },
177    stake_aggregator::StakeAggregator,
178    subscription_handler::SubscriptionHandler,
179    transaction_input_loader::TransactionInputLoader,
180    transaction_outputs::TransactionOutputs,
181    validator_tx_finalizer::ValidatorTxFinalizer,
182    verify_indexes::verify_indexes,
183};
184
185#[cfg(test)]
186#[path = "unit_tests/authority_tests.rs"]
187pub mod authority_tests;
188
189#[cfg(test)]
190#[path = "unit_tests/transaction_tests.rs"]
191pub mod transaction_tests;
192
193#[cfg(test)]
194#[path = "unit_tests/batch_transaction_tests.rs"]
195mod batch_transaction_tests;
196
197#[cfg(test)]
198#[path = "unit_tests/move_integration_tests.rs"]
199pub mod move_integration_tests;
200
201#[cfg(test)]
202#[path = "unit_tests/gas_tests.rs"]
203mod gas_tests;
204
205#[cfg(test)]
206#[path = "unit_tests/batch_verification_tests.rs"]
207mod batch_verification_tests;
208
209#[cfg(test)]
210#[path = "unit_tests/coin_deny_list_tests.rs"]
211mod coin_deny_list_tests;
212
213#[cfg(test)]
214#[path = "unit_tests/pre_execution_failure_tests.rs"]
215mod pre_execution_failure_tests;
216
217#[cfg(test)]
218#[path = "unit_tests/auth_unit_test_utils.rs"]
219pub mod auth_unit_test_utils;
220
221#[cfg(any(test, feature = "test-utils"))]
222pub mod authority_test_utils;
223
224pub mod authority_per_epoch_store;
225pub mod authority_per_epoch_store_pruner;
226
227pub mod authority_store_pruner;
228pub mod authority_store_tables;
229pub mod authority_store_types;
230pub mod epoch_start_configuration;
231pub mod shared_object_congestion_tracker;
232pub mod shared_object_version_manager;
233pub mod suggested_gas_price_calculator;
234#[cfg(any(test, feature = "test-utils"))]
235pub mod test_authority_builder;
236pub mod transaction_deferral;
237
238pub(crate) mod authority_store;
239pub mod backpressure;
240pub(crate) mod dropped_tx_status_cache;
241
242/// Prometheus metrics which can be displayed in Grafana, queried and alerted on
243pub struct AuthorityMetrics {
244    tx_orders: IntCounter,
245    total_certs: IntCounter,
246    total_cert_attempts: IntCounter,
247    total_effects: IntCounter,
248    pub shared_obj_tx: IntCounter,
249    sponsored_tx: IntCounter,
250    tx_already_processed: IntCounter,
251    num_input_objs: Histogram,
252    num_shared_objects: Histogram,
253    batch_size: Histogram,
254
255    authority_state_handle_transaction_latency: Histogram,
256
257    execute_certificate_latency_single_writer: Histogram,
258    execute_certificate_latency_shared_object: Histogram,
259
260    internal_execution_latency: Histogram,
261    /// Number of times the validator refused to report effects (signed or
262    /// unsigned, labeled by RPC surface) because it had previously signed
263    /// different effects for the same transaction.
264    signed_effects_equivocation_prevented: IntCounterVec,
265    execution_load_input_objects_latency: Histogram,
266    prepare_certificate_latency: Histogram,
267    commit_certificate_latency: Histogram,
268    epoch_end_db_snapshot_handover_latency: Histogram,
269    epoch_end_db_snapshots_skipped: IntCounter,
270
271    pub(crate) transaction_manager_num_enqueued_certificates: IntCounterVec,
272    pub(crate) transaction_manager_num_missing_objects: IntGauge,
273    pub(crate) transaction_manager_num_pending_certificates: IntGauge,
274    pub(crate) transaction_manager_num_executing_certificates: IntGauge,
275    pub(crate) transaction_manager_num_ready: IntGauge,
276    pub(crate) transaction_manager_object_cache_size: IntGauge,
277    pub(crate) transaction_manager_object_cache_hits: IntCounter,
278    pub(crate) transaction_manager_object_cache_misses: IntCounter,
279    pub(crate) transaction_manager_object_cache_evictions: IntCounter,
280    pub(crate) transaction_manager_package_cache_size: IntGauge,
281    pub(crate) transaction_manager_package_cache_hits: IntCounter,
282    pub(crate) transaction_manager_package_cache_misses: IntCounter,
283    pub(crate) transaction_manager_package_cache_evictions: IntCounter,
284    pub(crate) transaction_manager_transaction_queue_age_s: Histogram,
285
286    pub(crate) execution_driver_executed_transactions: IntCounter,
287    pub(crate) execution_driver_dispatch_queue: IntGauge,
288    pub(crate) execution_queueing_delay_s: Histogram,
289    pub(crate) prepare_cert_gas_latency_ratio: Histogram,
290    pub(crate) execution_gas_latency_ratio: Histogram,
291
292    pub(crate) skipped_consensus_txns: IntCounter,
293    pub(crate) skipped_consensus_txns_cache_hit: IntCounter,
294
295    pub(crate) authority_overload_status: IntGauge,
296    /// Percentage of transactions shed due to consensus queue length.
297    pub(crate) consensus_queue_load_shedding_percentage: IntGauge,
298    /// This authority's locally computed load shedding percentage, taken as the
299    /// max of its latency/rate-based, transaction-manager-queue-based, and
300    /// writeback-cache-backpressure signals.
301    pub(crate) local_post_consensus_load_shedding_percentage: IntGauge,
302
303    pub(crate) transaction_overload_sources: IntCounterVec,
304
305    // Post processing metrics
306    post_processing_total_events_emitted: IntCounter,
307    post_processing_total_tx_indexed: IntCounter,
308    post_processing_total_tx_had_event_processed: IntCounter,
309    post_processing_total_failures: IntCounter,
310
311    // Consensus handler metrics
312    pub consensus_handler_processed: IntCounterVec,
313    pub consensus_handler_transaction_sizes: HistogramVec,
314    pub consensus_handler_num_low_scoring_authorities: IntGauge,
315    pub consensus_handler_scores: IntGaugeVec,
316    pub consensus_handler_deferred_transactions: IntCounter,
317    pub consensus_handler_congested_transactions: IntCounter,
318    pub consensus_handler_cancelled_transactions: IntCounter,
319    /// Number of user transactions dropped during a consensus commit because
320    /// post-consensus conflict/lock validation rejected them. Distinct from
321    /// `consensus_handler_load_shedding_dropped_transactions`.
322    pub consensus_handler_validation_dropped_transactions: IntCounter,
323    /// Number of user transactions dropped during a consensus commit by
324    /// post-consensus load shedding, i.e. probabilistically rejected at the
325    /// quorum `consensus_handler_load_shedding_percentage` rate.
326    pub consensus_handler_load_shedding_dropped_transactions: IntCounter,
327    /// Stake-weighted quorum (2f+1) load shedding percentage enforced on user
328    /// transactions in the most recent consensus commit. This is the cluster
329    /// value actually applied post-consensus, as opposed to this authority's
330    /// own `authority_load_shedding_percentage`. 0 when the P-COOL flow is
331    /// disabled.
332    pub consensus_handler_load_shedding_percentage: IntGauge,
333    pub consensus_handler_max_object_costs: IntGaugeVec,
334    pub consensus_committed_subdags: IntCounterVec,
335    pub consensus_committed_messages: IntGaugeVec,
336    pub consensus_committed_user_transactions: IntGaugeVec,
337    pub consensus_handler_leader_round: IntGauge,
338    pub consensus_calculated_throughput: IntGauge,
339    pub consensus_calculated_throughput_profile: IntGauge,
340
341    pub validator_scoreboard_scores: IntGaugeVec,
342    pub invalid_misbehavior_reports_by_authority: IntGaugeVec,
343
344    pub limits_metrics: Arc<LimitsMetrics>,
345
346    /// bytecode verifier metrics for tracking timeouts
347    pub bytecode_verifier_metrics: Arc<BytecodeVerifierMetrics>,
348
349    /// Count of multisig signatures
350    pub multisig_sig_count: IntCounter,
351
352    // Tracks recent average txn queueing delay between when it is ready for execution
353    // until it starts executing.
354    pub execution_queueing_latency: LatencyObserver,
355
356    // Tracks the rate at which transactions become ready for execution in the
357    // scheduler. The need for the Mutex is that the tracker is updated in the
358    // scheduler and read in the overload_monitor. There should be low mutex
359    // contention because the update side is effectively single threaded and the
360    // read rate in overload_monitor is low. If the update side becomes
361    // multi-threaded, we can create one rate tracker per thread.
362    pub txn_ready_rate_tracker: Arc<Mutex<RateTracker>>,
363
364    // Tracks the rate of transactions starts execution in execution driver.
365    // Similar reason for using a Mutex here as to `txn_ready_rate_tracker`.
366    pub execution_rate_tracker: Arc<Mutex<RateTracker>>,
367}
368
369// Override default Prom buckets for positive numbers in 0-10M range
370const POSITIVE_INT_BUCKETS: &[f64] = &[
371    1., 2., 5., 7., 10., 20., 50., 70., 100., 200., 500., 700., 1000., 2000., 5000., 7000., 10000.,
372    20000., 50000., 70000., 100000., 200000., 500000., 700000., 1000000., 2000000., 5000000.,
373    7000000., 10000000.,
374];
375
376const LATENCY_SEC_BUCKETS: &[f64] = &[
377    0.0005, 0.001, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1., 2., 3., 4., 5., 6., 7., 8., 9.,
378    10., 20., 30., 60., 90.,
379];
380
381// Buckets for low latency samples. Starts from 10us.
382const LOW_LATENCY_SEC_BUCKETS: &[f64] = &[
383    0.00001, 0.00002, 0.00005, 0.0001, 0.0002, 0.0005, 0.001, 0.002, 0.005, 0.01, 0.02, 0.05, 0.1,
384    0.2, 0.5, 1., 2., 5., 10., 20., 50., 100.,
385];
386
387const GAS_LATENCY_RATIO_BUCKETS: &[f64] = &[
388    10.0, 50.0, 100.0, 200.0, 300.0, 400.0, 500.0, 600.0, 700.0, 800.0, 900.0, 1000.0, 2000.0,
389    3000.0, 4000.0, 5000.0, 6000.0, 7000.0, 8000.0, 9000.0, 10000.0, 50000.0, 100000.0, 1000000.0,
390];
391
392impl AuthorityMetrics {
393    pub fn new(registry: &prometheus_filtered::Registry) -> AuthorityMetrics {
394        let execute_certificate_latency = register_histogram_vec_with_registry!(
395            "authority_state_execute_certificate_latency",
396            "Latency of executing certificates, including waiting for inputs",
397            &["tx_type"],
398            LATENCY_SEC_BUCKETS.to_vec(),
399            registry;
400            MetricLevel::Info,
401        )
402        .unwrap();
403
404        let execute_certificate_latency_single_writer =
405            execute_certificate_latency.with_label_values(&[TX_TYPE_SINGLE_WRITER_TX]);
406        let execute_certificate_latency_shared_object =
407            execute_certificate_latency.with_label_values(&[TX_TYPE_SHARED_OBJ_TX]);
408
409        Self {
410            tx_orders: register_int_counter_with_registry!(
411                "total_transaction_orders",
412                "Total number of transaction orders",
413                registry;
414                MetricLevel::Warn,
415            )
416                .unwrap(),
417            total_certs: register_int_counter_with_registry!(
418                "total_transaction_certificates",
419                "Total number of transaction certificates handled",
420                registry;
421                MetricLevel::Warn,
422            )
423                .unwrap(),
424            total_cert_attempts: register_int_counter_with_registry!(
425                "total_handle_certificate_attempts",
426                "Number of calls to handle_certificate",
427                registry,
428            )
429                .unwrap(),
430            // total_effects == total transactions finished
431            total_effects: register_int_counter_with_registry!(
432                "total_transaction_effects",
433                "Total number of transaction effects produced",
434                registry;
435                MetricLevel::Warn,
436            )
437                .unwrap(),
438
439            shared_obj_tx: register_int_counter_with_registry!(
440                "num_shared_obj_tx",
441                "Number of transactions involving shared objects",
442                registry;
443                MetricLevel::Warn,
444            )
445                .unwrap(),
446
447            sponsored_tx: register_int_counter_with_registry!(
448                "num_sponsored_tx",
449                "Number of sponsored transactions",
450                registry,
451            )
452                .unwrap(),
453
454            tx_already_processed: register_int_counter_with_registry!(
455                "num_tx_already_processed",
456                "Number of transaction orders already processed previously",
457                registry,
458            )
459                .unwrap(),
460            num_input_objs: register_histogram_with_registry!(
461                "num_input_objects",
462                "Distribution of number of input TX objects per TX",
463                POSITIVE_INT_BUCKETS.to_vec(),
464                registry;
465                MetricLevel::Warn,
466            )
467                .unwrap(),
468            num_shared_objects: register_histogram_with_registry!(
469                "num_shared_objects",
470                "Number of shared input objects per TX",
471                POSITIVE_INT_BUCKETS.to_vec(),
472                registry,
473            )
474                .unwrap(),
475            batch_size: register_histogram_with_registry!(
476                "batch_size",
477                "Distribution of size of transaction batch",
478                POSITIVE_INT_BUCKETS.to_vec(),
479                registry,
480            )
481                .unwrap(),
482            authority_state_handle_transaction_latency: register_histogram_with_registry!(
483                "authority_state_handle_transaction_latency",
484                "Latency of handling transactions",
485                LATENCY_SEC_BUCKETS.to_vec(),
486                registry,
487            )
488                .unwrap(),
489            execute_certificate_latency_single_writer,
490            execute_certificate_latency_shared_object,
491            internal_execution_latency: register_histogram_with_registry!(
492                "authority_state_internal_execution_latency",
493                "Latency of actual certificate executions",
494                LATENCY_SEC_BUCKETS.to_vec(),
495                registry,
496            )
497                .unwrap(),
498            signed_effects_equivocation_prevented: register_int_counter_vec_with_registry!(
499                "authority_state_signed_effects_equivocation_prevented",
500                "Number of times the validator refused to report effects that differ from previously signed effects for the same transaction, by RPC surface",
501                &["surface"],
502                registry,
503            )
504            .unwrap(),
505            execution_load_input_objects_latency: register_histogram_with_registry!(
506                "authority_state_execution_load_input_objects_latency",
507                "Latency of loading input objects for execution",
508                LOW_LATENCY_SEC_BUCKETS.to_vec(),
509                registry,
510            )
511                .unwrap(),
512            prepare_certificate_latency: register_histogram_with_registry!(
513                "authority_state_prepare_certificate_latency",
514                "Latency of executing certificates, before committing the results",
515                LATENCY_SEC_BUCKETS.to_vec(),
516                registry,
517            )
518                .unwrap(),
519            commit_certificate_latency: register_histogram_with_registry!(
520                "authority_state_commit_certificate_latency",
521                "Latency of committing certificate execution results",
522                LATENCY_SEC_BUCKETS.to_vec(),
523                registry,
524            )
525                .unwrap(),
526            epoch_end_db_snapshot_handover_latency: register_histogram_with_registry!(
527                "epoch_end_db_snapshot_handover_latency",
528                "Time the epoch boundary waits for the consumer to take its database \
529                 snapshot of the perpetual store",
530                LATENCY_SEC_BUCKETS.to_vec(),
531                registry,
532            ).unwrap(),
533            epoch_end_db_snapshots_skipped: register_int_counter_with_registry!(
534                "epoch_end_db_snapshots_skipped",
535                "Epoch boundaries whose database snapshot was not taken, because the \
536                 consumer was busy or did not take it in time",
537                registry,
538            ).unwrap(),
539            transaction_manager_num_enqueued_certificates: register_int_counter_vec_with_registry!(
540                "transaction_manager_num_enqueued_certificates",
541                "Current number of certificates enqueued to TransactionManager",
542                &["result"],
543                registry,
544            )
545                .unwrap(),
546            transaction_manager_num_missing_objects: register_int_gauge_with_registry!(
547                "transaction_manager_num_missing_objects",
548                "Current number of missing objects in TransactionManager",
549                registry,
550            )
551                .unwrap(),
552            transaction_manager_num_pending_certificates: register_int_gauge_with_registry!(
553                "transaction_manager_num_pending_certificates",
554                "Number of certificates pending in TransactionManager, with at least 1 missing input object",
555                registry;
556                MetricLevel::Warn,
557            )
558                .unwrap(),
559            transaction_manager_num_executing_certificates: register_int_gauge_with_registry!(
560                "transaction_manager_num_executing_certificates",
561                "Number of executing certificates, including queued and actually running certificates",
562                registry;
563                MetricLevel::Warn,
564            )
565                .unwrap(),
566            transaction_manager_num_ready: register_int_gauge_with_registry!(
567                "transaction_manager_num_ready",
568                "Number of ready transactions in TransactionManager",
569                registry,
570            )
571                .unwrap(),
572            transaction_manager_object_cache_size: register_int_gauge_with_registry!(
573                "transaction_manager_object_cache_size",
574                "Current size of object-availability cache in TransactionManager",
575                registry,
576            )
577                .unwrap(),
578            transaction_manager_object_cache_hits: register_int_counter_with_registry!(
579                "transaction_manager_object_cache_hits",
580                "Number of object-availability cache hits in TransactionManager",
581                registry,
582            )
583                .unwrap(),
584            authority_overload_status: register_int_gauge_with_registry!(
585                "authority_overload_status",
586                "Whether authority is current experiencing overload and enters load shedding mode.",
587                registry;
588                MetricLevel::Warn,)
589                .unwrap(),
590            local_post_consensus_load_shedding_percentage: register_int_gauge_with_registry!(
591                "authority_load_shedding_percentage",
592                "This authority's locally computed load shedding percentage. In the P-COOL flow this is the value broadcast to peers, not necessarily the rate enforced (see consensus_handler_load_shedding_percentage).",
593                registry;
594                MetricLevel::Info,)
595                .unwrap(),
596            consensus_queue_load_shedding_percentage: register_int_gauge_with_registry!(
597                "consensus_queue_load_shedding_percentage",
598                "Percentage of transactions shed due to consensus queue length. Separate admission-control signal, not an input to authority_load_shedding_percentage.",
599                registry)
600                .unwrap(),
601            transaction_manager_object_cache_misses: register_int_counter_with_registry!(
602                "transaction_manager_object_cache_misses",
603                "Number of object-availability cache misses in TransactionManager",
604                registry,
605            )
606                .unwrap(),
607            transaction_manager_object_cache_evictions: register_int_counter_with_registry!(
608                "transaction_manager_object_cache_evictions",
609                "Number of object-availability cache evictions in TransactionManager",
610                registry,
611            )
612                .unwrap(),
613            transaction_manager_package_cache_size: register_int_gauge_with_registry!(
614                "transaction_manager_package_cache_size",
615                "Current size of package-availability cache in TransactionManager",
616                registry,
617            )
618                .unwrap(),
619            transaction_manager_package_cache_hits: register_int_counter_with_registry!(
620                "transaction_manager_package_cache_hits",
621                "Number of package-availability cache hits in TransactionManager",
622                registry,
623            )
624                .unwrap(),
625            transaction_manager_package_cache_misses: register_int_counter_with_registry!(
626                "transaction_manager_package_cache_misses",
627                "Number of package-availability cache misses in TransactionManager",
628                registry,
629            )
630                .unwrap(),
631            transaction_manager_package_cache_evictions: register_int_counter_with_registry!(
632                "transaction_manager_package_cache_evictions",
633                "Number of package-availability cache evictions in TransactionManager",
634                registry,
635            )
636                .unwrap(),
637            transaction_manager_transaction_queue_age_s: register_histogram_with_registry!(
638                "transaction_manager_transaction_queue_age_s",
639                "Time spent in waiting for transaction in the queue",
640                LATENCY_SEC_BUCKETS.to_vec(),
641                registry;
642                MetricLevel::Warn,
643            )
644                .unwrap(),
645            transaction_overload_sources: register_int_counter_vec_with_registry!(
646                "transaction_overload_sources",
647                "Number of times each source indicates transaction overload.",
648                &["source"],
649                registry)
650                .unwrap(),
651            execution_driver_executed_transactions: register_int_counter_with_registry!(
652                "execution_driver_executed_transactions",
653                "Cumulative number of transaction executed by execution driver",
654                registry;
655                MetricLevel::Warn,
656            )
657                .unwrap(),
658            execution_driver_dispatch_queue: register_int_gauge_with_registry!(
659                "execution_driver_dispatch_queue",
660                "Number of transaction pending in execution driver dispatch queue",
661                registry,
662            )
663                .unwrap(),
664            execution_queueing_delay_s: register_histogram_with_registry!(
665                "execution_queueing_delay_s",
666                "Queueing delay between a transaction is ready for execution until it starts executing.",
667                LATENCY_SEC_BUCKETS.to_vec(),
668                registry
669            )
670                .unwrap(),
671            prepare_cert_gas_latency_ratio: register_histogram_with_registry!(
672                "prepare_cert_gas_latency_ratio",
673                "The ratio of computation gas divided by VM execution latency.",
674                GAS_LATENCY_RATIO_BUCKETS.to_vec(),
675                registry
676            )
677                .unwrap(),
678            execution_gas_latency_ratio: register_histogram_with_registry!(
679                "execution_gas_latency_ratio",
680                "The ratio of computation gas divided by certificate execution latency, include committing certificate.",
681                GAS_LATENCY_RATIO_BUCKETS.to_vec(),
682                registry
683            )
684                .unwrap(),
685            skipped_consensus_txns: register_int_counter_with_registry!(
686                "skipped_consensus_txns",
687                "Total number of consensus transactions skipped",
688                registry,
689            )
690                .unwrap(),
691            skipped_consensus_txns_cache_hit: register_int_counter_with_registry!(
692                "skipped_consensus_txns_cache_hit",
693                "Total number of consensus transactions skipped because of local cache hit",
694                registry,
695            )
696                .unwrap(),
697            post_processing_total_events_emitted: register_int_counter_with_registry!(
698                "post_processing_total_events_emitted",
699                "Total number of events emitted in post processing",
700                registry,
701            )
702                .unwrap(),
703            post_processing_total_tx_indexed: register_int_counter_with_registry!(
704                "post_processing_total_tx_indexed",
705                "Total number of txes indexed in post processing",
706                registry,
707            )
708                .unwrap(),
709            post_processing_total_tx_had_event_processed: register_int_counter_with_registry!(
710                "post_processing_total_tx_had_event_processed",
711                "Total number of txes finished event processing in post processing",
712                registry,
713            )
714                .unwrap(),
715            post_processing_total_failures: register_int_counter_with_registry!(
716                "post_processing_total_failures",
717                "Total number of failure in post processing",
718                registry,
719            )
720                .unwrap(),
721            consensus_handler_processed: register_int_counter_vec_with_registry!(
722                "consensus_handler_processed",
723                "Number of transactions processed by consensus handler",
724                &["class"],
725                registry
726            ).unwrap(),
727            consensus_handler_transaction_sizes: register_histogram_vec_with_registry!(
728                "consensus_handler_transaction_sizes",
729                "Sizes of each type of transactions processed by consensus handler",
730                &["class"],
731                POSITIVE_INT_BUCKETS.to_vec(),
732                registry;
733                MetricLevel::Warn,
734            ).unwrap(),
735            consensus_handler_num_low_scoring_authorities: register_int_gauge_with_registry!(
736                "consensus_handler_num_low_scoring_authorities",
737                "Number of low scoring authorities based on reputation scores from consensus",
738                registry
739            ).unwrap(),
740            consensus_handler_scores: register_int_gauge_vec_with_registry!(
741                "consensus_handler_scores",
742                "scores from consensus for each authority",
743                &["authority"],
744                registry,
745            ).unwrap(),
746            validator_scoreboard_scores: register_int_gauge_vec_with_registry!(
747                "validator_scoreboard_scores",
748                "Per-authority validator scores published by the local Scoreboard after each consensus commit. Range [0, MAX_SCORE].",
749                &["authority"],
750                registry;
751                MetricLevel::Warn,
752            ).unwrap(),
753            invalid_misbehavior_reports_by_authority: register_int_gauge_vec_with_registry!(
754                "invalid_misbehavior_reports_by_authority",
755                "Cumulative count of invalid misbehavior reports received from each reporting authority in the current epoch. Bumped when a `MisbehaviorReport` consensus transaction fails sender/authority match or payload validation. Snapshot republished after each consensus commit.",
756                &["authority"],
757                registry;
758                MetricLevel::Warn,
759            ).unwrap(),
760            consensus_handler_deferred_transactions: register_int_counter_with_registry!(
761                "consensus_handler_deferred_transactions",
762                "Number of transactions deferred by consensus handler",
763                registry,
764            ).unwrap(),
765            consensus_handler_congested_transactions: register_int_counter_with_registry!(
766                "consensus_handler_congested_transactions",
767                "Number of transactions deferred by consensus handler due to congestion",
768                registry,
769            ).unwrap(),
770            consensus_handler_cancelled_transactions: register_int_counter_with_registry!(
771                "consensus_handler_cancelled_transactions",
772                "Number of transactions cancelled by consensus handler",
773                registry,
774            ).unwrap(),
775            consensus_handler_validation_dropped_transactions: register_int_counter_with_registry!(
776                "consensus_handler_validation_dropped_transactions",
777                "Number of UserTransactionV1 transactions dropped by post-consensus validation",
778                registry,
779            ).unwrap(),
780            consensus_handler_load_shedding_dropped_transactions: register_int_counter_with_registry!(
781                "consensus_handler_load_shedding_dropped_transactions",
782                "Number of user transactions dropped by post-consensus load shedding, based on the quorum load shedding percentage",
783                registry,
784            ).unwrap(),
785            consensus_handler_load_shedding_percentage: register_int_gauge_with_registry!(
786                "consensus_handler_load_shedding_percentage",
787                "Stake-weighted quorum (2f+1) load shedding percentage enforced on user transactions in the most recent consensus commit. 0 when the P-COOL flow is disabled.",
788                registry,
789            ).unwrap(),
790            consensus_handler_max_object_costs: register_int_gauge_vec_with_registry!(
791                "consensus_handler_max_congestion_control_object_costs",
792                "Max object costs for congestion control in the current consensus commit",
793                &["commit_type"],
794                registry,
795            ).unwrap(),
796            consensus_committed_subdags: register_int_counter_vec_with_registry!(
797                "consensus_committed_subdags",
798                "Number of committed subdags, sliced by leader",
799                &["authority"],
800                registry,
801            ).unwrap(),
802            consensus_committed_messages: register_int_gauge_vec_with_registry!(
803                "consensus_committed_messages",
804                "Total number of committed consensus messages, sliced by author",
805                &["authority"],
806                registry;
807                MetricLevel::Warn,
808            ).unwrap(),
809            consensus_committed_user_transactions: register_int_gauge_vec_with_registry!(
810                "consensus_committed_user_transactions",
811                "Number of committed user transactions, sliced by submitter",
812                &["authority"],
813                registry,
814            ).unwrap(),
815            consensus_handler_leader_round: register_int_gauge_with_registry!(
816                "consensus_handler_leader_round",
817                "The leader round of the current consensus output being processed in the consensus handler",
818                registry;
819                MetricLevel::Warn,
820            ).unwrap(),
821            limits_metrics: Arc::new(LimitsMetrics::new(registry)),
822            bytecode_verifier_metrics: Arc::new(BytecodeVerifierMetrics::new(registry)),
823            multisig_sig_count: register_int_counter_with_registry!(
824                "multisig_sig_count",
825                "Count of multisig signatures",
826                registry,
827            )
828                .unwrap(),
829            consensus_calculated_throughput: register_int_gauge_with_registry!(
830                "consensus_calculated_throughput",
831                "The calculated throughput from consensus output. Result is calculated based on unique transactions.",
832                registry,
833            ).unwrap(),
834            consensus_calculated_throughput_profile: register_int_gauge_with_registry!(
835                "consensus_calculated_throughput_profile",
836                "The current active calculated throughput profile",
837                registry
838            ).unwrap(),
839            execution_queueing_latency: LatencyObserver::new(),
840            txn_ready_rate_tracker: Arc::new(Mutex::new(RateTracker::new(Duration::from_secs(10)))),
841            execution_rate_tracker: Arc::new(Mutex::new(RateTracker::new(Duration::from_secs(10)))),
842        }
843    }
844
845    /// Reset metrics that contain `hostname` as one of the labels. This is
846    /// needed to avoid retaining metrics for long-gone committee members and
847    /// only exposing metrics for the committee in the current epoch.
848    pub fn reset_on_reconfigure(&self) {
849        self.consensus_committed_messages.reset();
850        self.consensus_handler_scores.reset();
851        self.validator_scoreboard_scores.reset();
852        self.invalid_misbehavior_reports_by_authority.reset();
853        self.consensus_committed_user_transactions.reset();
854    }
855}
856
857/// a Trait object for `Signer` that is:
858/// - Pin, i.e. confined to one place in memory (we don't want to copy private keys).
859/// - Sync, i.e. can be safely shared between threads.
860///
861/// Typically instantiated with Box::pin(keypair) where keypair is a `KeyPair`
862pub type StableSyncAuthoritySigner = Pin<Arc<dyn Signer<AuthoritySignature> + Send + Sync>>;
863
864/// Execution env contains the "environment" for the transaction to be executed
865/// in, that is, all the information necessary for execution that is not
866/// specified by the transaction itself.
867#[derive(Debug, Clone, Default)]
868pub struct ExecutionEnv {
869    /// The assigned version of each shared object for the transaction.
870    pub assigned_versions: AssignedVersions,
871    /// The expected digest of the effects of the transaction, if executing from
872    /// checkpoint or other sources where the effects are known in advance.
873    pub expected_effects_digest: Option<TransactionEffectsDigest>,
874}
875
876impl ExecutionEnv {
877    pub fn new() -> Self {
878        Default::default()
879    }
880
881    pub fn with_expected_effects_digest(
882        mut self,
883        expected_effects_digest: TransactionEffectsDigest,
884    ) -> Self {
885        self.expected_effects_digest = Some(expected_effects_digest);
886        self
887    }
888
889    pub fn with_assigned_versions(mut self, assigned_versions: AssignedVersions) -> Self {
890        if !assigned_versions.is_empty() {
891            self.assigned_versions = assigned_versions;
892        }
893        self
894    }
895}
896
897pub struct AuthorityState {
898    // Fixed size, static, identity of the authority
899    /// The name of this authority.
900    pub name: AuthorityName,
901    /// The signature key of the authority.
902    pub secret: StableSyncAuthoritySigner,
903
904    /// The database
905    input_loader: TransactionInputLoader,
906    execution_cache_trait_pointers: ExecutionCacheTraitPointers,
907
908    epoch_store: ArcSwap<AuthorityPerEpochStore>,
909
910    /// This lock denotes current 'execution epoch'.
911    /// Execution acquires read lock, checks transaction epoch and holds it
912    /// until all writes are complete. Reconfiguration acquires write lock,
913    /// changes the epoch and revert all transactions from previous epoch
914    /// that are executed but did not make into checkpoint.
915    execution_lock: RwLock<EpochId>,
916
917    pub indexes: Option<Arc<IndexStore>>,
918    pub grpc_indexes_store: Option<Arc<GrpcIndexesStore>>,
919
920    pub subscription_handler: Arc<SubscriptionHandler>,
921    pub checkpoint_store: Arc<CheckpointStore>,
922
923    committee_store: Arc<CommitteeStore>,
924
925    /// Schedules transaction execution.
926    execution_scheduler: Arc<ExecutionSchedulerWrapper>,
927
928    /// Shuts down the execution task. Used only in testing.
929    #[cfg_attr(not(test), expect(unused))]
930    tx_execution_shutdown: Mutex<Option<oneshot::Sender<()>>>,
931
932    pub metrics: Arc<AuthorityMetrics>,
933    /// The store pruner. The checkpoint executor uses it to nudge the pruner
934    /// after each checkpoint.
935    pruner: AuthorityStorePruner,
936    authority_per_epoch_pruner: AuthorityPerEpochStorePruner,
937    checkpoint_progress_tracker: Option<Arc<CheckpointProgressTracker>>,
938
939    pub config: NodeConfig,
940
941    /// Current overload status in this authority. Updated periodically.
942    pub overload_info: AuthorityOverloadInfo,
943
944    pub validator_tx_finalizer: Option<Arc<ValidatorTxFinalizer<NetworkAuthorityClient>>>,
945
946    /// The chain identifier is derived from the digest of the genesis
947    /// checkpoint.
948    chain_identifier: ChainIdentifier,
949
950    pub(crate) congestion_tracker: Arc<CongestionTracker>,
951
952    /// Traffic controller for IOTA core servers (json-rpc, validator service)
953    pub traffic_controller: Option<Arc<TrafficController>>,
954    /// Set when a consumer wants a database snapshot of the perpetual store at
955    /// each epoch boundary. See [`Self::hand_over_epoch_end_db_snapshot`].
956    epoch_end_db_snapshots: Option<EpochEndDbSnapshotHandle>,
957}
958
959/// The authority state encapsulates all state, drives execution, and ensures
960/// safety.
961///
962/// Note the authority operations can be accessed through a read ref (&) and do
963/// not require &mut. Internally a database is synchronized through a mutex
964/// lock.
965///
966/// Repeating valid commands should produce no changes and return no error.
967impl AuthorityState {
968    pub fn is_committee_validator(&self, epoch_store: &AuthorityPerEpochStore) -> bool {
969        epoch_store.committee().authority_exists(&self.name)
970    }
971
972    pub fn is_active_validator(&self, epoch_store: &AuthorityPerEpochStore) -> bool {
973        epoch_store
974            .active_validators()
975            .iter()
976            .any(|a| AuthorityName::from(a) == self.name)
977    }
978
979    pub fn is_fullnode(&self, epoch_store: &AuthorityPerEpochStore) -> bool {
980        !self.is_committee_validator(epoch_store)
981    }
982
983    pub fn committee_store(&self) -> &Arc<CommitteeStore> {
984        &self.committee_store
985    }
986
987    pub fn clone_committee_store(&self) -> Arc<CommitteeStore> {
988        self.committee_store.clone()
989    }
990
991    pub fn overload_config(&self) -> &AuthorityOverloadConfig {
992        &self.config.authority_overload_config
993    }
994
995    pub fn get_epoch_state_commitments(
996        &self,
997        epoch: EpochId,
998    ) -> IotaResult<Option<Vec<CheckpointCommitment>>> {
999        self.checkpoint_store.get_epoch_state_commitments(epoch)
1000    }
1001
1002    /// Runs deny list, input object validation, gas checks, coin deny list, and
1003    /// MoveAuthenticator checks. Returns the owned object refs for optional
1004    /// version validation. Does NOT acquire locks or sign the transaction.
1005    ///
1006    /// `deny_config` is the deny rule source to enforce, chosen per caller:
1007    /// the local config alone, the local config combined with the governance
1008    /// rules (admission), or the governance-derived active set alone
1009    /// (post-consensus) — the latter two when `deny_rule_governance` is
1010    /// enabled.
1011    ///
1012    /// `epoch_gated_coin_deny_list` selects how the coin deny list is read:
1013    /// `false` reads the latest value, so denials apply immediately - for
1014    /// validator-local admission (signing); `true` reads the value settled
1015    /// before the current epoch, which is deterministic across validators
1016    /// regardless of each validator's execution progress - required
1017    /// post-consensus, where the verdict decides whether the transaction
1018    /// stays in the committed set. The two read modes intentionally disagree
1019    /// about deny-list changes made in the current epoch, in both directions:
1020    /// - An entry added this epoch is enforced at admission right away, while the epoch-gated
1021    ///   layers enforce it only from the next epoch. Since execution and post-consensus must read
1022    ///   epoch-gated to stay deterministic, admission is the only layer that can react to a new
1023    ///   denial or global pause before the epoch boundary.
1024    /// - An entry removed this epoch is admitted right away but still denied by the epoch-gated
1025    ///   post-consensus read, so such transactions are sequenced by consensus and then
1026    ///   deterministically dropped (no execution, no gas charged) until the removal settles at the
1027    ///   next epoch boundary. The wasted consensus slot is accepted: post-consensus must handle
1028    ///   deterministic drops regardless (owned-object double-spend losers, for example), and
1029    ///   validators that skip admission can put such transactions into their blocks anyway, so no
1030    ///   admission policy can limit how many deterministically-dropped transactions reach
1031    ///   consensus.
1032    ///
1033    /// `verifier_limits_source` says where the metered bytecode verifier takes
1034    /// its limits for the packages the transaction publishes. Validator-local
1035    /// admission passes this validator's own `VerifierSigningConfig`;
1036    /// post-consensus validation passes the protocol config, because every
1037    /// validator must reach the same verdict there.
1038    #[instrument(level = "trace", skip_all, fields(tx_digest = ?transaction.digest()))]
1039    pub(crate) async fn handle_transaction_validation_checks(
1040        &self,
1041        transaction: &VerifiedTransaction,
1042        epoch_store: &Arc<AuthorityPerEpochStore>,
1043        deny_config: &dyn DenyRuleConfig,
1044        epoch_gated_coin_deny_list: bool,
1045        verifier_limits_source: VerifierLimitsSource<'_>,
1046    ) -> IotaResult<Vec<ObjectReference>> {
1047        let protocol_config = epoch_store.protocol_config();
1048        let reference_gas_price = epoch_store.reference_gas_price();
1049
1050        let epoch = epoch_store.epoch();
1051
1052        let tx = transaction.data().transaction();
1053
1054        // Note: the deny checks may do redundant package loads but:
1055        // - they only load packages when there is an active package deny map
1056        // - the loads are cached anyway
1057        iota_transaction_checks::deny::check_transaction_for_validation(
1058            tx,
1059            transaction.signatures(),
1060            &transaction.input_objects()?,
1061            &tx.receiving_objects(),
1062            deny_config,
1063            self.get_backing_package_store().as_ref(),
1064        )?;
1065
1066        // Load all transaction-related input objects including ones for every
1067        // `MoveAuthenticator`. Loading all objects eagerly means that any invalid
1068        // reference — missing object, wrong version, inaccessible object — causes a
1069        // pre-consensus rejection.
1070        let (tx_input_objects, tx_receiving_objects, per_authenticator_inputs) =
1071            self.read_objects_for_validation(transaction, epoch)?;
1072
1073        let move_authenticators = transaction.move_authenticators();
1074
1075        // Check the inputs for signing.
1076        // If there are `MoveAuthenticator` signatures, their input objects and the
1077        // account objects are also checked and must be provided.
1078        // It is also checked if there is enough gas to execute the transaction and its
1079        // authenticators.
1080        let (gas_status, tx_checked_input_objects, per_authenticator_checked_inputs) = self
1081            .check_transaction_inputs_for_validation(
1082                protocol_config,
1083                reference_gas_price,
1084                tx,
1085                tx_input_objects,
1086                &tx_receiving_objects,
1087                &move_authenticators,
1088                per_authenticator_inputs,
1089                verifier_limits_source,
1090            )?;
1091
1092        // Get the input objects for the authenticators, if there are
1093        // `MoveAuthenticator`s.
1094        let per_authenticator_checked_input_objects: Vec<_> = per_authenticator_checked_inputs
1095            .iter()
1096            .map(|i| &i.0)
1097            .collect();
1098
1099        // Move authenticators cannot use owned objects, so their inputs never
1100        // acquire owned-object locks.
1101        debug_assert!(
1102            per_authenticator_checked_input_objects
1103                .iter()
1104                .all(|objects| objects.inner().filter_owned_objects().is_empty()),
1105            "Move authenticator input objects must not contain owned objects"
1106        );
1107
1108        // The package holding each authenticate function is known only now that the
1109        // `AuthenticatorFunctionRef`s are loaded, so the deny-list check for it
1110        // stands apart from the one above. It must stay ahead of two things: the
1111        // filtering below, which drops authenticators that do not run
1112        // pre-consensus, so that every authenticator is covered; and the
1113        // authenticator execution itself, so that a denied package is never run.
1114        if protocol_config.deny_authenticator_packages() {
1115            iota_transaction_checks::deny::check_authenticator_packages(
1116                deny_config,
1117                per_authenticator_checked_inputs
1118                    .iter()
1119                    .map(|(_, authenticator_function_ref)| authenticator_function_ref),
1120                self.get_backing_package_store().as_ref(),
1121            )?;
1122        }
1123
1124        // Check if any of the sender, the transaction input objects, the receiving
1125        // objects and the authenticator input objects are in the coin deny
1126        // list, which would prevent the transaction from being signed.
1127        check_coin_deny_list_v1(
1128            tx.sender(),
1129            &tx_checked_input_objects,
1130            &tx_receiving_objects,
1131            &per_authenticator_checked_input_objects,
1132            &self.get_object_store(),
1133            epoch_gated_coin_deny_list.then_some(epoch),
1134        )?;
1135
1136        let (kind, signer, gas_data) = tx.execution_parts();
1137
1138        let (sender_authenticator_function_ref, sponsor_authenticator_function_ref) =
1139            extract_auth_fun_refs(signer, gas_data.owner, |address| {
1140                move_authenticators
1141                    .iter()
1142                    .zip(per_authenticator_checked_inputs.iter())
1143                    .find(|(move_authenticator, _)| move_authenticator.address() == address)
1144                    .map(|(_, (_, auth_fun_ref))| auth_fun_ref.clone())
1145            });
1146
1147        // Filter the authenticators and their checked inputs down to those that must
1148        // be executed pre-consensus. This is done *after* the deny-list check so
1149        // that all MoveAuthenticator input objects are covered by that check regardless
1150        // of deferral.
1151        let pre_consensus_move_authenticators =
1152            pre_consensus_move_authenticators(transaction, protocol_config);
1153        // Asserted before the zip below pairs them positionally; the two lists
1154        // come from independent computations.
1155        debug_assert_eq!(
1156            move_authenticators.len(),
1157            per_authenticator_checked_inputs.len(),
1158            "Move authenticators amount must match the number of checked authenticator inputs"
1159        );
1160        let (move_authenticators, per_authenticator_checked_inputs): (Vec<_>, Vec<_>) =
1161            move_authenticators
1162                .into_iter()
1163                .zip(per_authenticator_checked_inputs)
1164                .filter(|(a, _)| pre_consensus_move_authenticators.contains(a))
1165                .unzip();
1166        let per_authenticator_checked_input_objects: Vec<_> = per_authenticator_checked_inputs
1167            .iter()
1168            .map(|i| &i.0)
1169            .collect();
1170
1171        // If there are `MoveAuthenticator` signatures, execute them and check if they
1172        // all succeed.
1173        if !move_authenticators.is_empty() {
1174            let aggregated_authenticator_input_objects =
1175                iota_transaction_checks::aggregate_authenticator_input_objects(
1176                    &per_authenticator_checked_input_objects,
1177                )?;
1178
1179            let move_authenticators = move_authenticators
1180                .into_iter()
1181                .zip(per_authenticator_checked_inputs)
1182                .map(
1183                    |(
1184                        move_authenticator,
1185                        (authenticator_checked_input_objects, authenticator_function_ref),
1186                    )| {
1187                        (
1188                            move_authenticator.to_owned(),
1189                            authenticator_function_ref,
1190                            authenticator_checked_input_objects,
1191                        )
1192                    },
1193                )
1194                .collect();
1195
1196            // It is supposed that `MoveAuthenticator` availability is checked in
1197            // `SenderSignedTransaction::validity_check`.
1198
1199            // Serialize the Transaction for the auth context before decomposing.
1200            let tx_bytes = bcs::to_bytes(tx).expect("Transaction serialization cannot fail");
1201
1202            let (sender_auth_digest, sponsor_auth_digest) =
1203                transaction.data().compute_auth_digests()?;
1204
1205            let auth_context_data = AuthContextData {
1206                transaction_data_bytes: tx_bytes,
1207                sender_auth_digest,
1208                sponsor_auth_digest,
1209                sender_authenticator_function_ref,
1210                sponsor_authenticator_function_ref,
1211            };
1212
1213            // Execute the Move authenticators.
1214            let validation_result = epoch_store.executor().authenticate_transaction(
1215                self.get_backing_store().as_ref(),
1216                protocol_config,
1217                self.metrics.limits_metrics.clone(),
1218                &epoch_store.epoch_start_config().epoch_data().epoch_id(),
1219                epoch_store
1220                    .epoch_start_config()
1221                    .epoch_data()
1222                    .epoch_start_timestamp(),
1223                gas_data,
1224                gas_status,
1225                move_authenticators,
1226                aggregated_authenticator_input_objects,
1227                kind,
1228                signer,
1229                transaction.digest().to_owned(),
1230                auth_context_data,
1231                &mut None,
1232            );
1233
1234            if let Err(validation_error) = validation_result {
1235                return Err(IotaError::MoveAuthenticatorExecutionFailure {
1236                    error: validation_error.to_string(),
1237                });
1238            }
1239        }
1240
1241        Ok(tx_checked_input_objects.inner().filter_owned_objects())
1242    }
1243
1244    /// This is a private method and should be kept that way. It doesn't check
1245    /// whether the provided transaction is a system transaction, and hence
1246    /// can only be called internally.
1247    async fn handle_transaction_impl(
1248        &self,
1249        transaction: VerifiedTransaction,
1250        epoch_store: &Arc<AuthorityPerEpochStore>,
1251    ) -> IotaResult<VerifiedSignedTransaction> {
1252        // Ensure that validator cannot reconfigure while we are signing the tx
1253        let _execution_lock = self.execution_lock_for_signing()?;
1254
1255        let owned_objects = self
1256            .handle_transaction_validation_checks(
1257                &transaction,
1258                epoch_store,
1259                &self.config.transaction_deny_config,
1260                // Latest-value coin deny-list read: admission is validator-local,
1261                // and denials should take effect immediately. Unlike the P-COOL
1262                // submission path, no post-consensus re-check follows - this is
1263                // the only sender-side coin deny check in the certificate flow.
1264                false,
1265                VerifierLimitsSource::NodeConfig(&self.config.verifier_signing_config),
1266            )
1267            .await?;
1268
1269        let epoch = epoch_store.epoch();
1270        let signed_transaction =
1271            VerifiedSignedTransaction::new(epoch, transaction, self.name, &*self.secret);
1272
1273        // Check and write locks, to signed transaction, into the database
1274        // The call to self.set_transaction_lock checks the lock is not conflicting,
1275        // and returns ConflictingTransaction error in case there is a lock on a
1276        // different existing transaction.
1277        self.get_cache_writer().try_acquire_transaction_locks(
1278            epoch_store,
1279            &owned_objects,
1280            signed_transaction.clone(),
1281        )?;
1282
1283        Ok(signed_transaction)
1284    }
1285
1286    /// Initiate a new transaction.
1287    #[instrument(name = "handle_transaction", level = "trace", skip_all, fields(tx_digest = ?transaction.digest(), sender = transaction.data().transaction().gas_owner().to_string()
1288    ))]
1289    pub async fn handle_transaction(
1290        &self,
1291        epoch_store: &Arc<AuthorityPerEpochStore>,
1292        transaction: VerifiedTransaction,
1293    ) -> IotaResult<HandleTransactionResponse> {
1294        let tx_digest = *transaction.digest();
1295        debug!("handle_transaction");
1296
1297        // Ensure an idempotent answer.
1298        if let Some((_, status)) = self.get_transaction_status(&tx_digest, epoch_store)? {
1299            return Ok(HandleTransactionResponse { status });
1300        }
1301
1302        let _metrics_guard = self
1303            .metrics
1304            .authority_state_handle_transaction_latency
1305            .start_timer();
1306        self.metrics.tx_orders.inc();
1307
1308        let signed = self.handle_transaction_impl(transaction, epoch_store).await;
1309        match signed {
1310            Ok(s) => {
1311                if self.is_committee_validator(epoch_store) {
1312                    if let Some(validator_tx_finalizer) = &self.validator_tx_finalizer {
1313                        let tx = s.clone();
1314                        let validator_tx_finalizer = validator_tx_finalizer.clone();
1315                        let cache_reader = self.get_transaction_cache_reader().clone();
1316                        let epoch_store = epoch_store.clone();
1317                        spawn_monitored_task!(epoch_store.within_alive_epoch(
1318                            validator_tx_finalizer.track_signed_tx(cache_reader, &epoch_store, tx)
1319                        ));
1320                    }
1321                }
1322                Ok(HandleTransactionResponse {
1323                    status: TransactionStatus::Signed(s.into_inner().into_sig()),
1324                })
1325            }
1326            // It happens frequently that while we are checking the validity of the transaction, it
1327            // has just been executed.
1328            // In that case, we could still return Ok to avoid showing confusing errors.
1329            Err(err) => Ok(HandleTransactionResponse {
1330                status: self
1331                    .get_transaction_status(&tx_digest, epoch_store)?
1332                    .ok_or(err)?
1333                    .1,
1334            }),
1335        }
1336    }
1337
1338    pub fn check_system_overload_at_signing(&self) -> bool {
1339        self.config
1340            .authority_overload_config
1341            .check_system_overload_at_signing
1342    }
1343
1344    pub fn check_system_overload_at_execution(&self) -> bool {
1345        self.config
1346            .authority_overload_config
1347            .check_system_overload_at_execution
1348    }
1349
1350    /// Checks system overload conditions before accepting a transaction.
1351    ///
1352    /// In certificate-less (P-COOL) mode: only checks consensus
1353    /// queue overload, since execution-based overload will be handled
1354    /// post-consensus.
1355    ///
1356    /// In certificate mode: runs all checks — authority overload
1357    /// (execution latency), the execution scheduler (execution queue),
1358    /// consensus adapter (queue limit), and writeback cache backpressure.
1359    pub(crate) fn check_system_overload(
1360        &self,
1361        consensus_adapter: &Arc<ConsensusAdapter>,
1362        tx: &SenderSignedTransaction,
1363        do_authority_overload_check: bool,
1364        pcool_flow_enabled: bool,
1365    ) -> IotaResult {
1366        if pcool_flow_enabled {
1367            // Graduated shedding: 0% to 100% as consensus queue fills from soft
1368            // to hard limit.
1369            self.check_consensus_queue_graduated_limits(consensus_adapter, tx)
1370                .tap_err(|_| {
1371                    self.update_overload_metrics("consensus");
1372                })?;
1373
1374            // NOTE: graduated shedding at 100% already rejects everything at or above
1375            // `max_pending_transactions`, so the queue-length part of the check below
1376            // is redundant but harmless. But `check_consensus_overload()` should be
1377            // kept here because it also verifies that `submit_semaphore` has permits
1378            // (see `check_consensus_hard_limits` in consensus_adapter.rs), which is a
1379            // separate concurrency limit not covered by the graduated shedding.
1380            consensus_adapter.check_consensus_overload().tap_err(|_| {
1381                self.update_overload_metrics("consensus");
1382            })?;
1383        } else {
1384            if do_authority_overload_check {
1385                self.check_authority_overload(tx).tap_err(|_| {
1386                    self.update_overload_metrics("execution_queue");
1387                })?;
1388            }
1389            self.execution_scheduler
1390                .check_execution_overload(self.overload_config(), tx)
1391                .tap_err(|_| {
1392                    self.update_overload_metrics("execution_pending");
1393                })?;
1394            consensus_adapter.check_consensus_overload().tap_err(|_| {
1395                self.update_overload_metrics("consensus");
1396            })?;
1397
1398            let pending_tx_count = self
1399                .get_cache_commit()
1400                .approximate_pending_transaction_count();
1401            if pending_tx_count
1402                > self
1403                    .config
1404                    .execution_cache_config
1405                    .writeback_cache
1406                    .backpressure_threshold_for_rpc()
1407            {
1408                return Err(IotaError::ValidatorOverloadedRetryAfter {
1409                    retry_after_secs: 10,
1410                });
1411            }
1412        }
1413
1414        Ok(())
1415    }
1416
1417    /// Rejects `tx_data` via graduated shedding based on consensus queue
1418    /// length. Scales from 0% at the soft limit to 100% at
1419    /// `max_pending_transactions`. Returns `ValidatorOverloadedRetryAfter`
1420    /// for probabilistic rejection (shedding percentage < 100%, via
1421    /// `overload_monitor_accept_tx`) or `TooManyTransactionsPendingConsensus`
1422    /// for unconditional rejection (shedding percentage >= 100%). Updates
1423    /// `consensus_queue_load_shedding_percentage` metric.
1424    fn check_consensus_queue_graduated_limits(
1425        &self,
1426        consensus_adapter: &Arc<ConsensusAdapter>,
1427        tx: &SenderSignedTransaction,
1428    ) -> IotaResult {
1429        let num_inflight_txs = consensus_adapter.num_inflight_transactions() as usize;
1430
1431        let shedding_pct = compute_graduated_load_shedding_percentage(
1432            num_inflight_txs,
1433            consensus_adapter.max_pending_transactions(),
1434            consensus_adapter.graduated_load_shedding_soft_limit_pct(),
1435        );
1436
1437        self.metrics
1438            .consensus_queue_load_shedding_percentage
1439            .set(shedding_pct as i64);
1440
1441        if shedding_pct == 0 {
1442            return Ok(());
1443        }
1444
1445        // At/above the hard limit, rejection is unconditional (not
1446        // probabilistic), so the seed-rotation retry hint of
1447        // `ValidatorOverloadedRetryAfter` doesn't apply - return the
1448        // capacity-bound error instead.
1449        if shedding_pct >= 100 {
1450            return Err(IotaError::TooManyTransactionsPendingConsensus);
1451        }
1452
1453        overload_monitor_accept_tx(shedding_pct, tx.digest())
1454    }
1455
1456    fn check_authority_overload(&self, tx: &SenderSignedTransaction) -> IotaResult {
1457        if !self.overload_info.is_overload.load(Ordering::Relaxed) {
1458            return Ok(());
1459        }
1460
1461        let load_shedding_percentage = self
1462            .overload_info
1463            .local_load_shedding_percentage
1464            .load(Ordering::Relaxed);
1465        overload_monitor_accept_tx(load_shedding_percentage, tx.digest())
1466    }
1467
1468    fn update_overload_metrics(&self, source: &str) {
1469        self.metrics
1470            .transaction_overload_sources
1471            .with_label_values(&[source])
1472            .inc();
1473    }
1474
1475    /// Wait for a certificate to be executed.
1476    /// For consensus transactions, it needs to be sequenced by the consensus.
1477    /// For owned object transactions, this function will enqueue the
1478    /// transaction for execution.
1479    #[instrument(level = "trace", skip_all)]
1480    pub async fn wait_for_certificate_execution(
1481        &self,
1482        certificate: &VerifiedCertificate,
1483        epoch_store: &Arc<AuthorityPerEpochStore>,
1484    ) -> IotaResult<TransactionEffects> {
1485        let _metrics_guard = if certificate.contains_shared_object() {
1486            self.metrics
1487                .execute_certificate_latency_shared_object
1488                .start_timer()
1489        } else {
1490            self.metrics
1491                .execute_certificate_latency_single_writer
1492                .start_timer()
1493        };
1494        trace!("wait_for_certificate_execution");
1495
1496        self.metrics.total_cert_attempts.inc();
1497
1498        if !certificate.contains_shared_object() {
1499            // Shared object transactions need to be sequenced by the consensus before
1500            // enqueueing for execution, done in
1501            // AuthorityPerEpochStore::handle_consensus_transaction(). For owned
1502            // object transactions, they can be enqueued for execution immediately.
1503            self.execution_scheduler.enqueue(
1504                vec![(
1505                    Schedulable::Transaction(VerifiedExecutableTransaction::new_from_certificate(
1506                        certificate.clone(),
1507                    )),
1508                    ExecutionEnv::new(),
1509                )],
1510                epoch_store,
1511            );
1512        }
1513
1514        // tx could be reverted when epoch ends, so we must be careful not to return a
1515        // result here after the epoch ends.
1516        epoch_store
1517            .within_alive_epoch(self.notify_read_effects(
1518                "AuthorityState::wait_for_certificate_execution",
1519                certificate,
1520            ))
1521            .await
1522            .and_then(|r| r)
1523    }
1524
1525    /// Internal logic to execute a transaction.
1526    ///
1527    /// Guarantees that
1528    /// - If input objects are available, return no permanent failure.
1529    /// - Execution and output commit are atomic. i.e. outputs are only written to storage,
1530    /// on successful execution; crashed execution has no observable effect and
1531    /// can be retried.
1532    ///
1533    /// It is caller's responsibility to ensure input objects are available and
1534    /// locks are set. If this cannot be satisfied by the caller,
1535    /// `wait_for_certificate_execution()` should be called instead.
1536    ///
1537    /// Should only be called within iota-core.
1538    #[instrument(level = "trace", skip_all, fields(tx_digest = ?transaction.digest()))]
1539    pub fn try_execute_immediately(
1540        &self,
1541        transaction: &VerifiedExecutableTransaction,
1542        execution_env: ExecutionEnv,
1543        epoch_store: &Arc<AuthorityPerEpochStore>,
1544    ) -> IotaResult<(TransactionEffects, Option<ExecutionError>)> {
1545        let _scope = monitored_scope("Execution::try_execute_immediately");
1546        let _metrics_guard = self.metrics.internal_execution_latency.start_timer();
1547
1548        let tx_digest = transaction.digest();
1549
1550        // Acquire a lock to prevent concurrent executions of the same transaction.
1551        let tx_guard = epoch_store.acquire_tx_guard(transaction)?;
1552
1553        // The transaction could have been processed by a concurrent attempt of the
1554        // same transaction, so check if the effects have already been written.
1555        if let Some(effects) = self
1556            .get_transaction_cache_reader()
1557            .try_get_executed_effects(tx_digest)?
1558        {
1559            if let Some(expected_effects_digest_inner) = execution_env.expected_effects_digest {
1560                assert_eq!(
1561                    effects.digest(),
1562                    expected_effects_digest_inner,
1563                    "Unexpected effects digest for transaction {tx_digest}"
1564                );
1565            }
1566            tx_guard.release();
1567            return Ok((effects, None));
1568        }
1569
1570        let (tx_input_objects, per_authenticator_inputs) = self.read_objects_for_execution(
1571            tx_guard.as_lock_guard(),
1572            transaction,
1573            execution_env.assigned_versions,
1574            epoch_store,
1575        )?;
1576
1577        self.process_transaction(
1578            tx_guard,
1579            transaction,
1580            tx_input_objects,
1581            per_authenticator_inputs,
1582            execution_env.expected_effects_digest,
1583            epoch_store,
1584        )
1585        .tap_err(|e| info!(?tx_digest, "process_transaction failed: {e}"))
1586        .tap_ok(
1587            |(fx, _)| debug!(?tx_digest, fx_digest=?fx.digest(), "process_transaction succeeded"),
1588        )
1589    }
1590
1591    pub fn read_objects_for_execution(
1592        &self,
1593        tx_lock: &TxLockGuard,
1594        transaction: &VerifiedExecutableTransaction,
1595        assigned_shared_object_versions: AssignedVersions,
1596        epoch_store: &Arc<AuthorityPerEpochStore>,
1597    ) -> IotaResult<(InputObjects, Vec<(InputObjects, ObjectReadResult)>)> {
1598        let _scope = monitored_scope("Execution::load_input_objects");
1599        let _metrics_guard = self
1600            .metrics
1601            .execution_load_input_objects_latency
1602            .start_timer();
1603
1604        let input_objects = transaction.collect_all_input_object_kind_for_reading()?;
1605
1606        let input_objects = self.input_loader.read_objects_for_execution(
1607            &transaction.key(),
1608            tx_lock,
1609            &input_objects,
1610            &assigned_shared_object_versions,
1611            epoch_store.epoch(),
1612        )?;
1613
1614        transaction.split_input_objects_into_groups_for_reading(input_objects)
1615    }
1616
1617    /// Test only wrapper for `try_execute_immediately()` above, useful for
1618    /// checking errors if the pre-conditions are not satisfied, and
1619    /// executing change epoch transactions.
1620    pub fn try_execute_for_test(
1621        &self,
1622        certificate: &VerifiedCertificate,
1623        execution_env: ExecutionEnv,
1624    ) -> IotaResult<(VerifiedSignedTransactionEffects, Option<ExecutionError>)> {
1625        let epoch_store = self.epoch_store_for_testing();
1626        let (effects, execution_error_opt) = self.try_execute_immediately(
1627            &VerifiedExecutableTransaction::new_from_certificate(certificate.clone()),
1628            execution_env,
1629            &epoch_store,
1630        )?;
1631        let signed_effects = self.sign_effects(effects, &epoch_store)?;
1632        Ok((signed_effects, execution_error_opt))
1633    }
1634
1635    /// Non-fallible version of `try_execute_for_test()`.
1636    pub fn execute_for_test(
1637        &self,
1638        certificate: &VerifiedCertificate,
1639        execution_env: ExecutionEnv,
1640    ) -> (VerifiedSignedTransactionEffects, Option<ExecutionError>) {
1641        self.try_execute_for_test(certificate, execution_env)
1642            .expect("try_execute_for_test should not fail")
1643    }
1644
1645    pub async fn notify_read_effects(
1646        &self,
1647        task_name: &'static str,
1648        certificate: &VerifiedCertificate,
1649    ) -> IotaResult<TransactionEffects> {
1650        self.get_transaction_cache_reader()
1651            .try_notify_read_executed_effects(task_name, &[*certificate.digest()])
1652            .await
1653            .map(|mut r| r.pop().expect("must return correct number of effects"))
1654    }
1655
1656    fn check_owned_locks(&self, owned_object_refs: &[ObjectReference]) -> IotaResult {
1657        self.get_object_cache_reader()
1658            .try_check_owned_objects_are_live(owned_object_refs)
1659    }
1660
1661    /// This function captures the required state to debug a forked transaction.
1662    /// The dump is written to a file in dir `path`, with name prefixed by the
1663    /// transaction digest. NOTE: Since this info escapes the validator
1664    /// context, make sure not to leak any private info here
1665    pub(crate) fn debug_dump_transaction_state(
1666        &self,
1667        tx_digest: &TransactionDigest,
1668        effects: &TransactionEffects,
1669        expected_effects_digest: TransactionEffectsDigest,
1670        inner_temporary_store: &InnerTemporaryStore,
1671        transaction: &VerifiedExecutableTransaction,
1672        debug_dump_config: &StateDebugDumpConfig,
1673    ) -> IotaResult<PathBuf> {
1674        // Fall back to the OS temp directory if no dump directory is configured.
1675        // This is safe: dump files are named by transaction digest, so no collisions.
1676        let dump_dir = debug_dump_config
1677            .dump_file_directory
1678            .as_ref()
1679            .cloned()
1680            .unwrap_or(std::env::temp_dir());
1681        let epoch_store = self.load_epoch_store_one_call_per_task();
1682
1683        NodeStateDump::new(
1684            tx_digest,
1685            effects,
1686            expected_effects_digest,
1687            self.get_object_store().as_ref(),
1688            &epoch_store,
1689            inner_temporary_store,
1690            transaction,
1691        )?
1692        .write_to_file(&dump_dir)
1693        .map_err(|e| IotaError::FileIO(e.to_string()))
1694    }
1695
1696    #[instrument(name = "process_certificate", level = "trace", skip_all, fields(tx_digest = ?transaction.digest(), sender = ?transaction.data().transaction().gas_owner().to_string()))]
1697    pub(crate) fn process_transaction(
1698        &self,
1699        tx_guard: TxGuard,
1700        transaction: &VerifiedExecutableTransaction,
1701        tx_input_objects: InputObjects,
1702        per_authenticator_inputs: Vec<(InputObjects, ObjectReadResult)>,
1703        expected_effects_digest: Option<TransactionEffectsDigest>,
1704        epoch_store: &Arc<AuthorityPerEpochStore>,
1705    ) -> IotaResult<(TransactionEffects, Option<ExecutionError>)> {
1706        let process_transaction_start_time = tokio::time::Instant::now();
1707        let digest = *transaction.digest();
1708
1709        let _scope = monitored_scope("Execution::process_certificate");
1710
1711        fail_point_if!("correlated-crash-process-transaction", || {
1712            if iota_simulator::random::deterministic_probability_once(digest, 0.01) {
1713                iota_simulator::task::kill_current_node(None);
1714            }
1715        });
1716
1717        let execution_guard = self.execution_lock_for_executable_transaction(transaction);
1718        // Any caller that verifies the signatures on the transaction will have already
1719        // checked the epoch. But paths that don't verify sigs (e.g. execution
1720        // from checkpoint, reading from db) present the possibility of an epoch
1721        // mismatch. If this transaction is not finalized in previous epoch, then it's
1722        // invalid.
1723        let execution_guard = match execution_guard {
1724            Ok(execution_guard) => execution_guard,
1725            Err(err) => {
1726                tx_guard.release();
1727                return Err(err);
1728            }
1729        };
1730        // Since we obtain a reference to the epoch store before taking the execution
1731        // lock, it's possible that reconfiguration has happened and they no
1732        // longer match.
1733        if *execution_guard != epoch_store.epoch() {
1734            tx_guard.release();
1735            info!("The epoch of the execution_guard doesn't match the epoch store");
1736            return Err(IotaError::WrongEpoch {
1737                expected_epoch: epoch_store.epoch(),
1738                actual_epoch: *execution_guard,
1739            });
1740        }
1741
1742        // Errors originating from `execute_transaction` may be transient (failure to
1743        // read locks) or non-transient (transaction input is invalid, move vm
1744        // errors). However, all errors from this function occur before we have
1745        // written anything to the db, so we commit the tx guard and rely on the
1746        // client to retry the tx (if it was transient).
1747        let (inner_temporary_store, effects, execution_error_opt) = match self.execute_transaction(
1748            &execution_guard,
1749            transaction,
1750            tx_input_objects,
1751            per_authenticator_inputs,
1752            epoch_store,
1753        ) {
1754            Err(e) => {
1755                info!(name = ?self.name, ?digest, "Error preparing transaction: {e}");
1756                tx_guard.release();
1757                return Err(e);
1758            }
1759            Ok(res) => res,
1760        };
1761
1762        if let Some(expected_effects_digest) = expected_effects_digest {
1763            if effects.digest() != expected_effects_digest {
1764                // We dont want to mask the original error, so we log it and continue.
1765                match self.debug_dump_transaction_state(
1766                    &digest,
1767                    &effects,
1768                    expected_effects_digest,
1769                    &inner_temporary_store,
1770                    transaction,
1771                    &self.config.state_debug_dump_config,
1772                ) {
1773                    Ok(out_path) => {
1774                        info!(
1775                            "Dumped node state for transaction {} to {}",
1776                            digest,
1777                            out_path.as_path().display().to_string()
1778                        );
1779                    }
1780                    Err(e) => {
1781                        error!("Error dumping state for transaction {}: {e}", digest);
1782                    }
1783                }
1784                error!(
1785                    tx_digest = ?digest,
1786                    ?expected_effects_digest,
1787                    actual_effects = ?effects,
1788                    "fork detected!"
1789                );
1790                panic!(
1791                    "Transaction {} is expected to have effects digest {}, but got {}!",
1792                    digest,
1793                    expected_effects_digest,
1794                    effects.digest(),
1795                );
1796            }
1797        }
1798
1799        fail_point!("crash");
1800
1801        self.commit_transaction(
1802            transaction,
1803            inner_temporary_store,
1804            &effects,
1805            tx_guard,
1806            execution_guard,
1807            expected_effects_digest,
1808            epoch_store,
1809        )?;
1810
1811        let elapsed = process_transaction_start_time.elapsed().as_micros() as f64;
1812        if elapsed > 0.0 {
1813            self.metrics
1814                .execution_gas_latency_ratio
1815                .observe(effects.gas_cost_summary().computation_cost as f64 / elapsed);
1816        };
1817        Ok((effects, execution_error_opt))
1818    }
1819
1820    pub async fn reconfigure_traffic_control(
1821        &self,
1822        params: TrafficControlReconfigParams,
1823    ) -> Result<TrafficControlReconfigParams, IotaError> {
1824        if let Some(traffic_controller) = self.traffic_controller.as_ref() {
1825            traffic_controller.admin_reconfigure(params)
1826        } else {
1827            Err(IotaError::InvalidAdminRequest(
1828                "Traffic controller is not configured on this node".to_string(),
1829            ))
1830        }
1831    }
1832
1833    #[instrument(level = "trace", skip_all)]
1834    fn commit_transaction(
1835        &self,
1836        transaction: &VerifiedExecutableTransaction,
1837        inner_temporary_store: InnerTemporaryStore,
1838        effects: &TransactionEffects,
1839        tx_guard: TxGuard,
1840        _execution_guard: ExecutionLockReadGuard<'_>,
1841        expected_effects_digest: Option<TransactionEffectsDigest>,
1842        epoch_store: &Arc<AuthorityPerEpochStore>,
1843    ) -> IotaResult {
1844        let _scope: Option<iota_metrics::MonitoredScopeGuard> =
1845            monitored_scope("Execution::commit_certificate");
1846        let _metrics_guard = self.metrics.commit_certificate_latency.start_timer();
1847
1848        let tx_digest = transaction.digest();
1849        let input_object_count = inner_temporary_store.input_objects.len();
1850        let shared_object_count = effects.input_shared_objects().len();
1851
1852        let output_keys = inner_temporary_store.get_output_keys(effects);
1853
1854        // index transaction
1855        let _ = self
1856            .post_process_one_tx(transaction, effects, &inner_temporary_store, epoch_store)
1857            .tap_err(|e| {
1858                self.metrics.post_processing_total_failures.inc();
1859                error!(?tx_digest, "tx post processing failed: {e}");
1860            });
1861
1862        // The insertion to epoch_store is not atomic with the insertion to the
1863        // perpetual store. This is OK because we insert to the epoch store
1864        // first. And during lookups we always look up in the perpetual store first.
1865        epoch_store.insert_executed_in_epoch(tx_digest);
1866
1867        let key = transaction.key();
1868        if !matches!(key, TransactionKey::Digest(_)) {
1869            epoch_store.insert_tx_key(key, *tx_digest)?;
1870        }
1871
1872        // Allow testing what happens if we crash here.
1873        fail_point!("crash");
1874
1875        let transaction_outputs = TransactionOutputs::build_transaction_outputs(
1876            transaction.clone().into_unsigned(),
1877            effects.clone(),
1878            inner_temporary_store,
1879        );
1880        self.get_cache_writer()
1881            .try_write_transaction_outputs(epoch_store.epoch(), transaction_outputs.into())?;
1882
1883        self.report_failed_deny_rule_update_execution(
1884            transaction,
1885            effects,
1886            expected_effects_digest,
1887            epoch_store,
1888        );
1889
1890        if transaction.transaction().is_end_of_epoch_tx() {
1891            // At the end of epoch, since system packages may have been upgraded, force
1892            // reload them in the cache.
1893            self.get_object_cache_reader()
1894                .force_reload_system_packages(&BuiltInFramework::all_package_ids());
1895        }
1896
1897        // `commit_transaction()` finished, the tx is fully committed to the store.
1898        tx_guard.commit_tx();
1899
1900        match self.execution_scheduler.as_ref() {
1901            ExecutionSchedulerWrapper::ExecutionScheduler(_) => {}
1902            ExecutionSchedulerWrapper::TransactionManager(tm) => {
1903                // Notifies transaction manager about transaction and output objects committed.
1904                // This provides necessary information to transaction manager to start executing
1905                // additional ready transactions.
1906                tm.notify_commit(tx_digest, output_keys, epoch_store);
1907                // A transaction with a non-digest key can execute from a synced
1908                // checkpoint, in which case local randomness generation — the only
1909                // other caller of `notify_transaction_key` — never runs for that
1910                // round and would leave the env parked under its key forever. The
1911                // enqueue this triggers is filtered out as already executed.
1912                if let Some(key) = transaction.non_digest_key() {
1913                    tm.notify_transaction_key(epoch_store, key, *tx_digest);
1914                }
1915            }
1916        }
1917
1918        self.update_metrics(transaction, input_object_count, shared_object_count);
1919
1920        Ok(())
1921    }
1922
1923    fn update_metrics(
1924        &self,
1925        transaction: &VerifiedExecutableTransaction,
1926        input_object_count: usize,
1927        shared_object_count: usize,
1928    ) {
1929        // count signature by scheme, for multisig
1930        if transaction.has_multisig() {
1931            self.metrics.multisig_sig_count.inc();
1932        }
1933
1934        self.metrics.total_effects.inc();
1935        self.metrics.total_certs.inc();
1936
1937        if shared_object_count > 0 {
1938            self.metrics.shared_obj_tx.inc();
1939        }
1940
1941        if transaction.is_sponsored_tx() {
1942            self.metrics.sponsored_tx.inc();
1943        }
1944
1945        self.metrics
1946            .num_input_objs
1947            .observe(input_object_count as f64);
1948        self.metrics
1949            .num_shared_objects
1950            .observe(shared_object_count as f64);
1951        self.metrics
1952            .batch_size
1953            .observe(transaction.data().transaction().kind().num_commands() as f64);
1954    }
1955
1956    /// `execute_transaction()` validates the transaction input, and executes
1957    /// the transaction, returning effects, output objects, events, etc.
1958    ///
1959    /// It reads state from the db (both owned and shared locks), but it has no
1960    /// side effects.
1961    ///
1962    /// It can be generally understood that a failure of `execute_transaction`
1963    /// indicates a non-transient error, e.g. the transaction input is
1964    /// somehow invalid, the correct locks are not held, etc. However, this
1965    /// is not entirely true, as a transient db read error may also cause
1966    /// this function to fail.
1967    #[instrument(level = "trace", skip_all)]
1968    fn execute_transaction(
1969        &self,
1970        _execution_guard: &ExecutionLockReadGuard<'_>,
1971        transaction: &VerifiedExecutableTransaction,
1972        tx_input_objects: InputObjects,
1973        per_authenticator_inputs: Vec<(InputObjects, ObjectReadResult)>,
1974        epoch_store: &Arc<AuthorityPerEpochStore>,
1975    ) -> IotaResult<(
1976        InnerTemporaryStore,
1977        TransactionEffects,
1978        Option<ExecutionError>,
1979    )> {
1980        let _scope = monitored_scope("Execution::execute_certificate");
1981        let _metrics_guard = self.metrics.prepare_certificate_latency.start_timer();
1982        let prepare_transaction_start_time = tokio::time::Instant::now();
1983
1984        let protocol_config = epoch_store.protocol_config();
1985
1986        let reference_gas_price = epoch_store.reference_gas_price();
1987
1988        let epoch_id = epoch_store.epoch_start_config().epoch_data().epoch_id();
1989        let epoch_start_timestamp = epoch_store
1990            .epoch_start_config()
1991            .epoch_data()
1992            .epoch_start_timestamp();
1993
1994        let backing_store = self.get_backing_store().as_ref();
1995
1996        let tx_digest = *transaction.digest();
1997
1998        // TODO: We need to move this to a more appropriate place to avoid redundant
1999        // checks.
2000        let tx = transaction.data().transaction();
2001        tx.validity_check(protocol_config)?;
2002
2003        let (kind, signer, gas_data) = tx.execution_parts();
2004
2005        let move_authenticators = transaction.move_authenticators();
2006
2007        #[cfg_attr(not(any(msim, fail_points)), expect(unused_mut))]
2008        let (inner_temp_store, _, mut effects, _, execution_error_opt) = if move_authenticators
2009            .is_empty()
2010        {
2011            // No Move authentication required, proceed to execute the transaction directly.
2012
2013            // The cost of partially re-auditing a transaction before execution is
2014            // tolerated.
2015            let (tx_gas_status, tx_checked_input_objects) =
2016                iota_transaction_checks::check_certificate_input(
2017                    transaction,
2018                    tx_input_objects,
2019                    protocol_config,
2020                    reference_gas_price,
2021                )?;
2022
2023            let owned_object_refs = tx_checked_input_objects.inner().filter_owned_objects();
2024            self.check_owned_locks(&owned_object_refs)?;
2025            epoch_store.executor().execute_transaction_to_effects(
2026                backing_store,
2027                protocol_config,
2028                self.metrics.limits_metrics.clone(),
2029                // TODO: would be nice to pass the whole NodeConfig here, but it creates a
2030                // cyclic dependency w/ iota-adapter
2031                self.config
2032                    .expensive_safety_check_config
2033                    .enable_deep_per_tx_iota_conservation_check(),
2034                self.config.certificate_deny_config.certificate_deny_set(),
2035                &epoch_id,
2036                epoch_start_timestamp,
2037                tx_checked_input_objects,
2038                gas_data,
2039                tx_gas_status,
2040                kind,
2041                signer,
2042                tx_digest,
2043                &mut None,
2044            )
2045        } else {
2046            // One or more `MoveAuthenticator` signatures present — authenticate each and
2047            // then execute the transaction.
2048            // It is supposed that `MoveAuthenticator` availability is checked in
2049            // `SenderSignedTransaction::validity_check`.
2050
2051            debug_assert_eq!(
2052                move_authenticators.len(),
2053                per_authenticator_inputs.len(),
2054                "Move authenticators amount must match the number of authenticator inputs"
2055            );
2056
2057            // The first account that cannot be resolved fails the transaction before
2058            // any authenticator runs. The remaining inputs are still collected, since
2059            // for the failure effects, we have to charge gas for those inputs.
2060            let mut pre_execution_error = None;
2061            let mut per_authenticator_input_objects = Vec::with_capacity(move_authenticators.len());
2062            let mut function_refs = Vec::with_capacity(move_authenticators.len());
2063            for (move_authenticator, (authenticator_input_objects, account_object)) in
2064                move_authenticators.iter().zip(per_authenticator_inputs)
2065            {
2066                // The shape of the object to authenticate is decided by the
2067                // transaction bytes, and `validity_check` settles it in the
2068                // consensus handler, so it cannot be wrong here.
2069                let (account_object_id, account_object_version, account_object_digest) =
2070                    move_authenticator
2071                        .object_to_authenticate_components()
2072                        .expect(
2073                            "the object to authenticate is validated before consensus and cannot \
2074                                be invalid during execution",
2075                        );
2076
2077                // Unlike the shape above, resolving the account depends on state,
2078                // so post-consensus validation cannot keep this check. Here, an
2079                // account that fails it produces failure effects instead of a
2080                // panic.
2081                match self.check_move_account_for_execution(
2082                    account_object_id,
2083                    account_object_version,
2084                    account_object_digest,
2085                    account_object,
2086                    &move_authenticator.address(),
2087                    protocol_config,
2088                ) {
2089                    Ok(function_ref) => function_refs.push(function_ref),
2090                    Err(error) => {
2091                        if pre_execution_error.is_none() {
2092                            pre_execution_error = Some(error);
2093                        }
2094                    }
2095                }
2096
2097                per_authenticator_input_objects.push(authenticator_input_objects);
2098            }
2099
2100            // Serialize the Transaction for the auth context.
2101            let tx_bytes = bcs::to_bytes(tx).expect("Transaction serialization cannot fail");
2102
2103            let (sender_auth_digest, sponsor_auth_digest) =
2104                transaction.data().compute_auth_digests()?;
2105
2106            // Check the `MoveAuthenticator` input objects.
2107            // The `MoveAuthenticator` receiving objects are checked on the signing step.
2108            // `max_auth_gas` is used here as a Move authenticator gas budget until it is
2109            // not a part of the transaction data.
2110            let authenticator_gas_budget = protocol_config.max_auth_gas();
2111            let (
2112                gas_status,
2113                per_authenticator_checked_input_objects,
2114                authenticator_and_tx_checked_input_objects,
2115            ) = iota_transaction_checks::check_certificate_and_move_authenticator_input(
2116                transaction,
2117                tx_input_objects,
2118                per_authenticator_input_objects,
2119                authenticator_gas_budget,
2120                protocol_config,
2121                reference_gas_price,
2122            )?;
2123
2124            let owned_object_refs = authenticator_and_tx_checked_input_objects
2125                .inner()
2126                .filter_owned_objects();
2127            self.check_owned_locks(&owned_object_refs)?;
2128
2129            // With a pre-execution failure, no authenticator runs, the gas owner's
2130            // included, so the executor gets the failure in place of the list. That
2131            // changes neither
2132            // who pays nor the outcome: the gas owner is charged either way, as
2133            // when the first authenticator fails.
2134            let move_authenticators = if pre_execution_error.is_some() {
2135                Vec::new()
2136            } else {
2137                debug_assert_eq!(
2138                    move_authenticators.len(),
2139                    per_authenticator_checked_input_objects.len(),
2140                    "Move authenticators amount must match the number of checked authenticator inputs"
2141                );
2142                move_authenticators
2143                    .into_iter()
2144                    .zip(function_refs)
2145                    .zip(per_authenticator_checked_input_objects)
2146                    .map(
2147                        |(
2148                            (move_authenticator, function_ref),
2149                            authenticator_checked_input_objects,
2150                        )| {
2151                            MoveAuthenticatorForExecution {
2152                                authenticator: move_authenticator.to_owned(),
2153                                function_ref,
2154                                input_objects: authenticator_checked_input_objects,
2155                            }
2156                        },
2157                    )
2158                    .collect::<Vec<_>>()
2159            };
2160
2161            let (sender_authenticator_function_ref, sponsor_authenticator_function_ref) =
2162                extract_auth_fun_refs(signer, gas_data.owner, |address| {
2163                    move_authenticators
2164                        .iter()
2165                        .find(|a| a.authenticator.address() == address)
2166                        .map(|a| a.function_ref.authenticator_function_ref.clone())
2167                });
2168
2169            let auth_context_data = AuthContextData {
2170                transaction_data_bytes: tx_bytes,
2171                sender_auth_digest,
2172                sponsor_auth_digest,
2173                sender_authenticator_function_ref,
2174                sponsor_authenticator_function_ref,
2175            };
2176
2177            let authenticators = match pre_execution_error {
2178                Some(error) => MoveAuthenticatorsForExecution::ResolutionFailed(error),
2179                None => MoveAuthenticatorsForExecution::Resolved(move_authenticators),
2180            };
2181
2182            epoch_store
2183                .executor()
2184                .authenticate_then_execute_transaction_to_effects(
2185                    backing_store,
2186                    protocol_config,
2187                    self.metrics.limits_metrics.clone(),
2188                    self.config
2189                        .expensive_safety_check_config
2190                        .enable_deep_per_tx_iota_conservation_check(),
2191                    self.config.certificate_deny_config.certificate_deny_set(),
2192                    &epoch_id,
2193                    epoch_start_timestamp,
2194                    gas_data,
2195                    gas_status,
2196                    authenticators,
2197                    authenticator_and_tx_checked_input_objects,
2198                    kind,
2199                    signer,
2200                    tx_digest,
2201                    auth_context_data,
2202                    &mut None,
2203                )
2204        };
2205
2206        fail_point_if!("cp_execution_nondeterminism", || {
2207            #[cfg(msim)]
2208            self.create_fail_state(transaction, epoch_store, &mut effects);
2209        });
2210
2211        let elapsed = prepare_transaction_start_time.elapsed().as_micros() as f64;
2212        if elapsed > 0.0 {
2213            self.metrics
2214                .prepare_cert_gas_latency_ratio
2215                .observe(effects.gas_cost_summary().computation_cost as f64 / elapsed);
2216        }
2217
2218        Ok((inner_temp_store, effects, execution_error_opt.err()))
2219    }
2220
2221    pub fn prepare_transaction_for_benchmark(
2222        &self,
2223        transaction: &VerifiedExecutableTransaction,
2224        input_objects: InputObjects,
2225        epoch_store: &Arc<AuthorityPerEpochStore>,
2226    ) -> IotaResult<(
2227        InnerTemporaryStore,
2228        TransactionEffects,
2229        Option<ExecutionError>,
2230    )> {
2231        let lock = RwLock::new(epoch_store.epoch());
2232        let execution_guard = lock.try_read().unwrap();
2233
2234        self.execute_transaction(
2235            &execution_guard,
2236            transaction,
2237            input_objects,
2238            vec![],
2239            epoch_store,
2240        )
2241    }
2242
2243    /// Simulate a transaction without committing it.
2244    ///
2245    /// `checks` selects the Move VM semantics: `VmChecks::Enabled` runs the
2246    /// transaction as it would run on chain (a dry run), while
2247    /// `VmChecks::Disabled` relaxes the checks around entry functions and
2248    /// argument values (a dev inspect). Both report the per-command return
2249    /// values in [`SimulateTransactionResult::execution_result`].
2250    ///
2251    /// Under either `checks`, the simulation fills in whatever gas the
2252    /// transaction leaves unset, so that a caller with no gas to declare can
2253    /// leave all of it out: no gas payment mints a mock gas coin, whose ID is
2254    /// reported back in [`SimulateTransactionResult::mock_gas_id`]; a zero gas
2255    /// price becomes the epoch's reference gas price; and a zero gas budget
2256    /// becomes as much as the gas coins can back, up to
2257    /// [`max_tx_gas`](iota_protocol_config::ProtocolConfig::max_tx_gas).
2258    /// Anything the transaction does declare is metered as given, so a dry run
2259    /// still rejects the gas a validator would.
2260    ///
2261    /// Whatever the budget resolves to, the gas coins have to cover it, since
2262    /// execution reserves the whole budget from them before running any command
2263    /// and refunds it afterwards. A caller leaving the budget at zero to have
2264    /// the cost estimated therefore gets an estimate whatever its coins hold,
2265    /// but the reserved budget is off limits for the duration of the
2266    /// programmable transaction: a transaction that also pays out of its gas
2267    /// coin has to declare a budget leaving room for that, exactly as it would
2268    /// on chain. A balance too small to declare the minimum budget at all is
2269    /// rejected with [`UserInputError::GasBalanceTooLow`].
2270    pub fn simulate_transaction(
2271        &self,
2272        transaction: Transaction,
2273        checks: VmChecks,
2274    ) -> IotaResult<SimulateTransactionResult> {
2275        let epoch_store = self.load_epoch_store_one_call_per_task();
2276        self.simulate_transaction_in_epoch(&epoch_store, transaction, checks)
2277    }
2278
2279    /// Same as [`AuthorityState::simulate_transaction`], for callers that
2280    /// already hold an epoch store.
2281    ///
2282    /// Callers that derive gas parameters from an epoch, or resolve types
2283    /// against its executor once the simulation returns, should pass that same
2284    /// epoch store here so the whole operation observes one epoch.
2285    ///
2286    /// Nothing here checks that `epoch_store` is the current one — pinning a
2287    /// superseded epoch is the point, and is what
2288    /// [`AuthorityState::simulate_transaction`] does for the span of its own
2289    /// call. Keeping one across an unbounded period is the caller's problem:
2290    /// the simulation would run against that epoch's protocol config,
2291    /// executor, and reference gas price.
2292    #[instrument("simulate_tx", level = "trace", skip_all)]
2293    pub fn simulate_transaction_in_epoch(
2294        &self,
2295        epoch_store: &AuthorityPerEpochStore,
2296        transaction: Transaction,
2297        checks: VmChecks,
2298    ) -> IotaResult<SimulateTransactionResult> {
2299        if !self.is_fullnode(epoch_store) {
2300            return Err(IotaError::UnsupportedFeature {
2301                error: "simulate is only supported on fullnodes".to_string(),
2302            });
2303        }
2304
2305        self.simulate_transaction_inner(epoch_store, transaction, checks)
2306    }
2307
2308    /// Same as [`AuthorityState::simulate_transaction`], but runs on a
2309    /// validator too. Only the single-node benchmark, which has no fullnode
2310    /// to run against, needs this.
2311    pub fn simulate_transaction_for_benchmark(
2312        &self,
2313        transaction: Transaction,
2314        checks: VmChecks,
2315    ) -> IotaResult<SimulateTransactionResult> {
2316        let epoch_store = self.load_epoch_store_one_call_per_task();
2317        self.simulate_transaction_inner(&epoch_store, transaction, checks)
2318    }
2319
2320    #[instrument(level = "trace", skip_all)]
2321    fn simulate_transaction_inner(
2322        &self,
2323        epoch_store: &AuthorityPerEpochStore,
2324        mut transaction: Transaction,
2325        checks: VmChecks,
2326    ) -> IotaResult<SimulateTransactionResult> {
2327        if transaction.kind().is_system() {
2328            return Err(IotaError::UnsupportedFeature {
2329                error: "simulate does not support system transactions".to_string(),
2330            });
2331        }
2332
2333        transaction.check_serialized_size(epoch_store.protocol_config())?;
2334
2335        // Cheap validity checks for a transaction, including input size limits.
2336        // This does not check if gas objects are missing since we may create a
2337        // mock gas object. It checks for other transaction input validity.
2338        transaction.validity_check_no_gas_check(epoch_store.protocol_config())?;
2339
2340        // The full validity check caps the gas payment size alongside requiring a
2341        // gas payment at all, which a simulation relaxes so it can mock one. The cap
2342        // still applies, and is cheapest before any object is loaded.
2343        transaction.check_gas_payment_size(epoch_store.protocol_config())?;
2344
2345        let input_object_kinds = transaction.input_objects()?;
2346        let receiving_object_refs = transaction.receiving_objects();
2347
2348        // Since we need to simulate a validator signing the transaction, the first step
2349        // is to check if some transaction elements are denied.
2350        iota_transaction_checks::deny::check_transaction_for_validation(
2351            &transaction,
2352            &[],
2353            &input_object_kinds,
2354            &receiving_object_refs,
2355            &self.config.transaction_deny_config,
2356            self.get_backing_package_store().as_ref(),
2357        )?;
2358
2359        // Load input and receiving objects
2360        let (mut input_objects, receiving_objects) = self.input_loader.read_objects_for_signing(
2361            // We don't want to cache this transaction since it's a simulation.
2362            None,
2363            &input_object_kinds,
2364            &receiving_object_refs,
2365            epoch_store.epoch(),
2366        )?;
2367
2368        // Create a mock gas object if one was not provided
2369        let mock_gas_id = if transaction.gas().is_empty() {
2370            let mock_gas_object = mock_simulation_gas_coin(transaction.gas_data().owner);
2371            let mock_gas_object_ref = mock_gas_object.object_ref();
2372            transaction.gas_data_mut().objects = vec![mock_gas_object_ref];
2373            input_objects.push(ObjectReadResult::new_from_gas_object(&mock_gas_object));
2374            Some(mock_gas_object.id())
2375        } else {
2376            None
2377        };
2378
2379        let protocol_config = epoch_store.protocol_config();
2380
2381        iota_types::gas::fill_in_unset_simulation_gas(
2382            &mut transaction,
2383            &input_objects,
2384            epoch_store.reference_gas_price(),
2385            protocol_config,
2386        );
2387
2388        // `MoveAuthenticator`s are not supported in simulation, so we set the
2389        // `authenticator_gas_budget` to 0.
2390        let authenticator_gas_budget = 0;
2391
2392        // Checks enabled -> DRY-RUN, it means we are simulating a real TX
2393        // Checks disabled -> DEV-INSPECT, more relaxed Move VM checks
2394        let (gas_status, checked_input_objects) = if checks.enabled() {
2395            iota_transaction_checks::check_transaction_input(
2396                protocol_config,
2397                epoch_store.reference_gas_price(),
2398                &transaction,
2399                input_objects,
2400                &receiving_objects,
2401                &self.metrics.bytecode_verifier_metrics,
2402                VerifierLimitsSource::NodeConfig(&self.config.verifier_signing_config),
2403                authenticator_gas_budget,
2404            )?
2405        } else {
2406            // Execution smashes the gas coins and reserves the whole budget from them
2407            // before running any command, treating the input checks as having verified
2408            // that they are gas coins at all — so with those checks skipped here, this
2409            // has to stand in for them. With the checks enabled,
2410            // `check_transaction_input` covers it.
2411            iota_types::gas::check_gas_coins_cover_budget_in_simulation(
2412                &input_objects,
2413                transaction.gas(),
2414                transaction.gas_budget(),
2415            )?;
2416
2417            let checked_input_objects = iota_transaction_checks::check_simulation_input(
2418                protocol_config,
2419                transaction.kind(),
2420                input_objects,
2421                receiving_objects,
2422            )?;
2423            let gas_status = IotaGasStatus::new(
2424                transaction.gas_budget(),
2425                transaction.gas_price(),
2426                epoch_store.reference_gas_price(),
2427                protocol_config,
2428            )?;
2429
2430            (gas_status, checked_input_objects)
2431        };
2432
2433        // Create a new executor for the simulation
2434        let executor = iota_execution::executor(
2435            protocol_config,
2436            true, // silent
2437            None,
2438        )
2439        .expect("Creating an executor should not fail here");
2440
2441        // Execute the simulation
2442        let (kind, signer, gas_data) = transaction.execution_parts();
2443        let (inner_temp_store, _, effects, execution_result) = executor.dev_inspect_transaction(
2444            self.get_backing_store().as_ref(),
2445            protocol_config,
2446            self.metrics.limits_metrics.clone(),
2447            false, // expensive_checks
2448            self.config.certificate_deny_config.certificate_deny_set(),
2449            &epoch_store.epoch_start_config().epoch_data().epoch_id(),
2450            epoch_store
2451                .epoch_start_config()
2452                .epoch_data()
2453                .epoch_start_timestamp(),
2454            checked_input_objects,
2455            gas_data,
2456            gas_status,
2457            kind,
2458            signer,
2459            transaction.digest(),
2460            checks.disabled(),
2461        );
2462
2463        let mut input_objects = inner_temp_store.input_objects;
2464        iota_types::storage::extend_input_objects_with_loaded_runtime_objects(
2465            &mut input_objects,
2466            &effects,
2467            &inner_temp_store.loaded_runtime_objects,
2468            self.get_backing_store().as_object_store(),
2469        );
2470
2471        Ok(SimulateTransactionResult {
2472            input_objects,
2473            output_objects: inner_temp_store.written,
2474            events: effects.events_digest().map(|_| inner_temp_store.events),
2475            effects,
2476            execution_result,
2477            suggested_gas_price: self
2478                .congestion_tracker
2479                .get_prediction_suggested_gas_price(&transaction),
2480            mock_gas_id,
2481            gas_data: transaction.gas_data().clone(),
2482        })
2483    }
2484
2485    // Only used for testing because of how epoch store is loaded.
2486    pub fn reference_gas_price_for_testing(&self) -> Result<u64, anyhow::Error> {
2487        let epoch_store = self.epoch_store_for_testing();
2488        Ok(epoch_store.reference_gas_price())
2489    }
2490
2491    #[instrument(level = "trace", skip_all)]
2492    pub fn try_is_tx_already_executed(&self, digest: &TransactionDigest) -> IotaResult<bool> {
2493        self.get_transaction_cache_reader()
2494            .try_is_tx_already_executed(digest)
2495    }
2496
2497    /// Non-fallible version of `try_is_tx_already_executed`.
2498    pub fn is_tx_already_executed(&self, digest: &TransactionDigest) -> bool {
2499        self.try_is_tx_already_executed(digest)
2500            .expect("storage access failed")
2501    }
2502
2503    /// Indexes a transaction by updating various indexes in the `IndexStore`.
2504    #[instrument(level = "debug", skip_all, err)]
2505    fn index_tx(
2506        &self,
2507        indexes: &IndexStore,
2508        digest: &TransactionDigest,
2509        // TODO: index_tx really just need the transaction data here.
2510        transaction: &VerifiedExecutableTransaction,
2511        effects: &TransactionEffects,
2512        events: &TransactionEvents,
2513        timestamp_ms: u64,
2514        tx_coins: Option<TxCoins>,
2515        written: &WrittenObjects,
2516        inner_temporary_store: &InnerTemporaryStore,
2517        epoch_store: &Arc<AuthorityPerEpochStore>,
2518    ) -> IotaResult<u64> {
2519        let changes = self
2520            .process_object_index(effects, written, inner_temporary_store, epoch_store)
2521            .tap_err(|e| warn!(tx_digest=?digest, "Failed to process object index, index_tx is skipped: {e}"))?;
2522
2523        indexes.index_tx(
2524            transaction.data().transaction().sender(),
2525            transaction
2526                .data()
2527                .transaction()
2528                .input_objects()?
2529                .iter()
2530                .map(|o| o.object_id()),
2531            effects
2532                .all_changed_objects()
2533                .into_iter()
2534                .map(|(changed, _kind)| (*changed.reference(), *changed.owner())),
2535            transaction
2536                .data()
2537                .transaction()
2538                .move_calls()
2539                .into_iter()
2540                .map(|(package, module, function)| {
2541                    (*package, module.to_owned(), function.to_owned())
2542                }),
2543            events,
2544            changes,
2545            digest,
2546            timestamp_ms,
2547            tx_coins,
2548        )
2549    }
2550
2551    #[cfg(msim)]
2552    fn create_fail_state(
2553        &self,
2554        transaction: &VerifiedExecutableTransaction,
2555        epoch_store: &Arc<AuthorityPerEpochStore>,
2556        effects: &mut TransactionEffects,
2557    ) {
2558        use std::cell::RefCell;
2559
2560        use iota_types::effects::TransactionEffectsAPIForTesting;
2561        thread_local! {
2562            static FAIL_STATE: RefCell<(u64, HashSet<AuthorityName>)> = RefCell::new((0, HashSet::new()));
2563        }
2564        if !transaction.data().transaction().is_system_tx() {
2565            let committee = epoch_store.committee();
2566            let cur_stake = (**committee).weight(&self.name);
2567            if cur_stake > 0 {
2568                FAIL_STATE.with_borrow_mut(|fail_state| {
2569                    // let (&mut failing_stake, &mut failing_validators) = fail_state;
2570                    if fail_state.0 < committee.validity_threshold() {
2571                        fail_state.0 += cur_stake;
2572                        fail_state.1.insert(self.name);
2573                    }
2574
2575                    if fail_state.1.contains(&self.name) {
2576                        info!("cp_exec failing tx");
2577                        effects.gas_cost_summary_mut_for_testing().computation_cost += 1;
2578                    }
2579                });
2580            }
2581        }
2582    }
2583
2584    fn process_object_index(
2585        &self,
2586        effects: &TransactionEffects,
2587        written: &WrittenObjects,
2588        inner_temporary_store: &InnerTemporaryStore,
2589        epoch_store: &Arc<AuthorityPerEpochStore>,
2590    ) -> IotaResult<ObjectIndexChanges> {
2591        let mut layout_resolver =
2592            epoch_store
2593                .executor()
2594                .type_layout_resolver(Box::new(PackageStoreWithFallback::new(
2595                    inner_temporary_store,
2596                    self.get_backing_package_store(),
2597                )));
2598
2599        let modified_at_version = effects
2600            .modified_at_versions()
2601            .into_iter()
2602            .map(|modified| (*modified.object_id(), modified.version()))
2603            .collect::<HashMap<_, _>>();
2604
2605        let tx_digest = effects.transaction_digest();
2606        let mut deleted_owners = vec![];
2607        let mut deleted_dynamic_fields = vec![];
2608        for object_ref in effects.deleted().into_iter().chain(effects.wrapped()) {
2609            let old_version = modified_at_version.get(&object_ref.object_id).unwrap();
2610            // When we process the index, the latest object hasn't been written yet so
2611            // the old object must be present.
2612            match self.get_owner_at_version(&object_ref.object_id, *old_version).unwrap_or_else(
2613                |e| panic!("tx_digest={tx_digest}, error processing object owner index, cannot find owner for object {} at version {old_version:?}. Err: {e:?}", object_ref.object_id)
2614            ) {
2615                Owner::Address(addr) => deleted_owners.push((addr, object_ref.object_id)),
2616                Owner::Object(object_id) => {
2617                    deleted_dynamic_fields.push((object_id, object_ref.object_id))
2618                }
2619                _ => {}
2620            }
2621        }
2622
2623        let mut new_owners = vec![];
2624        let mut new_dynamic_fields = vec![];
2625
2626        for (changed, kind) in effects.all_changed_objects() {
2627            let (oref, owner) = (*changed.reference(), *changed.owner());
2628            let id = &oref.object_id;
2629            // For mutated objects, retrieve old owner and delete old index if there is a
2630            // owner change.
2631            if let WriteKind::Mutate = kind {
2632                let Some(old_version) = modified_at_version.get(id) else {
2633                    panic!(
2634                        "tx_digest={tx_digest}, error processing object owner index, cannot find modified at version for mutated object [{id}]."
2635                    );
2636                };
2637                // When we process the index, the latest object hasn't been written yet so
2638                // the old object must be present.
2639                let Some(old_object) = self
2640                    .get_object_store()
2641                    .try_get_object_by_key(id, *old_version)?
2642                else {
2643                    panic!(
2644                        "tx_digest={tx_digest}, error processing object owner index, cannot find owner for object {id} at version {old_version:?}"
2645                    );
2646                };
2647                if old_object.owner != owner {
2648                    match old_object.owner {
2649                        Owner::Address(addr) => {
2650                            deleted_owners.push((addr, *id));
2651                        }
2652                        Owner::Object(object_id) => deleted_dynamic_fields.push((object_id, *id)),
2653                        _ => {}
2654                    }
2655                }
2656            }
2657
2658            match owner {
2659                Owner::Address(addr) => {
2660                    // TODO: We can remove the object fetching after we added ObjectType to
2661                    // TransactionEffects
2662                    let new_object = written.get(id).unwrap_or_else(
2663                        || panic!("tx_digest={tx_digest}, error processing object owner index, written does not contain object {id}")
2664                    );
2665                    assert_eq!(
2666                        new_object.version(),
2667                        oref.version,
2668                        "tx_digest={} error processing object owner index, object {} from written has mismatched version. Actual: {}, expected: {}",
2669                        tx_digest,
2670                        id,
2671                        new_object.version(),
2672                        oref.version
2673                    );
2674
2675                    let object_type = new_object
2676                        .data
2677                        .opt_object_type()
2678                        .map(|ty| ObjectType::Struct(ty.clone()))
2679                        .unwrap_or(ObjectType::Package);
2680
2681                    new_owners.push((
2682                        (addr, *id),
2683                        ObjectInfo {
2684                            object_id: *id,
2685                            version: oref.version,
2686                            digest: oref.digest,
2687                            object_type,
2688                            owner,
2689                            previous_transaction: *effects.transaction_digest(),
2690                        },
2691                    ));
2692                }
2693                Owner::Object(owner) => {
2694                    let new_object = written.get(id).unwrap_or_else(
2695                        || panic!("tx_digest={tx_digest}, error processing object owner index, written does not contain object {id}")
2696                    );
2697                    assert_eq!(
2698                        new_object.version(),
2699                        oref.version,
2700                        "tx_digest={} error processing object owner index, object {} from written has mismatched version. Actual: {}, expected: {}",
2701                        tx_digest,
2702                        id,
2703                        new_object.version(),
2704                        oref.version
2705                    );
2706
2707                    let Some(df_info) = self
2708                        .try_create_dynamic_field_info(new_object, written, layout_resolver.as_mut())
2709                        .unwrap_or_else(|e| {
2710                            error!(
2711                                "try_create_dynamic_field_info should not fail, {}, new_object={}, new_object_type={}",
2712                                e,
2713                                new_object.id(),
2714                                ObjectType::from(new_object)
2715                            );
2716                            None
2717                        }
2718                        )
2719                    else {
2720                        // Skip indexing for non dynamic field objects.
2721                        continue;
2722                    };
2723                    new_dynamic_fields.push(((owner, *id), df_info))
2724                }
2725                _ => {}
2726            }
2727        }
2728
2729        Ok(ObjectIndexChanges {
2730            deleted_owners,
2731            deleted_dynamic_fields,
2732            new_owners,
2733            new_dynamic_fields,
2734        })
2735    }
2736
2737    fn try_create_dynamic_field_info(
2738        &self,
2739        o: &Object,
2740        written: &WrittenObjects,
2741        resolver: &mut dyn LayoutResolver,
2742    ) -> IotaResult<Option<DynamicFieldInfo>> {
2743        // Skip if not a move object
2744        let Some(move_object) = o.data.as_opt_struct().cloned() else {
2745            return Ok(None);
2746        };
2747
2748        // We only index dynamic field objects
2749        if !move_object.struct_tag().is_dynamic_field() {
2750            return Ok(None);
2751        }
2752
2753        let layout = match resolver.get_annotated_layout(move_object.struct_tag()) {
2754            Ok(annotated_layout) => annotated_layout.into_layout(),
2755            Err(e) => {
2756                error!(
2757                    "unable to load layout for type `{:?}`: {e}",
2758                    move_object.struct_tag()
2759                );
2760                return Ok(None);
2761            }
2762        };
2763
2764        let field =
2765            DFV::FieldVisitor::deserialize(move_object.contents(), &layout).map_err(|e| {
2766                IotaError::ObjectDeserialization {
2767                    error: e.to_string(),
2768                }
2769            })?;
2770
2771        let type_ = field.kind;
2772        let name_type: TypeTag = type_tag_core_to_sdk(&field.name_layout.into());
2773        let bcs_name = field.name_bytes.to_owned();
2774
2775        let name_value = BoundedVisitor::deserialize_value(field.name_bytes, field.name_layout)
2776            .map_err(|e| {
2777                warn!("{e}");
2778                IotaError::ObjectDeserialization {
2779                    error: e.to_string(),
2780                }
2781            })?;
2782
2783        let name = DynamicFieldName {
2784            type_tag: name_type,
2785            value: IotaMoveValue::from(name_value).to_json_value(),
2786        };
2787
2788        let value_metadata = field.value_metadata().map_err(|e| {
2789            warn!("{e}");
2790            IotaError::ObjectDeserialization {
2791                error: e.to_string(),
2792            }
2793        })?;
2794
2795        Ok(Some(match value_metadata {
2796            DFV::ValueMetadata::DynamicField(object_type) => DynamicFieldInfo {
2797                name,
2798                bcs_name,
2799                type_,
2800                object_type: object_type.to_canonical_string(/* with_prefix */ true),
2801                object_id: o.id(),
2802                version: o.version(),
2803                digest: o.digest(),
2804            },
2805
2806            DFV::ValueMetadata::DynamicObjectField(object_id) => {
2807                // Find the actual object from storage using the object id obtained from the
2808                // wrapper.
2809
2810                // Try to find the object in the written objects first.
2811                let (version, digest, object_type) = if let Some(object) = written.get(&object_id) {
2812                    (
2813                        object.version(),
2814                        object.digest(),
2815                        object.data.opt_object_type().unwrap().clone(),
2816                    )
2817                } else {
2818                    // If not found, try to find it in the database.
2819                    let object = self
2820                        .get_object_store()
2821                        .try_get_object_by_key(&object_id, o.version())?
2822                        .ok_or_else(|| UserInputError::ObjectNotFound {
2823                            object_id,
2824                            version: Some(o.version()),
2825                        })?;
2826                    let version = object.version();
2827                    let digest = object.digest();
2828                    let object_type = object.data.opt_object_type().unwrap().clone();
2829                    (version, digest, object_type)
2830                };
2831
2832                DynamicFieldInfo {
2833                    name,
2834                    bcs_name,
2835                    type_,
2836                    object_type: object_type.to_string(),
2837                    object_id,
2838                    version,
2839                    digest,
2840                }
2841            }
2842        }))
2843    }
2844
2845    #[instrument(level = "trace", skip_all, err)]
2846    fn post_process_one_tx(
2847        &self,
2848        transaction: &VerifiedExecutableTransaction,
2849        effects: &TransactionEffects,
2850        inner_temporary_store: &InnerTemporaryStore,
2851        epoch_store: &Arc<AuthorityPerEpochStore>,
2852    ) -> IotaResult {
2853        if self.indexes.is_none() {
2854            return Ok(());
2855        }
2856
2857        let _scope = monitored_scope("Execution::post_process_one_tx");
2858
2859        let tx_digest = transaction.digest();
2860        let timestamp_ms = Self::unixtime_now_ms();
2861        let events = &inner_temporary_store.events;
2862        let written = &inner_temporary_store.written;
2863        let tx_coins = self.fullnode_only_get_tx_coins_for_indexing(
2864            effects,
2865            inner_temporary_store,
2866            epoch_store,
2867        );
2868
2869        // Index tx
2870        if let Some(indexes) = &self.indexes {
2871            let _ = self
2872                .index_tx(
2873                    indexes.as_ref(),
2874                    tx_digest,
2875                    transaction,
2876                    effects,
2877                    events,
2878                    timestamp_ms,
2879                    tx_coins,
2880                    written,
2881                    inner_temporary_store,
2882                    epoch_store,
2883                )
2884                .tap_ok(|_| self.metrics.post_processing_total_tx_indexed.inc())
2885                .tap_err(|e| error!(?tx_digest, "Post processing - Couldn't index tx: {e}"))
2886                .expect("Indexing tx should not fail");
2887
2888            let effects: IotaTransactionBlockEffects = effects.clone().try_into()?;
2889            let events = self.make_transaction_block_events(
2890                events.clone(),
2891                *tx_digest,
2892                timestamp_ms,
2893                epoch_store,
2894                inner_temporary_store,
2895            )?;
2896            // Emit events
2897            self.subscription_handler
2898                .process_tx(transaction.data().transaction(), &effects, &events)
2899                .tap_ok(|_| {
2900                    self.metrics
2901                        .post_processing_total_tx_had_event_processed
2902                        .inc()
2903                })
2904                .tap_err(|e| {
2905                    warn!(
2906                        ?tx_digest,
2907                        "Post processing - Couldn't process events for tx: {}", e
2908                    )
2909                })?;
2910
2911            self.metrics
2912                .post_processing_total_events_emitted
2913                .inc_by(events.data.len() as u64);
2914        };
2915        Ok(())
2916    }
2917
2918    fn make_transaction_block_events(
2919        &self,
2920        transaction_events: TransactionEvents,
2921        digest: TransactionDigest,
2922        timestamp_ms: u64,
2923        epoch_store: &Arc<AuthorityPerEpochStore>,
2924        inner_temporary_store: &InnerTemporaryStore,
2925    ) -> IotaResult<IotaTransactionBlockEvents> {
2926        let mut layout_resolver =
2927            epoch_store
2928                .executor()
2929                .type_layout_resolver(Box::new(PackageStoreWithFallback::new(
2930                    inner_temporary_store,
2931                    self.get_backing_package_store(),
2932                )));
2933        IotaTransactionBlockEvents::try_from(
2934            transaction_events,
2935            digest,
2936            Some(timestamp_ms),
2937            layout_resolver.as_mut(),
2938        )
2939    }
2940
2941    pub fn unixtime_now_ms() -> u64 {
2942        let now = SystemTime::now()
2943            .duration_since(UNIX_EPOCH)
2944            .expect("Time went backwards")
2945            .as_millis();
2946        u64::try_from(now).expect("Travelling in time machine")
2947    }
2948
2949    #[instrument(level = "trace", skip_all)]
2950    pub async fn handle_transaction_info_request(
2951        &self,
2952        request: TransactionInfoRequest,
2953    ) -> IotaResult<TransactionInfoResponse> {
2954        let epoch_store = self.load_epoch_store_one_call_per_task();
2955        let (transaction, status) = self
2956            .get_transaction_status(&request.transaction_digest, &epoch_store)?
2957            .ok_or(IotaError::TransactionNotFound {
2958                digest: request.transaction_digest,
2959            })?;
2960        Ok(TransactionInfoResponse {
2961            transaction,
2962            status,
2963        })
2964    }
2965
2966    #[instrument(level = "trace", skip_all)]
2967    pub async fn handle_object_info_request(
2968        &self,
2969        request: ObjectInfoRequest,
2970    ) -> IotaResult<ObjectInfoResponse> {
2971        let epoch_store = self.load_epoch_store_one_call_per_task();
2972
2973        let requested_object_seq = match request.request_kind {
2974            ObjectInfoRequestKind::LatestObjectInfo => {
2975                self.try_get_object_or_tombstone(request.object_id)?
2976                    .ok_or_else(|| {
2977                        IotaError::from(UserInputError::ObjectNotFound {
2978                            object_id: request.object_id,
2979                            version: None,
2980                        })
2981                    })?
2982                    .version
2983            }
2984            ObjectInfoRequestKind::PastObjectInfoDebug(seq) => seq,
2985        };
2986
2987        let object = self
2988            .get_object_store()
2989            .try_get_object_by_key(&request.object_id, requested_object_seq)?
2990            .ok_or_else(|| {
2991                IotaError::from(UserInputError::ObjectNotFound {
2992                    object_id: request.object_id,
2993                    version: Some(requested_object_seq),
2994                })
2995            })?;
2996
2997        let layout = if let (LayoutGenerationOption::Generate, Some(move_obj)) =
2998            (request.generate_layout, object.data.as_opt_struct())
2999        {
3000            Some(into_struct_layout(
3001                epoch_store
3002                    .executor()
3003                    .type_layout_resolver(Box::new(self.get_backing_package_store().as_ref()))
3004                    .get_annotated_layout(move_obj.struct_tag())?,
3005            )?)
3006        } else {
3007            None
3008        };
3009
3010        let lock = if !object.is_address_owned() {
3011            // Only address owned objects have locks.
3012            None
3013        } else {
3014            self.get_transaction_lock(&object.object_ref(), &epoch_store)?
3015                .map(|s| s.into_inner())
3016        };
3017
3018        Ok(ObjectInfoResponse {
3019            object,
3020            layout,
3021            lock_for_debugging: lock,
3022        })
3023    }
3024
3025    #[instrument(level = "trace", skip_all)]
3026    pub fn handle_checkpoint_request(
3027        &self,
3028        request: &CheckpointRequest,
3029    ) -> IotaResult<CheckpointResponse> {
3030        let summary = if request.certified {
3031            let summary = match request.sequence_number {
3032                Some(seq) => self
3033                    .checkpoint_store
3034                    .get_checkpoint_by_sequence_number(seq)?,
3035                None => self.checkpoint_store.get_latest_certified_checkpoint()?,
3036            }
3037            .map(|v| v.into_inner());
3038            summary.map(CheckpointSummaryResponse::Certified)
3039        } else {
3040            let summary = match request.sequence_number {
3041                Some(seq) => self.checkpoint_store.get_locally_computed_checkpoint(seq)?,
3042                None => self
3043                    .checkpoint_store
3044                    .get_latest_locally_computed_checkpoint()?,
3045            };
3046            summary.map(CheckpointSummaryResponse::Pending)
3047        };
3048        let contents = match &summary {
3049            Some(s) => self
3050                .checkpoint_store
3051                .get_checkpoint_contents(&s.contents_digest())?,
3052            None => None,
3053        };
3054        Ok(CheckpointResponse {
3055            checkpoint: summary,
3056            contents,
3057        })
3058    }
3059
3060    fn check_protocol_version(
3061        supported_protocol_versions: SupportedProtocolVersions,
3062        current_version: ProtocolVersion,
3063    ) {
3064        info!("current protocol version is now {:?}", current_version);
3065        info!("supported versions are: {:?}", supported_protocol_versions);
3066        if !supported_protocol_versions.is_version_supported(current_version) {
3067            let msg = format!(
3068                "Unsupported protocol version. The network is at {current_version:?}, but this IotaNode only supports: {supported_protocol_versions:?}. Shutting down.",
3069            );
3070
3071            error!("{}", msg);
3072            eprintln!("{msg}");
3073
3074            #[cfg(not(msim))]
3075            std::process::exit(1);
3076
3077            #[cfg(msim)]
3078            iota_simulator::task::shutdown_current_node();
3079        }
3080    }
3081
3082    #[expect(clippy::disallowed_methods)] // allow unbounded_channel()
3083    pub async fn new(
3084        name: AuthorityName,
3085        secret: StableSyncAuthoritySigner,
3086        supported_protocol_versions: SupportedProtocolVersions,
3087        store: Arc<AuthorityStore>,
3088        execution_cache_trait_pointers: ExecutionCacheTraitPointers,
3089        epoch_store: Arc<AuthorityPerEpochStore>,
3090        committee_store: Arc<CommitteeStore>,
3091        indexes: Option<Arc<IndexStore>>,
3092        grpc_indexes_store: Option<Arc<GrpcIndexesStore>>,
3093        checkpoint_store: Arc<CheckpointStore>,
3094        prometheus_registry: &Registry,
3095        genesis_objects: &[Object],
3096        config: NodeConfig,
3097        validator_tx_finalizer: Option<Arc<ValidatorTxFinalizer<NetworkAuthorityClient>>>,
3098        chain_identifier: ChainIdentifier,
3099        checkpoint_progress_tracker: Option<Arc<CheckpointProgressTracker>>,
3100        policy_config: Option<PolicyConfig>,
3101        firewall_config: Option<RemoteFirewallConfig>,
3102        epoch_end_db_snapshots: Option<EpochEndDbSnapshotHandle>,
3103    ) -> Arc<Self> {
3104        Self::check_protocol_version(supported_protocol_versions, epoch_store.protocol_version());
3105
3106        let metrics = Arc::new(AuthorityMetrics::new(prometheus_registry));
3107        let (tx_ready_transactions, rx_ready_transactions) = unbounded_channel();
3108        let execution_scheduler = Arc::new(ExecutionSchedulerWrapper::new(
3109            execution_cache_trait_pointers.object_cache_reader.clone(),
3110            execution_cache_trait_pointers
3111                .transaction_cache_reader
3112                .clone(),
3113            tx_ready_transactions,
3114            &epoch_store,
3115            metrics.clone(),
3116        ));
3117        let (tx_execution_shutdown, rx_execution_shutdown) = oneshot::channel();
3118
3119        let authority_per_epoch_pruner = AuthorityPerEpochStorePruner::new(
3120            epoch_store.get_parent_path(),
3121            config
3122                .authority_store_pruning_config
3123                .num_latest_epoch_dbs_to_retain,
3124        )
3125        .await;
3126        let pruner = AuthorityStorePruner::new(
3127            store.perpetual_tables.clone(),
3128            checkpoint_store.clone(),
3129            grpc_indexes_store.clone(),
3130            indexes.clone(),
3131            config.authority_store_pruning_config.clone(),
3132            epoch_store.committee().authority_exists(&name),
3133            epoch_store.epoch_start_state().epoch_duration_ms(),
3134            prometheus_registry,
3135            checkpoint_progress_tracker.clone(),
3136        );
3137        let input_loader =
3138            TransactionInputLoader::new(execution_cache_trait_pointers.object_cache_reader.clone());
3139        let epoch = epoch_store.epoch();
3140        let rgp = epoch_store.reference_gas_price();
3141        let traffic_controller_metrics =
3142            Arc::new(TrafficControllerMetrics::new(prometheus_registry));
3143        let traffic_controller = policy_config.map(|policy_config| {
3144            Arc::new(TrafficController::init(
3145                policy_config,
3146                traffic_controller_metrics,
3147                firewall_config.clone(),
3148            ))
3149        });
3150        let state = Arc::new(AuthorityState {
3151            name,
3152            secret,
3153            execution_lock: RwLock::new(epoch),
3154            epoch_store: ArcSwap::new(epoch_store.clone()),
3155            input_loader,
3156            execution_cache_trait_pointers,
3157            indexes,
3158            grpc_indexes_store,
3159            subscription_handler: Arc::new(SubscriptionHandler::new(prometheus_registry)),
3160            checkpoint_store,
3161            committee_store,
3162            execution_scheduler,
3163            tx_execution_shutdown: Mutex::new(Some(tx_execution_shutdown)),
3164            metrics,
3165            pruner,
3166            authority_per_epoch_pruner,
3167            checkpoint_progress_tracker,
3168            config,
3169            overload_info: AuthorityOverloadInfo::default(),
3170            validator_tx_finalizer,
3171            chain_identifier,
3172            congestion_tracker: Arc::new(CongestionTracker::new(rgp)),
3173            traffic_controller,
3174            epoch_end_db_snapshots,
3175        });
3176
3177        // Start a task to execute ready transactions.
3178        let authority_state = Arc::downgrade(&state);
3179        spawn_monitored_task!(execution_process(
3180            authority_state,
3181            rx_ready_transactions,
3182            rx_execution_shutdown,
3183        ));
3184        // TODO: This doesn't belong to the constructor of AuthorityState.
3185        state
3186            .create_owner_index_if_empty(genesis_objects, &epoch_store)
3187            .expect("Error indexing genesis objects.");
3188
3189        state
3190    }
3191
3192    pub fn epoch_db_pruner(&self) -> &AuthorityPerEpochStorePruner {
3193        &self.authority_per_epoch_pruner
3194    }
3195
3196    // TODO: Consolidate our traits to reduce the number of methods here.
3197    pub fn get_object_cache_reader(&self) -> &Arc<dyn ObjectCacheRead> {
3198        &self.execution_cache_trait_pointers.object_cache_reader
3199    }
3200
3201    pub fn get_transaction_cache_reader(&self) -> &Arc<dyn TransactionCacheRead> {
3202        &self.execution_cache_trait_pointers.transaction_cache_reader
3203    }
3204
3205    pub fn get_cache_writer(&self) -> &Arc<dyn ExecutionCacheWrite> {
3206        &self.execution_cache_trait_pointers.cache_writer
3207    }
3208
3209    pub fn get_backing_store(&self) -> &Arc<dyn BackingStore + Send + Sync> {
3210        &self.execution_cache_trait_pointers.backing_store
3211    }
3212
3213    pub fn get_backing_package_store(&self) -> &Arc<dyn BackingPackageStore + Send + Sync> {
3214        &self.execution_cache_trait_pointers.backing_package_store
3215    }
3216
3217    pub fn get_object_store(&self) -> &Arc<dyn ObjectStore + Send + Sync> {
3218        &self.execution_cache_trait_pointers.object_store
3219    }
3220
3221    pub fn get_reconfig_api(&self) -> &Arc<dyn ExecutionCacheReconfigAPI> {
3222        &self.execution_cache_trait_pointers.reconfig_api
3223    }
3224
3225    pub fn get_global_state_hash_store(&self) -> &Arc<dyn GlobalStateHashStore> {
3226        &self.execution_cache_trait_pointers.global_state_hash_store
3227    }
3228
3229    pub fn get_checkpoint_cache(&self) -> &Arc<dyn CheckpointCache> {
3230        &self.execution_cache_trait_pointers.checkpoint_cache
3231    }
3232
3233    pub fn get_state_sync_store(&self) -> &Arc<dyn StateSyncAPI> {
3234        &self.execution_cache_trait_pointers.state_sync_store
3235    }
3236
3237    pub fn get_cache_commit(&self) -> &Arc<dyn ExecutionCacheCommit> {
3238        &self.execution_cache_trait_pointers.cache_commit
3239    }
3240
3241    pub fn database_for_testing(&self) -> Arc<AuthorityStore> {
3242        self.execution_cache_trait_pointers
3243            .testing_api
3244            .database_for_testing()
3245    }
3246
3247    pub async fn prune_checkpoints_for_eligible_epochs_for_testing(
3248        &self,
3249        config: NodeConfig,
3250        metrics: Arc<AuthorityStorePruningMetrics>,
3251    ) -> anyhow::Result<()> {
3252        AuthorityStorePruner::prune_checkpoints_for_eligible_epochs(
3253            &self.database_for_testing().perpetual_tables,
3254            &self.checkpoint_store,
3255            self.grpc_indexes_store.as_deref(),
3256            config.authority_store_pruning_config,
3257            metrics,
3258            EPOCH_DURATION_MS_FOR_TESTING,
3259            self.checkpoint_progress_tracker.as_ref(),
3260        )
3261        .await
3262    }
3263
3264    pub fn execution_scheduler(&self) -> &Arc<ExecutionSchedulerWrapper> {
3265        &self.execution_scheduler
3266    }
3267
3268    /// Whether this authority runs the `ExecutionScheduler` rather than the
3269    /// `TransactionManager`.
3270    pub fn uses_execution_scheduler(&self) -> bool {
3271        self.execution_scheduler.uses_execution_scheduler()
3272    }
3273
3274    fn create_owner_index_if_empty(
3275        &self,
3276        genesis_objects: &[Object],
3277        epoch_store: &Arc<AuthorityPerEpochStore>,
3278    ) -> IotaResult {
3279        let Some(index_store) = &self.indexes else {
3280            return Ok(());
3281        };
3282        if !index_store.is_empty() {
3283            return Ok(());
3284        }
3285
3286        let mut new_owners = vec![];
3287        let mut new_dynamic_fields = vec![];
3288        let mut layout_resolver = epoch_store
3289            .executor()
3290            .type_layout_resolver(Box::new(self.get_backing_package_store().as_ref()));
3291        for o in genesis_objects.iter() {
3292            match o.owner {
3293                Owner::Address(addr) => {
3294                    new_owners.push(((addr, o.id()), ObjectInfo::new(&o.object_ref(), o)))
3295                }
3296                Owner::Object(object_id) => {
3297                    let id = o.id();
3298                    let info = match self.try_create_dynamic_field_info(
3299                        o,
3300                        &BTreeMap::new(),
3301                        layout_resolver.as_mut(),
3302                    ) {
3303                        Ok(Some(info)) => info,
3304                        Ok(None) => continue,
3305                        Err(IotaError::UserInput {
3306                            error:
3307                                UserInputError::ObjectNotFound {
3308                                    object_id: not_found_id,
3309                                    version,
3310                                },
3311                        }) => {
3312                            warn!(
3313                                ?not_found_id,
3314                                ?version,
3315                                object_owner=?object_id,
3316                                field=?id,
3317                                "Skipping dynamic field: referenced genesis object not found"
3318                            );
3319                            continue;
3320                        }
3321                        Err(e) => return Err(e),
3322                    };
3323                    new_dynamic_fields.push(((object_id, id), info));
3324                }
3325                _ => {}
3326            }
3327        }
3328
3329        index_store.insert_genesis_objects(ObjectIndexChanges {
3330            deleted_owners: vec![],
3331            deleted_dynamic_fields: vec![],
3332            new_owners,
3333            new_dynamic_fields,
3334        })
3335    }
3336
3337    /// Attempts to acquire execution lock for an executable transaction.
3338    /// Returns the lock if the transaction is matching current executed epoch
3339    /// Returns None otherwise
3340    pub fn execution_lock_for_executable_transaction(
3341        &self,
3342        transaction: &VerifiedExecutableTransaction,
3343    ) -> IotaResult<ExecutionLockReadGuard<'_>> {
3344        let lock = self
3345            .execution_lock
3346            .try_read()
3347            .map_err(|_| IotaError::ValidatorHaltedAtEpochEnd)?;
3348        if *lock == transaction.auth_sig().epoch() {
3349            Ok(lock)
3350        } else {
3351            Err(IotaError::WrongEpoch {
3352                expected_epoch: *lock,
3353                actual_epoch: transaction.auth_sig().epoch(),
3354            })
3355        }
3356    }
3357
3358    /// Acquires the execution lock for the duration of a transaction signing
3359    /// request. This prevents reconfiguration from starting until we are
3360    /// finished handling the signing request. Otherwise, in-memory lock
3361    /// state could be cleared (by `ObjectLocks::clear_cached_locks`)
3362    /// while we are attempting to acquire locks for the transaction.
3363    pub fn execution_lock_for_signing(&self) -> IotaResult<ExecutionLockReadGuard<'_>> {
3364        self.execution_lock
3365            .try_read()
3366            .map_err(|_| IotaError::ValidatorHaltedAtEpochEnd)
3367    }
3368
3369    pub async fn execution_lock_for_reconfiguration(&self) -> ExecutionLockWriteGuard<'_> {
3370        self.execution_lock.write().await
3371    }
3372
3373    /// Reports a mirror that diverged from the object at the epoch boundary,
3374    /// where the two must agree. Reporting is the remedy: reconfiguration
3375    /// re-seeds the mirror from the object, so failing here would only pin
3376    /// the node to the diverged state. A missing object is fatal instead.
3377    /// Objects cannot be deleted, so the local store lost it and there is
3378    /// nothing to re-seed from. Nodes outside the closing committee are
3379    /// exempt. So is an epoch this node's consensus did not close. A
3380    /// checkpoint catch-up leaves the mirror legitimately behind until the
3381    /// re-seed.
3382    pub(crate) fn check_transaction_deny_rules_consistency(
3383        &self,
3384        cur_epoch_store: &AuthorityPerEpochStore,
3385        epoch_start_configuration: &EpochStartConfiguration,
3386    ) {
3387        if self.is_fullnode(cur_epoch_store) {
3388            return;
3389        }
3390        let Some(walked_deny_rules) = epoch_start_configuration.transaction_deny_rules_state()
3391        else {
3392            if cur_epoch_store
3393                .epoch_start_config()
3394                .transaction_deny_rules_obj_initial_shared_version()
3395                .is_some()
3396            {
3397                fatal!(
3398                    "TransactionDenyRules object existed in epoch {} but is missing from the \
3399                     state walked for the next epoch — the local store is corrupted; restore or \
3400                     state-sync before rejoining",
3401                    cur_epoch_store.epoch(),
3402                );
3403            }
3404            return;
3405        };
3406        // RejectAllTx proves this node's consensus processed every commit of
3407        // the epoch, so the mirror is complete. Otherwise the tail came from
3408        // synced checkpoints and the mirror's lag carries no signal.
3409        if cur_epoch_store
3410            .get_reconfig_state_read_lock_guard()
3411            .should_accept_tx()
3412        {
3413            info!(
3414                "skipping the deny-rule mirror comparison: consensus did not close epoch {} on \
3415                 this node",
3416                cur_epoch_store.epoch(),
3417            );
3418            return;
3419        }
3420        let mirrored_deny_rules = cur_epoch_store.get_mirrored_transaction_deny_rules();
3421        if *walked_deny_rules != *mirrored_deny_rules {
3422            debug_fatal!(
3423                "TransactionDenyRules object diverged from the mirrored state at the end of \
3424                 epoch {}; continuing from the object (walked: {walked_deny_rules:?}, mirrored: \
3425                 {mirrored_deny_rules:?})",
3426                cur_epoch_store.epoch(),
3427            );
3428            cur_epoch_store.metrics.deny_rule_mirror_divergence.set(1);
3429        }
3430    }
3431
3432    /// Reports a `TransactionDenyRulesUpdate` whose execution failed — an
3433    /// invariant violation, the update is built to exclude every expected
3434    /// failure. The object misses the delta until the epoch boundary re-seeds
3435    /// the mirror. Identification is by kind, so the report needs no tracking
3436    /// state and holds across restarts and replays.
3437    ///
3438    /// `expected_effects_digest` is `Some` when these effects were handed to
3439    /// this node with the transaction, which is the case while executing a
3440    /// certified checkpoint: the failure is then part of agreed history, so it
3441    /// is reported without asserting. Effects the node derived itself assert,
3442    /// because only then is the broken invariant its own.
3443    pub(crate) fn report_failed_deny_rule_update_execution(
3444        &self,
3445        transaction: &VerifiedExecutableTransaction,
3446        effects: &TransactionEffects,
3447        expected_effects_digest: Option<TransactionEffectsDigest>,
3448        epoch_store: &AuthorityPerEpochStore,
3449    ) {
3450        if !matches!(
3451            transaction.transaction().kind(),
3452            TransactionKind::TransactionDenyRulesUpdate(_)
3453        ) || effects.status().is_success()
3454        {
3455            return;
3456        }
3457        epoch_store
3458            .metrics
3459            .deny_rule_update_execution_failures
3460            .inc();
3461        if expected_effects_digest.is_some() {
3462            error!(
3463                digest = ?transaction.digest(),
3464                status = ?effects.status(),
3465                "TransactionDenyRulesUpdate failed execution; the object misses its delta until \
3466                 the epoch boundary re-seeds the mirror"
3467            );
3468            return;
3469        }
3470        debug_fatal!(
3471            "TransactionDenyRulesUpdate failed execution; the object misses its delta until the \
3472             epoch boundary re-seeds the mirror (digest: {:?}, status: {:?})",
3473            transaction.digest(),
3474            effects.status(),
3475        );
3476    }
3477
3478    #[instrument(level = "error", skip_all)]
3479    pub async fn reconfigure(
3480        &self,
3481        cur_epoch_store: &AuthorityPerEpochStore,
3482        supported_protocol_versions: SupportedProtocolVersions,
3483        new_committee: Committee,
3484        epoch_start_configuration: EpochStartConfiguration,
3485        state_hasher: Arc<GlobalStateHasher>,
3486        expensive_safety_check_config: &ExpensiveSafetyCheckConfig,
3487        epoch_supply_change: i64,
3488        epoch_last_checkpoint: CheckpointSequenceNumber,
3489    ) -> IotaResult<Arc<AuthorityPerEpochStore>> {
3490        Self::check_protocol_version(
3491            supported_protocol_versions,
3492            epoch_start_configuration
3493                .epoch_start_state()
3494                .protocol_version(),
3495        );
3496        self.metrics.reset_on_reconfigure();
3497        self.committee_store.insert_new_committee(&new_committee)?;
3498
3499        // Wait until no transactions are being executed.
3500        let mut execution_lock = self.execution_lock_for_reconfiguration().await;
3501
3502        // Terminate all epoch-specific tasks (those started with within_alive_epoch).
3503        cur_epoch_store.epoch_terminated().await;
3504
3505        let highest_locally_built_checkpoint_seq = self
3506            .checkpoint_store
3507            .get_latest_locally_computed_checkpoint()?
3508            .map(|c| c.sequence_number())
3509            .unwrap_or(0);
3510
3511        assert!(
3512            epoch_last_checkpoint >= highest_locally_built_checkpoint_seq,
3513            "expected {epoch_last_checkpoint} >= {highest_locally_built_checkpoint_seq}"
3514        );
3515
3516        // Safe to reconfigure now. No transactions are being executed,
3517        // and no epoch-specific tasks are running.
3518
3519        // TODO: revert_uncommitted_epoch_transactions will soon be unnecessary -
3520        // clear_state_end_of_epoch() can simply drop all uncommitted transactions
3521        self.revert_uncommitted_epoch_transactions(cur_epoch_store)
3522            .await?;
3523        self.get_reconfig_api()
3524            .clear_state_end_of_epoch(&execution_lock);
3525        self.check_system_consistency(
3526            cur_epoch_store,
3527            state_hasher,
3528            expensive_safety_check_config,
3529            epoch_supply_change,
3530        )?;
3531        self.check_transaction_deny_rules_consistency(cur_epoch_store, &epoch_start_configuration);
3532
3533        self.get_reconfig_api()
3534            .try_set_epoch_start_configuration(&epoch_start_configuration)?;
3535        self.hand_over_epoch_end_db_snapshot(cur_epoch_store.epoch())
3536            .await;
3537
3538        let new_epoch = new_committee.epoch;
3539        let new_epoch_store = self
3540            .reopen_epoch_db(
3541                cur_epoch_store,
3542                new_committee,
3543                epoch_start_configuration,
3544                expensive_safety_check_config,
3545                epoch_last_checkpoint,
3546            )
3547            .await?;
3548        assert_eq!(new_epoch_store.epoch(), new_epoch);
3549        match self.execution_scheduler.as_ref() {
3550            ExecutionSchedulerWrapper::ExecutionScheduler(_) => {}
3551            ExecutionSchedulerWrapper::TransactionManager(tm) => {
3552                tm.reconfigure(new_epoch);
3553            }
3554        }
3555        *execution_lock = new_epoch;
3556        // drop execution_lock after epoch store was updated
3557        // see also assert in AuthorityState::process_transaction
3558        // on the epoch store and execution lock epoch match
3559        Ok(new_epoch_store)
3560    }
3561
3562    /// Advance the epoch store to the next epoch for testing only.
3563    /// This only manually sets all the places where we have the epoch number.
3564    /// It doesn't properly reconfigure the node, hence should be only used for
3565    /// testing.
3566    pub async fn reconfigure_for_testing(&self) {
3567        self.reconfigure_for_testing_impl(None).await;
3568    }
3569
3570    /// Like [`Self::reconfigure_for_testing`], but the next epoch uses the
3571    /// given protocol config.
3572    pub async fn reconfigure_for_testing_with_protocol_config(
3573        &self,
3574        protocol_config: ProtocolConfig,
3575    ) {
3576        self.reconfigure_for_testing_impl(Some(protocol_config))
3577            .await;
3578    }
3579
3580    async fn reconfigure_for_testing_impl(&self, protocol_config: Option<ProtocolConfig>) {
3581        let mut execution_lock = self.execution_lock_for_reconfiguration().await;
3582        let epoch_store = self.epoch_store_for_testing().clone();
3583        // Default to the epoch store's config, whose override guard may have
3584        // been dropped. Read it under the lock so config and epoch store are
3585        // one snapshot.
3586        let protocol_config =
3587            protocol_config.unwrap_or_else(|| epoch_store.protocol_config().clone());
3588        let _guard =
3589            ProtocolConfig::apply_overrides_for_testing(move |_, _| protocol_config.clone());
3590        let new_epoch_store = epoch_store.new_at_next_epoch_for_testing(
3591            self.get_backing_package_store().clone(),
3592            &self.config.expensive_safety_check_config,
3593            self.checkpoint_store
3594                .get_epoch_last_checkpoint(epoch_store.epoch())
3595                .unwrap()
3596                .map(|c| c.sequence_number())
3597                .unwrap_or_default(),
3598        );
3599        let new_epoch = new_epoch_store.epoch();
3600        match self.execution_scheduler.as_ref() {
3601            ExecutionSchedulerWrapper::ExecutionScheduler(_) => {}
3602            ExecutionSchedulerWrapper::TransactionManager(tm) => {
3603                tm.reconfigure(new_epoch);
3604            }
3605        }
3606        self.epoch_store.store(new_epoch_store);
3607        epoch_store.epoch_terminated().await;
3608        *execution_lock = new_epoch;
3609    }
3610
3611    #[instrument(level = "error", skip_all)]
3612    fn check_system_consistency(
3613        &self,
3614        cur_epoch_store: &AuthorityPerEpochStore,
3615        state_hasher: Arc<GlobalStateHasher>,
3616        expensive_safety_check_config: &ExpensiveSafetyCheckConfig,
3617        epoch_supply_change: i64,
3618    ) -> IotaResult<()> {
3619        info!(
3620            "Performing iota conservation consistency check for epoch {}",
3621            cur_epoch_store.epoch()
3622        );
3623
3624        if cfg!(debug_assertions) {
3625            cur_epoch_store.check_all_executed_transactions_in_checkpoint();
3626        }
3627
3628        self.get_reconfig_api()
3629            .try_expensive_check_iota_conservation(cur_epoch_store, Some(epoch_supply_change))?;
3630
3631        // check for root state hash consistency with live object set
3632        if expensive_safety_check_config.enable_state_consistency_check() {
3633            info!(
3634                "Performing state consistency check for epoch {}",
3635                cur_epoch_store.epoch()
3636            );
3637            self.expensive_check_is_consistent_state(state_hasher, cur_epoch_store);
3638        }
3639
3640        if expensive_safety_check_config.enable_secondary_index_checks() {
3641            if let Some(indexes) = self.indexes.clone() {
3642                verify_indexes(self.get_global_state_hash_store().as_ref(), indexes)
3643                    .expect("secondary indexes are inconsistent");
3644            }
3645        }
3646
3647        Ok(())
3648    }
3649
3650    fn expensive_check_is_consistent_state(
3651        &self,
3652        state_hasher: Arc<GlobalStateHasher>,
3653        cur_epoch_store: &AuthorityPerEpochStore,
3654    ) {
3655        let live_object_set_hash = state_hasher.digest_live_object_set();
3656
3657        let root_state_hash: ECMHLiveObjectSetDigest = self
3658            .get_global_state_hash_store()
3659            .get_root_state_hash_for_epoch(cur_epoch_store.epoch())
3660            .expect("Retrieving root state hash cannot fail")
3661            .expect("Root state hash for epoch must exist")
3662            .1
3663            .digest()
3664            .into();
3665
3666        let is_inconsistent = root_state_hash != live_object_set_hash;
3667        if is_inconsistent {
3668            debug_fatal!(
3669                "Inconsistent state detected: root state hash: {:?}, live object set hash: {:?}",
3670                root_state_hash,
3671                live_object_set_hash
3672            );
3673        } else {
3674            info!("State consistency check passed");
3675        }
3676
3677        state_hasher.set_inconsistent_state(is_inconsistent);
3678    }
3679
3680    pub fn current_epoch_for_testing(&self) -> EpochId {
3681        self.epoch_store_for_testing().epoch()
3682    }
3683
3684    /// Lets the consumer take its database snapshot of the perpetual store as
3685    /// the epoch ends, when there is one. See
3686    /// [`EpochEndDbSnapshotHandle::hand_over`].
3687    #[instrument(level = "error", skip_all)]
3688    async fn hand_over_epoch_end_db_snapshot(&self, epoch: EpochId) {
3689        let Some(handle) = &self.epoch_end_db_snapshots else {
3690            return;
3691        };
3692        let _metrics_guard = self
3693            .metrics
3694            .epoch_end_db_snapshot_handover_latency
3695            .start_timer();
3696        if handle.hand_over(epoch).await == HandOver::Skipped {
3697            self.metrics.epoch_end_db_snapshots_skipped.inc();
3698        }
3699    }
3700
3701    /// Load the current epoch store. This can change during reconfiguration. To
3702    /// ensure that we never end up accessing different epoch stores in a
3703    /// single task, we need to make sure that this is called once per task.
3704    /// Each call needs to be carefully audited to ensure it is
3705    /// the case. This also means we should minimize the number of call-sites.
3706    /// Only call it when there is no way to obtain it from somewhere else.
3707    pub fn load_epoch_store_one_call_per_task(&self) -> Guard<Arc<AuthorityPerEpochStore>> {
3708        self.epoch_store.load()
3709    }
3710
3711    // Load the epoch store, should be used in tests only.
3712    pub fn epoch_store_for_testing(&self) -> Guard<Arc<AuthorityPerEpochStore>> {
3713        self.load_epoch_store_one_call_per_task()
3714    }
3715
3716    pub fn clone_committee_for_testing(&self) -> Committee {
3717        Committee::clone(self.epoch_store_for_testing().committee())
3718    }
3719
3720    #[instrument(level = "trace", skip_all)]
3721    pub fn try_get_object(&self, object_id: &ObjectId) -> IotaResult<Option<Object>> {
3722        self.get_object_store()
3723            .try_get_object(object_id)
3724            .map_err(Into::into)
3725    }
3726
3727    /// Non-fallible version of `try_get_object`.
3728    pub fn get_object(&self, object_id: &ObjectId) -> Option<Object> {
3729        self.try_get_object(object_id)
3730            .expect("storage access failed")
3731    }
3732
3733    pub fn get_iota_system_package_object_ref(&self) -> IotaResult<ObjectReference> {
3734        Ok(self
3735            .try_get_object(&ObjectId::SYSTEM)?
3736            .expect("system package should always exist")
3737            .object_ref())
3738    }
3739
3740    // This function is only used for testing.
3741    pub fn get_iota_system_state_object_for_testing(&self) -> IotaResult<IotaSystemState> {
3742        self.get_object_cache_reader()
3743            .try_get_iota_system_state_object_unsafe()
3744    }
3745
3746    #[instrument(level = "trace", skip_all)]
3747    pub fn get_checkpoint_by_sequence_number(
3748        &self,
3749        sequence_number: CheckpointSequenceNumber,
3750    ) -> IotaResult<Option<VerifiedCheckpoint>> {
3751        Ok(self
3752            .checkpoint_store
3753            .get_checkpoint_by_sequence_number(sequence_number)?)
3754    }
3755
3756    /// Wait for the given transactions to be included in a checkpoint.
3757    ///
3758    /// Returns a mapping from transaction digest to
3759    /// `(checkpoint_sequence_number, checkpoint_timestamp_ms)`.
3760    /// On timeout, returns partial results for any transactions that were
3761    /// already checkpointed.
3762    ///
3763    /// The wait survives epoch boundaries: a transaction in flight at a
3764    /// boundary may only be checkpointed in the next epoch, and still resolves
3765    /// here under the original deadline.
3766    pub async fn wait_for_checkpoint_inclusion(
3767        &self,
3768        digests: &[TransactionDigest],
3769        timeout: Duration,
3770    ) -> IotaResult<BTreeMap<TransactionDigest, (CheckpointSequenceNumber, u64)>> {
3771        let deadline = tokio::time::Instant::now() + timeout;
3772        let mut checkpoint_timestamp_cache = HashMap::<CheckpointSequenceNumber, u64>::new();
3773        let mut results = BTreeMap::new();
3774        let mut remaining = digests.to_vec();
3775        let mut epoch_store = self.load_epoch_store_one_call_per_task().clone();
3776
3777        loop {
3778            let wait = epoch_store.wait_for_transactions_in_checkpoint_with_timeout(
3779                &remaining,
3780                deadline.saturating_duration_since(tokio::time::Instant::now()),
3781                |seq| self.checkpoint_timestamp_ms_cached(seq, &mut checkpoint_timestamp_cache),
3782            );
3783            tokio::select! {
3784                wait_results = wait => {
3785                    for (digest, seq_and_ts) in remaining.iter().zip(wait_results?) {
3786                        if let Some(seq_and_ts) = seq_and_ts {
3787                            results.insert(*digest, seq_and_ts);
3788                        }
3789                    }
3790                    return Ok(results);
3791                }
3792                _ = epoch_store.wait_epoch_terminated() => {}
3793            }
3794
3795            // The epoch ended mid-wait, and this epoch store's notifications
3796            // can no longer fire: whatever is still uncheckpointed here is
3797            // checkpointed in the next epoch, on the next store. Cancelling
3798            // the wait may also have dropped notifications it had already
3799            // received, but the table write precedes each notification, so
3800            // re-reading the table recovers them.
3801            let found = match epoch_store.multi_get_transaction_checkpoint(&remaining) {
3802                Ok(found) => found,
3803                // The table handles were already released. They are released
3804                // long after the epoch's checkpoints are executed, so nothing
3805                // waited on here can still be checkpointed in the old epoch;
3806                // move on to the next store.
3807                Err(IotaError::EpochEnded(_)) => vec![None; remaining.len()],
3808                Err(err) => return Err(err),
3809            };
3810            let mut still_uncheckpointed = Vec::new();
3811            for (digest, found_seq) in remaining.iter().zip(found) {
3812                match found_seq {
3813                    Some(seq) => {
3814                        let ts = self
3815                            .checkpoint_timestamp_ms_cached(seq, &mut checkpoint_timestamp_cache);
3816                        results.insert(*digest, (seq, ts));
3817                    }
3818                    None => still_uncheckpointed.push(*digest),
3819                }
3820            }
3821            remaining = still_uncheckpointed;
3822            if remaining.is_empty() {
3823                return Ok(results);
3824            }
3825
3826            match self
3827                .wait_for_next_epoch_store(epoch_store.epoch(), deadline)
3828                .await
3829            {
3830                Some(next) => epoch_store = next,
3831                None => return Ok(results),
3832            }
3833        }
3834    }
3835
3836    /// Wait for the epoch store to be swapped to an epoch later than
3837    /// `prev_epoch`, returning `None` if `deadline` passes first.
3838    async fn wait_for_next_epoch_store(
3839        &self,
3840        prev_epoch: EpochId,
3841        deadline: tokio::time::Instant,
3842    ) -> Option<Arc<AuthorityPerEpochStore>> {
3843        // There is no notification for the epoch-store swap, and termination
3844        // and swap can come in either order (`reconfigure` terminates the old
3845        // epoch first, `reconfigure_for_testing` swaps first), so the swap is
3846        // polled at this interval.
3847        const EPOCH_STORE_SWAP_POLL_INTERVAL: Duration = Duration::from_millis(100);
3848
3849        loop {
3850            // Deliberately re-loaded on each poll; the one-call-per-task rule
3851            // guards against *unaware* mixing of epoch stores within a task.
3852            let current = self.load_epoch_store_one_call_per_task().clone();
3853            if current.epoch() > prev_epoch {
3854                return Some(current);
3855            }
3856            if tokio::time::Instant::now() >= deadline {
3857                return None;
3858            }
3859            tokio::time::sleep(EPOCH_STORE_SWAP_POLL_INTERVAL).await;
3860        }
3861    }
3862
3863    /// Resolve a checkpoint's timestamp, memoizing lookups in `cache` so
3864    /// multiple transactions in the same checkpoint trigger a single
3865    /// checkpoint summary lookup.
3866    fn checkpoint_timestamp_ms_cached(
3867        &self,
3868        seq: CheckpointSequenceNumber,
3869        cache: &mut HashMap<CheckpointSequenceNumber, u64>,
3870    ) -> u64 {
3871        *cache.entry(seq).or_insert_with(|| {
3872            self.get_checkpoint_by_sequence_number(seq)
3873                .ok()
3874                .flatten()
3875                .map(|c| c.timestamp_ms)
3876                .unwrap_or(0)
3877        })
3878    }
3879
3880    #[instrument(level = "trace", skip_all)]
3881    pub fn get_transaction_checkpoint_for_tests(
3882        &self,
3883        digest: &TransactionDigest,
3884        epoch_store: &AuthorityPerEpochStore,
3885    ) -> IotaResult<Option<VerifiedCheckpoint>> {
3886        let checkpoint = epoch_store.get_transaction_checkpoint(digest)?;
3887        let Some(checkpoint) = checkpoint else {
3888            return Ok(None);
3889        };
3890        let checkpoint = self
3891            .checkpoint_store
3892            .get_checkpoint_by_sequence_number(checkpoint)?;
3893        Ok(checkpoint)
3894    }
3895
3896    #[instrument(level = "trace", skip_all)]
3897    pub fn get_object_read(&self, object_id: &ObjectId) -> IotaResult<ObjectRead> {
3898        Ok(
3899            match self
3900                .get_object_cache_reader()
3901                .try_get_latest_object_or_tombstone(*object_id)?
3902            {
3903                Some((_, ObjectOrTombstone::Object(object))) => {
3904                    let layout = self.get_object_layout(&object)?;
3905                    ObjectRead::Exists(object.object_ref(), object, layout)
3906                }
3907                Some((_, ObjectOrTombstone::Tombstone(objref))) => ObjectRead::Deleted(objref),
3908                None => ObjectRead::NotExists(*object_id),
3909            },
3910        )
3911    }
3912
3913    /// Chain Identifier is the digest of the genesis checkpoint.
3914    pub fn get_chain_identifier(&self) -> ChainIdentifier {
3915        self.chain_identifier
3916    }
3917
3918    #[instrument(level = "trace", skip_all)]
3919    pub fn get_move_object<T>(&self, object_id: &ObjectId) -> IotaResult<T>
3920    where
3921        T: DeserializeOwned,
3922    {
3923        let o = self.get_object_read(object_id)?.into_object()?;
3924        if let Some(move_object) = o.data.as_opt_struct() {
3925            Ok(bcs::from_bytes(move_object.contents()).map_err(|e| {
3926                IotaError::ObjectDeserialization {
3927                    error: format!("{e}"),
3928                }
3929            })?)
3930        } else {
3931            Err(IotaError::ObjectDeserialization {
3932                error: format!("Provided object : [{object_id}] is not a Move object."),
3933            })
3934        }
3935    }
3936
3937    /// This function aims to serve rpc reads on past objects and
3938    /// we don't expect it to be called for other purposes.
3939    /// Depending on the object pruning policies that will be enforced in the
3940    /// future there is no software-level guarantee/SLA to retrieve an object
3941    /// with an old version even if it exists/existed.
3942    #[instrument(level = "trace", skip_all)]
3943    pub fn get_past_object_read(
3944        &self,
3945        object_id: &ObjectId,
3946        version: Version,
3947    ) -> IotaResult<PastObjectRead> {
3948        // Firstly we see if the object ever existed by getting its latest data
3949        let Some(obj_ref) = self
3950            .get_object_cache_reader()
3951            .try_get_latest_object_ref_or_tombstone(*object_id)?
3952        else {
3953            return Ok(PastObjectRead::ObjectNotExists(*object_id));
3954        };
3955
3956        if version > obj_ref.version {
3957            return Ok(PastObjectRead::VersionTooHigh {
3958                object_id: *object_id,
3959                asked_version: version,
3960                latest_version: obj_ref.version,
3961            });
3962        }
3963
3964        if version < obj_ref.version {
3965            // Read past objects
3966            return Ok(match self.read_object_at_version(object_id, version)? {
3967                Some((object, layout)) => {
3968                    let obj_ref = object.object_ref();
3969                    PastObjectRead::VersionFound(obj_ref, object, layout)
3970                }
3971
3972                None => PastObjectRead::VersionNotFound(*object_id, version),
3973            });
3974        }
3975
3976        if !obj_ref.digest.is_alive() {
3977            return Ok(PastObjectRead::ObjectDeleted(obj_ref));
3978        }
3979
3980        match self.read_object_at_version(object_id, obj_ref.version)? {
3981            Some((object, layout)) => Ok(PastObjectRead::VersionFound(obj_ref, object, layout)),
3982            None => {
3983                debug_fatal!(
3984                    "Object with in parent_entry is missing from object store, datastore is \
3985                     inconsistent",
3986                );
3987                Err(UserInputError::ObjectNotFound {
3988                    object_id: *object_id,
3989                    version: Some(obj_ref.version),
3990                }
3991                .into())
3992            }
3993        }
3994    }
3995
3996    #[instrument(level = "trace", skip_all)]
3997    fn read_object_at_version(
3998        &self,
3999        object_id: &ObjectId,
4000        version: Version,
4001    ) -> IotaResult<Option<(Object, Option<MoveStructLayout>)>> {
4002        let Some(object) = self
4003            .get_object_cache_reader()
4004            .try_get_object_by_key(object_id, version)?
4005        else {
4006            return Ok(None);
4007        };
4008
4009        let layout = self.get_object_layout(&object)?;
4010        Ok(Some((object, layout)))
4011    }
4012
4013    fn get_object_layout(&self, object: &Object) -> IotaResult<Option<MoveStructLayout>> {
4014        let layout = object
4015            .data
4016            .as_opt_struct()
4017            .map(|object| {
4018                into_struct_layout(
4019                    self.load_epoch_store_one_call_per_task()
4020                        .executor()
4021                        // TODO(cache) - must read through cache
4022                        .type_layout_resolver(Box::new(self.get_backing_package_store().as_ref()))
4023                        .get_annotated_layout(object.struct_tag())?,
4024                )
4025            })
4026            .transpose()?;
4027        Ok(layout)
4028    }
4029
4030    fn get_owner_at_version(&self, object_id: &ObjectId, version: Version) -> IotaResult<Owner> {
4031        self.get_object_store()
4032            .try_get_object_by_key(object_id, version)?
4033            .ok_or_else(|| {
4034                IotaError::from(UserInputError::ObjectNotFound {
4035                    object_id: *object_id,
4036                    version: Some(version),
4037                })
4038            })
4039            .map(|o| o.owner)
4040    }
4041
4042    #[instrument(level = "trace", skip_all)]
4043    pub fn get_owner_objects(
4044        &self,
4045        owner: Address,
4046        // If `Some`, the query will start from the next item after the specified cursor
4047        cursor: Option<ObjectId>,
4048        limit: usize,
4049        filter: Option<IotaObjectDataFilter>,
4050    ) -> IotaResult<Vec<ObjectInfo>> {
4051        if let Some(indexes) = &self.indexes {
4052            indexes.get_owner_objects(owner, cursor, limit, filter)
4053        } else {
4054            Err(IotaError::IndexStoreNotAvailable)
4055        }
4056    }
4057
4058    #[instrument(level = "trace", skip_all)]
4059    pub fn get_owned_coins_iterator_with_cursor(
4060        &self,
4061        owner: Address,
4062        // If `Some`, the query will start from the next item after the specified cursor
4063        cursor: (String, ObjectId),
4064        limit: usize,
4065        one_coin_type_only: bool,
4066    ) -> IotaResult<impl Iterator<Item = (String, ObjectId, CoinInfo)> + '_> {
4067        if let Some(indexes) = &self.indexes {
4068            indexes.get_owned_coins_iterator_with_cursor(owner, cursor, limit, one_coin_type_only)
4069        } else {
4070            Err(IotaError::IndexStoreNotAvailable)
4071        }
4072    }
4073
4074    #[instrument(level = "trace", skip_all)]
4075    pub fn get_owner_objects_iterator(
4076        &self,
4077        owner: Address,
4078        // If `Some`, the query will start from the next item after the specified cursor
4079        cursor: Option<ObjectId>,
4080        filter: Option<IotaObjectDataFilter>,
4081    ) -> IotaResult<impl Iterator<Item = ObjectInfo> + '_> {
4082        if let Some(indexes) = &self.indexes {
4083            indexes.get_owner_objects_iterator(owner, cursor, filter)
4084        } else {
4085            Err(IotaError::IndexStoreNotAvailable)
4086        }
4087    }
4088
4089    #[instrument(level = "trace", skip_all)]
4090    pub fn get_move_objects<T>(&self, owner: Address, tag: StructTag) -> IotaResult<Vec<T>>
4091    where
4092        T: DeserializeOwned,
4093    {
4094        let object_ids = self
4095            .get_owner_objects_iterator(owner, None, None)?
4096            .filter(|o| match &o.object_type {
4097                ObjectType::Struct(s) => *s == tag,
4098                ObjectType::Package => false,
4099            })
4100            .map(|info| ObjectKey(info.object_id, info.version))
4101            .collect::<Vec<_>>();
4102        let mut move_objects = vec![];
4103
4104        let objects = self
4105            .get_object_store()
4106            .try_multi_get_objects_by_key(&object_ids)?;
4107
4108        for (o, id) in objects.into_iter().zip(object_ids) {
4109            let object = o.ok_or_else(|| {
4110                IotaError::from(UserInputError::ObjectNotFound {
4111                    object_id: id.0,
4112                    version: Some(id.1),
4113                })
4114            })?;
4115            let move_object = object.data.as_opt_struct().ok_or_else(|| {
4116                IotaError::from(UserInputError::MovePackageAsObject { object_id: id.0 })
4117            })?;
4118            move_objects.push(bcs::from_bytes(move_object.contents()).map_err(|e| {
4119                IotaError::ObjectDeserialization {
4120                    error: format!("{e}"),
4121                }
4122            })?);
4123        }
4124        Ok(move_objects)
4125    }
4126
4127    #[instrument(level = "trace", skip_all)]
4128    pub fn get_dynamic_fields(
4129        &self,
4130        owner: ObjectId,
4131        // If `Some`, the query will start from the next item after the specified cursor
4132        cursor: Option<ObjectId>,
4133        limit: usize,
4134    ) -> IotaResult<Vec<(ObjectId, DynamicFieldInfo)>> {
4135        Ok(self
4136            .get_dynamic_fields_iterator(owner, cursor)?
4137            .take(limit)
4138            .collect::<Result<Vec<_>, _>>()?)
4139    }
4140
4141    fn get_dynamic_fields_iterator(
4142        &self,
4143        owner: ObjectId,
4144        // If `Some`, the query will start from the next item after the specified cursor
4145        cursor: Option<ObjectId>,
4146    ) -> IotaResult<impl Iterator<Item = Result<(ObjectId, DynamicFieldInfo), TypedStoreError>> + '_>
4147    {
4148        if let Some(indexes) = &self.indexes {
4149            indexes.get_dynamic_fields_iterator(owner, cursor)
4150        } else {
4151            Err(IotaError::IndexStoreNotAvailable)
4152        }
4153    }
4154
4155    #[instrument(level = "trace", skip_all)]
4156    pub fn get_dynamic_field_object_id(
4157        &self,
4158        owner: ObjectId,
4159        name_type: TypeTag,
4160        name_bcs_bytes: &[u8],
4161    ) -> IotaResult<Option<ObjectId>> {
4162        if let Some(indexes) = &self.indexes {
4163            indexes.get_dynamic_field_object_id(owner, name_type, name_bcs_bytes)
4164        } else {
4165            Err(IotaError::IndexStoreNotAvailable)
4166        }
4167    }
4168
4169    #[instrument(level = "trace", skip_all)]
4170    pub fn get_total_transaction_blocks(&self) -> IotaResult<u64> {
4171        Ok(self.get_indexes()?.next_sequence_number())
4172    }
4173
4174    #[instrument(level = "trace", skip_all)]
4175    pub async fn get_executed_transaction_and_effects(
4176        &self,
4177        digest: TransactionDigest,
4178        kv_store: Arc<TransactionKeyValueStore>,
4179    ) -> IotaResult<(TransactionEnvelope, TransactionEffects)> {
4180        let transaction = kv_store.get_tx(digest).await?;
4181        let effects = kv_store.get_fx_by_tx_digest(digest).await?;
4182        Ok((transaction, effects))
4183    }
4184
4185    #[instrument(level = "trace", skip_all)]
4186    pub fn multi_get_checkpoint_by_sequence_number(
4187        &self,
4188        sequence_numbers: &[CheckpointSequenceNumber],
4189    ) -> IotaResult<Vec<Option<VerifiedCheckpoint>>> {
4190        Ok(self
4191            .checkpoint_store
4192            .multi_get_checkpoint_by_sequence_number(sequence_numbers)?)
4193    }
4194
4195    #[instrument(level = "trace", skip_all)]
4196    pub fn get_transaction_events(
4197        &self,
4198        digest: &TransactionDigest,
4199    ) -> IotaResult<TransactionEvents> {
4200        self.get_transaction_cache_reader()
4201            .try_get_events(digest)?
4202            .ok_or(IotaError::TransactionEventsNotFound { digest: *digest })
4203    }
4204
4205    pub fn get_transaction_input_objects(
4206        &self,
4207        effects: &TransactionEffects,
4208    ) -> anyhow::Result<Vec<Object>> {
4209        iota_types::storage::get_transaction_input_objects(self.get_object_store(), effects)
4210            .map_err(Into::into)
4211    }
4212
4213    pub fn get_transaction_output_objects(
4214        &self,
4215        effects: &TransactionEffects,
4216    ) -> anyhow::Result<Vec<Object>> {
4217        iota_types::storage::get_transaction_output_objects(self.get_object_store(), effects)
4218            .map_err(Into::into)
4219    }
4220
4221    fn get_indexes(&self) -> IotaResult<Arc<IndexStore>> {
4222        match &self.indexes {
4223            Some(i) => Ok(i.clone()),
4224            None => Err(IotaError::UnsupportedFeature {
4225                error: "extended object indexing is not enabled on this server".into(),
4226            }),
4227        }
4228    }
4229
4230    pub async fn get_transactions_for_tests(
4231        self: &Arc<Self>,
4232        filter: Option<TransactionFilter>,
4233        cursor: Option<TransactionDigest>,
4234        limit: Option<usize>,
4235        reverse: bool,
4236    ) -> IotaResult<Vec<TransactionDigest>> {
4237        let metrics = KeyValueStoreMetrics::new_for_tests();
4238        let kv_store = Arc::new(TransactionKeyValueStore::new(
4239            "rocksdb",
4240            metrics,
4241            self.clone(),
4242        ));
4243        self.get_transactions(&kv_store, filter, cursor, limit, reverse)
4244            .await
4245    }
4246
4247    #[instrument(level = "trace", skip_all)]
4248    pub async fn get_transactions(
4249        &self,
4250        kv_store: &Arc<TransactionKeyValueStore>,
4251        filter: Option<TransactionFilter>,
4252        // If `Some`, the query will start from the next item after the specified cursor
4253        cursor: Option<TransactionDigest>,
4254        limit: Option<usize>,
4255        reverse: bool,
4256    ) -> IotaResult<Vec<TransactionDigest>> {
4257        if let Some(TransactionFilter::Checkpoint(sequence_number)) = filter {
4258            let checkpoint_contents = kv_store.get_checkpoint_contents(sequence_number).await?;
4259            let iter = checkpoint_contents.iter().map(|c| c.transaction);
4260            if reverse {
4261                let iter = iter
4262                    .rev()
4263                    .skip_while(|d| cursor.is_some() && Some(*d) != cursor)
4264                    .skip(usize::from(cursor.is_some()));
4265                return Ok(iter.take(limit.unwrap_or(usize::MAX)).collect());
4266            } else {
4267                let iter = iter
4268                    .skip_while(|d| cursor.is_some() && Some(*d) != cursor)
4269                    .skip(usize::from(cursor.is_some()));
4270                return Ok(iter.take(limit.unwrap_or(usize::MAX)).collect());
4271            }
4272        }
4273        self.get_indexes()?
4274            .get_transactions(filter, cursor, limit, reverse)
4275    }
4276
4277    pub fn get_checkpoint_store(&self) -> &Arc<CheckpointStore> {
4278        &self.checkpoint_store
4279    }
4280
4281    /// The store pruner; the checkpoint executor uses it to nudge the pruner
4282    /// after each checkpoint.
4283    pub fn pruner(&self) -> &AuthorityStorePruner {
4284        &self.pruner
4285    }
4286
4287    pub fn get_latest_checkpoint_sequence_number(&self) -> IotaResult<CheckpointSequenceNumber> {
4288        self.get_checkpoint_store()
4289            .get_highest_executed_checkpoint_seq_number()?
4290            .ok_or(IotaError::UserInput {
4291                error: UserInputError::LatestCheckpointSequenceNumberNotFound,
4292            })
4293    }
4294
4295    #[cfg(msim)]
4296    pub fn get_highest_pruned_checkpoint_for_testing(
4297        &self,
4298    ) -> IotaResult<CheckpointSequenceNumber> {
4299        self.database_for_testing()
4300            .perpetual_tables
4301            .get_highest_pruned_checkpoint()
4302            .map(|c| c.unwrap_or(0))
4303            .map_err(Into::into)
4304    }
4305
4306    #[instrument(level = "trace", skip_all)]
4307    pub fn get_checkpoint_summary_by_sequence_number(
4308        &self,
4309        sequence_number: CheckpointSequenceNumber,
4310    ) -> IotaResult<CheckpointSummary> {
4311        let verified_checkpoint = self
4312            .get_checkpoint_store()
4313            .get_checkpoint_by_sequence_number(sequence_number)?;
4314        match verified_checkpoint {
4315            Some(verified_checkpoint) => Ok(verified_checkpoint.into_inner().into_data()),
4316            None => Err(IotaError::UserInput {
4317                error: UserInputError::VerifiedCheckpointNotFound(sequence_number),
4318            }),
4319        }
4320    }
4321
4322    #[instrument(level = "trace", skip_all)]
4323    pub fn get_checkpoint_summary_by_digest(
4324        &self,
4325        digest: CheckpointDigest,
4326    ) -> IotaResult<CheckpointSummary> {
4327        let verified_checkpoint = self
4328            .get_checkpoint_store()
4329            .get_checkpoint_by_digest(&digest)?;
4330        match verified_checkpoint {
4331            Some(verified_checkpoint) => Ok(verified_checkpoint.into_inner().into_data()),
4332            None => Err(IotaError::UserInput {
4333                error: UserInputError::VerifiedCheckpointDigestNotFound(Base58::encode(digest)),
4334            }),
4335        }
4336    }
4337
4338    #[instrument(level = "trace", skip_all)]
4339    pub fn find_publish_txn_digest(&self, package_id: ObjectId) -> IotaResult<TransactionDigest> {
4340        if package_id.is_system_package() {
4341            return self.find_genesis_txn_digest();
4342        }
4343        Ok(self
4344            .get_object_read(&package_id)?
4345            .into_object()?
4346            .previous_transaction)
4347    }
4348
4349    #[instrument(level = "trace", skip_all)]
4350    pub fn find_genesis_txn_digest(&self) -> IotaResult<TransactionDigest> {
4351        let summary = self
4352            .get_verified_checkpoint_by_sequence_number(0)?
4353            .into_message();
4354        let content = self.get_checkpoint_contents(summary.contents_digest)?;
4355        let genesis_transaction = content.enumerate_transactions(&summary).next();
4356        Ok(genesis_transaction
4357            .ok_or(IotaError::UserInput {
4358                error: UserInputError::GenesisTransactionNotFound,
4359            })?
4360            .1
4361            .transaction)
4362    }
4363
4364    #[instrument(level = "trace", skip_all)]
4365    pub fn get_verified_checkpoint_by_sequence_number(
4366        &self,
4367        sequence_number: CheckpointSequenceNumber,
4368    ) -> IotaResult<VerifiedCheckpoint> {
4369        let verified_checkpoint = self
4370            .get_checkpoint_store()
4371            .get_checkpoint_by_sequence_number(sequence_number)?;
4372        match verified_checkpoint {
4373            Some(verified_checkpoint) => Ok(verified_checkpoint),
4374            None => Err(IotaError::UserInput {
4375                error: UserInputError::VerifiedCheckpointNotFound(sequence_number),
4376            }),
4377        }
4378    }
4379
4380    #[instrument(level = "trace", skip_all)]
4381    pub fn get_verified_checkpoint_summary_by_digest(
4382        &self,
4383        digest: CheckpointDigest,
4384    ) -> IotaResult<VerifiedCheckpoint> {
4385        let verified_checkpoint = self
4386            .get_checkpoint_store()
4387            .get_checkpoint_by_digest(&digest)?;
4388        match verified_checkpoint {
4389            Some(verified_checkpoint) => Ok(verified_checkpoint),
4390            None => Err(IotaError::UserInput {
4391                error: UserInputError::VerifiedCheckpointDigestNotFound(Base58::encode(digest)),
4392            }),
4393        }
4394    }
4395
4396    #[instrument(level = "trace", skip_all)]
4397    pub fn get_checkpoint_contents(
4398        &self,
4399        digest: CheckpointContentsDigest,
4400    ) -> IotaResult<CheckpointContents> {
4401        self.get_checkpoint_store()
4402            .get_checkpoint_contents(&digest)?
4403            .ok_or(IotaError::UserInput {
4404                error: UserInputError::CheckpointContentsNotFound(digest),
4405            })
4406    }
4407
4408    #[instrument(level = "trace", skip_all)]
4409    pub fn get_checkpoint_contents_by_sequence_number(
4410        &self,
4411        sequence_number: CheckpointSequenceNumber,
4412    ) -> IotaResult<CheckpointContents> {
4413        let verified_checkpoint = self
4414            .get_checkpoint_store()
4415            .get_checkpoint_by_sequence_number(sequence_number)?;
4416        match verified_checkpoint {
4417            Some(verified_checkpoint) => {
4418                let contents_digest = verified_checkpoint.into_inner().contents_digest;
4419                self.get_checkpoint_contents(contents_digest)
4420            }
4421            None => Err(IotaError::UserInput {
4422                error: UserInputError::VerifiedCheckpointNotFound(sequence_number),
4423            }),
4424        }
4425    }
4426
4427    #[instrument(level = "trace", skip_all)]
4428    pub async fn query_events(
4429        &self,
4430        kv_store: &Arc<TransactionKeyValueStore>,
4431        query: EventFilter,
4432        // If `Some`, the query will start from the next item after the specified cursor
4433        cursor: Option<EventID>,
4434        limit: usize,
4435        descending: bool,
4436    ) -> IotaResult<Vec<IotaEvent>> {
4437        let index_store = self.get_indexes()?;
4438
4439        // Get the tx_num from tx_digest
4440        let (tx_num, event_num) = if let Some(cursor) = cursor.as_ref() {
4441            let tx_seq = index_store.get_transaction_seq(&cursor.tx_digest)?.ok_or(
4442                IotaError::TransactionNotFound {
4443                    digest: cursor.tx_digest,
4444                },
4445            )?;
4446            (tx_seq, cursor.event_seq as usize)
4447        } else if descending {
4448            (u64::MAX, usize::MAX)
4449        } else {
4450            (0, 0)
4451        };
4452
4453        let limit = limit + 1;
4454        let mut event_keys = match query {
4455            EventFilter::All(filters) => {
4456                if filters.is_empty() {
4457                    index_store.all_events(tx_num, event_num, limit, descending)?
4458                } else {
4459                    return Err(IotaError::UserInput {
4460                        error: UserInputError::Unsupported(
4461                            "This query type does not currently support filter combinations"
4462                                .to_string(),
4463                        ),
4464                    });
4465                }
4466            }
4467            EventFilter::Transaction(digest) => {
4468                index_store.events_by_transaction(&digest, tx_num, event_num, limit, descending)?
4469            }
4470            EventFilter::MoveModule { package, module } => {
4471                let module_id = ModuleId::new(
4472                    AccountAddress::new(package.into_bytes()),
4473                    move_core_types::identifier::Identifier::new(module.as_str()).unwrap(),
4474                );
4475                index_store.events_by_module_id(&module_id, tx_num, event_num, limit, descending)?
4476            }
4477            EventFilter::MoveEventType(struct_name) => index_store
4478                .events_by_move_event_struct_name(
4479                    &struct_name,
4480                    tx_num,
4481                    event_num,
4482                    limit,
4483                    descending,
4484                )?,
4485            EventFilter::Sender(sender) => {
4486                index_store.events_by_sender(&sender, tx_num, event_num, limit, descending)?
4487            }
4488            EventFilter::TimeRange {
4489                start_time,
4490                end_time,
4491            } => index_store
4492                .event_iterator(start_time, end_time, tx_num, event_num, limit, descending)?,
4493            EventFilter::MoveEventModule { package, module } => index_store
4494                .events_by_move_event_module(
4495                    &ModuleId::new(
4496                        AccountAddress::new(package.into_bytes()),
4497                        move_core_types::identifier::Identifier::new(module.as_str()).unwrap(),
4498                    ),
4499                    tx_num,
4500                    event_num,
4501                    limit,
4502                    descending,
4503                )?,
4504            // not using "_ =>" because we want to make sure we remember to add new variants here
4505            EventFilter::Package(_)
4506            | EventFilter::MoveEventField { .. }
4507            | EventFilter::Any(_)
4508            | EventFilter::And(_, _)
4509            | EventFilter::Or(_, _) => {
4510                return Err(IotaError::UserInput {
4511                    error: UserInputError::Unsupported(
4512                        "This query type is not supported by the full node.".to_string(),
4513                    ),
4514                });
4515            }
4516        };
4517
4518        // skip one event if exclusive cursor is provided,
4519        // otherwise truncate to the original limit.
4520        if cursor.is_some() {
4521            if !event_keys.is_empty() {
4522                event_keys.remove(0);
4523            }
4524        } else {
4525            event_keys.truncate(limit - 1);
4526        }
4527
4528        // get the unique set of digests from the event_keys
4529        let transaction_digests = event_keys
4530            .iter()
4531            .map(|(_, digest, _, _)| *digest)
4532            .collect::<HashSet<_>>()
4533            .into_iter()
4534            .collect::<Vec<_>>();
4535
4536        let events = kv_store
4537            .multi_get_events_by_tx_digests(&transaction_digests)
4538            .await?;
4539
4540        let events_map: HashMap<_, _> =
4541            transaction_digests.iter().zip(events.into_iter()).collect();
4542
4543        let stored_events = event_keys
4544            .into_iter()
4545            .map(|k| {
4546                (
4547                    k,
4548                    events_map
4549                        .get(&k.1)
4550                        .expect("fetched digest is missing")
4551                        .clone()
4552                        .and_then(|e| e.get(k.2).cloned()),
4553                )
4554            })
4555            .map(
4556                |((_event_digest, tx_digest, event_seq, timestamp), event)| {
4557                    event
4558                        .map(|e| (e, tx_digest, event_seq, timestamp))
4559                        .ok_or(IotaError::TransactionEventsNotFound { digest: tx_digest })
4560                },
4561            )
4562            .collect::<Result<Vec<_>, _>>()?;
4563
4564        let epoch_store = self.load_epoch_store_one_call_per_task();
4565        let backing_store = self.get_backing_package_store().as_ref();
4566        let mut layout_resolver = epoch_store
4567            .executor()
4568            .type_layout_resolver(Box::new(backing_store));
4569        let mut events = vec![];
4570        for (e, tx_digest, event_seq, timestamp) in stored_events.into_iter() {
4571            events.push(IotaEvent::try_from(
4572                e.clone(),
4573                tx_digest,
4574                event_seq as u64,
4575                Some(timestamp),
4576                layout_resolver.get_annotated_layout(&e.struct_tag)?,
4577            )?)
4578        }
4579        Ok(events)
4580    }
4581
4582    pub fn insert_genesis_object(&self, object: Object) {
4583        self.get_reconfig_api()
4584            .try_insert_genesis_object(object)
4585            .expect("Cannot insert genesis object")
4586    }
4587
4588    pub fn insert_genesis_objects(&self, objects: &[Object]) {
4589        for o in objects {
4590            self.insert_genesis_object(o.clone());
4591        }
4592    }
4593
4594    /// Make a status response for a transaction
4595    #[instrument(level = "trace", skip_all)]
4596    pub fn get_transaction_status(
4597        &self,
4598        transaction_digest: &TransactionDigest,
4599        epoch_store: &Arc<AuthorityPerEpochStore>,
4600    ) -> IotaResult<Option<(SenderSignedTransaction, TransactionStatus)>> {
4601        // TODO: In the case of read path, we should not have to re-sign the effects.
4602        if let Some(effects) =
4603            self.get_signed_effects_and_maybe_resign(transaction_digest, epoch_store)?
4604        {
4605            if let Some(transaction) = self
4606                .get_transaction_cache_reader()
4607                .try_get_transaction_block(transaction_digest)?
4608            {
4609                let cert_sig = epoch_store.get_transaction_cert_sig(transaction_digest)?;
4610                let events = if effects.events_digest().is_some() {
4611                    self.get_transaction_events(effects.transaction_digest())?
4612                } else {
4613                    TransactionEvents::default()
4614                };
4615                return Ok(Some((
4616                    (*transaction).clone().into_message(),
4617                    TransactionStatus::Executed(cert_sig, effects.into_inner(), events),
4618                )));
4619            } else {
4620                // The read of effects and read of transaction are not atomic. It's possible
4621                // that we reverted the transaction (during epoch change) in
4622                // between the above two reads, and we end up having effects but
4623                // not transaction. In this case, we just fall through.
4624                debug!(tx_digest=?transaction_digest, "Signed effects exist but no transaction found");
4625            }
4626        }
4627        if let Some(signed) = epoch_store.get_signed_transaction(transaction_digest)? {
4628            self.metrics.tx_already_processed.inc();
4629            let (transaction, sig) = signed.into_inner().into_data_and_sig();
4630            Ok(Some((transaction, TransactionStatus::Signed(sig))))
4631        } else {
4632            Ok(None)
4633        }
4634    }
4635
4636    /// Get the signed effects of the given transaction. If the effects was
4637    /// signed in a previous epoch, re-sign it so that the caller is able to
4638    /// form a cert of the effects in the current epoch.
4639    #[instrument(level = "trace", skip_all)]
4640    pub fn get_signed_effects_and_maybe_resign(
4641        &self,
4642        transaction_digest: &TransactionDigest,
4643        epoch_store: &Arc<AuthorityPerEpochStore>,
4644    ) -> IotaResult<Option<VerifiedSignedTransactionEffects>> {
4645        let effects = self
4646            .get_transaction_cache_reader()
4647            .try_get_executed_effects(transaction_digest)?;
4648        match effects {
4649            Some(effects) => {
4650                // If the transaction was executed in previous epochs, the validator will
4651                // re-sign the effects with new current epoch so that a client is always able to
4652                // obtain an effects certificate at the current epoch.
4653                //
4654                // Why is this necessary? Consider the following case:
4655                // - assume there are 4 validators
4656                // - Quorum driver gets 2 signed effects before reconfig halt
4657                // - The tx makes it into final checkpoint.
4658                // - 2 validators go away and are replaced in the new epoch.
4659                // - The new epoch begins.
4660                // - The quorum driver cannot complete the partial effects cert from the previous
4661                //   epoch, because it may not be able to reach either of the 2 former validators.
4662                // - But, if the 2 validators that stayed are willing to re-sign the effects in the
4663                //   new epoch, the QD can make a new effects cert and return it to the client.
4664                //
4665                // This is a considered a short-term workaround. Eventually, Quorum Driver
4666                // should be able to return either an effects certificate, -or-
4667                // a proof of inclusion in a checkpoint. In the case above, the
4668                // Quorum Driver would return a proof of inclusion in the final
4669                // checkpoint, and this code would no longer be necessary.
4670                if effects.epoch() != epoch_store.epoch() {
4671                    debug!(
4672                        tx_digest=?transaction_digest,
4673                        effects_epoch=?effects.epoch(),
4674                        epoch=?epoch_store.epoch(),
4675                        "Re-signing the effects with the current epoch"
4676                    );
4677                }
4678                Ok(Some(self.sign_effects(effects, epoch_store)?))
4679            }
4680            None => Ok(None),
4681        }
4682    }
4683
4684    /// A client aggregating effects signatures towards a quorum assumes
4685    /// finality once it collects 2f+1 of them, so within an epoch this
4686    /// validator must never assert two different effects for the same
4687    /// transaction on any RPC surface, signed or unsigned. Executed effects
4688    /// can change across a restart if an uncommitted transaction is
4689    /// re-executed with divergent results (e.g. by a new binary), so every
4690    /// effects-reporting path calls this before returning effects, and
4691    /// refuses to contradict a signature that may already be in a client's
4692    /// hands.
4693    pub fn check_effects_against_previously_signed(
4694        &self,
4695        epoch_store: &AuthorityPerEpochStore,
4696        tx_digest: &TransactionDigest,
4697        effects_digest: &TransactionEffectsDigest,
4698        surface: &'static str,
4699    ) -> IotaResult<()> {
4700        if let Some(previously_signed_digest) = epoch_store.get_signed_effects_digest(tx_digest)? {
4701            if previously_signed_digest != *effects_digest {
4702                self.metrics
4703                    .signed_effects_equivocation_prevented
4704                    .with_label_values(&[surface])
4705                    .inc();
4706                error!(
4707                    ?tx_digest,
4708                    ?previously_signed_digest,
4709                    executed_digest = ?effects_digest,
4710                    surface,
4711                    "refusing to report effects that differ from previously signed effects"
4712                );
4713                return Err(IotaError::GenericAuthority {
4714                    error: format!(
4715                        "Refusing to report effects for transaction {tx_digest}: effects digest \
4716                         {effects_digest} differs from previously signed effects digest \
4717                         {previously_signed_digest}"
4718                    ),
4719                });
4720            }
4721        }
4722        Ok(())
4723    }
4724
4725    #[instrument(level = "trace", skip_all)]
4726    pub(crate) fn sign_effects(
4727        &self,
4728        effects: TransactionEffects,
4729        epoch_store: &Arc<AuthorityPerEpochStore>,
4730    ) -> IotaResult<VerifiedSignedTransactionEffects> {
4731        let tx_digest = *effects.transaction_digest();
4732
4733        self.check_effects_against_previously_signed(
4734            epoch_store,
4735            &tx_digest,
4736            &effects.digest(),
4737            "sign_effects",
4738        )?;
4739
4740        let signed_effects = match epoch_store.get_effects_signature(&tx_digest)? {
4741            Some(sig) => {
4742                debug_assert!(sig.epoch == epoch_store.epoch());
4743                SignedTransactionEffects::new_from_data_and_sig(effects, sig)
4744            }
4745            _ => {
4746                let sig = AuthoritySignInfo::new(
4747                    epoch_store.epoch(),
4748                    &effects,
4749                    Intent::iota_app(IntentScope::TransactionEffects),
4750                    self.name,
4751                    &*self.secret,
4752                );
4753
4754                let effects = SignedTransactionEffects::new_from_data_and_sig(effects, sig.clone());
4755
4756                epoch_store.insert_effects_digest_and_signature(
4757                    &tx_digest,
4758                    effects.digest(),
4759                    &sig,
4760                )?;
4761
4762                effects
4763            }
4764        };
4765
4766        Ok(VerifiedSignedTransactionEffects::new_unchecked(
4767            signed_effects,
4768        ))
4769    }
4770
4771    // Returns coin objects for indexing for fullnode if indexing is enabled.
4772    #[instrument(level = "trace", skip_all)]
4773    fn fullnode_only_get_tx_coins_for_indexing(
4774        &self,
4775        effects: &TransactionEffects,
4776        inner_temporary_store: &InnerTemporaryStore,
4777        epoch_store: &Arc<AuthorityPerEpochStore>,
4778    ) -> Option<TxCoins> {
4779        if self.indexes.is_none() || self.is_committee_validator(epoch_store) {
4780            return None;
4781        }
4782        let written_coin_objects = inner_temporary_store
4783            .written
4784            .iter()
4785            .filter_map(|(k, v)| {
4786                if v.is_coin() {
4787                    Some((*k, v.clone()))
4788                } else {
4789                    None
4790                }
4791            })
4792            .collect();
4793        let mut input_coin_objects = inner_temporary_store
4794            .input_objects
4795            .iter()
4796            .filter_map(|(k, v)| {
4797                if v.is_coin() {
4798                    Some((*k, v.clone()))
4799                } else {
4800                    None
4801                }
4802            })
4803            .collect::<ObjectMap>();
4804
4805        // Check for receiving objects that were actually used and modified during
4806        // execution. Their updated version will already showup in
4807        // "written_coins" but their input isn't included in the set of input
4808        // objects in a inner_temporary_store.
4809        for modified in effects.modified_at_versions() {
4810            let (object_id, version) = (*modified.object_id(), modified.version());
4811            if inner_temporary_store
4812                .loaded_runtime_objects
4813                .contains_key(&object_id)
4814            {
4815                if let Some(object) = self
4816                    .get_object_store()
4817                    .get_object_by_key(&object_id, version)
4818                {
4819                    if object.is_coin() {
4820                        input_coin_objects.insert(object_id, object);
4821                    }
4822                }
4823            }
4824        }
4825
4826        Some((input_coin_objects, written_coin_objects))
4827    }
4828
4829    /// Get the transaction envelope that currently locks the given object, if
4830    /// any. Since object locks are only valid for one epoch, we also need
4831    /// the epoch_id in the query. Returns UserInputError::ObjectNotFound if
4832    /// no lock records for the given object can be found.
4833    /// Returns UserInputError::ObjectVersionUnavailableForConsumption if the
4834    /// object record is at a different version.
4835    /// Returns Some(VerifiedEnvelope) if the given ObjectReference is locked by
4836    /// a certain transaction. Returns None if the a lock record is
4837    /// initialized for the given ObjectReference but not yet locked by any
4838    /// transaction,     or cannot find the transaction in transaction
4839    /// table, because of data race etc.
4840    #[instrument(level = "trace", skip_all)]
4841    pub fn get_transaction_lock(
4842        &self,
4843        object_ref: &ObjectReference,
4844        epoch_store: &AuthorityPerEpochStore,
4845    ) -> IotaResult<Option<VerifiedSignedTransaction>> {
4846        let lock_info = self
4847            .get_object_cache_reader()
4848            .try_get_lock(*object_ref, epoch_store)?;
4849        let lock_info = match lock_info {
4850            ObjectLockStatus::LockedAtDifferentVersion { locked_ref } => {
4851                return Err(UserInputError::ObjectVersionUnavailableForConsumption {
4852                    provided_obj_ref: *object_ref,
4853                    current_version: locked_ref.version,
4854                }
4855                .into());
4856            }
4857            ObjectLockStatus::Initialized => {
4858                return Ok(None);
4859            }
4860            ObjectLockStatus::LockedToTx { locked_by_tx } => locked_by_tx,
4861        };
4862
4863        epoch_store.get_signed_transaction(&lock_info)
4864    }
4865
4866    pub fn try_get_objects(&self, objects: &[ObjectId]) -> IotaResult<Vec<Option<Object>>> {
4867        self.get_object_cache_reader().try_get_objects(objects)
4868    }
4869
4870    /// Non-fallible version of `try_get_objects`.
4871    pub fn get_objects(&self, objects: &[ObjectId]) -> Vec<Option<Object>> {
4872        self.try_get_objects(objects)
4873            .expect("storage access failed")
4874    }
4875
4876    pub fn try_get_object_or_tombstone(
4877        &self,
4878        object_id: ObjectId,
4879    ) -> IotaResult<Option<ObjectReference>> {
4880        self.get_object_cache_reader()
4881            .try_get_latest_object_ref_or_tombstone(object_id)
4882    }
4883
4884    /// Non-fallible version of `try_get_object_or_tombstone`.
4885    pub fn get_object_or_tombstone(&self, object_id: ObjectId) -> Option<ObjectReference> {
4886        self.try_get_object_or_tombstone(object_id)
4887            .expect("storage access failed")
4888    }
4889
4890    /// Ordinarily, protocol upgrades occur when 2f + 1 + (f *
4891    /// ProtocolConfig::buffer_stake_for_protocol_upgrade_bps) vote for the
4892    /// upgrade.
4893    ///
4894    /// This method can be used to dynamic adjust the amount of buffer. If set
4895    /// to 0, the upgrade will go through with only 2f+1 votes.
4896    ///
4897    /// IMPORTANT: If this is used, it must be used on >=2f+1 validators (all
4898    /// should have the same value), or you risk halting the chain.
4899    pub fn set_override_protocol_upgrade_buffer_stake(
4900        &self,
4901        expected_epoch: EpochId,
4902        buffer_stake_bps: u64,
4903    ) -> IotaResult {
4904        let epoch_store = self.load_epoch_store_one_call_per_task();
4905        let actual_epoch = epoch_store.epoch();
4906        if actual_epoch != expected_epoch {
4907            return Err(IotaError::WrongEpoch {
4908                expected_epoch,
4909                actual_epoch,
4910            });
4911        }
4912
4913        epoch_store.set_override_protocol_upgrade_buffer_stake(buffer_stake_bps)
4914    }
4915
4916    pub fn clear_override_protocol_upgrade_buffer_stake(
4917        &self,
4918        expected_epoch: EpochId,
4919    ) -> IotaResult {
4920        let epoch_store = self.load_epoch_store_one_call_per_task();
4921        let actual_epoch = epoch_store.epoch();
4922        if actual_epoch != expected_epoch {
4923            return Err(IotaError::WrongEpoch {
4924                expected_epoch,
4925                actual_epoch,
4926            });
4927        }
4928
4929        epoch_store.clear_override_protocol_upgrade_buffer_stake()
4930    }
4931
4932    /// Get the set of system packages that are compiled in to this build, if
4933    /// those packages are compatible with the current versions of those
4934    /// packages on-chain.
4935    pub async fn get_available_system_packages(
4936        &self,
4937        protocol_config: &ProtocolConfig,
4938    ) -> Vec<ObjectReference> {
4939        let mut results = vec![];
4940
4941        let system_packages = BuiltInFramework::iter_system_packages();
4942
4943        // Add extra framework packages during simtest
4944        #[cfg(msim)]
4945        let extra_packages = framework_injection::get_extra_packages(self.name);
4946        #[cfg(msim)]
4947        let system_packages = {
4948            let mut packages: Vec<_> = system_packages.collect();
4949            packages.extend(extra_packages.iter());
4950            packages
4951        };
4952
4953        for system_package in system_packages {
4954            let modules = system_package.modules().to_vec();
4955            // In simtests, we could override the current built-in framework packages.
4956            #[cfg(msim)]
4957            let modules = framework_injection::get_override_modules(&system_package.id, self.name)
4958                .unwrap_or(modules);
4959
4960            let Some(obj_ref) = iota_framework::compare_system_package(
4961                &self.get_object_store(),
4962                &system_package.id,
4963                &modules,
4964                system_package.dependencies.to_vec(),
4965                protocol_config,
4966            )
4967            .await
4968            else {
4969                return vec![];
4970            };
4971            results.push(obj_ref);
4972        }
4973
4974        results
4975    }
4976
4977    /// Return the new versions, module bytes, and dependencies for the packages
4978    /// that have been committed to for a framework upgrade, in
4979    /// `system_packages`.  Loads the module contents from the binary, and
4980    /// performs the following checks:
4981    ///
4982    /// - Whether its contents matches what is on-chain already, in which case no upgrade is
4983    ///   required, and its contents are omitted from the output.
4984    /// - Whether the contents in the binary can form a package whose digest matches the input,
4985    ///   meaning the framework will be upgraded, and this authority can satisfy that upgrade, in
4986    ///   which case the contents are included in the output.
4987    ///
4988    /// If a needed version of the framework can't be loaded, the binary does
4989    /// not contain the bytes for that framework ID, or the resulting
4990    /// package fails the digest check, `None` is returned indicating that
4991    /// this authority cannot run the upgrade that the network voted on.
4992    ///
4993    /// All object lookups are pinned to the versions in `system_packages`
4994    /// instead of using the latest versions, so that the result is
4995    /// deterministic even if the change epoch transaction that performs the
4996    /// upgrade has already been executed locally (e.g. via state sync). In
4997    /// that case the reconstructed change epoch transaction is byte-identical
4998    /// to the executed one, and the caller detects it as already executed.
4999    async fn get_system_package_bytes(
5000        &self,
5001        system_packages: Vec<ObjectReference>,
5002        binary_config: &BinaryConfig,
5003    ) -> Option<Vec<SystemPackage>> {
5004        let object_store = self.get_object_cache_reader();
5005
5006        let mut res = Vec::with_capacity(system_packages.len());
5007        for system_package_ref in system_packages {
5008            if object_store
5009                .get_object_by_key(&system_package_ref.object_id, system_package_ref.version)
5010                .is_some_and(|object| object.object_ref() == system_package_ref)
5011            {
5012                // Skip this one because it doesn't need to be upgraded.
5013                info!(
5014                    "Framework {} does not need updating",
5015                    system_package_ref.object_id
5016                );
5017                continue;
5018            }
5019
5020            // The digest in `system_package_ref` commits to a package built on top of the
5021            // predecessor version's `previous_transaction` (see `compare_system_package`),
5022            // so it must be re-derived from that version. A ref at
5023            // `Version::OBJECT_START` is a freshly created package with no predecessor.
5024            let prev_transaction = if system_package_ref.version == Version::OBJECT_START {
5025                TransactionDigest::GENESIS_MARKER
5026            } else {
5027                let prev_version = system_package_ref
5028                    .version
5029                    .previous()
5030                    .expect("version is greater than Version::OBJECT_START");
5031                let Some(prev_object) =
5032                    object_store.get_object_by_key(&system_package_ref.object_id, prev_version)
5033                else {
5034                    error!(
5035                        "Framework {} not available locally at version {prev_version:?}, cannot \
5036                         derive upgrade to {system_package_ref:?}",
5037                        system_package_ref.object_id
5038                    );
5039                    return None;
5040                };
5041                prev_object.previous_transaction
5042            };
5043
5044            #[cfg(msim)]
5045            let FrameworkSystemPackage {
5046                id: _,
5047                bytes,
5048                dependencies,
5049            } = framework_injection::get_override_system_package(
5050                &system_package_ref.object_id,
5051                self.name,
5052            )
5053            .unwrap_or_else(|| {
5054                BuiltInFramework::get_package_by_id(&system_package_ref.object_id).clone()
5055            });
5056
5057            #[cfg(not(msim))]
5058            let FrameworkSystemPackage {
5059                id: _,
5060                bytes,
5061                dependencies,
5062            } = BuiltInFramework::get_package_by_id(&system_package_ref.object_id).clone();
5063
5064            let modules: Vec<_> = bytes
5065                .iter()
5066                .map(|m| CompiledModule::deserialize_with_config(m, binary_config).unwrap())
5067                .collect();
5068
5069            let new_object = Object::new_system_package(
5070                &modules,
5071                system_package_ref.version,
5072                dependencies.clone(),
5073                prev_transaction,
5074            );
5075
5076            let new_ref = new_object.object_ref();
5077            if new_ref != system_package_ref {
5078                debug_fatal!(
5079                    "Framework mismatch -- binary: {new_ref:?}\n  upgrade: {system_package_ref:?}"
5080                );
5081                return None;
5082            }
5083
5084            res.push(SystemPackage {
5085                version: system_package_ref.version,
5086                modules: bytes,
5087                dependencies,
5088            });
5089        }
5090
5091        Some(res)
5092    }
5093
5094    /// Returns the new protocol version and system packages that the network
5095    /// has voted to upgrade to. If the proposed protocol version is not
5096    /// supported, None is returned.
5097    fn is_protocol_version_supported_v1(
5098        proposed_protocol_version: ProtocolVersion,
5099        committee: &Committee,
5100        capabilities: Vec<AuthorityCapabilitiesV1>,
5101        mut buffer_stake_bps: u64,
5102    ) -> Option<(ProtocolVersion, Digest, Vec<ObjectReference>)> {
5103        if buffer_stake_bps > 10000 {
5104            warn!("clamping buffer_stake_bps to 10000");
5105            buffer_stake_bps = 10000;
5106        }
5107
5108        // For each validator, gather the protocol version and system packages that it
5109        // would like to upgrade to in the next epoch.
5110        let mut desired_upgrades: Vec<_> = capabilities
5111            .into_iter()
5112            .filter_map(|mut cap| {
5113                // A validator that lists no packages is voting against any change at all.
5114                if cap.available_system_packages.is_empty() {
5115                    return None;
5116                }
5117
5118                cap.available_system_packages.sort();
5119
5120                info!(
5121                    "validator {:?} supports {:?} with system packages: {:?}",
5122                    cap.authority.concise(),
5123                    cap.supported_protocol_versions,
5124                    cap.available_system_packages,
5125                );
5126
5127                // A validator that only supports the current protocol version is also voting
5128                // against any change, because framework upgrades always require a protocol
5129                // version bump.
5130                cap.supported_protocol_versions
5131                    .get_version_digest(proposed_protocol_version)
5132                    .map(|digest| (digest, cap.available_system_packages, cap.authority))
5133            })
5134            .collect();
5135
5136        // There can only be one set of votes that have a majority, find one if it
5137        // exists.
5138        desired_upgrades.sort();
5139        desired_upgrades
5140            .into_iter()
5141            .chunk_by(|(digest, packages, _authority)| (*digest, packages.clone()))
5142            .into_iter()
5143            .find_map(|((digest, packages), group)| {
5144                // should have been filtered out earlier.
5145                assert!(!packages.is_empty());
5146
5147                let mut stake_aggregator: StakeAggregator<(), true> =
5148                    StakeAggregator::new(Arc::new(committee.clone()));
5149
5150                for (_, _, authority) in group {
5151                    stake_aggregator.insert_generic(authority, ());
5152                }
5153
5154                let total_votes = stake_aggregator.total_votes();
5155                let quorum_threshold = committee.quorum_threshold();
5156                let effective_threshold = committee.effective_threshold(buffer_stake_bps);
5157
5158                info!(
5159                    protocol_config_digest = ?digest,
5160                    ?total_votes,
5161                    ?quorum_threshold,
5162                    ?buffer_stake_bps,
5163                    ?effective_threshold,
5164                    ?proposed_protocol_version,
5165                    ?packages,
5166                    "support for upgrade"
5167                );
5168
5169                let has_support = total_votes >= effective_threshold;
5170                has_support.then_some((proposed_protocol_version, digest, packages))
5171            })
5172    }
5173
5174    /// Selects the highest supported protocol version and system packages that
5175    /// the network has voted to upgrade to. If no upgrade is supported,
5176    /// returns the current protocol version and system packages.
5177    fn choose_protocol_version_and_system_packages_v1(
5178        current_protocol_version: ProtocolVersion,
5179        current_protocol_digest: Digest,
5180        committee: &Committee,
5181        capabilities: Vec<AuthorityCapabilitiesV1>,
5182        buffer_stake_bps: u64,
5183    ) -> (ProtocolVersion, Digest, Vec<ObjectReference>) {
5184        let mut next_protocol_version = current_protocol_version;
5185        let mut system_packages = vec![];
5186        let mut protocol_version_digest = current_protocol_digest;
5187
5188        // Finds the highest supported protocol version and system packages by
5189        // incrementing the proposed protocol version by one until no further
5190        // upgrades are supported.
5191        while let Some((version, digest, packages)) = Self::is_protocol_version_supported_v1(
5192            next_protocol_version + 1,
5193            committee,
5194            capabilities.clone(),
5195            buffer_stake_bps,
5196        ) {
5197            next_protocol_version = version;
5198            protocol_version_digest = digest;
5199            system_packages = packages;
5200        }
5201
5202        (
5203            next_protocol_version,
5204            protocol_version_digest,
5205            system_packages,
5206        )
5207    }
5208
5209    /// Returns the indices of validators that support the given protocol
5210    /// version and digest. This includes both committee and non-committee
5211    /// validators based on their capabilities. Uses active validators
5212    /// instead of committee indices.
5213    fn get_validators_supporting_protocol_version(
5214        target_protocol_version: ProtocolVersion,
5215        target_digest: Digest,
5216        active_validators: &[AggregateAuthorityPublicKey],
5217        capabilities: &[AuthorityCapabilitiesV1],
5218    ) -> Vec<u64> {
5219        let mut eligible_validators = Vec::new();
5220
5221        for capability in capabilities {
5222            // Check if this validator supports the target protocol version and digest
5223            if let Some(digest) = capability
5224                .supported_protocol_versions
5225                .get_version_digest(target_protocol_version)
5226            {
5227                if digest == target_digest {
5228                    // Find the validator's index in the active validators list
5229                    if let Some(index) = active_validators
5230                        .iter()
5231                        .position(|name| AuthorityName::from(name) == capability.authority)
5232                    {
5233                        eligible_validators.push(index as u64);
5234                    }
5235                }
5236            }
5237        }
5238
5239        // Sort indices for deterministic behavior
5240        eligible_validators.sort();
5241        eligible_validators
5242    }
5243
5244    /// Calculates the sum of weights for eligible validators that are part of
5245    /// the committee. Takes the indices from
5246    /// get_validators_supporting_protocol_version and maps them back
5247    /// to committee members to get their weights.
5248    fn calculate_eligible_validators_weight(
5249        eligible_validator_indices: &[u64],
5250        active_validators: &[AggregateAuthorityPublicKey],
5251        committee: &Committee,
5252    ) -> u64 {
5253        let mut total_weight = 0u64;
5254
5255        for &index in eligible_validator_indices {
5256            let authority_pubkey = &active_validators[index as usize];
5257            // Check if this validator is in the committee and get their weight
5258            if let Some((_, weight)) = committee
5259                .members()
5260                .find(|(name, _)| *name == AuthorityName::from(authority_pubkey))
5261            {
5262                total_weight += weight;
5263            }
5264        }
5265
5266        total_weight
5267    }
5268
5269    /// Creates and execute the advance epoch transaction to effects without
5270    /// committing it to the database. The effects of the change epoch tx
5271    /// are only written to the database after a certified checkpoint has been
5272    /// formed and executed by CheckpointExecutor.
5273    ///
5274    /// When a framework upgraded has been decided on, but the validator does
5275    /// not have the new versions of the packages locally, the validator
5276    /// cannot form the ChangeEpochTx. In this case it returns Err,
5277    /// indicating that the checkpoint builder should give up trying to make the
5278    /// final checkpoint. As long as the network is able to create a certified
5279    /// checkpoint (which should be ensured by the capabilities vote), it
5280    /// will arrive via state sync and be executed by CheckpointExecutor.
5281    #[instrument(level = "error", skip_all)]
5282    pub async fn create_and_execute_advance_epoch_tx(
5283        &self,
5284        epoch_store: &Arc<AuthorityPerEpochStore>,
5285        gas_cost_summary: &GasCostSummary,
5286        checkpoint: CheckpointSequenceNumber,
5287        epoch_start_timestamp_ms: CheckpointTimestamp,
5288        scores: Vec<u64>,
5289    ) -> CheckpointBuilderResult<(
5290        IotaSystemState,
5291        Option<SystemEpochInfoEvent>,
5292        TransactionEffects,
5293    )> {
5294        let mut txns = Vec::new();
5295
5296        // Create the TransactionDenyRules object once: the epoch-start
5297        // configuration is identical on every validator, so the whole
5298        // committee injects (or skips) the kind together. If this epoch
5299        // change falls into safe mode the creation is dropped with it, the
5300        // object stays absent, and the next epoch end injects it again.
5301        if epoch_store
5302            .protocol_config()
5303            .deny_rule_governance_on_chain()
5304            && epoch_store
5305                .epoch_start_config()
5306                .transaction_deny_rules_obj_initial_shared_version()
5307                .is_none()
5308        {
5309            txns.push(EndOfEpochTransactionKind::TransactionDenyRulesCreate);
5310        }
5311
5312        let next_epoch = epoch_store.epoch() + 1;
5313
5314        let buffer_stake_bps = epoch_store.get_effective_buffer_stake_bps();
5315        let authority_capabilities = epoch_store
5316            .get_capabilities_v1()
5317            .expect("read capabilities from db cannot fail");
5318        let (next_epoch_protocol_version, next_epoch_protocol_digest, next_epoch_system_packages) =
5319            Self::choose_protocol_version_and_system_packages_v1(
5320                epoch_store.protocol_version(),
5321                SupportedProtocolVersionsWithHashes::protocol_config_digest(
5322                    epoch_store.protocol_config(),
5323                ),
5324                epoch_store.committee(),
5325                authority_capabilities.clone(),
5326                buffer_stake_bps,
5327            );
5328
5329        // since system packages are created during the current epoch, they should abide
5330        // by the rules of the current epoch, including the current epoch's max
5331        // Move binary format version
5332        let config = epoch_store.protocol_config();
5333        let binary_config = to_binary_config(config, None);
5334        let Some(next_epoch_system_package_bytes) = self
5335            .get_system_package_bytes(next_epoch_system_packages.clone(), &binary_config)
5336            .await
5337        else {
5338            debug_fatal!(
5339                "upgraded system packages {:?} are not locally available, cannot create \
5340                ChangeEpochTx. validator binary must be upgraded to the correct version!",
5341                next_epoch_system_packages
5342            );
5343            // the checkpoint builder will keep retrying forever when it hits this error.
5344            // Eventually, one of two things will happen:
5345            // - The operator will upgrade this binary to one that has the new packages locally, and
5346            //   this function will succeed.
5347            // - The final checkpoint will be certified by other validators, we will receive it via
5348            //   state sync, and execute it. This will upgrade the framework packages, reconfigure,
5349            //   and most likely shut down in the new epoch (this validator likely doesn't support
5350            //   the new protocol version, or else it should have had the packages.)
5351            return Err(CheckpointBuilderError::SystemPackagesMissing);
5352        };
5353
5354        // Use ChangeEpochV3 or ChangeEpochV4 when the feature flags are enabled and
5355        // ChangeEpochV2 requirements are met
5356        if config.select_committee_from_eligible_validators() {
5357            // Get the list of eligible validators that support the target protocol version
5358            let active_validators = epoch_store.epoch_start_state().get_active_validators();
5359
5360            let mut eligible_active_validators = (0..active_validators.len() as u64).collect();
5361
5362            // Use validators supporting the target protocol version as eligible validators
5363            // in the next version if select_committee_supporting_next_epoch_version feature
5364            // flag is set to true.
5365            if config.select_committee_supporting_next_epoch_version() {
5366                eligible_active_validators = Self::get_validators_supporting_protocol_version(
5367                    next_epoch_protocol_version,
5368                    next_epoch_protocol_digest,
5369                    &active_validators,
5370                    &authority_capabilities,
5371                );
5372
5373                // Calculate the total weight of eligible validators in the committee
5374                let eligible_validators_weight = Self::calculate_eligible_validators_weight(
5375                    &eligible_active_validators,
5376                    &active_validators,
5377                    epoch_store.committee(),
5378                );
5379
5380                // Safety check: ensure eligible validators have enough stake
5381                // Use the same effective threshold calculation that was used to decide the
5382                // protocol version
5383                let committee = epoch_store.committee();
5384                let effective_threshold = committee.effective_threshold(buffer_stake_bps);
5385
5386                if eligible_validators_weight < effective_threshold {
5387                    error!(
5388                        "Eligible validators weight {eligible_validators_weight} is less than effective threshold {effective_threshold}. \
5389                        This could indicate a bug in validator selection logic or inconsistency with protocol version decision.",
5390                    );
5391                    // Pass all active validator indices as eligible validators
5392                    // to perform selection among all of them.
5393                    eligible_active_validators = (0..active_validators.len() as u64).collect();
5394                }
5395            }
5396
5397            // Use ChangeEpochV4 when the pass_validator_scores_to_advance_epoch feature
5398            // flag is enabled.
5399            if config.pass_validator_scores_to_advance_epoch() {
5400                txns.push(EndOfEpochTransactionKind::new_change_epoch_v4(
5401                    next_epoch,
5402                    next_epoch_protocol_version.as_u64(),
5403                    gas_cost_summary.storage_cost,
5404                    gas_cost_summary.computation_cost,
5405                    gas_cost_summary.computation_cost_burned,
5406                    gas_cost_summary.storage_rebate,
5407                    gas_cost_summary.non_refundable_storage_fee,
5408                    epoch_start_timestamp_ms,
5409                    next_epoch_system_package_bytes,
5410                    eligible_active_validators,
5411                    scores,
5412                    config.adjust_rewards_by_score(),
5413                ));
5414            } else {
5415                txns.push(EndOfEpochTransactionKind::new_change_epoch_v3(
5416                    next_epoch,
5417                    next_epoch_protocol_version.as_u64(),
5418                    gas_cost_summary.storage_cost,
5419                    gas_cost_summary.computation_cost,
5420                    gas_cost_summary.computation_cost_burned,
5421                    gas_cost_summary.storage_rebate,
5422                    gas_cost_summary.non_refundable_storage_fee,
5423                    epoch_start_timestamp_ms,
5424                    next_epoch_system_package_bytes,
5425                    eligible_active_validators,
5426                ));
5427            }
5428        } else if config.protocol_defined_base_fee()
5429            && config.max_committee_members_count_as_option().is_some()
5430        {
5431            txns.push(EndOfEpochTransactionKind::new_change_epoch_v2(
5432                next_epoch,
5433                next_epoch_protocol_version.as_u64(),
5434                gas_cost_summary.storage_cost,
5435                gas_cost_summary.computation_cost,
5436                gas_cost_summary.computation_cost_burned,
5437                gas_cost_summary.storage_rebate,
5438                gas_cost_summary.non_refundable_storage_fee,
5439                epoch_start_timestamp_ms,
5440                next_epoch_system_package_bytes,
5441            ));
5442        } else {
5443            txns.push(EndOfEpochTransactionKind::new_change_epoch(
5444                next_epoch,
5445                next_epoch_protocol_version.as_u64(),
5446                gas_cost_summary.storage_cost,
5447                gas_cost_summary.computation_cost,
5448                gas_cost_summary.storage_rebate,
5449                gas_cost_summary.non_refundable_storage_fee,
5450                epoch_start_timestamp_ms,
5451                next_epoch_system_package_bytes,
5452            ));
5453        }
5454
5455        let tx = VerifiedTransaction::new_end_of_epoch_transaction(txns);
5456
5457        let executable_tx = VerifiedExecutableTransaction::new_from_checkpoint(
5458            tx.clone(),
5459            epoch_store.epoch(),
5460            checkpoint,
5461        );
5462
5463        let tx_digest = executable_tx.digest();
5464
5465        info!(
5466            ?next_epoch,
5467            ?next_epoch_protocol_version,
5468            ?next_epoch_system_packages,
5469            computation_cost=?gas_cost_summary.computation_cost,
5470            computation_cost_burned=?gas_cost_summary.computation_cost_burned,
5471            storage_cost=?gas_cost_summary.storage_cost,
5472            storage_rebate=?gas_cost_summary.storage_rebate,
5473            non_refundable_storage_fee=?gas_cost_summary.non_refundable_storage_fee,
5474            ?tx_digest,
5475            "Creating advance epoch transaction"
5476        );
5477
5478        fail_point_async!("change_epoch_tx_delay");
5479        let tx_lock = epoch_store.acquire_tx_lock(tx_digest);
5480
5481        // The tx could have been executed by state sync already - if so simply return
5482        // an error. The checkpoint builder will shortly be terminated by
5483        // reconfiguration anyway.
5484        if self
5485            .get_transaction_cache_reader()
5486            .try_is_tx_already_executed(tx_digest)?
5487        {
5488            warn!("change epoch tx has already been executed via state sync");
5489            return Err(CheckpointBuilderError::ChangeEpochTxAlreadyExecuted);
5490        }
5491
5492        let execution_guard = self.execution_lock_for_executable_transaction(&executable_tx)?;
5493
5494        // We must manually assign the shared object versions to the transaction before
5495        // executing it. This is because we do not sequence end-of-epoch
5496        // transactions through consensus.
5497        let assigned_versions = epoch_store.assign_shared_object_versions_idempotent(
5498            self.get_object_cache_reader().as_ref(),
5499            std::iter::once(&Schedulable::Transaction(&executable_tx)),
5500        )?;
5501
5502        assert_eq!(assigned_versions.0.len(), 1);
5503        let assigned_versions = assigned_versions.0.into_iter().next().unwrap().1;
5504
5505        let (input_objects, _) = self.read_objects_for_execution(
5506            &tx_lock,
5507            &executable_tx,
5508            assigned_versions,
5509            epoch_store,
5510        )?;
5511
5512        let (temporary_store, effects, _execution_error_opt) = self.execute_transaction(
5513            &execution_guard,
5514            &executable_tx,
5515            input_objects,
5516            vec![],
5517            epoch_store,
5518        )?;
5519        let system_obj = get_iota_system_state(&temporary_store.written)
5520            .expect("change epoch tx must write to system object");
5521        // Find the SystemEpochInfoEvent emitted by the advance_epoch transaction.
5522        let system_epoch_info_event = temporary_store
5523            .events
5524            .0
5525            .into_iter()
5526            .find(|event| event.is_system_epoch_info_event())
5527            .map(SystemEpochInfoEvent::from);
5528        // The system epoch info event can be `None` in case if the `advance_epoch`
5529        // Move function call failed and was executed in the safe mode.
5530        assert!(system_epoch_info_event.is_some() || system_obj.safe_mode());
5531
5532        // We must write tx and effects to the state sync tables so that state sync is
5533        // able to deliver to the transaction to CheckpointExecutor after it is
5534        // included in a certified checkpoint.
5535        self.get_state_sync_store()
5536            .try_insert_transaction_and_effects(&tx, &effects)?;
5537
5538        info!(
5539            "Effects summary of the change epoch transaction: {:?}",
5540            effects.summary_for_debug()
5541        );
5542        epoch_store.record_checkpoint_builder_is_safe_mode_metric(system_obj.safe_mode());
5543        // The change epoch transaction cannot fail to execute.
5544        assert!(effects.status().is_success());
5545        Ok((system_obj, system_epoch_info_event, effects))
5546    }
5547
5548    /// This function is called at the very end of the epoch.
5549    /// This step is required before updating new epoch in the db and calling
5550    /// reopen_epoch_db.
5551    #[instrument(level = "error", skip_all)]
5552    async fn revert_uncommitted_epoch_transactions(
5553        &self,
5554        epoch_store: &AuthorityPerEpochStore,
5555    ) -> IotaResult {
5556        {
5557            let state = epoch_store.get_reconfig_state_write_lock_guard();
5558            if state.should_accept_user_certs() {
5559                // Need to change this so that consensus adapter do not accept certificates from
5560                // user. This can happen if our local validator did not initiate
5561                // epoch change locally, but 2f+1 nodes already concluded the
5562                // epoch.
5563                //
5564                // This lock is essentially a barrier (in the certificate mode only) for
5565                // `epoch_store.pending_consensus_certificates` table we are reading on the line
5566                // after this block
5567                epoch_store.close_user_certs(state);
5568            }
5569            // lock is dropped here
5570        }
5571
5572        // In the P-COOL flow, the list of pending consensus certificates is
5573        // always empty, so the reverting below is only for the certificate mode.
5574        if !epoch_store.protocol_config().enable_pcool_flow() {
5575            let pending_certificates = epoch_store.pending_consensus_certificates();
5576            info!(
5577                "Reverting {} locally executed transactions that was not included in the epoch: \
5578                    {:?}",
5579                pending_certificates.len(),
5580                pending_certificates,
5581            );
5582            for digest in pending_certificates {
5583                if epoch_store.is_transaction_executed_in_checkpoint(&digest)? {
5584                    info!(
5585                        "Not reverting pending consensus transaction {:?} - it was included in \
5586                            checkpoint",
5587                        digest
5588                    );
5589                    continue;
5590                }
5591                info!("Reverting {:?} at the end of epoch", digest);
5592                epoch_store.revert_executed_transaction(&digest)?;
5593                self.get_reconfig_api().try_revert_state_update(&digest)?;
5594            }
5595            info!("All uncommitted local transactions reverted");
5596        } else {
5597            info!("P-COOL mode: skipping revert of uncommitted epoch transactions");
5598        }
5599
5600        Ok(())
5601    }
5602
5603    #[instrument(level = "error", skip_all)]
5604    async fn reopen_epoch_db(
5605        &self,
5606        cur_epoch_store: &AuthorityPerEpochStore,
5607        new_committee: Committee,
5608        epoch_start_configuration: EpochStartConfiguration,
5609        expensive_safety_check_config: &ExpensiveSafetyCheckConfig,
5610        epoch_last_checkpoint: CheckpointSequenceNumber,
5611    ) -> IotaResult<Arc<AuthorityPerEpochStore>> {
5612        let new_epoch = new_committee.epoch;
5613        info!(new_epoch = ?new_epoch, "re-opening AuthorityEpochTables for new epoch");
5614        assert_eq!(
5615            epoch_start_configuration.epoch_start_state().epoch(),
5616            new_committee.epoch
5617        );
5618        fail_point!("before-open-new-epoch-store");
5619        let new_epoch_store = cur_epoch_store.new_at_next_epoch(
5620            self.name,
5621            new_committee,
5622            epoch_start_configuration,
5623            self.get_backing_package_store().clone(),
5624            expensive_safety_check_config,
5625            epoch_last_checkpoint,
5626        )?;
5627        self.epoch_store.store(new_epoch_store.clone());
5628        Ok(new_epoch_store)
5629    }
5630
5631    /// Resolves the account's `AuthenticatorFunctionRef` on the execution path,
5632    /// where the certificate has already passed validation before consensus.
5633    ///
5634    /// A deleted or cancelled account object is not an error here: its version
5635    /// is returned so execution can proceed and surface the proper effect
5636    /// (e.g. `InputObjectDeleted` or a shared-object congestion cancellation).
5637    ///
5638    /// # Errors
5639    ///
5640    /// Any failure the check finds with the transaction, as the error the
5641    /// transaction fails with, with the original error as the source:
5642    /// [`ExecutionErrorKind::AuthenticatorFunctionNotFound`] when the account
5643    /// has no authenticator function field or the field cannot be read, and
5644    /// [`ExecutionErrorKind::AccountNotSharedObject`] when the account is not a
5645    /// shared object.
5646    ///
5647    /// # Panics
5648    ///
5649    /// When the store cannot be read: that says nothing about the transaction,
5650    /// so nobody can be charged for it, and effects written here would differ
5651    /// from the other validators'.
5652    fn check_move_account_for_execution(
5653        &self,
5654        auth_account_object_id: ObjectId,
5655        auth_account_object_seq_number: Option<Version>,
5656        auth_account_object_digest: Option<ObjectDigest>,
5657        account_object: ObjectReadResult,
5658        signer: &Address,
5659        protocol_config: &ProtocolConfig,
5660    ) -> Result<AuthenticatorFunctionRefForExecution, ExecutionError> {
5661        match self.check_move_account(
5662            auth_account_object_id,
5663            auth_account_object_seq_number,
5664            auth_account_object_digest,
5665            account_object,
5666            signer,
5667            true,
5668            protocol_config,
5669        ) {
5670            Ok(function_ref) => Ok(function_ref),
5671            Err(MoveAccountCheckError::Input(input_error)) => {
5672                let status = match &input_error {
5673                    UserInputError::MoveAuthenticatorNotFound {
5674                        account_object_id, ..
5675                    }
5676                    | UserInputError::InvalidAuthenticatorFunctionRefField { account_object_id } => {
5677                        ExecutionErrorKind::AuthenticatorFunctionNotFound {
5678                            object_id: *account_object_id,
5679                        }
5680                    }
5681                    // Since protocol version 36, only a shared account named with
5682                    // an immutable or owned reference reaches the digest
5683                    // comparison. The fix is the same as for the other two: name
5684                    // the account as a shared object.
5685                    UserInputError::AccountObjectNotSupported { object_id }
5686                    | UserInputError::ImmutableAccountObjectNotSupported { object_id }
5687                    | UserInputError::InvalidAccountObjectDigest { object_id, .. } => {
5688                        ExecutionErrorKind::AccountNotSharedObject {
5689                            object_id: *object_id,
5690                        }
5691                    }
5692                    // The check produces nothing else today. A failure added to
5693                    // it gets this status until it is given one of its own.
5694                    _ => ExecutionErrorKind::FunctionNotFound,
5695                };
5696                Err(ExecutionError::new_with_source(status, input_error))
5697            }
5698            Err(MoveAccountCheckError::Storage(storage_error)) => {
5699                panic!("failed to read the store while checking a Move account: {storage_error}")
5700            }
5701        }
5702    }
5703
5704    /// Resolves the account's `AuthenticatorFunctionRef` on the validation
5705    /// (signing) path, rejecting the transaction when the account object was
5706    /// deleted or belongs to a cancelled transaction.
5707    fn check_move_account_for_validation(
5708        &self,
5709        auth_account_object_id: ObjectId,
5710        auth_account_object_seq_number: Option<Version>,
5711        auth_account_object_digest: Option<ObjectDigest>,
5712        account_object: ObjectReadResult,
5713        signer: &Address,
5714        protocol_config: &ProtocolConfig,
5715    ) -> IotaResult<AuthenticatorFunctionRefForExecution> {
5716        self.check_move_account(
5717            auth_account_object_id,
5718            auth_account_object_seq_number,
5719            auth_account_object_digest,
5720            account_object,
5721            signer,
5722            false,
5723            protocol_config,
5724        )
5725        .map_err(IotaError::from)
5726    }
5727
5728    /// Checks whether `authenticator` unlocks a valid Move account and returns
5729    /// the account-related `AuthenticatorFunctionRef`. Where the protocol
5730    /// config requires it, the account object must be shared, so that a
5731    /// transaction carrying a `MoveAuthenticator` always has a shared input
5732    /// and is ordered by consensus. When `is_execution` is
5733    /// set, a deleted or cancelled account object yields its version instead of
5734    /// an error, so execution can proceed to the proper effect. Prefer the
5735    /// `check_move_account_for_execution` / `check_move_account_for_validation`
5736    /// wrappers over calling this directly.
5737    ///
5738    /// # Errors
5739    ///
5740    /// [`MoveAccountCheckError::Input`] when the account fails a check, and
5741    /// [`MoveAccountCheckError::Storage`] when the store cannot be read. The
5742    /// error type is what lets the execution path distinguish the two, so a
5743    /// new kind of failure here has to be added to it.
5744    fn check_move_account(
5745        &self,
5746        auth_account_object_id: ObjectId,
5747        auth_account_object_seq_number: Option<Version>,
5748        auth_account_object_digest: Option<ObjectDigest>,
5749        account_object: ObjectReadResult,
5750        signer: &Address,
5751        is_execution: bool,
5752        protocol_config: &ProtocolConfig,
5753    ) -> Result<AuthenticatorFunctionRefForExecution, MoveAccountCheckError> {
5754        let auth_account_object_seq_number = match (&account_object.object, is_execution) {
5755            // In any case, if the account object is loaded, we can check its version and digest.
5756            // Then we return the version of the account object to be used for reading the
5757            // authenticator function ref dynamic field.
5758            (ObjectReadResultKind::Object(object), _) => {
5759                let account_object_addr = Address::from(auth_account_object_id);
5760                fp_ensure!(
5761                    signer == &account_object_addr,
5762                    UserInputError::IncorrectUserSignature {
5763                        error: format!("Move authenticator is trying to unlock {account_object_addr:?}, but given signer address is {signer:?}")
5764                    }
5765                    .into()
5766                );
5767
5768                if protocol_config.reject_immutable_account_objects() {
5769                    fp_ensure!(
5770                        !object.is_immutable(),
5771                        UserInputError::ImmutableAccountObjectNotSupported {
5772                            object_id: auth_account_object_id
5773                        }
5774                        .into()
5775                    );
5776                }
5777
5778                fp_ensure!(
5779                    object.is_shared() || object.is_immutable(),
5780                    UserInputError::AccountObjectNotSupported {
5781                        object_id: auth_account_object_id
5782                    }
5783                    .into()
5784                );
5785
5786                let auth_account_object_seq_number =
5787                    if let Some(expected_version) = auth_account_object_seq_number {
5788                        let account_object_version = object.version();
5789
5790                        fp_ensure!(
5791                            account_object_version == expected_version,
5792                            UserInputError::AccountObjectVersionMismatch {
5793                                object_id: auth_account_object_id,
5794                                expected_version,
5795                                actual_version: account_object_version,
5796                            }
5797                            .into()
5798                        );
5799
5800                        expected_version
5801                    } else {
5802                        object.version()
5803                    };
5804
5805                if let Some(expected_digest) = auth_account_object_digest {
5806                    let account_object_digest = object.digest();
5807                    fp_ensure!(
5808                        account_object_digest == expected_digest,
5809                        UserInputError::InvalidAccountObjectDigest {
5810                            object_id: auth_account_object_id,
5811                            expected_digest,
5812                            actual_digest: account_object_digest,
5813                        }
5814                        .into()
5815                    );
5816                }
5817
5818                Ok(auth_account_object_seq_number)
5819            }
5820            // If the account object is not loaded because it was deleted, we return the error in
5821            // the case in which we are not executing the transaction right after.
5822            (ObjectReadResultKind::DeletedSharedObject(version, digest), false) => {
5823                Err(UserInputError::AccountObjectDeleted {
5824                    account_id: account_object.id(),
5825                    account_version: *version,
5826                    transaction_digest: *digest,
5827                })
5828            }
5829            // If the account object is not loaded because the transaction was canceled, we return
5830            // the error in the case in which we are not executing the transaction right
5831            // after.
5832            (ObjectReadResultKind::CancelledTransactionObject(version), false) => {
5833                Err(UserInputError::AccountObjectInCanceledTransaction {
5834                    account_id: account_object.id(),
5835                    account_version: *version,
5836                })
5837            }
5838            // If the account object is not loaded because it was deleted, we return the version in
5839            // the case in which we are executing the transaction right after.
5840            // This version is used to read the authenticator function ref dynamic field because it
5841            // is greater than the version of the child dynamic field.
5842            (ObjectReadResultKind::DeletedSharedObject(version, _), true) => Ok(*version),
5843            // If the account object is not loaded because the transaction was canceled, we return
5844            // the version in the case in which we are executing the transaction right
5845            // after. This version is used to read the authenticator function ref
5846            // dynamic field because it is greater than the version of the child dynamic
5847            // field.
5848            (ObjectReadResultKind::CancelledTransactionObject(version), true) => Ok(*version),
5849        }?;
5850
5851        let authenticator_function_ref_field_id =
5852            derive_authenticator_function_ref_v1_dynamic_field_id(auth_account_object_id)?;
5853
5854        let authenticator_function_ref_field = self
5855            .get_object_cache_reader()
5856            .try_find_object_lt_or_eq_version(
5857                authenticator_function_ref_field_id,
5858                auth_account_object_seq_number,
5859            )
5860            .map_err(MoveAccountCheckError::Storage)?;
5861
5862        if let Some(authenticator_function_ref_field_obj) = authenticator_function_ref_field {
5863            Ok(authenticator_function_ref_v1_from_dynamic_field_object(
5864                auth_account_object_id,
5865                &authenticator_function_ref_field_obj,
5866            )?)
5867        } else {
5868            Err(UserInputError::MoveAuthenticatorNotFound {
5869                authenticator_function_ref_id: authenticator_function_ref_field_id,
5870                account_object_id: auth_account_object_id,
5871                account_object_version: auth_account_object_seq_number,
5872            }
5873            .into())
5874        }
5875    }
5876
5877    #[allow(clippy::type_complexity)]
5878    fn read_objects_for_validation(
5879        &self,
5880        transaction: &VerifiedTransaction,
5881        epoch: u64,
5882    ) -> IotaResult<(
5883        InputObjects,
5884        ReceivingObjects,
5885        Vec<(InputObjects, ObjectReadResult)>,
5886    )> {
5887        let (input_objects, tx_receiving_objects) = self.input_loader.read_objects_for_signing(
5888            Some(transaction.digest()),
5889            &transaction.collect_all_input_object_kind_for_reading()?,
5890            &transaction.data().transaction().receiving_objects(),
5891            epoch,
5892        )?;
5893
5894        transaction
5895            .split_input_objects_into_groups_for_reading(input_objects)
5896            .map(|(tx_input_objects, per_authenticator_inputs)| {
5897                (
5898                    tx_input_objects,
5899                    tx_receiving_objects,
5900                    per_authenticator_inputs,
5901                )
5902            })
5903    }
5904
5905    #[allow(clippy::type_complexity)]
5906    fn check_transaction_inputs_for_validation(
5907        &self,
5908        protocol_config: &ProtocolConfig,
5909        reference_gas_price: u64,
5910        tx: &Transaction,
5911        tx_input_objects: InputObjects,
5912        tx_receiving_objects: &ReceivingObjects,
5913        move_authenticators: &Vec<&MoveAuthenticator>,
5914        per_authenticator_inputs: Vec<(InputObjects, ObjectReadResult)>,
5915        verifier_limits_source: VerifierLimitsSource<'_>,
5916    ) -> IotaResult<(
5917        IotaGasStatus,
5918        CheckedInputObjects,
5919        Vec<(CheckedInputObjects, AuthenticatorFunctionRef)>,
5920    )> {
5921        let authenticator_gas_budget = if move_authenticators.is_empty() {
5922            0
5923        } else {
5924            // `max_auth_gas` is used here as a Move authenticator gas budget until it is
5925            // not a part of the transaction data.
5926            protocol_config.max_auth_gas()
5927        };
5928
5929        debug_assert_eq!(
5930            move_authenticators.len(),
5931            per_authenticator_inputs.len(),
5932            "Move authenticators amount must match the number of authenticator inputs"
5933        );
5934
5935        let per_authenticator_checked_inputs = move_authenticators
5936            .iter()
5937            .zip(per_authenticator_inputs)
5938            .map(
5939                |(move_authenticator, (authenticator_input_objects, account_object))| {
5940                    // Check basic `object_to_authenticate` preconditions and get its components.
5941                    let (
5942                        auth_account_object_id,
5943                        auth_account_object_seq_number,
5944                        auth_account_object_digest,
5945                    ) = move_authenticator.object_to_authenticate_components()?;
5946
5947                    let signer = move_authenticator.address();
5948
5949                    // Make sure the signer is a Move account.
5950                    let AuthenticatorFunctionRefForExecution {
5951                        authenticator_function_ref,
5952                        ..
5953                    } = self.check_move_account_for_validation(
5954                        auth_account_object_id,
5955                        auth_account_object_seq_number,
5956                        auth_account_object_digest,
5957                        account_object,
5958                        &signer,
5959                        protocol_config,
5960                    )?;
5961
5962                    // Check the MoveAuthenticator input objects.
5963                    let authenticator_checked_input_objects =
5964                        iota_transaction_checks::check_move_authenticator_input_for_validation(
5965                            authenticator_input_objects,
5966                        )?;
5967
5968                    Ok((
5969                        authenticator_checked_input_objects,
5970                        authenticator_function_ref,
5971                    ))
5972                },
5973            )
5974            .collect::<IotaResult<Vec<_>>>()?;
5975
5976        // Check the transaction inputs.
5977        let (gas_status, tx_checked_input_objects) =
5978            iota_transaction_checks::check_transaction_input(
5979                protocol_config,
5980                reference_gas_price,
5981                tx,
5982                tx_input_objects,
5983                tx_receiving_objects,
5984                &self.metrics.bytecode_verifier_metrics,
5985                verifier_limits_source,
5986                authenticator_gas_budget,
5987            )?;
5988
5989        Ok((
5990            gas_status,
5991            tx_checked_input_objects,
5992            per_authenticator_checked_inputs,
5993        ))
5994    }
5995
5996    #[cfg(test)]
5997    pub(crate) fn iter_live_object_set_for_testing(
5998        &self,
5999    ) -> impl Iterator<Item = authority_store_tables::LiveObject> + '_ {
6000        self.get_global_state_hash_store()
6001            .iter_cached_live_object_set_for_testing()
6002    }
6003
6004    #[cfg(test)]
6005    pub(crate) fn shutdown_execution_for_test(&self) {
6006        self.tx_execution_shutdown
6007            .lock()
6008            .take()
6009            .unwrap()
6010            .send(())
6011            .unwrap();
6012    }
6013
6014    /// NOTE: this function is only to be used for fuzzing and testing. Never
6015    /// use in prod
6016    pub async fn insert_objects_unsafe_for_testing_only(&self, objects: &[Object]) {
6017        self.get_reconfig_api().bulk_insert_genesis_objects(objects);
6018        self.get_object_cache_reader()
6019            .force_reload_system_packages(&BuiltInFramework::all_package_ids());
6020        self.get_reconfig_api()
6021            .clear_state_end_of_epoch(&self.execution_lock_for_reconfiguration().await);
6022    }
6023}
6024
6025pub struct RandomnessRoundReceiver {
6026    authority_state: Arc<AuthorityState>,
6027    randomness_rx: mpsc::Receiver<(EpochId, RandomnessRound, Vec<u8>)>,
6028}
6029
6030impl RandomnessRoundReceiver {
6031    pub fn spawn(
6032        authority_state: Arc<AuthorityState>,
6033        randomness_rx: mpsc::Receiver<(EpochId, RandomnessRound, Vec<u8>)>,
6034    ) -> JoinHandle<()> {
6035        let rrr = RandomnessRoundReceiver {
6036            authority_state,
6037            randomness_rx,
6038        };
6039        spawn_monitored_task!(rrr.run())
6040    }
6041
6042    async fn run(mut self) {
6043        info!("RandomnessRoundReceiver event loop started");
6044
6045        loop {
6046            tokio::select! {
6047                maybe_recv = self.randomness_rx.recv() => {
6048                    if let Some((epoch, round, bytes)) = maybe_recv {
6049                        self.handle_new_randomness(epoch, round, bytes).await;
6050                    } else {
6051                        break;
6052                    }
6053                },
6054            }
6055        }
6056
6057        info!("RandomnessRoundReceiver event loop ended");
6058    }
6059
6060    #[instrument(level = "debug", skip_all, fields(?epoch, ?round))]
6061    async fn handle_new_randomness(&self, epoch: EpochId, round: RandomnessRound, bytes: Vec<u8>) {
6062        fail_point_async!("randomness-delay");
6063
6064        let epoch_store = self.authority_state.load_epoch_store_one_call_per_task();
6065        if epoch_store.epoch() != epoch {
6066            warn!(
6067                "dropping randomness for epoch {epoch}, round {round}, because we are in epoch {}",
6068                epoch_store.epoch()
6069            );
6070            return;
6071        }
6072        let key = TransactionKey::RandomnessRound(epoch, round);
6073        let transaction = VerifiedTransaction::new_randomness_state_update(
6074            epoch,
6075            round,
6076            bytes,
6077            epoch_store
6078                .epoch_start_config()
6079                .randomness_obj_initial_shared_version(),
6080        );
6081        debug!(
6082            "created randomness state update transaction with digest: {:?}",
6083            transaction.digest()
6084        );
6085        let transaction = VerifiedExecutableTransaction::new_system(transaction, epoch);
6086        let digest = *transaction.digest();
6087
6088        // Randomness state updates contain the full bls signature for the random round,
6089        // which cannot necessarily be reconstructed again later. Therefore we must
6090        // immediately persist this transaction. If we crash before its outputs
6091        // are committed, this ensures we will be able to re-execute it.
6092        self.authority_state
6093            .get_cache_commit()
6094            .persist_transaction(&transaction);
6095
6096        // Notify the scheduler that the transaction key now has a known digest
6097        if epoch_store.insert_tx_key(key, digest).is_err() {
6098            warn!("epoch ended while handling new randomness");
6099        }
6100
6101        // TODO: delete this when transaction manager is deleted
6102        match self.authority_state.execution_scheduler().as_ref() {
6103            ExecutionSchedulerWrapper::ExecutionScheduler(_) => {}
6104            ExecutionSchedulerWrapper::TransactionManager(manager) => {
6105                // Notifies transaction manager about transaction and output objects
6106                // committed. This provides necessary information to transaction manager
6107                // to start executing additional ready transactions.
6108                manager.notify_transaction_key(&epoch_store, key, digest);
6109            }
6110        }
6111
6112        let authority_state = self.authority_state.clone();
6113        spawn_monitored_task!(async move {
6114            // Wait for transaction execution in a separate task, to avoid deadlock in case
6115            // of out-of-order randomness generation. (Each
6116            // RandomnessStateUpdate depends on the output of the
6117            // RandomnessStateUpdate from the previous round.)
6118            //
6119            // We set a very long timeout so that in case this gets stuck for some reason,
6120            // the validator will eventually crash rather than continuing in a
6121            // zombie mode.
6122            const RANDOMNESS_STATE_UPDATE_EXECUTION_TIMEOUT: Duration = Duration::from_secs(300);
6123            let result = tokio::time::timeout(
6124                RANDOMNESS_STATE_UPDATE_EXECUTION_TIMEOUT,
6125                authority_state
6126                    .get_transaction_cache_reader()
6127                    .try_notify_read_executed_effects(
6128                        "RandomnessRoundReceiver::notify_read_executed_effects_first",
6129                        &[digest],
6130                    ),
6131            )
6132            .await;
6133            let result = match result {
6134                Ok(result) => result,
6135                Err(_) => {
6136                    if cfg!(debug_assertions) {
6137                        // Crash on randomness update execution timeout in debug builds.
6138                        panic!(
6139                            "randomness state update transaction execution timed out at epoch {epoch}, round {round}"
6140                        );
6141                    }
6142                    warn!(
6143                        "randomness state update transaction execution timed out at epoch {epoch}, round {round}"
6144                    );
6145                    // Continue waiting as long as necessary in non-debug builds.
6146                    authority_state
6147                        .get_transaction_cache_reader()
6148                        .try_notify_read_executed_effects(
6149                            "RandomnessRoundReceiver::notify_read_executed_effects_second",
6150                            &[digest],
6151                        )
6152                        .await
6153                }
6154            };
6155
6156            let mut effects = result.unwrap_or_else(|_| panic!("failed to get effects for randomness state update transaction at epoch {epoch}, round {round}"));
6157            let effects = effects.pop().expect("should return effects");
6158            if *effects.status() != ExecutionStatus::Success {
6159                fatal!(
6160                    "failed to execute randomness state update transaction at epoch {epoch}, round {round}: {effects:?}"
6161                );
6162            }
6163            debug!(
6164                "successfully executed randomness state update transaction at epoch {epoch}, round {round}"
6165            );
6166        });
6167    }
6168}
6169
6170#[async_trait]
6171impl TransactionKeyValueStoreTrait for AuthorityState {
6172    async fn multi_get(
6173        &self,
6174        transaction_keys: &[TransactionDigest],
6175        effects_keys: &[TransactionDigest],
6176    ) -> IotaResult<KVStoreTransactionData> {
6177        let txns = if !transaction_keys.is_empty() {
6178            self.get_transaction_cache_reader()
6179                .try_multi_get_transaction_blocks(transaction_keys)?
6180                .into_iter()
6181                .map(|t| t.map(|t| (*t).clone().into_inner()))
6182                .collect()
6183        } else {
6184            vec![]
6185        };
6186
6187        let fx = if !effects_keys.is_empty() {
6188            self.get_transaction_cache_reader()
6189                .try_multi_get_executed_effects(effects_keys)?
6190        } else {
6191            vec![]
6192        };
6193
6194        Ok((txns, fx))
6195    }
6196
6197    async fn multi_get_checkpoints(
6198        &self,
6199        checkpoint_summaries: &[CheckpointSequenceNumber],
6200        checkpoint_contents: &[CheckpointSequenceNumber],
6201        checkpoint_summaries_by_digest: &[CheckpointDigest],
6202    ) -> IotaResult<(
6203        Vec<Option<CertifiedCheckpointSummary>>,
6204        Vec<Option<CheckpointContents>>,
6205        Vec<Option<CertifiedCheckpointSummary>>,
6206    )> {
6207        // TODO: use multi-get methods if it ever becomes important (unlikely)
6208        let mut summaries = Vec::with_capacity(checkpoint_summaries.len());
6209        let store = self.get_checkpoint_store();
6210        for seq in checkpoint_summaries {
6211            let checkpoint = store
6212                .get_checkpoint_by_sequence_number(*seq)?
6213                .map(|c| c.into_inner());
6214
6215            summaries.push(checkpoint);
6216        }
6217
6218        let mut contents = Vec::with_capacity(checkpoint_contents.len());
6219        for seq in checkpoint_contents {
6220            let checkpoint = store
6221                .get_checkpoint_by_sequence_number(*seq)?
6222                .and_then(|summary| {
6223                    store
6224                        .get_checkpoint_contents(&summary.contents_digest)
6225                        .expect("db read cannot fail")
6226                });
6227            contents.push(checkpoint);
6228        }
6229
6230        let mut summaries_by_digest = Vec::with_capacity(checkpoint_summaries_by_digest.len());
6231        for digest in checkpoint_summaries_by_digest {
6232            let checkpoint = store
6233                .get_checkpoint_by_digest(digest)?
6234                .map(|c| c.into_inner());
6235            summaries_by_digest.push(checkpoint);
6236        }
6237
6238        Ok((summaries, contents, summaries_by_digest))
6239    }
6240
6241    async fn get_transaction_perpetual_checkpoint(
6242        &self,
6243        digest: TransactionDigest,
6244    ) -> IotaResult<Option<CheckpointSequenceNumber>> {
6245        self.get_checkpoint_cache()
6246            .try_get_transaction_perpetual_checkpoint(&digest)
6247            .map(|res| res.map(|(_epoch, checkpoint)| checkpoint))
6248    }
6249
6250    async fn get_object(
6251        &self,
6252        object_id: ObjectId,
6253        version: VersionNumber,
6254    ) -> IotaResult<Option<Object>> {
6255        self.get_object_cache_reader()
6256            .try_get_object_by_key(&object_id, version)
6257    }
6258
6259    #[instrument(skip_all)]
6260    async fn multi_get_objects(
6261        &self,
6262        object_keys: &[ObjectKey],
6263    ) -> IotaResult<Vec<Option<Object>>> {
6264        Ok(self
6265            .get_object_cache_reader()
6266            .multi_get_objects_by_key(object_keys))
6267    }
6268
6269    async fn multi_get_transactions_perpetual_checkpoints(
6270        &self,
6271        digests: &[TransactionDigest],
6272    ) -> IotaResult<Vec<Option<CheckpointSequenceNumber>>> {
6273        let res = self
6274            .get_checkpoint_cache()
6275            .try_multi_get_transactions_perpetual_checkpoints(digests)?;
6276
6277        Ok(res
6278            .into_iter()
6279            .map(|maybe| maybe.map(|(_epoch, checkpoint)| checkpoint))
6280            .collect())
6281    }
6282
6283    #[instrument(skip(self, digests), fields(digests = digests.iter().map(|d| d.to_string()).collect::<Vec<String>>().join(", ")))]
6284    async fn multi_get_events_by_tx_digests(
6285        &self,
6286        digests: &[TransactionDigest],
6287    ) -> IotaResult<Vec<Option<TransactionEvents>>> {
6288        if digests.is_empty() {
6289            return Ok(vec![]);
6290        }
6291
6292        Ok(self
6293            .get_transaction_cache_reader()
6294            .multi_get_events(digests))
6295    }
6296}
6297
6298#[cfg(msim)]
6299pub mod framework_injection {
6300    use std::{
6301        cell::RefCell,
6302        collections::{BTreeMap, BTreeSet},
6303    };
6304
6305    use iota_framework::{BuiltInFramework, SystemPackage};
6306    use iota_sdk_types::ObjectId;
6307    use iota_types::base_types::AuthorityName;
6308    use move_binary_format::CompiledModule;
6309
6310    type FrameworkOverrideConfig = BTreeMap<ObjectId, PackageOverrideConfig>;
6311
6312    // Thread local cache because all simtests run in a single unique thread.
6313    thread_local! {
6314        static OVERRIDE: RefCell<FrameworkOverrideConfig> = RefCell::new(FrameworkOverrideConfig::default());
6315    }
6316
6317    type Framework = Vec<CompiledModule>;
6318
6319    pub type PackageUpgradeCallback =
6320        Box<dyn Fn(AuthorityName) -> Option<Framework> + Send + Sync + 'static>;
6321
6322    enum PackageOverrideConfig {
6323        Global(Framework),
6324        PerValidator(PackageUpgradeCallback),
6325    }
6326
6327    fn compiled_modules_to_bytes(modules: &[CompiledModule]) -> Vec<Vec<u8>> {
6328        modules
6329            .iter()
6330            .map(|m| {
6331                let mut buf = Vec::new();
6332                m.serialize_with_version(m.version, &mut buf).unwrap();
6333                buf
6334            })
6335            .collect()
6336    }
6337
6338    pub fn set_override(package_id: ObjectId, modules: Vec<CompiledModule>) {
6339        OVERRIDE.with(|bs| {
6340            bs.borrow_mut()
6341                .insert(package_id, PackageOverrideConfig::Global(modules))
6342        });
6343    }
6344
6345    pub fn set_override_cb(package_id: ObjectId, func: PackageUpgradeCallback) {
6346        OVERRIDE.with(|bs| {
6347            bs.borrow_mut()
6348                .insert(package_id, PackageOverrideConfig::PerValidator(func))
6349        });
6350    }
6351
6352    pub fn get_override_bytes(package_id: &ObjectId, name: AuthorityName) -> Option<Vec<Vec<u8>>> {
6353        OVERRIDE.with(|cfg| {
6354            cfg.borrow().get(package_id).and_then(|entry| match entry {
6355                PackageOverrideConfig::Global(framework) => {
6356                    Some(compiled_modules_to_bytes(framework))
6357                }
6358                PackageOverrideConfig::PerValidator(func) => {
6359                    func(name).map(|fw| compiled_modules_to_bytes(&fw))
6360                }
6361            })
6362        })
6363    }
6364
6365    pub fn get_override_modules(
6366        package_id: &ObjectId,
6367        name: AuthorityName,
6368    ) -> Option<Vec<CompiledModule>> {
6369        OVERRIDE.with(|cfg| {
6370            cfg.borrow().get(package_id).and_then(|entry| match entry {
6371                PackageOverrideConfig::Global(framework) => Some(framework.clone()),
6372                PackageOverrideConfig::PerValidator(func) => func(name),
6373            })
6374        })
6375    }
6376
6377    pub fn get_override_system_package(
6378        package_id: &ObjectId,
6379        name: AuthorityName,
6380    ) -> Option<SystemPackage> {
6381        let bytes = get_override_bytes(package_id, name)?;
6382        // A built-in package keeps the dependencies it declares. Anything else is a
6383        // package being added -- including one at a system address that is not built in
6384        // yet -- and is assumed to depend on all existing system packages.
6385        let dependencies = BuiltInFramework::try_get_package_by_id(package_id)
6386            .map_or_else(BuiltInFramework::all_package_ids, |package| {
6387                package.dependencies.to_vec()
6388            });
6389        Some(SystemPackage {
6390            id: *package_id,
6391            bytes,
6392            dependencies,
6393        })
6394    }
6395
6396    pub fn get_extra_packages(name: AuthorityName) -> Vec<SystemPackage> {
6397        let built_in = BTreeSet::from_iter(BuiltInFramework::all_package_ids());
6398        let extra: Vec<ObjectId> = OVERRIDE.with(|cfg| {
6399            cfg.borrow()
6400                .keys()
6401                .filter_map(|package| (!built_in.contains(package)).then_some(*package))
6402                .collect()
6403        });
6404
6405        extra
6406            .into_iter()
6407            .map(|package| SystemPackage {
6408                id: package,
6409                bytes: get_override_bytes(&package, name).unwrap(),
6410                dependencies: BuiltInFramework::all_package_ids(),
6411            })
6412            .collect()
6413    }
6414}
6415
6416#[derive(Debug, Serialize, Deserialize, Clone)]
6417pub struct ObjDumpFormat {
6418    pub id: ObjectId,
6419    pub version: VersionNumber,
6420    pub digest: ObjectDigest,
6421    pub object: Object,
6422}
6423
6424impl ObjDumpFormat {
6425    fn new(object: Object) -> Self {
6426        let oref = object.object_ref();
6427        Self {
6428            id: oref.object_id,
6429            version: oref.version,
6430            digest: oref.digest,
6431            object,
6432        }
6433    }
6434}
6435
6436#[derive(Debug, Serialize, Deserialize, Clone)]
6437pub struct NodeStateDump {
6438    pub tx_digest: TransactionDigest,
6439    pub sender_signed_data: SenderSignedTransaction,
6440    pub executed_epoch: u64,
6441    pub reference_gas_price: u64,
6442    pub protocol_version: u64,
6443    pub epoch_start_timestamp_ms: u64,
6444    pub computed_effects: TransactionEffects,
6445    pub expected_effects_digest: TransactionEffectsDigest,
6446    pub relevant_system_packages: Vec<ObjDumpFormat>,
6447    pub shared_objects: Vec<ObjDumpFormat>,
6448    pub loaded_child_objects: Vec<ObjDumpFormat>,
6449    pub modified_at_versions: Vec<ObjDumpFormat>,
6450    pub runtime_reads: Vec<ObjDumpFormat>,
6451    pub input_objects: Vec<ObjDumpFormat>,
6452}
6453
6454impl NodeStateDump {
6455    pub fn new(
6456        tx_digest: &TransactionDigest,
6457        effects: &TransactionEffects,
6458        expected_effects_digest: TransactionEffectsDigest,
6459        object_store: &dyn ObjectStore,
6460        epoch_store: &Arc<AuthorityPerEpochStore>,
6461        inner_temporary_store: &InnerTemporaryStore,
6462        transaction: &VerifiedExecutableTransaction,
6463    ) -> IotaResult<Self> {
6464        // Epoch info
6465        let executed_epoch = epoch_store.epoch();
6466        let reference_gas_price = epoch_store.reference_gas_price();
6467        let epoch_start_config = epoch_store.epoch_start_config();
6468        let protocol_version = epoch_store.protocol_version().as_u64();
6469        let epoch_start_timestamp_ms = epoch_start_config.epoch_data().epoch_start_timestamp();
6470
6471        // Record all system packages at this version
6472        let mut relevant_system_packages = Vec::new();
6473        for sys_package_id in BuiltInFramework::all_package_ids() {
6474            if let Some(w) = object_store.try_get_object(&sys_package_id)? {
6475                relevant_system_packages.push(ObjDumpFormat::new(w))
6476            }
6477        }
6478
6479        // Record all the shared objects
6480        let mut shared_objects = Vec::new();
6481        for kind in effects.input_shared_objects() {
6482            match kind {
6483                InputSharedObject::Mutate(obj_ref) | InputSharedObject::ReadOnly(obj_ref) => {
6484                    if let Some(w) =
6485                        object_store.try_get_object_by_key(&obj_ref.object_id, obj_ref.version)?
6486                    {
6487                        shared_objects.push(ObjDumpFormat::new(w))
6488                    }
6489                }
6490                InputSharedObject::ReadDeleted(..)
6491                | InputSharedObject::MutateDeleted(..)
6492                // TODO: consider record congested objects.
6493                | InputSharedObject::Canceled(..) => (),
6494                _ => unimplemented!(
6495                    "a new InputSharedObject enum variant was added and needs to be handled"
6496                ),
6497            }
6498        }
6499
6500        // Record all loaded child objects
6501        // Child objects which are read but not mutated are not tracked anywhere else
6502        let mut loaded_child_objects = Vec::new();
6503        for (id, meta) in &inner_temporary_store.loaded_runtime_objects {
6504            if let Some(w) = object_store.try_get_object_by_key(id, meta.version)? {
6505                loaded_child_objects.push(ObjDumpFormat::new(w))
6506            }
6507        }
6508
6509        // Record all modified objects
6510        let mut modified_at_versions = Vec::new();
6511        for modified in effects.modified_at_versions() {
6512            let (id, ver) = (modified.object_id(), modified.version());
6513            if let Some(w) = object_store.try_get_object_by_key(id, ver)? {
6514                modified_at_versions.push(ObjDumpFormat::new(w))
6515            }
6516        }
6517
6518        // Packages read at runtime, which were not previously loaded into the temoorary
6519        // store Some packages may be fetched at runtime and wont show up in
6520        // input objects
6521        let mut runtime_reads = Vec::new();
6522        for obj in inner_temporary_store
6523            .runtime_packages_loaded_from_db
6524            .values()
6525        {
6526            runtime_reads.push(ObjDumpFormat::new(obj.object().clone()));
6527        }
6528
6529        // All other input objects should already be in `inner_temporary_store.objects`
6530
6531        Ok(Self {
6532            tx_digest: *tx_digest,
6533            executed_epoch,
6534            reference_gas_price,
6535            epoch_start_timestamp_ms,
6536            protocol_version,
6537            relevant_system_packages,
6538            shared_objects,
6539            loaded_child_objects,
6540            modified_at_versions,
6541            runtime_reads,
6542            sender_signed_data: transaction.clone().into_message(),
6543            input_objects: inner_temporary_store
6544                .input_objects
6545                .values()
6546                .map(|o| ObjDumpFormat::new(o.clone()))
6547                .collect(),
6548            computed_effects: effects.clone(),
6549            expected_effects_digest,
6550        })
6551    }
6552
6553    pub fn all_objects(&self) -> Vec<ObjDumpFormat> {
6554        let mut objects = Vec::new();
6555        objects.extend(self.relevant_system_packages.clone());
6556        objects.extend(self.shared_objects.clone());
6557        objects.extend(self.loaded_child_objects.clone());
6558        objects.extend(self.modified_at_versions.clone());
6559        objects.extend(self.runtime_reads.clone());
6560        objects.extend(self.input_objects.clone());
6561        objects
6562    }
6563
6564    pub fn write_to_file(&self, path: &Path) -> Result<PathBuf, anyhow::Error> {
6565        let file_name = format!(
6566            "{}_{}_NODE_DUMP.json",
6567            self.tx_digest,
6568            AuthorityState::unixtime_now_ms()
6569        );
6570        let mut path = path.to_path_buf();
6571        path.push(&file_name);
6572        let mut file = File::create(path.clone())?;
6573        file.write_all(serde_json::to_string_pretty(self)?.as_bytes())?;
6574        Ok(path)
6575    }
6576
6577    pub fn read_from_file(path: &PathBuf) -> Result<Self, anyhow::Error> {
6578        let file = File::open(path)?;
6579        serde_json::from_reader(file).map_err(|e| anyhow::anyhow!(e))
6580    }
6581}
6582
6583/// The reason [`AuthorityState::check_move_account`] failed.
6584/// Execution treats the two cases differently, because only
6585/// the first one is about the transaction.
6586enum MoveAccountCheckError {
6587    /// The account fails a check.
6588    Input(UserInputError),
6589
6590    /// This validator failed to read its store.
6591    Storage(IotaError),
6592}
6593
6594impl From<UserInputError> for MoveAccountCheckError {
6595    fn from(error: UserInputError) -> Self {
6596        Self::Input(error)
6597    }
6598}
6599
6600impl From<MoveAccountCheckError> for IotaError {
6601    fn from(error: MoveAccountCheckError) -> Self {
6602        match error {
6603            MoveAccountCheckError::Input(error) => IotaError::UserInput { error },
6604            MoveAccountCheckError::Storage(error) => error,
6605        }
6606    }
6607}
6608
6609/// Returns the [`MoveAuthenticator`]s to execute during the pre-consensus
6610/// phase.
6611///
6612/// When `pre_consensus_sponsor_only_move_authentication` is enabled:
6613/// - For sponsored transactions: only the sponsor's [`MoveAuthenticator`] is returned (empty if the
6614///   sponsor does not use one).
6615/// - For non-sponsored transactions: all [`MoveAuthenticator`]s are returned (currently only the
6616///   sender's).
6617///
6618/// When the flag is not set, all [`MoveAuthenticator`]s are returned for
6619/// compatibility.
6620fn pre_consensus_move_authenticators<'a>(
6621    tx: &'a VerifiedTransaction,
6622    protocol_config: &ProtocolConfig,
6623) -> Vec<&'a MoveAuthenticator> {
6624    if protocol_config.pre_consensus_sponsor_only_move_authentication() {
6625        if tx.transaction().is_sponsored_tx() {
6626            if let Some(sponsor_move_authenticator) = tx.sponsor_move_authenticator() {
6627                vec![sponsor_move_authenticator]
6628            } else {
6629                vec![]
6630            }
6631        } else {
6632            tx.move_authenticators()
6633        }
6634    } else {
6635        tx.move_authenticators()
6636    }
6637}