1use 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
242pub 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 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 pub(crate) consensus_queue_load_shedding_percentage: IntGauge,
298 pub(crate) local_post_consensus_load_shedding_percentage: IntGauge,
302
303 pub(crate) transaction_overload_sources: IntCounterVec,
304
305 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 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 pub consensus_handler_validation_dropped_transactions: IntCounter,
323 pub consensus_handler_load_shedding_dropped_transactions: IntCounter,
327 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 pub bytecode_verifier_metrics: Arc<BytecodeVerifierMetrics>,
348
349 pub multisig_sig_count: IntCounter,
351
352 pub execution_queueing_latency: LatencyObserver,
355
356 pub txn_ready_rate_tracker: Arc<Mutex<RateTracker>>,
363
364 pub execution_rate_tracker: Arc<Mutex<RateTracker>>,
367}
368
369const 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
381const 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: 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 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
857pub type StableSyncAuthoritySigner = Pin<Arc<dyn Signer<AuthoritySignature> + Send + Sync>>;
863
864#[derive(Debug, Clone, Default)]
868pub struct ExecutionEnv {
869 pub assigned_versions: AssignedVersions,
871 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 pub name: AuthorityName,
901 pub secret: StableSyncAuthoritySigner,
903
904 input_loader: TransactionInputLoader,
906 execution_cache_trait_pointers: ExecutionCacheTraitPointers,
907
908 epoch_store: ArcSwap<AuthorityPerEpochStore>,
909
910 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 execution_scheduler: Arc<ExecutionSchedulerWrapper>,
927
928 #[cfg_attr(not(test), expect(unused))]
930 tx_execution_shutdown: Mutex<Option<oneshot::Sender<()>>>,
931
932 pub metrics: Arc<AuthorityMetrics>,
933 pruner: AuthorityStorePruner,
936 authority_per_epoch_pruner: AuthorityPerEpochStorePruner,
937 checkpoint_progress_tracker: Option<Arc<CheckpointProgressTracker>>,
938
939 pub config: NodeConfig,
940
941 pub overload_info: AuthorityOverloadInfo,
943
944 pub validator_tx_finalizer: Option<Arc<ValidatorTxFinalizer<NetworkAuthorityClient>>>,
945
946 chain_identifier: ChainIdentifier,
949
950 pub(crate) congestion_tracker: Arc<CongestionTracker>,
951
952 pub traffic_controller: Option<Arc<TrafficController>>,
954 epoch_end_db_snapshots: Option<EpochEndDbSnapshotHandle>,
957}
958
959impl 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 #[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 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 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 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 let per_authenticator_checked_input_objects: Vec<_> = per_authenticator_checked_inputs
1095 .iter()
1096 .map(|i| &i.0)
1097 .collect();
1098
1099 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 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_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 let pre_consensus_move_authenticators =
1152 pre_consensus_move_authenticators(transaction, protocol_config);
1153 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 !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 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 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 async fn handle_transaction_impl(
1248 &self,
1249 transaction: VerifiedTransaction,
1250 epoch_store: &Arc<AuthorityPerEpochStore>,
1251 ) -> IotaResult<VerifiedSignedTransaction> {
1252 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 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 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 #[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 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 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 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 self.check_consensus_queue_graduated_limits(consensus_adapter, tx)
1370 .tap_err(|_| {
1371 self.update_overload_metrics("consensus");
1372 })?;
1373
1374 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 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 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 #[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 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 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 #[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 let tx_guard = epoch_store.acquire_tx_guard(transaction)?;
1552
1553 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 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 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 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 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 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 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 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 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 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 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 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 self.get_object_cache_reader()
1894 .force_reload_system_packages(&BuiltInFramework::all_package_ids());
1895 }
1896
1897 tx_guard.commit_tx();
1899
1900 match self.execution_scheduler.as_ref() {
1901 ExecutionSchedulerWrapper::ExecutionScheduler(_) => {}
1902 ExecutionSchedulerWrapper::TransactionManager(tm) => {
1903 tm.notify_commit(tx_digest, output_keys, epoch_store);
1907 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 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 #[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 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 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 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 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 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 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 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 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 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 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 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 #[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 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 transaction.validity_check_no_gas_check(epoch_store.protocol_config())?;
2339
2340 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 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 let (mut input_objects, receiving_objects) = self.input_loader.read_objects_for_signing(
2361 None,
2363 &input_object_kinds,
2364 &receiving_object_refs,
2365 epoch_store.epoch(),
2366 )?;
2367
2368 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 let authenticator_gas_budget = 0;
2391
2392 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 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 let executor = iota_execution::executor(
2435 protocol_config,
2436 true, None,
2438 )
2439 .expect("Creating an executor should not fail here");
2440
2441 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, 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 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 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 #[instrument(level = "debug", skip_all, err)]
2505 fn index_tx(
2506 &self,
2507 indexes: &IndexStore,
2508 digest: &TransactionDigest,
2509 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 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 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 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 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 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 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 let Some(move_object) = o.data.as_opt_struct().cloned() else {
2745 return Ok(None);
2746 };
2747
2748 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(true),
2801 object_id: o.id(),
2802 version: o.version(),
2803 digest: o.digest(),
2804 },
2805
2806 DFV::ValueMetadata::DynamicObjectField(object_id) => {
2807 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 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 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 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 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)] 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 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 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 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 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 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 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 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 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 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 let mut execution_lock = self.execution_lock_for_reconfiguration().await;
3501
3502 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 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 Ok(new_epoch_store)
3560 }
3561
3562 pub async fn reconfigure_for_testing(&self) {
3567 self.reconfigure_for_testing_impl(None).await;
3568 }
3569
3570 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 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 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 #[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 pub fn load_epoch_store_one_call_per_task(&self) -> Guard<Arc<AuthorityPerEpochStore>> {
3708 self.epoch_store.load()
3709 }
3710
3711 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 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 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 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 let found = match epoch_store.multi_get_transaction_checkpoint(&remaining) {
3802 Ok(found) => found,
3803 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 async fn wait_for_next_epoch_store(
3839 &self,
3840 prev_epoch: EpochId,
3841 deadline: tokio::time::Instant,
3842 ) -> Option<Arc<AuthorityPerEpochStore>> {
3843 const EPOCH_STORE_SWAP_POLL_INTERVAL: Duration = Duration::from_millis(100);
3848
3849 loop {
3850 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 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 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 #[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 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 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 .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 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 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 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 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 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 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 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 cursor: Option<EventID>,
4434 limit: usize,
4435 descending: bool,
4436 ) -> IotaResult<Vec<IotaEvent>> {
4437 let index_store = self.get_indexes()?;
4438
4439 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 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 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 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 #[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 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 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 #[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 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 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 #[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 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 #[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 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 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 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 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 #[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 #[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 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 info!(
5014 "Framework {} does not need updating",
5015 system_package_ref.object_id
5016 );
5017 continue;
5018 }
5019
5020 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 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 let mut desired_upgrades: Vec<_> = capabilities
5111 .into_iter()
5112 .filter_map(|mut cap| {
5113 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 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 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 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 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 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 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 if let Some(digest) = capability
5224 .supported_protocol_versions
5225 .get_version_digest(target_protocol_version)
5226 {
5227 if digest == target_digest {
5228 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 eligible_validators.sort();
5241 eligible_validators
5242 }
5243
5244 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 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 #[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 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 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 return Err(CheckpointBuilderError::SystemPackagesMissing);
5352 };
5353
5354 if config.select_committee_from_eligible_validators() {
5357 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 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 let eligible_validators_weight = Self::calculate_eligible_validators_weight(
5375 &eligible_active_validators,
5376 &active_validators,
5377 epoch_store.committee(),
5378 );
5379
5380 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 eligible_active_validators = (0..active_validators.len() as u64).collect();
5394 }
5395 }
5396
5397 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 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 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 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 assert!(system_epoch_info_event.is_some() || system_obj.safe_mode());
5531
5532 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 assert!(effects.status().is_success());
5545 Ok((system_obj, system_epoch_info_event, effects))
5546 }
5547
5548 #[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 epoch_store.close_user_certs(state);
5568 }
5569 }
5571
5572 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 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 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 _ => 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 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 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 (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 (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 (ObjectReadResultKind::CancelledTransactionObject(version), false) => {
5833 Err(UserInputError::AccountObjectInCanceledTransaction {
5834 account_id: account_object.id(),
5835 account_version: *version,
5836 })
5837 }
5838 (ObjectReadResultKind::DeletedSharedObject(version, _), true) => Ok(*version),
5843 (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 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 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 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 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 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 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 self.authority_state
6093 .get_cache_commit()
6094 .persist_transaction(&transaction);
6095
6096 if epoch_store.insert_tx_key(key, digest).is_err() {
6098 warn!("epoch ended while handling new randomness");
6099 }
6100
6101 match self.authority_state.execution_scheduler().as_ref() {
6103 ExecutionSchedulerWrapper::ExecutionScheduler(_) => {}
6104 ExecutionSchedulerWrapper::TransactionManager(manager) => {
6105 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 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 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 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 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! {
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 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 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 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 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 | InputSharedObject::Canceled(..) => (),
6494 _ => unimplemented!(
6495 "a new InputSharedObject enum variant was added and needs to be handled"
6496 ),
6497 }
6498 }
6499
6500 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 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 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 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
6583enum MoveAccountCheckError {
6587 Input(UserInputError),
6589
6590 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
6609fn 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}