Skip to main content

iota_replay/
replay.rs

1// Copyright (c) Mysten Labs, Inc.
2// Modifications Copyright (c) 2024 IOTA Stiftung
3// SPDX-License-Identifier: Apache-2.0
4
5use 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// TODO: add persistent cache. But perf is good enough already.
75
76#[derive(Debug, Serialize, Deserialize)]
77pub struct ExecutionSandboxState {
78    /// Information describing the transaction
79    pub transaction_info: OnChainTransactionInfo,
80    /// All the objects that are required for the execution of the transaction
81    pub required_objects: Vec<Object>,
82    /// Temporary store from executing this locally in
83    /// `execute_transaction_to_effects`
84    #[serde(skip)]
85    pub local_exec_temporary_store: Option<InnerTemporaryStore>,
86    /// Effects from executing this locally in `execute_transaction_to_effects`
87    pub local_exec_effects: IotaTransactionBlockEffects,
88    /// Status from executing this locally in `execute_transaction_to_effects`
89    #[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    /// Utility to diff effects in a human readable format
110    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    /// Protocol version at this point
134    pub protocol_version: u64,
135    /// The first epoch that uses this protocol version
136    pub epoch_start: u64,
137    /// The last epoch that uses this protocol version
138    pub epoch_end: u64,
139    /// The first checkpoint in this protocol v ersion
140    pub checkpoint_start: Option<u64>,
141    /// The last checkpoint in this protocol version
142    pub checkpoint_end: Option<u64>,
143    /// The transaction which triggered this epoch change
144    pub epoch_change_tx: TransactionDigest,
145}
146
147#[derive(Clone)]
148pub struct Storage {
149    /// These are objects at the frontier of the execution's view
150    /// They might not be the latest object currently but they are the latest
151    /// objects for the TX at the time it was run
152    /// This store cannot be shared between runners
153    pub live_objects_store: Arc<Mutex<BTreeMap<ObjectId, Object>>>,
154
155    /// Package cache and object version cache can be shared between runners
156    /// Non system packages are immutable so we can cache these
157    pub package_cache: Arc<Mutex<BTreeMap<ObjectId, Object>>>,
158    /// Object contents are frozen at their versions so we can cache these
159    /// We must place system packages here as well
160    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    // For a given protocol version, what TX created it, and what is the valid range of epochs
229    // at this protocol version.
230    pub protocol_version_epoch_table: BTreeMap<u64, ProtocolVersionSummary>,
231    // For a given protocol version, the mapping valid sequence numbers for each framework package
232    pub protocol_version_system_package_table: BTreeMap<u64, BTreeMap<ObjectId, Version>>,
233    // The current protocol version for this execution
234    pub current_protocol_version: u64,
235    // All state is contained here
236    pub storage: Storage,
237    // Debug events
238    pub exec_store_events: Arc<Mutex<Vec<ExecutionStoreEvent>>>,
239    // Debug events
240    pub metrics: Arc<LimitsMetrics>,
241    // Used for fetching data from the network or remote store
242    pub fetcher: Fetchers,
243
244    // One can optionally override the executor version
245    // -1 implies use latest version
246    pub executor_version: Option<i64>,
247    // One can optionally override the protocol version
248    // -1 implies use latest version
249    // None implies use the protocol version at the time of execution
250    pub protocol_version: Option<i64>,
251    // Whether or not to enable the gas profiler, the PathBuf contains either a user specified
252    // filepath or the default current directory and name format for the profile output
253    pub enable_profiler: Option<PathBuf>,
254    pub config_and_versions: Option<Vec<(ObjectId, Version)>>,
255    // Retry policies due to RPC errors
256    pub num_retries_for_timeout: u32,
257    pub sleep_period_for_timeout: std::time::Duration,
258}
259
260impl LocalExec {
261    /// Wrapper around fetcher in case we want to add more functionality
262    /// Such as fetching from local DB from snapshot
263    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    /// Wrapper around fetcher in case we want to add more functionality
286    /// Such as fetching from local DB from snapshot
287    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        // Get the child objects loaded
315        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    /// This captures the state of the network at a given point in time and
358    /// populates prptocol version tables including which system packages to
359    /// fetch If this function is called across epoch boundaries, the info
360    /// might be stale. But it should only be called once per epoch.
361    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        // Use a throwaway metrics registry for local execution.
382        let registry = prometheus_filtered::Registry::new();
383        let metrics = Arc::new(LimitsMetrics::new(&registry));
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            // TODO: make these configurable
397            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        // Use a throwaway metrics registry for local execution.
411        let registry = prometheus_filtered::Registry::new();
412        let metrics = Arc::new(LimitsMetrics::new(&registry));
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            // TODO: make these configurable
440            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        // Backfill the store
456        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        // Download latest version of all packages that are not system packages
494        // This is okay since the versions can never change
495        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            // We dont always want the latest in store
508            // self.storage.store.insert(o_ref.object_id, obj.clone());
509            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    // TODO: remove this after `futures::executor::block_on` is removed.
526    #[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    // TODO: remove this after `futures::executor::block_on` is removed.
571    #[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            // info!("Downloading latest object {object_id}");
578            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            // This is a child object which was not found in the store (e.g., due to exists
657            // check before creating the dynamic field).
658            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        // Get all the TXs at this checkpoint
688        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        // Initialize the state necessary for execution
732        // Get the input objects
733        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        // At this point we have all the objects needed for replay
748
749        // This assumes we already initialized the protocol version table
750        // `protocol_version_epoch_table`
751        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        // We could probably cache the executor per protocol config
759        let executor = get_executor(
760            ov,
761            protocol_config,
762            expensive_safety_check_config,
763            self.enable_profiler.clone(),
764        );
765
766        // All prep done
767        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            // Standard path: no MoveAuthenticator
793            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            // MoveAuthenticator path: split input objects and run authentication
811            // before PTB execution, matching the production flow.
812            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    /// Must be called after `init_for_execution`
990    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    /// Executes a transaction with the state specified in `pre_run_sandbox`
1017    /// This is useful for executing a transaction with a specific state
1018    /// However if the state in invalid, the behavior is undefined.
1019    pub async fn certificate_execute_with_sandbox_state(
1020        pre_run_sandbox: &ExecutionSandboxState,
1021    ) -> Result<ExecutionSandboxState, ReplayEngineError> {
1022        // These cannot be changed and are inherited from the sandbox state
1023        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        // TODO: This will not work for deleted shared objects. We need to persist that
1037        // information in the sandbox. TODO: A lot of the following code is
1038        // replicated in several places. We should introduce a few traits and
1039        // make them shared so that we don't have to fix one by one when we have major
1040        // execution layer changes.
1041        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            // Standard path: no MoveAuthenticator
1052            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            // MoveAuthenticator path: read all objects (tx + auth), split, and
1081            // run authentication before PTB execution.
1082            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, // We dont capture it for cert exec run
1227            local_exec_effects: effects,
1228            local_exec_status: Some(exec_res),
1229        })
1230    }
1231
1232    /// Must be called after `init_for_execution`
1233    /// This executes from
1234    /// `iota_core::authority::AuthorityState::try_execute_immediately`
1235    pub async fn certificate_execute(
1236        &mut self,
1237        tx_digest: &TransactionDigest,
1238        expensive_safety_check_config: ExpensiveSafetyCheckConfig,
1239    ) -> Result<ExecutionSandboxState, ReplayEngineError> {
1240        // Use the lighterweight execution engine to get the pre-run state
1241        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    /// Must be called after `init_for_execution`
1248    /// This executes from
1249    /// `iota_adapter::execution_engine::execute_transaction_to_effects`
1250    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    /// This is the only function which accesses the network during execution
1307    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            // Check if its a system package because we must've downloaded all
1323            // TODO: Will return this check once we can download completely for
1324            // other networks assert!(
1325            //     !self.system_package_ids().contains(obj_id),
1326            //     "All system packages should be downloaded already"
1327            // );
1328        } 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    /// Must be called after `populate_protocol_version_tables`
1368    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        // Exception for Genesis: Protocol version 1 at epoch 0
1400        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        // Somehow the genesis TX did not emit any event, but we know it was the start
1407        // of version 1 So we need to manually add this range
1408        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        // This is the final tx digest for the epoch change. We need this to track the
1418        // final checkpoint
1419        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                // Same range
1427                continue;
1428            }
1429
1430            // Change in prot version
1431            // Find the last checkpoint
1432            curr_checkpoint = self
1433                .fetcher
1434                .get_transaction(&event.id.tx_digest)
1435                .await?
1436                .checkpoint;
1437            // Insert the last range
1438            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        // Insert the last range
1457        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        // Naive impl but works for now
1481        // Can improve with range algos & data structures
1482        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        // This can be more efficient but small footprint so okay for now
1498        // Table is sorted from earliest to latest
1499        for (
1500            prot_ver,
1501            ProtocolVersionSummary {
1502                epoch_change_tx: tx_digest,
1503                ..
1504            },
1505        ) in self.protocol_version_epoch_table.clone()
1506        {
1507            // Use the previous versions protocol version table
1508            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                // Oldest appears first in list, so reverse
1522                for ver in versions.iter().rev() {
1523                    if ver.1 == tx_digest {
1524                        // Found the version for this protocol version
1525                        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        // Extract all the transactions which created or mutated this object
1552        while !system_package_objs.is_empty() {
1553            // For the given object and its version, record the transaction which upgraded
1554            // or created it
1555            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            // Next round
1568            // Get the previous version of each object if exists
1569            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                    // This happens when the RPC server prunes older object
1584                    // Replays in the current protocol version will work but old ones might not
1585                    // as we cannot fetch the package
1586                    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                    // This happens when the RPC server prunes older object
1594                    // Replays in the current protocol version will work but old ones might not
1595                    // as we cannot fetch the package
1596                    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                    // This happens when the RPC server prunes older object
1619                    // Replays in the current protocol version will work but old ones might not
1620                    // as we cannot fetch the package
1621                    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                // NB: the version of the deny list object doesn't matter
1735                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        // Fetch full transaction content
1758        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        // Download the objects at the version right before the execution of this TX
1779        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        // Extract the epoch start timestamp
1809        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            // Find the protocol version for this epoch
1829            // This assumes we already initialized the protocol version table
1830            // `protocol_version_epoch_table`
1831            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        // Config objects don't show up in the node state dump so they need to be
1853        // provided.
1854        let config_objects = self.add_config_objects_if_needed(effects.status());
1855
1856        // Fetch full transaction content
1857        // let tx_info = self.fetcher.get_transaction(tx_digest).await?;
1858
1859        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        // Download the objects at the version right before the execution of this TX
1870        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        // Extract the epoch start timestamp
1899        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        // Download the input objects
1935        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        // for deleted shared objects, we need to look at the transaction dependencies
1941        // to find the correct transaction dependency for a deleted shared
1942        // object.
1943        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                    // We already downloaded
1971                    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        // Download the imm and owned objects
1992        let mut in_objs = self.multi_download_and_store(&imm_owned_inputs).await?;
1993
1994        // For packages, download latest if non framework
1995        // If framework, download relevant for the current protocol version
1996        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        // Add shared objects
2004        in_objs.extend(shared_inputs);
2005
2006        // TODO(Zhe): Account for cancelled transaction assigned version here, and
2007        // tests.
2008        let resolved_input_objs = tx_info
2009            .input_objects
2010            .iter()
2011            .flat_map(|kind| match kind {
2012                InputObjectKind::MovePackage(i) => {
2013                    // Okay to unwrap since we downloaded it
2014                    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                    // we already downloaded
2046                    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    /// Given the OnChainTransactionInfo, download and store the input objects,
2072    /// and other info necessary for execution
2073    async fn initialize_execution_env_state(
2074        &mut self,
2075        tx_info: &OnChainTransactionInfo,
2076    ) -> Result<InputObjects, ReplayEngineError> {
2077        // We need this for other activities in this session
2078        self.current_protocol_version = tx_info.protocol_version.as_u64();
2079
2080        // Download the objects at the version right before the execution of this TX
2081        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        // Download shared objects at the version right before the execution of this TX
2091        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        // Download gas (although this should already be in cache from modified at
2098        // versions?)
2099        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        // Fetch the input objects we know from the raw transaction
2107        let input_objs = self
2108            .resolve_download_input_objects(tx_info, deleted_shared_refs)
2109            .await?;
2110
2111        // Fetch the receiving objects
2112        self.multi_download_and_store(&tx_info.receiving_objs)
2113            .await?;
2114
2115        // Fetch specified config objects if any
2116        self.multi_download_and_store(&tx_info.config_objects)
2117            .await?;
2118
2119        // Prep the object runtime for dynamic fields
2120        // Download the child objects accessed at the version right before the execution
2121        // of this TX
2122        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        // If the transaction uses MoveAuthenticators, download the authenticator
2127        // function ref dynamic field objects so they are available during execution.
2128        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            // Get account object version from the already-downloaded objects
2138            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
2157/// Loads the `AuthenticatorFunctionRefForExecution` from a dynamic field on the
2158/// account object. This is a simplified version of `check_move_account()` in
2159/// `authority.rs` — we skip validation since we trust on-chain state.
2160fn 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
2183// <---------------------  Implement necessary traits for LocalExec to work with
2184// exec engine ----------------------->
2185
2186impl BackingPackageStore for LocalExec {
2187    /// In this case we might need to download a dependency package which was
2188    /// not present in the modified at versions list because packages are
2189    /// immutable
2190    fn get_package_object(&self, package_id: &ObjectId) -> IotaResult<Option<PackageObject>> {
2191        fn inner(self_: &LocalExec, package_id: &ObjectId) -> IotaResult<Option<Object>> {
2192            // If package not present fetch it from the network
2193            self_
2194                .get_or_download_object(package_id, true /* we expect a Move package */)
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    /// This uses `get_object`, which does not download from the network
2212    /// Hence all objects must be in store already
2213    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    /// This fetches a module which must already be present in the store
2309    /// We do not download
2310    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        // Recording event here will be double-counting since its already recorded in
2332        // the get_module fn
2333        (**self).get_module(module_id)
2334    }
2335}
2336
2337impl ObjectStore for LocalExec {
2338    /// The object must be present in store by normal process we used to
2339    /// backfill store in init We dont download if not present
2340    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    /// The object must be present in store by normal process we used to
2362    /// backfill store in init We dont download if not present
2363    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        // Recording event here will be double-counting since its already recorded in
2401        // the get_module fn
2402        (**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        // Recording event here will be double-counting since its already recorded in
2411        // the get_module fn
2412        (**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
2434// <--------------------- Util functions ----------------------->
2435
2436pub 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}