1use std::{
6 collections::{BTreeMap, HashSet},
7 path::PathBuf,
8 sync::{Arc, Mutex},
9};
10
11use futures::executor::block_on;
12use iota_config::node::ExpensiveSafetyCheckConfig;
13use iota_core::authority::NodeStateDump;
14use iota_execution::Executor;
15use iota_framework::BuiltInFramework;
16use iota_json_rpc_types::{
17 IotaExecutionStatus, IotaTransactionBlockEffects, IotaTransactionBlockEffectsAPI,
18};
19use iota_protocol_config::{Chain, ProtocolConfig};
20use iota_sdk::{IotaClient, IotaClientBuilder};
21use iota_sdk_types::{
22 GasPayment, MoveAuthenticator, ObjectDigest, ObjectId, ObjectReference, Owner,
23 SenderSignedTransaction, TransactionDigest, TransactionKind, Version,
24};
25use iota_types::{
26 IOTA_DENY_LIST_OBJECT_ID,
27 account_abstraction::authenticator_function::{
28 AuthenticatorFunctionRefForExecution,
29 authenticator_function_ref_v1_from_dynamic_field_object,
30 derive_authenticator_function_ref_v1_dynamic_field_id, extract_auth_fun_refs,
31 },
32 auth_context::AuthContextData,
33 base_types::VersionNumber,
34 committee::EpochId,
35 error::{ExecutionError, IotaError, IotaResult},
36 executable_transaction::VerifiedExecutableTransaction,
37 execution::SharedInput,
38 gas::IotaGasStatus,
39 in_memory_storage::InMemoryStorage,
40 inner_temporary_store::InnerTemporaryStore,
41 message_envelope::Message,
42 metrics::LimitsMetrics,
43 move_authenticator::MoveAuthenticatorExt,
44 object::Object,
45 storage::{
46 BackingPackageStore, ChildObjectResolver, ObjectStore, PackageObject, get_module,
47 get_module_by_id,
48 },
49 transaction::{
50 CheckedInputObjects, InputObjectKind, InputObjects, ObjectReadResult, ObjectReadResultKind,
51 SenderSignedTransactionAPI, TransactionAPI, TransactionEnvelope, VerifiedTransaction,
52 },
53};
54use move_binary_format::CompiledModule;
55use move_bytecode_utils::module_cache::GetModule;
56use move_core_types::{language_storage::ModuleId, resolver::ModuleResolver};
57use prometheus_filtered::Registry;
58use serde::{Deserialize, Serialize};
59use similar::{ChangeTag, TextDiff};
60use tracing::{error, info, trace, warn};
61
62use crate::{
63 chain_from_chain_id,
64 data_fetcher::{
65 DataFetcher, Fetchers, NodeStateDumpFetcher, RemoteFetcher, extract_epoch_and_version,
66 },
67 displays::{
68 Pretty,
69 transaction_displays::{FullPTB, transform_command_results_to_annotated},
70 },
71 types::*,
72};
73
74#[derive(Debug, Serialize, Deserialize)]
77pub struct ExecutionSandboxState {
78 pub transaction_info: OnChainTransactionInfo,
80 pub required_objects: Vec<Object>,
82 #[serde(skip)]
85 pub local_exec_temporary_store: Option<InnerTemporaryStore>,
86 pub local_exec_effects: IotaTransactionBlockEffects,
88 #[serde(skip)]
90 pub local_exec_status: Option<Result<(), ExecutionError>>,
91}
92
93impl ExecutionSandboxState {
94 pub fn check_effects(&self) -> Result<(), ReplayEngineError> {
95 if self.transaction_info.effects != self.local_exec_effects {
96 error!("Replay tool forked {}", self.transaction_info.tx_digest);
97 let diff = self.diff_effects();
98 println!("{diff}");
99 return Err(ReplayEngineError::EffectsForked {
100 digest: self.transaction_info.tx_digest,
101 diff: format!("\n{diff}"),
102 on_chain: Box::new(self.transaction_info.effects.clone()),
103 local: Box::new(self.local_exec_effects.clone()),
104 });
105 }
106 Ok(())
107 }
108
109 pub fn diff_effects(&self) -> String {
111 let eff1 = &self.transaction_info.effects;
112 let eff2 = &self.local_exec_effects;
113 let on_chain_str = format!("{eff1:#?}");
114 let local_chain_str = format!("{eff2:#?}");
115 let mut res = vec![];
116
117 let diff = TextDiff::from_lines(&on_chain_str, &local_chain_str);
118 for change in diff.iter_all_changes() {
119 let sign = match change.tag() {
120 ChangeTag::Delete => "---",
121 ChangeTag::Insert => "+++",
122 ChangeTag::Equal => " ",
123 };
124 res.push(format!("{sign}{change}"));
125 }
126
127 res.join("")
128 }
129}
130
131#[derive(Debug, Clone, PartialEq, Eq)]
132pub struct ProtocolVersionSummary {
133 pub protocol_version: u64,
135 pub epoch_start: u64,
137 pub epoch_end: u64,
139 pub checkpoint_start: Option<u64>,
141 pub checkpoint_end: Option<u64>,
143 pub epoch_change_tx: TransactionDigest,
145}
146
147#[derive(Clone)]
148pub struct Storage {
149 pub live_objects_store: Arc<Mutex<BTreeMap<ObjectId, Object>>>,
154
155 pub package_cache: Arc<Mutex<BTreeMap<ObjectId, Object>>>,
158 pub object_version_cache: Arc<Mutex<BTreeMap<(ObjectId, Version), Object>>>,
161}
162
163impl std::fmt::Display for Storage {
164 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
165 writeln!(f, "Live object store")?;
166 for (id, obj) in self
167 .live_objects_store
168 .lock()
169 .expect("Unable to lock")
170 .iter()
171 {
172 writeln!(f, "{}: {:?}", id, obj.object_ref())?;
173 }
174 writeln!(f, "Package cache")?;
175 for (id, obj) in self.package_cache.lock().expect("Unable to lock").iter() {
176 writeln!(f, "{}: {:?}", id, obj.object_ref())?;
177 }
178 writeln!(f, "Object version cache")?;
179 for id in self
180 .object_version_cache
181 .lock()
182 .expect("Unable to lock")
183 .keys()
184 {
185 writeln!(f, "{}: {}", id.0, id.1)?;
186 }
187
188 write!(f, "")
189 }
190}
191
192impl Storage {
193 pub fn default() -> Self {
194 Self {
195 live_objects_store: Arc::new(Mutex::new(BTreeMap::new())),
196 package_cache: Arc::new(Mutex::new(BTreeMap::new())),
197 object_version_cache: Arc::new(Mutex::new(BTreeMap::new())),
198 }
199 }
200
201 pub fn all_objects(&self) -> Vec<Object> {
202 self.live_objects_store
203 .lock()
204 .expect("Unable to lock")
205 .values()
206 .cloned()
207 .chain(
208 self.package_cache
209 .lock()
210 .expect("Unable to lock")
211 .values()
212 .cloned(),
213 )
214 .chain(
215 self.object_version_cache
216 .lock()
217 .expect("Unable to lock")
218 .values()
219 .cloned(),
220 )
221 .collect::<Vec<_>>()
222 }
223}
224
225#[derive(Clone)]
226pub struct LocalExec {
227 pub client: Option<IotaClient>,
228 pub protocol_version_epoch_table: BTreeMap<u64, ProtocolVersionSummary>,
231 pub protocol_version_system_package_table: BTreeMap<u64, BTreeMap<ObjectId, Version>>,
233 pub current_protocol_version: u64,
235 pub storage: Storage,
237 pub exec_store_events: Arc<Mutex<Vec<ExecutionStoreEvent>>>,
239 pub metrics: Arc<LimitsMetrics>,
241 pub fetcher: Fetchers,
243
244 pub executor_version: Option<i64>,
247 pub protocol_version: Option<i64>,
251 pub enable_profiler: Option<PathBuf>,
254 pub config_and_versions: Option<Vec<(ObjectId, Version)>>,
255 pub num_retries_for_timeout: u32,
257 pub sleep_period_for_timeout: std::time::Duration,
258}
259
260impl LocalExec {
261 pub async fn multi_download(
264 &self,
265 objs: &[(ObjectId, Version)],
266 ) -> Result<Vec<Object>, ReplayEngineError> {
267 let mut num_retries_for_timeout = self.num_retries_for_timeout as i64;
268 while num_retries_for_timeout >= 0 {
269 match self.fetcher.multi_get_versioned(objs).await {
270 Ok(objs) => return Ok(objs),
271 Err(ReplayEngineError::IotaRpcRequestTimeout) => {
272 warn!(
273 "RPC request timed out. Retries left {}. Sleeping for {}s",
274 num_retries_for_timeout,
275 self.sleep_period_for_timeout.as_secs()
276 );
277 num_retries_for_timeout -= 1;
278 tokio::time::sleep(self.sleep_period_for_timeout).await;
279 }
280 Err(e) => return Err(e),
281 }
282 }
283 Err(ReplayEngineError::IotaRpcRequestTimeout)
284 }
285 pub async fn multi_download_latest(
288 &self,
289 objs: &[ObjectId],
290 ) -> Result<Vec<Object>, ReplayEngineError> {
291 let mut num_retries_for_timeout = self.num_retries_for_timeout as i64;
292 while num_retries_for_timeout >= 0 {
293 match self.fetcher.multi_get_latest(objs).await {
294 Ok(objs) => return Ok(objs),
295 Err(ReplayEngineError::IotaRpcRequestTimeout) => {
296 warn!(
297 "RPC request timed out. Retries left {}. Sleeping for {}s",
298 num_retries_for_timeout,
299 self.sleep_period_for_timeout.as_secs()
300 );
301 num_retries_for_timeout -= 1;
302 tokio::time::sleep(self.sleep_period_for_timeout).await;
303 }
304 Err(e) => return Err(e),
305 }
306 }
307 Err(ReplayEngineError::IotaRpcRequestTimeout)
308 }
309
310 pub async fn fetch_loaded_child_refs(
311 &self,
312 tx_digest: &TransactionDigest,
313 ) -> Result<Vec<(ObjectId, Version)>, ReplayEngineError> {
314 self.fetcher.get_loaded_child_objects(tx_digest).await
316 }
317
318 pub async fn new_from_fn_url(http_url: &str) -> Result<Self, ReplayEngineError> {
319 Self::new_for_remote(
320 IotaClientBuilder::default()
321 .request_timeout(RPC_TIMEOUT_ERR_SLEEP_RETRY_PERIOD)
322 .max_concurrent_requests(MAX_CONCURRENT_REQUESTS)
323 .build(http_url)
324 .await?,
325 None,
326 )
327 .await
328 }
329
330 pub async fn replay_with_network_config(
331 rpc_url: String,
332 tx_digest: TransactionDigest,
333 expensive_safety_check_config: ExpensiveSafetyCheckConfig,
334 use_authority: bool,
335 executor_version: Option<i64>,
336 protocol_version: Option<i64>,
337 enable_profiler: Option<PathBuf>,
338 config_and_versions: Option<Vec<(ObjectId, Version)>>,
339 ) -> Result<ExecutionSandboxState, ReplayEngineError> {
340 info!("Using RPC URL: {}", rpc_url);
341 LocalExec::new_from_fn_url(&rpc_url)
342 .await?
343 .init_for_execution()
344 .await?
345 .execute_transaction(
346 &tx_digest,
347 expensive_safety_check_config,
348 use_authority,
349 executor_version,
350 protocol_version,
351 enable_profiler,
352 config_and_versions,
353 )
354 .await
355 }
356
357 pub async fn init_for_execution(mut self) -> Result<Self, ReplayEngineError> {
362 self.populate_protocol_version_tables().await?;
363 tokio::task::yield_now().await;
364 Ok(self)
365 }
366
367 pub async fn reset_for_new_execution_with_client(self) -> Result<Self, ReplayEngineError> {
368 Self::new_for_remote(
369 self.client.expect("Remote client not initialized"),
370 Some(self.fetcher.into_remote()),
371 )
372 .await?
373 .init_for_execution()
374 .await
375 }
376
377 pub async fn new_for_remote(
378 client: IotaClient,
379 remote_fetcher: Option<RemoteFetcher>,
380 ) -> Result<Self, ReplayEngineError> {
381 let registry = prometheus_filtered::Registry::new();
383 let metrics = Arc::new(LimitsMetrics::new(®istry));
384
385 let fetcher = remote_fetcher.unwrap_or(RemoteFetcher::new(client.clone()));
386
387 Ok(Self {
388 client: Some(client),
389 protocol_version_epoch_table: BTreeMap::new(),
390 protocol_version_system_package_table: BTreeMap::new(),
391 current_protocol_version: 0,
392 exec_store_events: Arc::new(Mutex::new(Vec::new())),
393 metrics,
394 storage: Storage::default(),
395 fetcher: Fetchers::Remote(fetcher),
396 num_retries_for_timeout: RPC_TIMEOUT_ERR_NUM_RETRIES,
398 sleep_period_for_timeout: RPC_TIMEOUT_ERR_SLEEP_RETRY_PERIOD,
399 executor_version: None,
400 protocol_version: None,
401 enable_profiler: None,
402 config_and_versions: None,
403 })
404 }
405
406 pub async fn new_for_state_dump(
407 path: &str,
408 backup_rpc_url: Option<String>,
409 ) -> Result<Self, ReplayEngineError> {
410 let registry = prometheus_filtered::Registry::new();
412 let metrics = Arc::new(LimitsMetrics::new(®istry));
413
414 let state = NodeStateDump::read_from_file(&PathBuf::from(path))?;
415 let current_protocol_version = state.protocol_version;
416 let fetcher = match backup_rpc_url {
417 Some(url) => NodeStateDumpFetcher::new(
418 state,
419 Some(RemoteFetcher::new(
420 IotaClientBuilder::default()
421 .request_timeout(RPC_TIMEOUT_ERR_SLEEP_RETRY_PERIOD)
422 .max_concurrent_requests(MAX_CONCURRENT_REQUESTS)
423 .build(url)
424 .await?,
425 )),
426 ),
427 None => NodeStateDumpFetcher::new(state, None),
428 };
429
430 Ok(Self {
431 client: None,
432 protocol_version_epoch_table: BTreeMap::new(),
433 protocol_version_system_package_table: BTreeMap::new(),
434 current_protocol_version,
435 exec_store_events: Arc::new(Mutex::new(Vec::new())),
436 metrics,
437 storage: Storage::default(),
438 fetcher: Fetchers::NodeStateDump(fetcher),
439 num_retries_for_timeout: RPC_TIMEOUT_ERR_NUM_RETRIES,
441 sleep_period_for_timeout: RPC_TIMEOUT_ERR_SLEEP_RETRY_PERIOD,
442 executor_version: None,
443 protocol_version: None,
444 enable_profiler: None,
445 config_and_versions: None,
446 })
447 }
448
449 pub async fn multi_download_and_store(
450 &mut self,
451 objs: &[(ObjectId, Version)],
452 ) -> Result<Vec<Object>, ReplayEngineError> {
453 let objs = self.multi_download(objs).await?;
454
455 for obj in objs.iter() {
457 let o_ref = obj.object_ref();
458 self.storage
459 .live_objects_store
460 .lock()
461 .expect("Can't lock")
462 .insert(o_ref.object_id, obj.clone());
463 self.storage
464 .object_version_cache
465 .lock()
466 .expect("Cannot lock")
467 .insert((o_ref.object_id, o_ref.version), obj.clone());
468 if obj.is_package() {
469 self.storage
470 .package_cache
471 .lock()
472 .expect("Cannot lock")
473 .insert(o_ref.object_id, obj.clone());
474 }
475 }
476 tokio::task::yield_now().await;
477 Ok(objs)
478 }
479
480 pub async fn multi_download_relevant_packages_and_store(
481 &mut self,
482 objs: Vec<ObjectId>,
483 protocol_version: u64,
484 ) -> Result<Vec<Object>, ReplayEngineError> {
485 let syst_packages_objs = if self.protocol_version.is_some_and(|i| i < 0) {
486 BuiltInFramework::genesis_objects().collect()
487 } else {
488 let syst_packages =
489 self.system_package_versions_for_protocol_version(protocol_version)?;
490 self.multi_download(&syst_packages).await?
491 };
492
493 let non_system_package_objs: Vec<_> = objs
496 .into_iter()
497 .filter(|o| !Self::system_package_ids(self.current_protocol_version).contains(o))
498 .collect();
499 let objs = self
500 .multi_download_latest(&non_system_package_objs)
501 .await?
502 .into_iter()
503 .chain(syst_packages_objs);
504
505 for obj in objs.clone() {
506 let o_ref = obj.object_ref();
507 self.storage
510 .object_version_cache
511 .lock()
512 .expect("Cannot lock")
513 .insert((o_ref.object_id, o_ref.version), obj.clone());
514 if obj.is_package() {
515 self.storage
516 .package_cache
517 .lock()
518 .expect("Cannot lock")
519 .insert(o_ref.object_id, obj.clone());
520 }
521 }
522 Ok(objs.collect())
523 }
524
525 #[expect(clippy::disallowed_methods)]
527 pub fn download_object(
528 &self,
529 object_id: &ObjectId,
530 version: Version,
531 ) -> Result<Object, ReplayEngineError> {
532 if self
533 .storage
534 .object_version_cache
535 .lock()
536 .expect("Cannot lock")
537 .contains_key(&(*object_id, version))
538 {
539 return Ok(self
540 .storage
541 .object_version_cache
542 .lock()
543 .expect("Cannot lock")
544 .get(&(*object_id, version))
545 .ok_or(ReplayEngineError::InternalCacheInvariantViolation {
546 id: *object_id,
547 version: Some(version),
548 })?
549 .clone());
550 }
551
552 let o = block_on(self.multi_download(&[(*object_id, version)])).map(|mut q| {
553 q.pop().unwrap_or_else(|| {
554 panic!(
555 "Downloaded obj response cannot be empty {:?}",
556 (*object_id, version)
557 )
558 })
559 })?;
560
561 let o_ref = o.object_ref();
562 self.storage
563 .object_version_cache
564 .lock()
565 .expect("Cannot lock")
566 .insert((o_ref.object_id, o_ref.version), o.clone());
567 Ok(o)
568 }
569
570 #[expect(clippy::disallowed_methods)]
572 pub fn download_latest_object(
573 &self,
574 object_id: &ObjectId,
575 ) -> Result<Option<Object>, ReplayEngineError> {
576 let resp = block_on({
577 self.multi_download_latest(&[*object_id])
579 })
580 .map(|mut q| {
581 q.pop()
582 .unwrap_or_else(|| panic!("Downloaded obj response cannot be empty {}", *object_id))
583 });
584
585 match resp {
586 Ok(v) => Ok(Some(v)),
587 Err(ReplayEngineError::ObjectNotExist { id }) => {
588 error!(
589 "Could not find object {id} on RPC server. It might have been pruned, deleted, or never existed."
590 );
591 Ok(None)
592 }
593 Err(ReplayEngineError::ObjectDeleted {
594 id,
595 version,
596 digest,
597 }) => {
598 error!("Object {id} {version} {digest} was deleted on RPC server.");
599 Ok(None)
600 }
601 Err(err) => Err(ReplayEngineError::IotaRpcError {
602 err: err.to_string(),
603 }),
604 }
605 }
606
607 #[expect(clippy::disallowed_methods)]
608 pub fn download_object_by_upper_bound(
609 &self,
610 object_id: &ObjectId,
611 version_upper_bound: VersionNumber,
612 ) -> Result<Option<Object>, ReplayEngineError> {
613 let local_object = self
614 .storage
615 .live_objects_store
616 .lock()
617 .expect("Can't lock")
618 .get(object_id)
619 .cloned();
620 if local_object.is_some() {
621 return Ok(local_object);
622 }
623 let response = block_on({
624 self.fetcher
625 .get_child_object(object_id, version_upper_bound)
626 });
627 match response {
628 Ok(object) => {
629 let obj_ref = object.object_ref();
630 self.storage
631 .live_objects_store
632 .lock()
633 .expect("Can't lock")
634 .insert(*object_id, object.clone());
635 self.storage
636 .object_version_cache
637 .lock()
638 .expect("Can't lock")
639 .insert((obj_ref.object_id, obj_ref.version), object.clone());
640 Ok(Some(object))
641 }
642 Err(ReplayEngineError::ObjectNotExist { id }) => {
643 error!(
644 "Could not find child object {id} on RPC server. It might have been pruned, deleted, or never existed."
645 );
646 Ok(None)
647 }
648 Err(ReplayEngineError::ObjectDeleted {
649 id,
650 version,
651 digest,
652 }) => {
653 error!("Object {id} {version} {digest} was deleted on RPC server.");
654 Ok(None)
655 }
656 Err(ReplayEngineError::ObjectVersionNotFound { id, version }) => {
659 info!(
660 "Object {id} {version} not found on RPC server -- this may have been pruned or never existed."
661 );
662 Ok(None)
663 }
664 Err(err) => Err(ReplayEngineError::IotaRpcError {
665 err: err.to_string(),
666 }),
667 }
668 }
669
670 pub async fn get_checkpoint_txs(
671 &self,
672 checkpoint_id: u64,
673 ) -> Result<Vec<TransactionDigest>, ReplayEngineError> {
674 self.fetcher
675 .get_checkpoint_txs(checkpoint_id)
676 .await
677 .map_err(|e| ReplayEngineError::IotaRpcError { err: e.to_string() })
678 }
679
680 pub async fn execute_all_in_checkpoints(
681 &mut self,
682 checkpoint_ids: &[u64],
683 expensive_safety_check_config: &ExpensiveSafetyCheckConfig,
684 terminate_early: bool,
685 use_authority: bool,
686 ) -> Result<(u64, u64), ReplayEngineError> {
687 let mut txs = Vec::new();
689 for checkpoint_id in checkpoint_ids {
690 txs.extend(self.get_checkpoint_txs(*checkpoint_id).await?);
691 }
692 let num = txs.len();
693 let mut succeeded = 0;
694 for tx in txs {
695 match self
696 .execute_transaction(
697 &tx,
698 expensive_safety_check_config.clone(),
699 use_authority,
700 None,
701 None,
702 None,
703 None,
704 )
705 .await
706 .map(|q| q.check_effects())
707 {
708 Err(e) | Ok(Err(e)) => {
709 if terminate_early {
710 return Err(e);
711 }
712 error!("Error executing tx: {}, {:#?}", tx, e);
713 continue;
714 }
715 _ => (),
716 }
717
718 succeeded += 1;
719 }
720 Ok((succeeded, num as u64))
721 }
722
723 pub async fn execution_engine_execute_with_tx_info_impl(
724 &mut self,
725 tx_info: &OnChainTransactionInfo,
726 override_transaction_kind: Option<TransactionKind>,
727 expensive_safety_check_config: ExpensiveSafetyCheckConfig,
728 ) -> Result<ExecutionSandboxState, ReplayEngineError> {
729 let tx_digest = &tx_info.tx_digest;
730
731 let input_objects = self.initialize_execution_env_state(tx_info).await?;
734 let unique_shared_object_ids: HashSet<_> = input_objects
735 .filter_shared_objects()
736 .iter()
737 .map(|s| match s {
738 SharedInput::Existing(obj_ref) => obj_ref.object_id,
739 SharedInput::Deleted((id, _, _, _)) => *id,
740 SharedInput::Cancelled((id, _)) => *id,
741 })
742 .collect();
743 assert_eq!(
744 unique_shared_object_ids.len(),
745 tx_info.shared_object_refs.len()
746 );
747 let protocol_config =
752 &ProtocolConfig::get_for_version(tx_info.protocol_version, tx_info.chain);
753
754 let metrics = self.metrics.clone();
755
756 let ov = self.executor_version;
757
758 let executor = get_executor(
760 ov,
761 protocol_config,
762 expensive_safety_check_config,
763 self.enable_profiler.clone(),
764 );
765
766 let expensive_checks = true;
768 let transaction_kind = override_transaction_kind.unwrap_or(tx_info.kind.clone());
769 let certificate_deny_set = HashSet::new();
770 let gas_status = if tx_info.kind.is_system() {
771 IotaGasStatus::new_unmetered()
772 } else {
773 IotaGasStatus::new(
774 tx_info.gas_budget,
775 tx_info.gas_price,
776 tx_info.reference_gas_price,
777 protocol_config,
778 )
779 .expect("Failed to create gas status")
780 };
781 let gas_data = GasPayment {
782 objects: tx_info.gas.clone(),
783 owner: tx_info.gas_owner.unwrap_or(tx_info.sender),
784 price: tx_info.gas_price,
785 budget: tx_info.gas_budget,
786 };
787
788 let move_authenticators = tx_info.sender_signed_data.move_authenticators();
789
790 let (inner_store, gas_status, effects, _timings, result) = if move_authenticators.is_empty()
791 {
792 executor.execute_transaction_to_effects(
794 &self,
795 protocol_config,
796 metrics.clone(),
797 expensive_checks,
798 &certificate_deny_set,
799 &tx_info.executed_epoch,
800 tx_info.epoch_start_timestamp,
801 CheckedInputObjects::new_for_replay(input_objects.clone()),
802 gas_data,
803 gas_status,
804 transaction_kind.clone(),
805 tx_info.sender,
806 *tx_digest,
807 &mut None,
808 )
809 } else {
810 let (_, per_authenticator_inputs) = tx_info
813 .sender_signed_data
814 .split_input_objects_into_groups_for_reading(input_objects.clone())
815 .map_err(|e| ReplayEngineError::GeneralError { err: e.to_string() })?;
816
817 debug_assert_eq!(
818 move_authenticators.len(),
819 per_authenticator_inputs.len(),
820 "Move authenticators amount must match the number of authenticator inputs"
821 );
822
823 let move_authenticators = move_authenticators
824 .into_iter()
825 .zip(per_authenticator_inputs)
826 .map(
827 |(move_authenticator, (authenticator_inputs, account_object))| {
828 let account_version = match &account_object.object {
829 ObjectReadResultKind::Object(obj) => obj.version(),
830 _ => {
831 return Err(ReplayEngineError::GeneralError {
832 err: format!(
833 "Account object {} is not available",
834 account_object.id()
835 ),
836 });
837 }
838 };
839
840 let authenticator_function_ref = load_authenticator_function_ref(
841 move_authenticator,
842 account_version,
843 |id| {
844 self.storage
845 .live_objects_store
846 .lock()
847 .expect("Can't lock")
848 .get(id)
849 .cloned()
850 },
851 )?;
852
853 Ok((
854 move_authenticator.to_owned(),
855 authenticator_function_ref,
856 CheckedInputObjects::new_for_replay(authenticator_inputs),
857 ))
858 },
859 )
860 .collect::<Result<Vec<_>, ReplayEngineError>>()?;
861
862 let (sender_auth_digest, sponsor_auth_digest) = tx_info
863 .sender_signed_data
864 .compute_auth_digests()
865 .map_err(IotaError::from)?;
866
867 let (sender_authenticator_function_ref, sponsor_authenticator_function_ref) =
868 extract_auth_fun_refs(tx_info.sender, gas_data.owner, |address| {
869 move_authenticators
870 .iter()
871 .find(|t| t.0.address() == address)
872 .map(|t| t.1.authenticator_function_ref.clone())
873 });
874
875 let auth_context_data = AuthContextData {
876 transaction_data_bytes: bcs::to_bytes(tx_info.sender_signed_data.transaction())
877 .expect("Transaction serialization cannot fail"),
878 sender_auth_digest,
879 sponsor_auth_digest,
880 sender_authenticator_function_ref,
881 sponsor_authenticator_function_ref,
882 };
883
884 executor.authenticate_then_execute_transaction_to_effects(
885 &self,
886 protocol_config,
887 metrics.clone(),
888 expensive_checks,
889 &certificate_deny_set,
890 &tx_info.executed_epoch,
891 tx_info.epoch_start_timestamp,
892 gas_data,
893 gas_status,
894 move_authenticators,
895 CheckedInputObjects::new_for_replay(input_objects.clone()),
896 transaction_kind.clone(),
897 tx_info.sender,
898 *tx_digest,
899 auth_context_data,
900 &mut None,
901 )
902 };
903
904 if let Err(err) = self.pretty_print_for_tracing(
905 &gas_status,
906 &executor,
907 tx_info,
908 &transaction_kind,
909 protocol_config,
910 metrics,
911 expensive_checks,
912 input_objects,
913 ) {
914 error!("Failed to pretty print for tracing: {:?}", err);
915 }
916
917 let all_required_objects = self.storage.all_objects();
918
919 let effects =
920 IotaTransactionBlockEffects::try_from(effects).map_err(ReplayEngineError::from)?;
921
922 Ok(ExecutionSandboxState {
923 transaction_info: tx_info.clone(),
924 required_objects: all_required_objects,
925 local_exec_temporary_store: Some(inner_store),
926 local_exec_effects: effects,
927 local_exec_status: Some(result),
928 })
929 }
930
931 fn pretty_print_for_tracing(
932 &self,
933 gas_status: &IotaGasStatus,
934 executor: &Arc<dyn Executor + Send + Sync>,
935 tx_info: &OnChainTransactionInfo,
936 transaction_kind: &TransactionKind,
937 protocol_config: &ProtocolConfig,
938 metrics: Arc<LimitsMetrics>,
939 expensive_checks: bool,
940 input_objects: InputObjects,
941 ) -> anyhow::Result<()> {
942 trace!(target: "replay_gas_info", "{}", Pretty(gas_status));
943
944 let skip_checks = true;
945 let gas_data = GasPayment {
946 objects: tx_info.gas.clone(),
947 owner: tx_info.gas_owner.unwrap_or(tx_info.sender),
948 price: tx_info.gas_price,
949 budget: tx_info.gas_budget,
950 };
951 if let TransactionKind::Programmable(pt) = transaction_kind {
952 trace!(
953 target: "replay_ptb_info",
954 "{}",
955 Pretty(&FullPTB {
956 ptb: pt.clone(),
957 results: transform_command_results_to_annotated(
958 executor,
959 &self.clone(),
960 executor.dev_inspect_transaction(
961 &self,
962 protocol_config,
963 metrics,
964 expensive_checks,
965 &HashSet::new(),
966 &tx_info.executed_epoch,
967 tx_info.epoch_start_timestamp,
968 CheckedInputObjects::new_for_replay(input_objects),
969 gas_data,
970 IotaGasStatus::new(
971 tx_info.gas_budget,
972 tx_info.gas_price,
973 tx_info.reference_gas_price,
974 protocol_config,
975 )?,
976 transaction_kind.clone(),
977 tx_info.sender,
978 tx_info.sender_signed_data.digest(),
979 skip_checks,
980 )
981 .3
982 .unwrap_or_default(),
983 )?,
984 }));
985 }
986 Ok(())
987 }
988
989 pub async fn execution_engine_execute_impl(
991 &mut self,
992 tx_digest: &TransactionDigest,
993 expensive_safety_check_config: ExpensiveSafetyCheckConfig,
994 ) -> Result<ExecutionSandboxState, ReplayEngineError> {
995 if self.is_remote_replay() {
996 assert!(
997 !self.protocol_version_system_package_table.is_empty()
998 || !self.protocol_version_epoch_table.is_empty(),
999 "Required tables not populated. Must call `init_for_execution` before executing transactions"
1000 );
1001 }
1002
1003 let tx_info = if self.is_remote_replay() {
1004 self.resolve_tx_components(tx_digest).await?
1005 } else {
1006 self.resolve_tx_components_from_dump(tx_digest).await?
1007 };
1008 self.execution_engine_execute_with_tx_info_impl(
1009 &tx_info,
1010 None,
1011 expensive_safety_check_config,
1012 )
1013 .await
1014 }
1015
1016 pub async fn certificate_execute_with_sandbox_state(
1020 pre_run_sandbox: &ExecutionSandboxState,
1021 ) -> Result<ExecutionSandboxState, ReplayEngineError> {
1022 let executed_epoch = pre_run_sandbox.transaction_info.executed_epoch;
1024 let reference_gas_price = pre_run_sandbox.transaction_info.reference_gas_price;
1025 let epoch_start_timestamp = pre_run_sandbox.transaction_info.epoch_start_timestamp;
1026 let protocol_config = ProtocolConfig::get_for_version(
1027 pre_run_sandbox.transaction_info.protocol_version,
1028 pre_run_sandbox.transaction_info.chain,
1029 );
1030 let required_objects = pre_run_sandbox.required_objects.clone();
1031 let store = InMemoryStorage::new(required_objects.clone());
1032
1033 let transaction =
1034 TransactionEnvelope::new(pre_run_sandbox.transaction_info.sender_signed_data.clone());
1035
1036 let executable = VerifiedExecutableTransaction::new_from_quorum_execution(
1042 VerifiedTransaction::new_unchecked(transaction),
1043 executed_epoch,
1044 );
1045 let sender_signed_data = &pre_run_sandbox.transaction_info.sender_signed_data;
1046 let executor = iota_execution::executor(&protocol_config, true, None).unwrap();
1047
1048 let move_authenticators = sender_signed_data.move_authenticators();
1049
1050 let (_, _, effects, _timings, exec_res) = if move_authenticators.is_empty() {
1051 let input_objects = store.read_input_objects_for_transaction(
1053 &TransactionEnvelope::new(sender_signed_data.clone()),
1054 );
1055 let (gas_status, input_objects) = iota_transaction_checks::check_certificate_input(
1056 &executable,
1057 input_objects,
1058 &protocol_config,
1059 reference_gas_price,
1060 )
1061 .unwrap();
1062 let (kind, signer, gas_data) = executable.transaction().execution_parts();
1063 executor.execute_transaction_to_effects(
1064 &store,
1065 &protocol_config,
1066 Arc::new(LimitsMetrics::new(&Registry::new())),
1067 true,
1068 &HashSet::new(),
1069 &executed_epoch,
1070 epoch_start_timestamp,
1071 input_objects,
1072 gas_data,
1073 gas_status,
1074 kind,
1075 signer,
1076 *executable.digest(),
1077 &mut None,
1078 )
1079 } else {
1080 let all_input_object_kinds = sender_signed_data
1083 .collect_all_input_object_kind_for_reading()
1084 .unwrap();
1085 let all_input_objects: InputObjects = all_input_object_kinds
1086 .into_iter()
1087 .map(|kind| {
1088 let id = kind.object_id();
1089 let obj = store
1090 .get_object(&id)
1091 .expect("Object must be in store")
1092 .clone();
1093 ObjectReadResult::new(kind, obj.into())
1094 })
1095 .collect::<Vec<_>>()
1096 .into();
1097
1098 let (tx_input_objects, per_authenticator_inputs) = sender_signed_data
1099 .split_input_objects_into_groups_for_reading(all_input_objects)
1100 .unwrap();
1101
1102 debug_assert_eq!(
1103 move_authenticators.len(),
1104 per_authenticator_inputs.len(),
1105 "Move authenticators amount must match the number of authenticator inputs"
1106 );
1107
1108 let per_authenticator_inputs = move_authenticators
1109 .iter()
1110 .zip(per_authenticator_inputs)
1111 .map(
1112 |(move_authenticator, (authenticator_inputs, account_object))| {
1113 let account_version = match &account_object.object {
1114 ObjectReadResultKind::Object(obj) => obj.version(),
1115 _ => {
1116 return Err(ReplayEngineError::GeneralError {
1117 err: format!(
1118 "Account object {} is not available",
1119 account_object.id()
1120 ),
1121 });
1122 }
1123 };
1124
1125 let authenticator_function_ref = load_authenticator_function_ref(
1126 move_authenticator,
1127 account_version,
1128 |id| store.get_object(id).cloned(),
1129 )
1130 .unwrap();
1131
1132 Ok((authenticator_inputs, authenticator_function_ref))
1133 },
1134 )
1135 .collect::<Result<Vec<_>, ReplayEngineError>>()?;
1136
1137 let per_authenticator_input_objects = per_authenticator_inputs
1138 .iter()
1139 .map(|(authenticator_input_objects, _)| authenticator_input_objects.clone())
1140 .collect::<Vec<_>>();
1141
1142 let authenticator_gas_budget = protocol_config.max_auth_gas();
1143 let (gas_status, per_authenticator_checked_input_objects, union_checked_input_objects) =
1144 iota_transaction_checks::check_certificate_and_move_authenticator_input(
1145 &executable,
1146 tx_input_objects,
1147 per_authenticator_input_objects,
1148 authenticator_gas_budget,
1149 &protocol_config,
1150 reference_gas_price,
1151 )
1152 .unwrap();
1153
1154 debug_assert_eq!(
1155 move_authenticators.len(),
1156 per_authenticator_checked_input_objects.len(),
1157 "Move authenticators amount must match the number of checked authenticator inputs"
1158 );
1159
1160 let move_authenticators = move_authenticators
1161 .into_iter()
1162 .zip(per_authenticator_inputs)
1163 .zip(per_authenticator_checked_input_objects)
1164 .map(
1165 |(
1166 (move_authenticator, (_, authenticator_function_ref_for_execution)),
1167 authenticator_checked_input_objects,
1168 )| {
1169 (
1170 move_authenticator.to_owned(),
1171 authenticator_function_ref_for_execution,
1172 authenticator_checked_input_objects,
1173 )
1174 },
1175 )
1176 .collect::<Vec<_>>();
1177
1178 let (kind, signer, gas_data) = executable.transaction().execution_parts();
1179 let (sender_auth_digest, sponsor_auth_digest) = sender_signed_data
1180 .compute_auth_digests()
1181 .map_err(IotaError::from)?;
1182
1183 let (sender_authenticator_function_ref, sponsor_authenticator_function_ref) =
1184 extract_auth_fun_refs(signer, gas_data.owner, |address| {
1185 move_authenticators
1186 .iter()
1187 .find(|t| t.0.address() == address)
1188 .map(|t| t.1.authenticator_function_ref.clone())
1189 });
1190
1191 let auth_context_data = AuthContextData {
1192 transaction_data_bytes: bcs::to_bytes(sender_signed_data.transaction())
1193 .expect("Transaction serialization cannot fail"),
1194 sender_auth_digest,
1195 sponsor_auth_digest,
1196 sender_authenticator_function_ref,
1197 sponsor_authenticator_function_ref,
1198 };
1199
1200 executor.authenticate_then_execute_transaction_to_effects(
1201 &store,
1202 &protocol_config,
1203 Arc::new(LimitsMetrics::new(&Registry::new())),
1204 true,
1205 &HashSet::new(),
1206 &executed_epoch,
1207 epoch_start_timestamp,
1208 gas_data,
1209 gas_status,
1210 move_authenticators,
1211 union_checked_input_objects,
1212 kind,
1213 signer,
1214 *executable.digest(),
1215 auth_context_data,
1216 &mut None,
1217 )
1218 };
1219
1220 let effects =
1221 IotaTransactionBlockEffects::try_from(effects).map_err(ReplayEngineError::from)?;
1222
1223 Ok(ExecutionSandboxState {
1224 transaction_info: pre_run_sandbox.transaction_info.clone(),
1225 required_objects,
1226 local_exec_temporary_store: None, local_exec_effects: effects,
1228 local_exec_status: Some(exec_res),
1229 })
1230 }
1231
1232 pub async fn certificate_execute(
1236 &mut self,
1237 tx_digest: &TransactionDigest,
1238 expensive_safety_check_config: ExpensiveSafetyCheckConfig,
1239 ) -> Result<ExecutionSandboxState, ReplayEngineError> {
1240 let pre_run_sandbox = self
1242 .execution_engine_execute_impl(tx_digest, expensive_safety_check_config)
1243 .await?;
1244 Self::certificate_execute_with_sandbox_state(&pre_run_sandbox).await
1245 }
1246
1247 pub async fn execution_engine_execute(
1251 &mut self,
1252 tx_digest: &TransactionDigest,
1253 expensive_safety_check_config: ExpensiveSafetyCheckConfig,
1254 ) -> Result<ExecutionSandboxState, ReplayEngineError> {
1255 let sandbox_state = self
1256 .execution_engine_execute_impl(tx_digest, expensive_safety_check_config)
1257 .await?;
1258
1259 Ok(sandbox_state)
1260 }
1261
1262 pub async fn execute_state_dump(
1263 &mut self,
1264 expensive_safety_check_config: ExpensiveSafetyCheckConfig,
1265 ) -> Result<(ExecutionSandboxState, NodeStateDump), ReplayEngineError> {
1266 assert!(!self.is_remote_replay());
1267
1268 let d = match self.fetcher.clone() {
1269 Fetchers::NodeStateDump(d) => d,
1270 _ => panic!("Invalid fetcher for state dump"),
1271 };
1272 let tx_digest = d.node_state_dump.clone().tx_digest;
1273 let sandbox_state = self
1274 .execution_engine_execute_impl(&tx_digest, expensive_safety_check_config)
1275 .await?;
1276
1277 Ok((sandbox_state, d.node_state_dump))
1278 }
1279
1280 pub async fn execute_transaction(
1281 &mut self,
1282 tx_digest: &TransactionDigest,
1283 expensive_safety_check_config: ExpensiveSafetyCheckConfig,
1284 use_authority: bool,
1285 executor_version: Option<i64>,
1286 protocol_version: Option<i64>,
1287 enable_profiler: Option<PathBuf>,
1288 config_and_versions: Option<Vec<(ObjectId, Version)>>,
1289 ) -> Result<ExecutionSandboxState, ReplayEngineError> {
1290 self.executor_version = executor_version;
1291 self.protocol_version = protocol_version;
1292 self.enable_profiler = enable_profiler;
1293 self.config_and_versions = config_and_versions;
1294 if use_authority {
1295 self.certificate_execute(tx_digest, expensive_safety_check_config.clone())
1296 .await
1297 } else {
1298 self.execution_engine_execute(tx_digest, expensive_safety_check_config)
1299 .await
1300 }
1301 }
1302 fn system_package_ids(_protocol_version: u64) -> Vec<ObjectId> {
1303 BuiltInFramework::all_package_ids()
1304 }
1305
1306 pub fn get_or_download_object(
1308 &self,
1309 obj_id: &ObjectId,
1310 package_expected: bool,
1311 ) -> Result<Option<Object>, ReplayEngineError> {
1312 if package_expected {
1313 if let Some(obj) = self
1314 .storage
1315 .package_cache
1316 .lock()
1317 .expect("Cannot lock")
1318 .get(obj_id)
1319 {
1320 return Ok(Some(obj.clone()));
1321 };
1322 } else if let Some(obj) = self
1329 .storage
1330 .live_objects_store
1331 .lock()
1332 .expect("Can't lock")
1333 .get(obj_id)
1334 {
1335 return Ok(Some(obj.clone()));
1336 }
1337
1338 let Some(o) = self.download_latest_object(obj_id)? else {
1339 return Ok(None);
1340 };
1341
1342 if o.is_package() {
1343 assert!(
1344 package_expected,
1345 "Did not expect package but downloaded object is a package: {obj_id}"
1346 );
1347
1348 self.storage
1349 .package_cache
1350 .lock()
1351 .expect("Cannot lock")
1352 .insert(*obj_id, o.clone());
1353 }
1354 let o_ref = o.object_ref();
1355 self.storage
1356 .object_version_cache
1357 .lock()
1358 .expect("Cannot lock")
1359 .insert((o_ref.object_id, o_ref.version), o.clone());
1360 Ok(Some(o))
1361 }
1362
1363 pub fn is_remote_replay(&self) -> bool {
1364 matches!(self.fetcher, Fetchers::Remote(_))
1365 }
1366
1367 pub fn system_package_versions_for_protocol_version(
1369 &self,
1370 protocol_version: u64,
1371 ) -> Result<Vec<(ObjectId, Version)>, ReplayEngineError> {
1372 match &self.fetcher {
1373 Fetchers::Remote(_) => Ok(self
1374 .protocol_version_system_package_table
1375 .get(&protocol_version)
1376 .ok_or(ReplayEngineError::FrameworkObjectVersionTableNotPopulated {
1377 protocol_version,
1378 })?
1379 .clone()
1380 .into_iter()
1381 .collect()),
1382
1383 Fetchers::NodeStateDump(d) => Ok(d
1384 .node_state_dump
1385 .relevant_system_packages
1386 .iter()
1387 .map(|w| (w.id, w.version, w.digest))
1388 .map(|q| (q.0, q.1))
1389 .collect()),
1390 }
1391 }
1392
1393 pub async fn protocol_ver_to_epoch_map(
1394 &self,
1395 ) -> Result<BTreeMap<u64, ProtocolVersionSummary>, ReplayEngineError> {
1396 let mut range_map = BTreeMap::new();
1397 let epoch_change_events = self.fetcher.get_epoch_change_events(false).await?;
1398
1399 let mut tx_digest = *self
1401 .fetcher
1402 .get_checkpoint_txs(0)
1403 .await?
1404 .first()
1405 .expect("Genesis TX must be in first checkpoint");
1406 let (mut start_epoch, mut start_protocol_version, mut start_checkpoint) =
1409 (0, 1, Some(0u64));
1410
1411 let (mut curr_epoch, mut curr_protocol_version, mut curr_checkpoint) =
1412 (start_epoch, start_protocol_version, start_checkpoint);
1413
1414 (start_epoch, start_protocol_version, start_checkpoint) =
1415 (curr_epoch, curr_protocol_version, curr_checkpoint);
1416
1417 let mut end_epoch_tx_digest = tx_digest;
1420
1421 for event in epoch_change_events {
1422 (curr_epoch, curr_protocol_version) = extract_epoch_and_version(event.clone())?;
1423 end_epoch_tx_digest = event.id.tx_digest;
1424
1425 if start_protocol_version == curr_protocol_version {
1426 continue;
1428 }
1429
1430 curr_checkpoint = self
1433 .fetcher
1434 .get_transaction(&event.id.tx_digest)
1435 .await?
1436 .checkpoint;
1437 range_map.insert(
1439 start_protocol_version,
1440 ProtocolVersionSummary {
1441 protocol_version: start_protocol_version,
1442 epoch_start: start_epoch,
1443 epoch_end: curr_epoch - 1,
1444 checkpoint_start: start_checkpoint,
1445 checkpoint_end: curr_checkpoint.map(|x| x - 1),
1446 epoch_change_tx: tx_digest,
1447 },
1448 );
1449
1450 start_epoch = curr_epoch;
1451 start_protocol_version = curr_protocol_version;
1452 tx_digest = event.id.tx_digest;
1453 start_checkpoint = curr_checkpoint;
1454 }
1455
1456 range_map.insert(
1458 curr_protocol_version,
1459 ProtocolVersionSummary {
1460 protocol_version: curr_protocol_version,
1461 epoch_start: start_epoch,
1462 epoch_end: curr_epoch,
1463 checkpoint_start: curr_checkpoint,
1464 checkpoint_end: self
1465 .fetcher
1466 .get_transaction(&end_epoch_tx_digest)
1467 .await?
1468 .checkpoint,
1469 epoch_change_tx: tx_digest,
1470 },
1471 );
1472
1473 Ok(range_map)
1474 }
1475
1476 pub fn protocol_version_for_epoch(
1477 epoch: u64,
1478 mp: &BTreeMap<u64, (TransactionDigest, u64, u64)>,
1479 ) -> u64 {
1480 let mut version = 1;
1483 for (k, v) in mp.iter().rev() {
1484 if v.1 <= epoch {
1485 version = *k;
1486 break;
1487 }
1488 }
1489 version
1490 }
1491
1492 pub async fn populate_protocol_version_tables(&mut self) -> Result<(), ReplayEngineError> {
1493 self.protocol_version_epoch_table = self.protocol_ver_to_epoch_map().await?;
1494
1495 let system_package_revisions = self.system_package_versions().await?;
1496
1497 for (
1500 prot_ver,
1501 ProtocolVersionSummary {
1502 epoch_change_tx: tx_digest,
1503 ..
1504 },
1505 ) in self.protocol_version_epoch_table.clone()
1506 {
1507 let mut working = if prot_ver <= 1 {
1509 BTreeMap::new()
1510 } else {
1511 self.protocol_version_system_package_table
1512 .iter()
1513 .rev()
1514 .find(|(ver, _)| **ver <= prot_ver)
1515 .expect("Prev entry must exist")
1516 .1
1517 .clone()
1518 };
1519
1520 for (id, versions) in system_package_revisions.iter() {
1521 for ver in versions.iter().rev() {
1523 if ver.1 == tx_digest {
1524 working.insert(*id, ver.0);
1526 break;
1527 }
1528 }
1529 }
1530 self.protocol_version_system_package_table
1531 .insert(prot_ver, working);
1532 }
1533 Ok(())
1534 }
1535
1536 pub async fn system_package_versions(
1537 &self,
1538 ) -> Result<BTreeMap<ObjectId, Vec<(Version, TransactionDigest)>>, ReplayEngineError> {
1539 let system_package_ids = Self::system_package_ids(
1540 *self
1541 .protocol_version_epoch_table
1542 .keys()
1543 .peekable()
1544 .last()
1545 .expect("Protocol version epoch table not populated"),
1546 );
1547 let mut system_package_objs = self.multi_download_latest(&system_package_ids).await?;
1548
1549 let mut mapping = BTreeMap::new();
1550
1551 while !system_package_objs.is_empty() {
1553 let previous_txs: Vec<_> = system_package_objs
1556 .iter()
1557 .map(|o| (o.object_ref(), o.previous_transaction))
1558 .collect();
1559
1560 previous_txs.iter().for_each(|(object_ref, tx)| {
1561 mapping
1562 .entry(object_ref.object_id)
1563 .or_insert(vec![])
1564 .push((object_ref.version, *tx));
1565 });
1566
1567 let previous_ver_refs: Vec<_> = previous_txs
1570 .iter()
1571 .filter_map(|(q, _)| {
1572 let prev_ver = q.version - 1;
1573 if prev_ver == 0 {
1574 None
1575 } else {
1576 Some((q.object_id, prev_ver))
1577 }
1578 })
1579 .collect();
1580 system_package_objs = match self.multi_download(&previous_ver_refs).await {
1581 Ok(packages) => packages,
1582 Err(ReplayEngineError::ObjectNotExist { id }) => {
1583 warn!(
1587 "Object {} does not exist on RPC server. This might be due to pruning. Historical replays might not work",
1588 id
1589 );
1590 break;
1591 }
1592 Err(ReplayEngineError::ObjectVersionNotFound { id, version }) => {
1593 warn!(
1597 "Object {} at version {} does not exist on RPC server. This might be due to pruning. Historical replays might not work",
1598 id, version
1599 );
1600 break;
1601 }
1602 Err(ReplayEngineError::ObjectVersionTooHigh {
1603 id,
1604 asked_version,
1605 latest_version,
1606 }) => {
1607 warn!(
1608 "Object {} at version {} does not exist on RPC server. Latest version is {}. This might be due to pruning. Historical replays might not work",
1609 id, asked_version, latest_version
1610 );
1611 break;
1612 }
1613 Err(ReplayEngineError::ObjectDeleted {
1614 id,
1615 version,
1616 digest,
1617 }) => {
1618 warn!(
1622 "Object {} at version {} digest {} deleted from RPC server. This might be due to pruning. Historical replays might not work",
1623 id, version, digest
1624 );
1625 break;
1626 }
1627 Err(e) => return Err(e),
1628 };
1629 }
1630 Ok(mapping)
1631 }
1632
1633 pub async fn get_protocol_config(
1634 &self,
1635 epoch_id: EpochId,
1636 chain: Chain,
1637 ) -> Result<ProtocolConfig, ReplayEngineError> {
1638 match self.protocol_version {
1639 Some(x) if x < 0 => Ok(ProtocolConfig::get_for_max_version_UNSAFE()),
1640 Some(v) => Ok(ProtocolConfig::get_for_version((v as u64).into(), chain)),
1641 None => self
1642 .protocol_version_epoch_table
1643 .iter()
1644 .rev()
1645 .find(|(_, rg)| epoch_id >= rg.epoch_start)
1646 .map(|(p, _rg)| Ok(ProtocolConfig::get_for_version((*p).into(), chain)))
1647 .unwrap_or_else(|| {
1648 Err(ReplayEngineError::ProtocolVersionNotFound { epoch: epoch_id })
1649 }),
1650 }
1651 }
1652
1653 pub async fn checkpoints_for_epoch(
1654 &self,
1655 epoch_id: u64,
1656 ) -> Result<(u64, u64), ReplayEngineError> {
1657 let epoch_change_events = self
1658 .fetcher
1659 .get_epoch_change_events(true)
1660 .await?
1661 .into_iter()
1662 .collect::<Vec<_>>();
1663 let (start_checkpoint, start_epoch_idx) = if epoch_id == 0 {
1664 (0, 1)
1665 } else {
1666 let idx = epoch_change_events
1667 .iter()
1668 .position(|ev| match extract_epoch_and_version(ev.clone()) {
1669 Ok((epoch, _)) => epoch == epoch_id,
1670 Err(_) => false,
1671 })
1672 .ok_or(ReplayEngineError::EventNotFound { epoch: epoch_id })?;
1673 let epoch_change_tx = epoch_change_events[idx].id.tx_digest;
1674 (
1675 self.fetcher
1676 .get_transaction(&epoch_change_tx)
1677 .await?
1678 .checkpoint
1679 .unwrap_or_else(|| {
1680 panic!(
1681 "Checkpoint for transaction {epoch_change_tx} not present. Could be due to pruning"
1682 )
1683 }),
1684 idx,
1685 )
1686 };
1687
1688 let next_epoch_change_tx = epoch_change_events
1689 .get(start_epoch_idx + 1)
1690 .map(|v| v.id.tx_digest)
1691 .ok_or(ReplayEngineError::UnableToDetermineCheckpoint { epoch: epoch_id })?;
1692
1693 let next_epoch_checkpoint = self
1694 .fetcher
1695 .get_transaction(&next_epoch_change_tx)
1696 .await?
1697 .checkpoint
1698 .unwrap_or_else(|| {
1699 panic!(
1700 "Checkpoint for transaction {next_epoch_change_tx} not present. Could be due to pruning"
1701 )
1702 });
1703
1704 Ok((start_checkpoint, next_epoch_checkpoint - 1))
1705 }
1706
1707 pub async fn get_epoch_start_timestamp_and_rgp(
1708 &self,
1709 epoch_id: u64,
1710 tx_digest: &TransactionDigest,
1711 ) -> Result<(u64, u64), ReplayEngineError> {
1712 if epoch_id == 0 {
1713 return Err(ReplayEngineError::TransactionNotSupported {
1714 digest: *tx_digest,
1715 reason: "Transactions from epoch 0 not supported".to_string(),
1716 });
1717 }
1718 self.fetcher
1719 .get_epoch_start_timestamp_and_rgp(epoch_id)
1720 .await
1721 }
1722
1723 fn add_config_objects_if_needed(
1724 &self,
1725 status: &IotaExecutionStatus,
1726 ) -> Vec<(ObjectId, Version)> {
1727 match parse_effect_error_for_denied_coins(status) {
1728 Some(coin_type) => {
1729 let Some(mut config_id_and_version) = self.config_and_versions.clone() else {
1730 panic!(
1731 "Need to specify the config object ID and version for '{coin_type}' in order to replay this transaction"
1732 );
1733 };
1734 if !config_id_and_version
1736 .iter()
1737 .any(|(id, _)| id == &IOTA_DENY_LIST_OBJECT_ID)
1738 {
1739 let deny_list_oid_version = self.download_latest_object(&IOTA_DENY_LIST_OBJECT_ID)
1740 .ok()
1741 .flatten()
1742 .expect("Unable to download the deny list object for a transaction that requires it")
1743 .version();
1744 config_id_and_version.push((IOTA_DENY_LIST_OBJECT_ID, deny_list_oid_version));
1745 }
1746 config_id_and_version
1747 }
1748 None => vec![],
1749 }
1750 }
1751
1752 async fn resolve_tx_components(
1753 &self,
1754 tx_digest: &TransactionDigest,
1755 ) -> Result<OnChainTransactionInfo, ReplayEngineError> {
1756 assert!(self.is_remote_replay());
1757 let tx_info = self.fetcher.get_transaction(tx_digest).await?;
1759 let sender = match tx_info.clone().transaction.unwrap().data {
1760 iota_json_rpc_types::IotaTransactionBlockData::V1(tx) => tx.sender,
1761 };
1762 let IotaTransactionBlockEffects::V1(effects) = tx_info.clone().effects.unwrap();
1763
1764 let config_objects = self.add_config_objects_if_needed(effects.status());
1765
1766 let raw_tx_bytes = tx_info.clone().raw_transaction;
1767 let orig_tx: SenderSignedTransaction = bcs::from_bytes(&raw_tx_bytes).unwrap();
1768 let input_objs = orig_tx
1769 .collect_all_input_object_kind_for_reading()
1770 .map_err(|e| match e {
1771 IotaError::UserInput { error } => ReplayEngineError::UserInputError { err: error },
1772 other => ReplayEngineError::GeneralError {
1773 err: other.to_string(),
1774 },
1775 })?;
1776 let tx_kind_orig = orig_tx.transaction().kind();
1777
1778 let modified_at_versions: Vec<(ObjectId, Version)> = effects.modified_at_versions();
1780
1781 let shared_object_refs: Vec<ObjectReference> = effects
1782 .shared_objects()
1783 .iter()
1784 .map(|so_ref| {
1785 if so_ref.digest == ObjectDigest::OBJECT_DELETED {
1786 unimplemented!(
1787 "Replay of deleted shared object transactions is not supported yet"
1788 );
1789 } else {
1790 *so_ref
1791 }
1792 })
1793 .collect();
1794 let gas_data = match tx_info.clone().transaction.unwrap().data {
1795 iota_json_rpc_types::IotaTransactionBlockData::V1(tx) => tx.gas_data,
1796 };
1797 let gas_object_refs = gas_data.payment;
1798 let receiving_objs = orig_tx
1799 .transaction()
1800 .receiving_objects()
1801 .into_iter()
1802 .map(|obj_ref| (obj_ref.object_id, obj_ref.version))
1803 .collect();
1804
1805 let epoch_id = effects.executed_epoch;
1806 let chain = chain_from_chain_id(self.fetcher.get_chain_id().await?.as_str());
1807
1808 let (epoch_start_timestamp, reference_gas_price) = self
1810 .get_epoch_start_timestamp_and_rgp(epoch_id, tx_digest)
1811 .await?;
1812
1813 Ok(OnChainTransactionInfo {
1814 kind: tx_kind_orig.clone(),
1815 sender,
1816 modified_at_versions,
1817 input_objects: input_objs,
1818 shared_object_refs,
1819 gas: gas_object_refs,
1820 gas_owner: (gas_data.owner != sender).then_some(gas_data.owner),
1821 gas_budget: gas_data.budget,
1822 gas_price: gas_data.price,
1823 executed_epoch: epoch_id,
1824 dependencies: effects.dependencies().to_vec(),
1825 effects: IotaTransactionBlockEffects::V1(effects),
1826 receiving_objs,
1827 config_objects,
1828 protocol_version: self.get_protocol_config(epoch_id, chain).await?.version,
1832 tx_digest: *tx_digest,
1833 epoch_start_timestamp,
1834 sender_signed_data: orig_tx.clone(),
1835 reference_gas_price,
1836 chain,
1837 })
1838 }
1839
1840 async fn resolve_tx_components_from_dump(
1841 &self,
1842 tx_digest: &TransactionDigest,
1843 ) -> Result<OnChainTransactionInfo, ReplayEngineError> {
1844 assert!(!self.is_remote_replay());
1845
1846 let dp = self.fetcher.as_node_state_dump();
1847
1848 let sender = dp.node_state_dump.sender_signed_data.transaction().sender();
1849 let orig_tx = dp.node_state_dump.sender_signed_data.clone();
1850 let effects = dp.node_state_dump.computed_effects.clone();
1851 let effects = IotaTransactionBlockEffects::try_from(effects).unwrap();
1852 let config_objects = self.add_config_objects_if_needed(effects.status());
1855
1856 let input_objs = orig_tx
1860 .collect_all_input_object_kind_for_reading()
1861 .map_err(|e| match e {
1862 IotaError::UserInput { error } => ReplayEngineError::UserInputError { err: error },
1863 other => ReplayEngineError::GeneralError {
1864 err: other.to_string(),
1865 },
1866 })?;
1867 let tx_kind_orig = orig_tx.transaction().kind();
1868
1869 let modified_at_versions: Vec<(ObjectId, Version)> = effects.modified_at_versions();
1871
1872 let shared_object_refs: Vec<ObjectReference> = effects
1873 .shared_objects()
1874 .iter()
1875 .map(|so_ref| {
1876 if so_ref.digest == ObjectDigest::OBJECT_DELETED {
1877 unimplemented!(
1878 "Replay of deleted shared object transactions is not supported yet"
1879 );
1880 } else {
1881 *so_ref
1882 }
1883 })
1884 .collect();
1885 let receiving_objs = orig_tx
1886 .transaction()
1887 .receiving_objects()
1888 .into_iter()
1889 .map(|obj_ref| (obj_ref.object_id, obj_ref.version))
1890 .collect();
1891
1892 let epoch_id = dp.node_state_dump.executed_epoch;
1893
1894 let chain = chain_from_chain_id(self.fetcher.get_chain_id().await?.as_str());
1895
1896 let protocol_config =
1897 ProtocolConfig::get_for_version(dp.node_state_dump.protocol_version.into(), chain);
1898 let (epoch_start_timestamp, reference_gas_price) = self
1900 .get_epoch_start_timestamp_and_rgp(epoch_id, tx_digest)
1901 .await?;
1902 let gas_data = orig_tx.transaction().gas_data();
1903 let gas_object_refs: Vec<_> = gas_data.clone().objects;
1904
1905 Ok(OnChainTransactionInfo {
1906 kind: tx_kind_orig.clone(),
1907 sender,
1908 modified_at_versions,
1909 input_objects: input_objs,
1910 shared_object_refs,
1911 gas: gas_object_refs,
1912 gas_owner: (gas_data.owner != sender).then_some(gas_data.owner),
1913 gas_budget: gas_data.budget,
1914 gas_price: gas_data.price,
1915 executed_epoch: epoch_id,
1916 dependencies: effects.dependencies().to_vec(),
1917 effects,
1918 receiving_objs,
1919 config_objects,
1920 protocol_version: protocol_config.version,
1921 tx_digest: *tx_digest,
1922 epoch_start_timestamp,
1923 sender_signed_data: orig_tx.clone(),
1924 reference_gas_price,
1925 chain,
1926 })
1927 }
1928
1929 async fn resolve_download_input_objects(
1930 &mut self,
1931 tx_info: &OnChainTransactionInfo,
1932 deleted_shared_objects: Vec<ObjectReference>,
1933 ) -> Result<InputObjects, ReplayEngineError> {
1934 let mut package_inputs = vec![];
1936 let mut imm_owned_inputs = vec![];
1937 let mut shared_inputs = vec![];
1938 let mut deleted_shared_info_map = BTreeMap::new();
1939
1940 if !deleted_shared_objects.is_empty() {
1944 for tx_digest in tx_info.dependencies.iter() {
1945 let tx_info = self.resolve_tx_components(tx_digest).await?;
1946 for obj_ref in tx_info.shared_object_refs.iter() {
1947 deleted_shared_info_map
1948 .insert(obj_ref.object_id, (tx_info.tx_digest, obj_ref.version));
1949 }
1950 }
1951 }
1952
1953 tx_info
1954 .input_objects
1955 .iter()
1956 .map(|kind| match kind {
1957 InputObjectKind::MovePackage(i) => {
1958 package_inputs.push(*i);
1959 Ok(())
1960 }
1961 InputObjectKind::ImmOrOwnedMoveObject(o_ref) => {
1962 imm_owned_inputs.push((o_ref.object_id, o_ref.version));
1963 Ok(())
1964 }
1965 InputObjectKind::SharedMoveObject {
1966 id,
1967 initial_shared_version: _,
1968 mutable: _,
1969 } if !deleted_shared_info_map.contains_key(id) => {
1970 if let Some(o) = self
1972 .storage
1973 .live_objects_store
1974 .lock()
1975 .expect("Can't lock")
1976 .get(id)
1977 {
1978 shared_inputs.push(o.clone());
1979 Ok(())
1980 } else {
1981 Err(ReplayEngineError::InternalCacheInvariantViolation {
1982 id: *id,
1983 version: None,
1984 })
1985 }
1986 }
1987 _ => Ok(()),
1988 })
1989 .collect::<Result<Vec<_>, _>>()?;
1990
1991 let mut in_objs = self.multi_download_and_store(&imm_owned_inputs).await?;
1993
1994 in_objs.extend(
1997 self.multi_download_relevant_packages_and_store(
1998 package_inputs,
1999 tx_info.protocol_version.as_u64(),
2000 )
2001 .await?,
2002 );
2003 in_objs.extend(shared_inputs);
2005
2006 let resolved_input_objs = tx_info
2009 .input_objects
2010 .iter()
2011 .flat_map(|kind| match kind {
2012 InputObjectKind::MovePackage(i) => {
2013 Some(ObjectReadResult::new(
2015 *kind,
2016 self.storage
2017 .package_cache
2018 .lock()
2019 .expect("Cannot lock")
2020 .get(i)
2021 .unwrap_or(
2022 &self
2023 .download_latest_object(i)
2024 .expect("Object download failed")
2025 .expect("Object not found on chain"),
2026 )
2027 .clone()
2028 .into(),
2029 ))
2030 }
2031 InputObjectKind::ImmOrOwnedMoveObject(o_ref) => Some(ObjectReadResult::new(
2032 *kind,
2033 self.storage
2034 .object_version_cache
2035 .lock()
2036 .expect("Cannot lock")
2037 .get(&(o_ref.object_id, o_ref.version))
2038 .unwrap()
2039 .clone()
2040 .into(),
2041 )),
2042 InputObjectKind::SharedMoveObject { id, .. }
2043 if !deleted_shared_info_map.contains_key(id) =>
2044 {
2045 Some(ObjectReadResult::new(
2047 *kind,
2048 self.storage
2049 .live_objects_store
2050 .lock()
2051 .expect("Can't lock")
2052 .get(id)
2053 .unwrap()
2054 .clone()
2055 .into(),
2056 ))
2057 }
2058 InputObjectKind::SharedMoveObject { id, .. } => {
2059 let (digest, version) = deleted_shared_info_map.get(id).unwrap();
2060 Some(ObjectReadResult::new(
2061 *kind,
2062 ObjectReadResultKind::DeletedSharedObject(*version, *digest),
2063 ))
2064 }
2065 })
2066 .collect();
2067
2068 Ok(InputObjects::new(resolved_input_objs))
2069 }
2070
2071 async fn initialize_execution_env_state(
2074 &mut self,
2075 tx_info: &OnChainTransactionInfo,
2076 ) -> Result<InputObjects, ReplayEngineError> {
2077 self.current_protocol_version = tx_info.protocol_version.as_u64();
2079
2080 self.multi_download_and_store(&tx_info.modified_at_versions)
2082 .await?;
2083
2084 let (shared_refs, deleted_shared_refs): (Vec<ObjectReference>, Vec<ObjectReference>) =
2085 tx_info
2086 .shared_object_refs
2087 .iter()
2088 .partition(|r| r.digest != ObjectDigest::OBJECT_DELETED);
2089
2090 let shared_refs: Vec<_> = shared_refs
2092 .iter()
2093 .map(|r| (r.object_id, r.version))
2094 .collect();
2095 self.multi_download_and_store(&shared_refs).await?;
2096
2097 let gas_refs: Vec<_> = tx_info
2100 .gas
2101 .iter()
2102 .filter_map(|w| (w.object_id != ObjectId::ZERO).then_some((w.object_id, w.version)))
2103 .collect();
2104 self.multi_download_and_store(&gas_refs).await?;
2105
2106 let input_objs = self
2108 .resolve_download_input_objects(tx_info, deleted_shared_refs)
2109 .await?;
2110
2111 self.multi_download_and_store(&tx_info.receiving_objs)
2113 .await?;
2114
2115 self.multi_download_and_store(&tx_info.config_objects)
2117 .await?;
2118
2119 let loaded_child_refs = self.fetch_loaded_child_refs(&tx_info.tx_digest).await?;
2123 self.multi_download_and_store(&loaded_child_refs).await?;
2124 tokio::task::yield_now().await;
2125
2126 for move_authenticator in tx_info.sender_signed_data.move_authenticators() {
2129 let (account_object_id, _, _) = move_authenticator
2130 .object_to_authenticate_components()
2131 .map_err(|e| ReplayEngineError::GeneralError { err: e.to_string() })?;
2132
2133 let authenticator_function_ref_field_id =
2134 derive_authenticator_function_ref_v1_dynamic_field_id(account_object_id)
2135 .map_err(|e| ReplayEngineError::GeneralError { err: e.to_string() })?;
2136
2137 let account_object_version = self
2139 .storage
2140 .live_objects_store
2141 .lock()
2142 .expect("Can't lock")
2143 .get(&account_object_id)
2144 .map(|obj| obj.version())
2145 .expect("Account object should have been downloaded as part of input objects");
2146
2147 self.download_object_by_upper_bound(
2148 &authenticator_function_ref_field_id,
2149 account_object_version,
2150 )?;
2151 }
2152
2153 Ok(input_objs)
2154 }
2155}
2156
2157fn load_authenticator_function_ref(
2161 move_authenticator: &MoveAuthenticator,
2162 account_object_version: Version,
2163 get_object: impl Fn(&ObjectId) -> Option<Object>,
2164) -> Result<AuthenticatorFunctionRefForExecution, ReplayEngineError> {
2165 let (account_object_id, _, _) = move_authenticator
2166 .object_to_authenticate_components()
2167 .map_err(|e| ReplayEngineError::GeneralError { err: e.to_string() })?;
2168
2169 let field_id = derive_authenticator_function_ref_v1_dynamic_field_id(account_object_id)
2170 .map_err(|e| ReplayEngineError::GeneralError { err: e.to_string() })?;
2171
2172 let field_obj = get_object(&field_id).ok_or_else(|| ReplayEngineError::GeneralError {
2173 err: format!(
2174 "Authenticator function ref dynamic field {field_id} not found in storage \
2175 for account object {account_object_id} at version {account_object_version}"
2176 ),
2177 })?;
2178
2179 authenticator_function_ref_v1_from_dynamic_field_object(account_object_id, &field_obj)
2180 .map_err(|e| ReplayEngineError::GeneralError { err: e.to_string() })
2181}
2182
2183impl BackingPackageStore for LocalExec {
2187 fn get_package_object(&self, package_id: &ObjectId) -> IotaResult<Option<PackageObject>> {
2191 fn inner(self_: &LocalExec, package_id: &ObjectId) -> IotaResult<Option<Object>> {
2192 self_
2194 .get_or_download_object(package_id, true )
2195 .map_err(|e| IotaError::Storage(e.to_string()))
2196 }
2197
2198 let res = inner(self, package_id);
2199 self.exec_store_events
2200 .lock()
2201 .expect("Unable to lock events list")
2202 .push(ExecutionStoreEvent::BackingPackageGetPackageObject {
2203 package_id: *package_id,
2204 result: res.clone(),
2205 });
2206 res.map(|o| o.map(PackageObject::new))
2207 }
2208}
2209
2210impl ChildObjectResolver for LocalExec {
2211 fn read_child_object(
2214 &self,
2215 parent: &ObjectId,
2216 child: &ObjectId,
2217 child_version_upper_bound: Version,
2218 ) -> IotaResult<Option<Object>> {
2219 fn inner(
2220 self_: &LocalExec,
2221 parent: &ObjectId,
2222 child: &ObjectId,
2223 child_version_upper_bound: Version,
2224 ) -> IotaResult<Option<Object>> {
2225 let child_object =
2226 match self_.download_object_by_upper_bound(child, child_version_upper_bound)? {
2227 None => return Ok(None),
2228 Some(o) => o,
2229 };
2230 let child_version = child_object.version();
2231 if child_object.version() > child_version_upper_bound {
2232 return Err(IotaError::Unknown(format!(
2233 "Invariant Violation. Replay loaded child_object {child} at version \
2234 {child_version} but expected the version to be <= {child_version_upper_bound}"
2235 )));
2236 }
2237 let parent = *parent;
2238 if child_object.owner != Owner::Object(parent) {
2239 return Err(IotaError::InvalidChildObjectAccess {
2240 object: *child,
2241 given_parent: parent,
2242 actual_owner: child_object.owner,
2243 });
2244 }
2245 Ok(Some(child_object))
2246 }
2247
2248 let res = inner(self, parent, child, child_version_upper_bound);
2249 self.exec_store_events
2250 .lock()
2251 .expect("Unable to lock events list")
2252 .push(
2253 ExecutionStoreEvent::ChildObjectResolverStoreReadChildObject {
2254 parent: *parent,
2255 child: *child,
2256 result: res.clone(),
2257 },
2258 );
2259 res
2260 }
2261
2262 fn get_object_received_at_version(
2263 &self,
2264 owner: &ObjectId,
2265 receiving_object_id: &ObjectId,
2266 receive_object_at_version: Version,
2267 _epoch_id: EpochId,
2268 ) -> IotaResult<Option<Object>> {
2269 fn inner(
2270 self_: &LocalExec,
2271 owner: &ObjectId,
2272 receiving_object_id: &ObjectId,
2273 receive_object_at_version: Version,
2274 ) -> IotaResult<Option<Object>> {
2275 let recv_object = match self_.try_get_object(receiving_object_id)? {
2276 None => return Ok(None),
2277 Some(o) => o,
2278 };
2279 if recv_object.version() != receive_object_at_version {
2280 return Err(IotaError::Unknown(format!(
2281 "Invariant Violation. Replay loaded child_object {receiving_object_id} at version \
2282 {receive_object_at_version} but expected the version to be == {receive_object_at_version}"
2283 )));
2284 }
2285 if recv_object.owner != Owner::Address((*owner).into()) {
2286 return Ok(None);
2287 }
2288 Ok(Some(recv_object))
2289 }
2290
2291 let res = inner(self, owner, receiving_object_id, receive_object_at_version);
2292 self.exec_store_events
2293 .lock()
2294 .expect("Unable to lock events list")
2295 .push(ExecutionStoreEvent::ReceiveObject {
2296 owner: *owner,
2297 receive: *receiving_object_id,
2298 receive_at_version: receive_object_at_version,
2299 result: res.clone(),
2300 });
2301 res
2302 }
2303}
2304
2305impl ModuleResolver for LocalExec {
2306 type Error = IotaError;
2307
2308 fn get_module(&self, module_id: &ModuleId) -> IotaResult<Option<Vec<u8>>> {
2311 fn inner(self_: &LocalExec, module_id: &ModuleId) -> IotaResult<Option<Vec<u8>>> {
2312 get_module(self_, module_id)
2313 }
2314
2315 let res = inner(self, module_id);
2316 self.exec_store_events
2317 .lock()
2318 .expect("Unable to lock events list")
2319 .push(ExecutionStoreEvent::ModuleResolverGetModule {
2320 module_id: module_id.clone(),
2321 result: res.clone(),
2322 });
2323 res
2324 }
2325}
2326
2327impl ModuleResolver for &mut LocalExec {
2328 type Error = IotaError;
2329
2330 fn get_module(&self, module_id: &ModuleId) -> IotaResult<Option<Vec<u8>>> {
2331 (**self).get_module(module_id)
2334 }
2335}
2336
2337impl ObjectStore for LocalExec {
2338 fn try_get_object(
2341 &self,
2342 object_id: &ObjectId,
2343 ) -> iota_types::storage::error::Result<Option<Object>> {
2344 let res = self
2345 .storage
2346 .live_objects_store
2347 .lock()
2348 .expect("Can't lock")
2349 .get(object_id)
2350 .cloned();
2351 self.exec_store_events
2352 .lock()
2353 .expect("Unable to lock events list")
2354 .push(ExecutionStoreEvent::ObjectStoreGetObject {
2355 object_id: *object_id,
2356 result: Ok(res.clone()),
2357 });
2358 Ok(res)
2359 }
2360
2361 fn try_get_object_by_key(
2364 &self,
2365 object_id: &ObjectId,
2366 version: VersionNumber,
2367 ) -> iota_types::storage::error::Result<Option<Object>> {
2368 let res = self
2369 .storage
2370 .live_objects_store
2371 .lock()
2372 .expect("Can't lock")
2373 .get(object_id)
2374 .and_then(|obj| {
2375 if obj.version() == version {
2376 Some(obj.clone())
2377 } else {
2378 None
2379 }
2380 });
2381
2382 self.exec_store_events
2383 .lock()
2384 .expect("Unable to lock events list")
2385 .push(ExecutionStoreEvent::ObjectStoreGetObjectByKey {
2386 object_id: *object_id,
2387 version,
2388 result: Ok(res.clone()),
2389 });
2390
2391 Ok(res)
2392 }
2393}
2394
2395impl ObjectStore for &mut LocalExec {
2396 fn try_get_object(
2397 &self,
2398 object_id: &ObjectId,
2399 ) -> iota_types::storage::error::Result<Option<Object>> {
2400 (**self).try_get_object(object_id)
2403 }
2404
2405 fn try_get_object_by_key(
2406 &self,
2407 object_id: &ObjectId,
2408 version: VersionNumber,
2409 ) -> iota_types::storage::error::Result<Option<Object>> {
2410 (**self).try_get_object_by_key(object_id, version)
2413 }
2414}
2415
2416impl GetModule for LocalExec {
2417 type Error = IotaError;
2418 type Item = CompiledModule;
2419
2420 fn get_module_by_id(&self, id: &ModuleId) -> IotaResult<Option<Self::Item>> {
2421 let res = get_module_by_id(self, id);
2422
2423 self.exec_store_events
2424 .lock()
2425 .expect("Unable to lock events list")
2426 .push(ExecutionStoreEvent::GetModuleGetModuleByModuleId {
2427 id: id.clone(),
2428 result: res.clone(),
2429 });
2430 res
2431 }
2432}
2433
2434pub fn get_executor(
2437 executor_version_override: Option<i64>,
2438 protocol_config: &ProtocolConfig,
2439 _expensive_safety_check_config: ExpensiveSafetyCheckConfig,
2440 enable_profiler: Option<PathBuf>,
2441) -> Arc<dyn Executor + Send + Sync> {
2442 let protocol_config = executor_version_override
2443 .map(|q| {
2444 let ver = if q < 0 {
2445 ProtocolConfig::get_for_max_version_UNSAFE().execution_version()
2446 } else {
2447 q as u64
2448 };
2449
2450 let mut c = protocol_config.clone();
2451 c.set_execution_version_for_testing(ver);
2452 c
2453 })
2454 .unwrap_or(protocol_config.clone());
2455
2456 let silent = true;
2457 iota_execution::executor(&protocol_config, silent, enable_profiler)
2458 .expect("Creating an executor should not fail here")
2459}
2460
2461fn parse_effect_error_for_denied_coins(status: &IotaExecutionStatus) -> Option<String> {
2462 let IotaExecutionStatus::Failure { error } = status else {
2463 return None;
2464 };
2465 parse_denied_error_string(error)
2466}
2467
2468fn parse_denied_error_string(error: &str) -> Option<String> {
2469 let regulated_regex = regex::Regex::new(
2470 r#"CoinTypeGlobalPause.*?"(.*?)"|AddressDeniedForCoin.*coin_type:.*?"(.*?)""#,
2471 )
2472 .unwrap();
2473
2474 let caps = regulated_regex.captures(error)?;
2475 Some(caps.get(1).or(caps.get(2))?.as_str().to_string())
2476}
2477
2478#[cfg(test)]
2479mod tests {
2480 use super::parse_denied_error_string;
2481 #[test]
2482 fn test_regex_regulated_coin_errors() {
2483 let test_bank = vec![
2484 "CoinTypeGlobalPause { coin_type: \"39a572c071784c280ee8ee8c683477e059d1381abc4366f9a58ffac3f350a254::rcoin::RCOIN\" }",
2485 "AddressDeniedForCoin { address: B, coin_type: \"39a572c071784c280ee8ee8c683477e059d1381abc4366f9a58ffac3f350a254::rcoin::RCOIN\" }",
2486 ];
2487 let expected_string =
2488 "39a572c071784c280ee8ee8c683477e059d1381abc4366f9a58ffac3f350a254::rcoin::RCOIN";
2489
2490 for test in &test_bank {
2491 assert!(parse_denied_error_string(test).unwrap() == expected_string);
2492 }
2493 }
2494}