Skip to main content

sui_replay/
replay.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4use crate::chain_from_chain_id;
5use crate::{
6    data_fetcher::{
7        DataFetcher, Fetchers, NodeStateDumpFetcher, RemoteFetcher, extract_epoch_and_version,
8    },
9    displays::{
10        Pretty,
11        transaction_displays::{FullPTB, transform_command_results_to_annotated},
12    },
13    types::*,
14};
15use futures::executor::block_on;
16use move_binary_format::CompiledModule;
17use move_bytecode_utils::module_cache::GetModule;
18use move_core_types::resolver::SerializedPackage;
19use move_core_types::{
20    account_address::AccountAddress, language_storage::ModuleId, resolver::ModuleResolver,
21};
22use prometheus::Registry;
23use serde::{Deserialize, Serialize};
24use similar::{ChangeTag, TextDiff};
25use std::{
26    collections::{BTreeMap, HashSet},
27    path::PathBuf,
28    sync::Arc,
29    sync::Mutex,
30};
31use sui_config::node::ExpensiveSafetyCheckConfig;
32use sui_core::authority::NodeStateDump;
33use sui_execution::Executor;
34use sui_framework::BuiltInFramework;
35use sui_json_rpc_types::{
36    SuiExecutionStatus, SuiTransactionBlockEffects, SuiTransactionBlockEffectsAPI,
37};
38use sui_protocol_config::{Chain, ProtocolConfig};
39use sui_sdk::{SuiClient, SuiClientBuilder};
40use sui_types::SUI_DENY_LIST_OBJECT_ID;
41use sui_types::error::SuiErrorKind;
42use sui_types::execution_params::{
43    ExecutionOrEarlyError, FundsWithdrawStatus, get_early_execution_error,
44};
45use sui_types::in_memory_storage::InMemoryStorage;
46use sui_types::message_envelope::Message;
47use sui_types::storage::{PackageObject, get_module, get_package};
48use sui_types::transaction::GasData;
49use sui_types::transaction::TransactionKind::ProgrammableTransaction;
50use sui_types::{
51    DEEPBOOK_PACKAGE_ID,
52    base_types::{ObjectID, ObjectRef, SequenceNumber, VersionNumber},
53    committee::EpochId,
54    digests::{ObjectDigest, TransactionDigest},
55    error::{ExecutionError, SuiError, SuiResult},
56    executable_transaction::VerifiedExecutableTransaction,
57    gas::SuiGasStatus,
58    inner_temporary_store::InnerTemporaryStore,
59    metrics::ExecutionMetrics,
60    object::{Object, Owner},
61    storage::get_module_by_id,
62    storage::{BackingPackageStore, ObjectStore, ParentSync, RuntimeObjectResolver},
63    transaction::{
64        CheckedInputObjects, InputObjectKind, InputObjects, ObjectReadResult, ObjectReadResultKind,
65        SenderSignedData, Transaction, TransactionDataAPI, TransactionKind, VerifiedTransaction,
66    },
67};
68use tracing::{error, info, trace, warn};
69
70// TODO: add persistent cache. But perf is good enough already.
71
72#[derive(Debug, Serialize, Deserialize)]
73pub struct ExecutionSandboxState {
74    /// Information describing the transaction
75    pub transaction_info: OnChainTransactionInfo,
76    /// All the objects that are required for the execution of the transaction
77    pub required_objects: Vec<Object>,
78    /// Temporary store from executing this locally in `execute_transaction_to_effects`
79    #[serde(skip)]
80    pub local_exec_temporary_store: Option<InnerTemporaryStore>,
81    /// Effects from executing this locally in `execute_transaction_to_effects`
82    pub local_exec_effects: SuiTransactionBlockEffects,
83    /// Status from executing this locally in `execute_transaction_to_effects`
84    #[serde(skip)]
85    pub local_exec_status: Option<Result<(), ExecutionError>>,
86}
87
88impl ExecutionSandboxState {
89    #[allow(clippy::result_large_err)]
90    pub fn check_effects(&self) -> Result<(), ReplayEngineError> {
91        let SuiTransactionBlockEffects::V1(mut local_effects) = self.local_exec_effects.clone();
92        let SuiTransactionBlockEffects::V1(on_chain_effects) =
93            self.transaction_info.effects.clone();
94
95        // Handle backwards compatibility with the new `abort_error` field in
96        // `SuiTransactionBlockEffects`
97        if on_chain_effects.abort_error.is_none() {
98            local_effects.abort_error = None;
99        }
100        let local_effects = SuiTransactionBlockEffects::V1(local_effects);
101        let on_chain_effects = SuiTransactionBlockEffects::V1(on_chain_effects);
102
103        if on_chain_effects != local_effects {
104            error!("Replay tool forked {}", self.transaction_info.tx_digest);
105            let diff = Self::diff_effects(&on_chain_effects, &local_effects);
106            println!("{}", diff);
107            return Err(ReplayEngineError::EffectsForked {
108                digest: self.transaction_info.tx_digest,
109                diff: format!("\n{}", diff),
110                on_chain: Box::new(on_chain_effects),
111                local: Box::new(local_effects),
112            });
113        }
114        Ok(())
115    }
116
117    /// Utility to diff effects in a human readable format
118    pub fn diff_effects(
119        on_chain_effects: &SuiTransactionBlockEffects,
120        local_effects: &SuiTransactionBlockEffects,
121    ) -> String {
122        let on_chain_str = format!("{:#?}", on_chain_effects);
123        let local_chain_str = format!("{:#?}", local_effects);
124        let mut res = vec![];
125
126        let diff = TextDiff::from_lines(&on_chain_str, &local_chain_str);
127        for change in diff.iter_all_changes() {
128            let sign = match change.tag() {
129                ChangeTag::Delete => "---",
130                ChangeTag::Insert => "+++",
131                ChangeTag::Equal => "   ",
132            };
133            res.push(format!("{}{}", sign, change));
134        }
135
136        res.join("")
137    }
138}
139
140#[derive(Debug, Clone, PartialEq, Eq)]
141pub struct ProtocolVersionSummary {
142    /// Protocol version at this point
143    pub protocol_version: u64,
144    /// The first epoch that uses this protocol version
145    pub epoch_start: u64,
146    /// The last epoch that uses this protocol version
147    pub epoch_end: u64,
148    /// The first checkpoint in this protocol v ersion
149    pub checkpoint_start: Option<u64>,
150    /// The last checkpoint in this protocol version
151    pub checkpoint_end: Option<u64>,
152    /// The transaction which triggered this epoch change
153    pub epoch_change_tx: TransactionDigest,
154}
155
156#[derive(Clone)]
157pub struct Storage {
158    /// These are objects at the frontier of the execution's view
159    /// They might not be the latest object currently but they are the latest objects
160    /// for the TX at the time it was run
161    /// This store cannot be shared between runners
162    pub live_objects_store: Arc<Mutex<BTreeMap<ObjectID, Object>>>,
163
164    /// Package cache and object version cache can be shared between runners
165    /// Non system packages are immutable so we can cache these
166    pub package_cache: Arc<Mutex<BTreeMap<ObjectID, Object>>>,
167    /// Object contents are frozen at their versions so we can cache these
168    /// We must place system packages here as well
169    pub object_version_cache: Arc<Mutex<BTreeMap<(ObjectID, SequenceNumber), Object>>>,
170}
171
172impl std::fmt::Display for Storage {
173    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
174        writeln!(f, "Live object store")?;
175        for (id, obj) in self
176            .live_objects_store
177            .lock()
178            .expect("Unable to lock")
179            .iter()
180        {
181            writeln!(f, "{}: {:?}", id, obj.compute_object_reference())?;
182        }
183        writeln!(f, "Package cache")?;
184        for (id, obj) in self.package_cache.lock().expect("Unable to lock").iter() {
185            writeln!(f, "{}: {:?}", id, obj.compute_object_reference())?;
186        }
187        writeln!(f, "Object version cache")?;
188        for (id, _) in self
189            .object_version_cache
190            .lock()
191            .expect("Unable to lock")
192            .iter()
193        {
194            writeln!(f, "{}: {}", id.0, id.1)?;
195        }
196
197        write!(f, "")
198    }
199}
200
201impl Storage {
202    pub fn default() -> Self {
203        Self {
204            live_objects_store: Arc::new(Mutex::new(BTreeMap::new())),
205            package_cache: Arc::new(Mutex::new(BTreeMap::new())),
206            object_version_cache: Arc::new(Mutex::new(BTreeMap::new())),
207        }
208    }
209
210    pub fn all_objects(&self) -> Vec<Object> {
211        self.live_objects_store
212            .lock()
213            .expect("Unable to lock")
214            .values()
215            .cloned()
216            .chain(
217                self.package_cache
218                    .lock()
219                    .expect("Unable to lock")
220                    .values()
221                    .cloned(),
222            )
223            .chain(
224                self.object_version_cache
225                    .lock()
226                    .expect("Unable to lock")
227                    .values()
228                    .cloned(),
229            )
230            .collect::<Vec<_>>()
231    }
232}
233
234#[derive(Clone)]
235pub struct LocalExec {
236    pub client: Option<SuiClient>,
237    // For a given protocol version, what TX created it, and what is the valid range of epochs
238    // at this protocol version.
239    pub protocol_version_epoch_table: BTreeMap<u64, ProtocolVersionSummary>,
240    // For a given protocol version, the mapping valid sequence numbers for each framework package
241    pub protocol_version_system_package_table: BTreeMap<u64, BTreeMap<ObjectID, SequenceNumber>>,
242    // The current protocol version for this execution
243    pub current_protocol_version: u64,
244    // All state is contained here
245    pub storage: Storage,
246    // Debug events
247    pub exec_store_events: Arc<Mutex<Vec<ExecutionStoreEvent>>>,
248    // Debug events
249    pub metrics: Arc<ExecutionMetrics>,
250    // Used for fetching data from the network or remote store
251    pub fetcher: Fetchers,
252
253    // One can optionally override the executor version
254    // -1 implies use latest version
255    pub executor_version: Option<i64>,
256    // One can optionally override the protocol version
257    // -1 implies use latest version
258    // None implies use the protocol version at the time of execution
259    pub protocol_version: Option<i64>,
260    pub config_and_versions: Option<Vec<(ObjectID, SequenceNumber)>>,
261    // Retry policies due to RPC errors
262    pub num_retries_for_timeout: u32,
263    pub sleep_period_for_timeout: std::time::Duration,
264}
265
266impl LocalExec {
267    /// Wrapper around fetcher in case we want to add more functionality
268    /// Such as fetching from local DB from snapshot
269    pub async fn multi_download(
270        &self,
271        objs: &[(ObjectID, SequenceNumber)],
272    ) -> Result<Vec<Object>, ReplayEngineError> {
273        let mut num_retries_for_timeout = self.num_retries_for_timeout as i64;
274        while num_retries_for_timeout >= 0 {
275            match self.fetcher.multi_get_versioned(objs).await {
276                Ok(objs) => return Ok(objs),
277                Err(ReplayEngineError::SuiRpcRequestTimeout) => {
278                    warn!(
279                        "RPC request timed out. Retries left {}. Sleeping for {}s",
280                        num_retries_for_timeout,
281                        self.sleep_period_for_timeout.as_secs()
282                    );
283                    num_retries_for_timeout -= 1;
284                    tokio::time::sleep(self.sleep_period_for_timeout).await;
285                }
286                Err(e) => return Err(e),
287            }
288        }
289        Err(ReplayEngineError::SuiRpcRequestTimeout)
290    }
291    /// Wrapper around fetcher in case we want to add more functionality
292    /// Such as fetching from local DB from snapshot
293    pub async fn multi_download_latest(
294        &self,
295        objs: &[ObjectID],
296    ) -> Result<Vec<Object>, ReplayEngineError> {
297        let mut num_retries_for_timeout = self.num_retries_for_timeout as i64;
298        while num_retries_for_timeout >= 0 {
299            match self.fetcher.multi_get_latest(objs).await {
300                Ok(objs) => return Ok(objs),
301                Err(ReplayEngineError::SuiRpcRequestTimeout) => {
302                    warn!(
303                        "RPC request timed out. Retries left {}. Sleeping for {}s",
304                        num_retries_for_timeout,
305                        self.sleep_period_for_timeout.as_secs()
306                    );
307                    num_retries_for_timeout -= 1;
308                    tokio::time::sleep(self.sleep_period_for_timeout).await;
309                }
310                Err(e) => return Err(e),
311            }
312        }
313        Err(ReplayEngineError::SuiRpcRequestTimeout)
314    }
315
316    pub async fn fetch_loaded_child_refs(
317        &self,
318        tx_digest: &TransactionDigest,
319    ) -> Result<Vec<(ObjectID, SequenceNumber)>, ReplayEngineError> {
320        // Get the child objects loaded
321        self.fetcher.get_loaded_child_objects(tx_digest).await
322    }
323
324    pub async fn new_from_fn_url(http_url: &str) -> Result<Self, ReplayEngineError> {
325        Self::new_for_remote(
326            SuiClientBuilder::default()
327                .request_timeout(RPC_TIMEOUT_ERR_SLEEP_RETRY_PERIOD)
328                .max_concurrent_requests(MAX_CONCURRENT_REQUESTS)
329                .build(http_url)
330                .await?,
331            None,
332        )
333        .await
334    }
335
336    pub async fn replay_with_network_config(
337        rpc_url: String,
338        tx_digest: TransactionDigest,
339        expensive_safety_check_config: ExpensiveSafetyCheckConfig,
340        use_authority: bool,
341        executor_version: Option<i64>,
342        protocol_version: Option<i64>,
343        config_and_versions: Option<Vec<(ObjectID, SequenceNumber)>>,
344    ) -> Result<ExecutionSandboxState, ReplayEngineError> {
345        info!("Using RPC URL: {}", rpc_url);
346        LocalExec::new_from_fn_url(&rpc_url)
347            .await?
348            .init_for_execution()
349            .await?
350            .execute_transaction(
351                &tx_digest,
352                expensive_safety_check_config,
353                use_authority,
354                executor_version,
355                protocol_version,
356                config_and_versions,
357            )
358            .await
359    }
360
361    /// This captures the state of the network at a given point in time and populates
362    /// prptocol version tables including which system packages to fetch
363    /// If this function is called across epoch boundaries, the info might be stale.
364    /// But it should only be called once per epoch.
365    pub async fn init_for_execution(mut self) -> Result<Self, ReplayEngineError> {
366        self.populate_protocol_version_tables().await?;
367        tokio::task::yield_now().await;
368        Ok(self)
369    }
370
371    pub async fn reset_for_new_execution_with_client(self) -> Result<Self, ReplayEngineError> {
372        Self::new_for_remote(
373            self.client.expect("Remote client not initialized"),
374            Some(self.fetcher.into_remote()),
375        )
376        .await?
377        .init_for_execution()
378        .await
379    }
380
381    pub async fn new_for_remote(
382        client: SuiClient,
383        remote_fetcher: Option<RemoteFetcher>,
384    ) -> Result<Self, ReplayEngineError> {
385        // Use a throwaway metrics registry for local execution.
386        let registry = prometheus::Registry::new();
387        let metrics = Arc::new(ExecutionMetrics::new(&registry));
388
389        let fetcher = remote_fetcher.unwrap_or(RemoteFetcher::new(client.clone()));
390
391        Ok(Self {
392            client: Some(client),
393            protocol_version_epoch_table: BTreeMap::new(),
394            protocol_version_system_package_table: BTreeMap::new(),
395            current_protocol_version: 0,
396            exec_store_events: Arc::new(Mutex::new(Vec::new())),
397            metrics,
398            storage: Storage::default(),
399            fetcher: Fetchers::Remote(fetcher),
400            // TODO: make these configurable
401            num_retries_for_timeout: RPC_TIMEOUT_ERR_NUM_RETRIES,
402            sleep_period_for_timeout: RPC_TIMEOUT_ERR_SLEEP_RETRY_PERIOD,
403            executor_version: None,
404            protocol_version: None,
405            config_and_versions: None,
406        })
407    }
408
409    pub async fn new_for_state_dump(
410        path: &str,
411        backup_rpc_url: Option<String>,
412    ) -> Result<Self, ReplayEngineError> {
413        // Use a throwaway metrics registry for local execution.
414        let registry = prometheus::Registry::new();
415        let metrics = Arc::new(ExecutionMetrics::new(&registry));
416
417        let state = NodeStateDump::read_from_file(&PathBuf::from(path))?;
418        let current_protocol_version = state.protocol_version;
419        let fetcher = match backup_rpc_url {
420            Some(url) => NodeStateDumpFetcher::new(
421                state,
422                Some(RemoteFetcher::new(
423                    SuiClientBuilder::default()
424                        .request_timeout(RPC_TIMEOUT_ERR_SLEEP_RETRY_PERIOD)
425                        .max_concurrent_requests(MAX_CONCURRENT_REQUESTS)
426                        .build(url)
427                        .await?,
428                )),
429            ),
430            None => NodeStateDumpFetcher::new(state, None),
431        };
432
433        Ok(Self {
434            client: None,
435            protocol_version_epoch_table: BTreeMap::new(),
436            protocol_version_system_package_table: BTreeMap::new(),
437            current_protocol_version,
438            exec_store_events: Arc::new(Mutex::new(Vec::new())),
439            metrics,
440            storage: Storage::default(),
441            fetcher: Fetchers::NodeStateDump(fetcher),
442            // TODO: make these configurable
443            num_retries_for_timeout: RPC_TIMEOUT_ERR_NUM_RETRIES,
444            sleep_period_for_timeout: RPC_TIMEOUT_ERR_SLEEP_RETRY_PERIOD,
445            executor_version: None,
446            protocol_version: None,
447            config_and_versions: None,
448        })
449    }
450
451    pub async fn multi_download_and_store(
452        &mut self,
453        objs: &[(ObjectID, SequenceNumber)],
454    ) -> Result<Vec<Object>, ReplayEngineError> {
455        let objs = self.multi_download(objs).await?;
456
457        // Backfill the store
458        for obj in objs.iter() {
459            let o_ref = obj.compute_object_reference();
460            self.storage
461                .live_objects_store
462                .lock()
463                .expect("Can't lock")
464                .insert(o_ref.0, obj.clone());
465            self.storage
466                .object_version_cache
467                .lock()
468                .expect("Cannot lock")
469                .insert((o_ref.0, o_ref.1), obj.clone());
470            if obj.is_package() {
471                self.storage
472                    .package_cache
473                    .lock()
474                    .expect("Cannot lock")
475                    .insert(o_ref.0, obj.clone());
476            }
477        }
478        tokio::task::yield_now().await;
479        Ok(objs)
480    }
481
482    pub async fn multi_download_relevant_packages_and_store(
483        &mut self,
484        objs: Vec<ObjectID>,
485        protocol_version: u64,
486    ) -> Result<Vec<Object>, ReplayEngineError> {
487        let syst_packages_objs = if self.protocol_version.is_some_and(|i| i < 0) {
488            BuiltInFramework::genesis_objects().collect()
489        } else {
490            let syst_packages =
491                self.system_package_versions_for_protocol_version(protocol_version)?;
492            self.multi_download(&syst_packages).await?
493        };
494
495        // Download latest version of all packages that are not system packages
496        // This is okay since the versions can never change
497        let non_system_package_objs: Vec<_> = objs
498            .into_iter()
499            .filter(|o| !Self::system_package_ids(self.current_protocol_version).contains(o))
500            .collect();
501        let objs = self
502            .multi_download_latest(&non_system_package_objs)
503            .await?
504            .into_iter()
505            .chain(syst_packages_objs);
506
507        for obj in objs.clone() {
508            let o_ref = obj.compute_object_reference();
509            // We dont always want the latest in store
510            //self.storage.store.insert(o_ref.0, obj.clone());
511            self.storage
512                .object_version_cache
513                .lock()
514                .expect("Cannot lock")
515                .insert((o_ref.0, o_ref.1), obj.clone());
516            if obj.is_package() {
517                self.storage
518                    .package_cache
519                    .lock()
520                    .expect("Cannot lock")
521                    .insert(o_ref.0, obj.clone());
522            }
523        }
524        Ok(objs.collect())
525    }
526
527    // TODO: remove this after `futures::executor::block_on` is removed.
528    #[allow(clippy::disallowed_methods, clippy::result_large_err)]
529    pub fn download_object(
530        &self,
531        object_id: &ObjectID,
532        version: SequenceNumber,
533    ) -> Result<Object, ReplayEngineError> {
534        if self
535            .storage
536            .object_version_cache
537            .lock()
538            .expect("Cannot lock")
539            .contains_key(&(*object_id, version))
540        {
541            return Ok(self
542                .storage
543                .object_version_cache
544                .lock()
545                .expect("Cannot lock")
546                .get(&(*object_id, version))
547                .ok_or(ReplayEngineError::InternalCacheInvariantViolation {
548                    id: *object_id,
549                    version: Some(version),
550                })?
551                .clone());
552        }
553
554        let o = block_on(self.multi_download(&[(*object_id, version)])).map(|mut q| {
555            q.pop().unwrap_or_else(|| {
556                panic!(
557                    "Downloaded obj response cannot be empty {:?}",
558                    (*object_id, version)
559                )
560            })
561        })?;
562
563        let o_ref = o.compute_object_reference();
564        self.storage
565            .object_version_cache
566            .lock()
567            .expect("Cannot lock")
568            .insert((o_ref.0, o_ref.1), o.clone());
569        Ok(o)
570    }
571
572    // TODO: remove this after `futures::executor::block_on` is removed.
573    #[allow(clippy::disallowed_methods, clippy::result_large_err)]
574    pub fn download_latest_object(
575        &self,
576        object_id: &ObjectID,
577    ) -> Result<Option<Object>, ReplayEngineError> {
578        let resp = block_on({
579            //info!("Downloading latest object {object_id}");
580            self.multi_download_latest(std::slice::from_ref(object_id))
581        })
582        .map(|mut q| {
583            q.pop()
584                .unwrap_or_else(|| panic!("Downloaded obj response cannot be empty {}", *object_id))
585        });
586
587        match resp {
588            Ok(v) => Ok(Some(v)),
589            Err(ReplayEngineError::ObjectNotExist { id }) => {
590                error!(
591                    "Could not find object {id} on RPC server. It might have been pruned, deleted, or never existed."
592                );
593                Ok(None)
594            }
595            Err(ReplayEngineError::ObjectDeleted {
596                id,
597                version,
598                digest,
599            }) => {
600                error!("Object {id} {version} {digest} was deleted on RPC server.");
601                Ok(None)
602            }
603            Err(err) => Err(ReplayEngineError::SuiRpcError {
604                err: err.to_string(),
605            }),
606        }
607    }
608
609    #[allow(clippy::disallowed_methods, clippy::result_large_err)]
610    pub fn download_object_by_upper_bound(
611        &self,
612        object_id: &ObjectID,
613        version_upper_bound: VersionNumber,
614    ) -> Result<Option<Object>, ReplayEngineError> {
615        let local_object = self
616            .storage
617            .live_objects_store
618            .lock()
619            .expect("Can't lock")
620            .get(object_id)
621            .cloned();
622        if local_object.is_some() {
623            return Ok(local_object);
624        }
625        let response = block_on({
626            self.fetcher
627                .get_child_object(object_id, version_upper_bound)
628        });
629        match response {
630            Ok(object) => {
631                let obj_ref = object.compute_object_reference();
632                self.storage
633                    .live_objects_store
634                    .lock()
635                    .expect("Can't lock")
636                    .insert(*object_id, object.clone());
637                self.storage
638                    .object_version_cache
639                    .lock()
640                    .expect("Can't lock")
641                    .insert((obj_ref.0, obj_ref.1), object.clone());
642                Ok(Some(object))
643            }
644            Err(ReplayEngineError::ObjectNotExist { id }) => {
645                error!(
646                    "Could not find child object {id} on RPC server. It might have been pruned, deleted, or never existed."
647                );
648                Ok(None)
649            }
650            Err(ReplayEngineError::ObjectDeleted {
651                id,
652                version,
653                digest,
654            }) => {
655                error!("Object {id} {version} {digest} was deleted on RPC server.");
656                Ok(None)
657            }
658            // This is a child object which was not found in the store (e.g., due to exists
659            // check before creating the dynamic field).
660            Err(ReplayEngineError::ObjectVersionNotFound { id, version }) => {
661                info!(
662                    "Object {id} {version} not found on RPC server -- this may have been pruned or never existed."
663                );
664                Ok(None)
665            }
666            Err(err) => Err(ReplayEngineError::SuiRpcError {
667                err: err.to_string(),
668            }),
669        }
670    }
671
672    pub async fn get_checkpoint_txs(
673        &self,
674        checkpoint_id: u64,
675    ) -> Result<Vec<TransactionDigest>, ReplayEngineError> {
676        self.fetcher
677            .get_checkpoint_txs(checkpoint_id)
678            .await
679            .map_err(|e| ReplayEngineError::SuiRpcError { err: e.to_string() })
680    }
681
682    pub async fn execute_all_in_checkpoints(
683        &mut self,
684        checkpoint_ids: &[u64],
685        expensive_safety_check_config: &ExpensiveSafetyCheckConfig,
686        terminate_early: bool,
687        use_authority: bool,
688    ) -> Result<(u64, u64), ReplayEngineError> {
689        // Get all the TXs at this checkpoint
690        let mut txs = Vec::new();
691        for checkpoint_id in checkpoint_ids {
692            txs.extend(self.get_checkpoint_txs(*checkpoint_id).await?);
693        }
694        let num = txs.len();
695        let mut succeeded = 0;
696        for tx in txs {
697            match self
698                .execute_transaction(
699                    &tx,
700                    expensive_safety_check_config.clone(),
701                    use_authority,
702                    None,
703                    None,
704                    None,
705                )
706                .await
707                .map(|q| q.check_effects())
708            {
709                Err(e) | Ok(Err(e)) => {
710                    if terminate_early {
711                        return Err(e);
712                    }
713                    error!("Error executing tx: {},  {:#?}", tx, e);
714                    continue;
715                }
716                _ => (),
717            }
718
719            succeeded += 1;
720        }
721        Ok((succeeded, num as u64))
722    }
723
724    pub async fn execution_engine_execute_with_tx_info_impl(
725        &mut self,
726        tx_info: &OnChainTransactionInfo,
727        override_transaction_kind: Option<TransactionKind>,
728        expensive_safety_check_config: ExpensiveSafetyCheckConfig,
729    ) -> Result<ExecutionSandboxState, ReplayEngineError> {
730        let tx_digest = &tx_info.tx_digest;
731        // Before protocol version 16, the generation of effects depends on the wrapped tombstones.
732        // It is not possible to retrieve such data for replay.
733        if tx_info.protocol_version.as_u64() < 16 {
734            warn!(
735                "Protocol version ({:?}) too old: {}, skipping transaction",
736                tx_info.protocol_version, tx_digest
737            );
738            return Err(ReplayEngineError::TransactionNotSupported {
739                digest: *tx_digest,
740                reason: "Protocol version too old".to_string(),
741            });
742        }
743        // Initialize the state necessary for execution
744        // Get the input objects
745        let input_objects = self.initialize_execution_env_state(tx_info).await?;
746        assert_eq!(
747            &input_objects.filter_shared_objects().len(),
748            &tx_info.shared_object_refs.len()
749        );
750        // At this point we have all the objects needed for replay
751
752        // This assumes we already initialized the protocol version table `protocol_version_epoch_table`
753        let protocol_config =
754            &ProtocolConfig::get_for_version(tx_info.protocol_version, tx_info.chain);
755
756        let metrics = self.metrics.clone();
757
758        let ov = self.executor_version;
759
760        // We could probably cache the executor per protocol config
761        let executor = get_executor(ov, protocol_config, expensive_safety_check_config);
762
763        // All prep done
764        let expensive_checks = true;
765        let transaction_kind = override_transaction_kind.unwrap_or(tx_info.kind.clone());
766        let gas_status = if tx_info.kind.is_system_tx() {
767            SuiGasStatus::new_unmetered()
768        } else {
769            SuiGasStatus::new(
770                tx_info.gas_budget,
771                tx_info.gas_price,
772                tx_info.reference_gas_price,
773                protocol_config,
774            )
775            .expect("Failed to create gas status")
776        };
777        let gas_data = GasData {
778            payment: tx_info.gas.clone(),
779            owner: tx_info.gas_owner.unwrap_or(tx_info.sender),
780            price: tx_info.gas_price,
781            budget: tx_info.gas_budget,
782        };
783        let checked_input_objects = CheckedInputObjects::new_for_replay(input_objects.clone());
784        let early_execution_error = get_early_execution_error(
785            tx_digest,
786            &checked_input_objects,
787            &HashSet::new(),
788            // TODO(address-balances): Support balance withdraw status for replay
789            &FundsWithdrawStatus::MaybeSufficient,
790        );
791        let execution_params = match early_execution_error {
792            None => ExecutionOrEarlyError::ok(None),
793            Some(errors) => ExecutionOrEarlyError::failed(errors, None),
794        };
795        let (inner_store, gas_status, effects, _timings, result) = executor
796            .execute_transaction_to_effects_and_execution_error(
797                &self,
798                protocol_config,
799                metrics.clone(),
800                expensive_checks,
801                execution_params,
802                &tx_info.executed_epoch,
803                tx_info.epoch_start_timestamp,
804                checked_input_objects,
805                std::collections::BTreeMap::new(),
806                gas_data,
807                gas_status,
808                transaction_kind.clone(),
809                None, // compat_args
810                tx_info.sender,
811                *tx_digest,
812                &mut None,
813            );
814
815        if let Err(err) = self.pretty_print_for_tracing(
816            &gas_status,
817            &executor,
818            tx_info,
819            &transaction_kind,
820            protocol_config,
821            metrics,
822            expensive_checks,
823            input_objects.clone(),
824        ) {
825            error!("Failed to pretty print for tracing: {:?}", err);
826        }
827
828        let all_required_objects = self.storage.all_objects();
829
830        let effects =
831            SuiTransactionBlockEffects::try_from(effects).map_err(ReplayEngineError::from)?;
832
833        Ok(ExecutionSandboxState {
834            transaction_info: tx_info.clone(),
835            required_objects: all_required_objects,
836            local_exec_temporary_store: Some(inner_store),
837            local_exec_effects: effects,
838            local_exec_status: Some(result),
839        })
840    }
841
842    fn pretty_print_for_tracing(
843        &self,
844        gas_status: &SuiGasStatus,
845        executor: &Arc<dyn Executor + Send + Sync>,
846        tx_info: &OnChainTransactionInfo,
847        transaction_kind: &TransactionKind,
848        protocol_config: &ProtocolConfig,
849        metrics: Arc<ExecutionMetrics>,
850        expensive_checks: bool,
851        input_objects: InputObjects,
852    ) -> anyhow::Result<()> {
853        trace!(target: "replay_gas_info", "{}", Pretty(gas_status));
854
855        let skip_checks = true;
856        let gas_data = GasData {
857            payment: tx_info.gas.clone(),
858            owner: tx_info.gas_owner.unwrap_or(tx_info.sender),
859            price: tx_info.gas_price,
860            budget: tx_info.gas_budget,
861        };
862        let checked_input_objects = CheckedInputObjects::new_for_replay(input_objects.clone());
863        let early_execution_error = get_early_execution_error(
864            &tx_info.tx_digest,
865            &checked_input_objects,
866            &HashSet::new(),
867            // TODO(address-balances): Support balance withdraw status for replay
868            &FundsWithdrawStatus::MaybeSufficient,
869        );
870        let execution_params = match early_execution_error {
871            None => ExecutionOrEarlyError::ok(None),
872            Some(errors) => ExecutionOrEarlyError::failed(errors, None),
873        };
874        if let ProgrammableTransaction(pt) = transaction_kind {
875            trace!(
876                target: "replay_ptb_info",
877                "{}",
878                Pretty(&FullPTB {
879                    ptb: pt.clone(),
880                    results: transform_command_results_to_annotated(
881                        protocol_config,
882                        executor,
883                        &self.clone(),
884                        executor.dev_inspect_transaction(
885                            &self,
886                            protocol_config,
887                            metrics,
888                            expensive_checks,
889                            execution_params,
890                            &tx_info.executed_epoch,
891                            tx_info.epoch_start_timestamp,
892                            CheckedInputObjects::new_for_replay(input_objects),
893                            gas_data,
894                            SuiGasStatus::new(
895                                tx_info.gas_budget,
896                                tx_info.gas_price,
897                                tx_info.reference_gas_price,
898                                protocol_config,
899                            )?,
900                            transaction_kind.clone(),
901                            None, // compat_args
902                            tx_info.sender,
903                            tx_info.sender_signed_data.digest(),
904                            skip_checks,
905                        )
906                        .3
907                        .unwrap_or_default(),
908                    )?,
909            }));
910        }
911        Ok(())
912    }
913
914    /// Must be called after `init_for_execution`
915    #[allow(clippy::result_large_err)]
916    pub async fn execution_engine_execute_impl(
917        &mut self,
918        tx_digest: &TransactionDigest,
919        expensive_safety_check_config: ExpensiveSafetyCheckConfig,
920    ) -> Result<ExecutionSandboxState, ReplayEngineError> {
921        if self.is_remote_replay() {
922            assert!(
923                !self.protocol_version_system_package_table.is_empty()
924                    || !self.protocol_version_epoch_table.is_empty(),
925                "Required tables not populated. Must call `init_for_execution` before executing transactions"
926            );
927        }
928
929        let tx_info = if self.is_remote_replay() {
930            self.resolve_tx_components(tx_digest).await?
931        } else {
932            self.resolve_tx_components_from_dump(tx_digest).await?
933        };
934        self.execution_engine_execute_with_tx_info_impl(
935            &tx_info,
936            None,
937            expensive_safety_check_config,
938        )
939        .await
940    }
941
942    /// Executes a transaction with the state specified in `pre_run_sandbox`
943    /// This is useful for executing a transaction with a specific state
944    /// However if the state in invalid, the behavior is undefined.
945    #[allow(clippy::result_large_err)]
946    pub async fn certificate_execute_with_sandbox_state(
947        pre_run_sandbox: &ExecutionSandboxState,
948    ) -> Result<ExecutionSandboxState, ReplayEngineError> {
949        // These cannot be changed and are inherited from the sandbox state
950        let executed_epoch = pre_run_sandbox.transaction_info.executed_epoch;
951        let reference_gas_price = pre_run_sandbox.transaction_info.reference_gas_price;
952        let epoch_start_timestamp = pre_run_sandbox.transaction_info.epoch_start_timestamp;
953        let protocol_config = ProtocolConfig::get_for_version(
954            pre_run_sandbox.transaction_info.protocol_version,
955            pre_run_sandbox.transaction_info.chain,
956        );
957        let required_objects = pre_run_sandbox.required_objects.clone();
958        let store = InMemoryStorage::new(required_objects.clone());
959
960        let transaction =
961            Transaction::new(pre_run_sandbox.transaction_info.sender_signed_data.clone());
962
963        // TODO: This will not work for deleted shared objects. We need to persist that information in the sandbox.
964        // TODO: A lot of the following code is replicated in several places. We should introduce a few
965        // traits and make them shared so that we don't have to fix one by one when we have major execution
966        // layer changes.
967        let input_objects = store.read_input_objects_for_transaction(&transaction);
968        let executable = VerifiedExecutableTransaction::new_from_consensus(
969            VerifiedTransaction::new_unchecked(transaction),
970            executed_epoch,
971        );
972        let (gas_status, input_objects) = sui_transaction_checks::check_certificate_input(
973            &executable,
974            input_objects,
975            &protocol_config,
976            reference_gas_price,
977        )
978        .unwrap();
979        let (kind, signer, gas_data) = executable.transaction_data().execution_parts();
980        let executor = sui_execution::executor(&protocol_config, true).unwrap();
981        let early_execution_error = get_early_execution_error(
982            executable.digest(),
983            &input_objects,
984            &HashSet::new(),
985            // TODO(address-balances): Support balance withdraw status for replay
986            &FundsWithdrawStatus::MaybeSufficient,
987        );
988        let execution_params = match early_execution_error {
989            None => ExecutionOrEarlyError::ok(None),
990            Some(errors) => ExecutionOrEarlyError::failed(errors, None),
991        };
992        let (_, _, effects, _timings, exec_res) = executor
993            .execute_transaction_to_effects_and_execution_error(
994                &store,
995                &protocol_config,
996                Arc::new(ExecutionMetrics::new(&Registry::new())),
997                true,
998                execution_params,
999                &executed_epoch,
1000                epoch_start_timestamp,
1001                input_objects,
1002                std::collections::BTreeMap::new(),
1003                gas_data,
1004                gas_status,
1005                kind,
1006                None, // compat_args
1007                signer,
1008                *executable.digest(),
1009                &mut None,
1010            );
1011
1012        let effects =
1013            SuiTransactionBlockEffects::try_from(effects).map_err(ReplayEngineError::from)?;
1014
1015        Ok(ExecutionSandboxState {
1016            transaction_info: pre_run_sandbox.transaction_info.clone(),
1017            required_objects,
1018            local_exec_temporary_store: None, // We dont capture it for cert exec run
1019            local_exec_effects: effects,
1020            local_exec_status: Some(exec_res),
1021        })
1022    }
1023
1024    /// Must be called after `init_for_execution`
1025    /// This executes from `sui_core::authority::AuthorityState::try_execute_immediately`
1026    #[allow(clippy::result_large_err)]
1027    pub async fn certificate_execute(
1028        &mut self,
1029        tx_digest: &TransactionDigest,
1030        expensive_safety_check_config: ExpensiveSafetyCheckConfig,
1031    ) -> Result<ExecutionSandboxState, ReplayEngineError> {
1032        // Use the lighterweight execution engine to get the pre-run state
1033        let pre_run_sandbox = self
1034            .execution_engine_execute_impl(tx_digest, expensive_safety_check_config)
1035            .await?;
1036        Self::certificate_execute_with_sandbox_state(&pre_run_sandbox).await
1037    }
1038
1039    /// Must be called after `init_for_execution`
1040    /// This executes from `sui_adapter::execution_engine::execute_transaction_to_effects`
1041    #[allow(clippy::result_large_err)]
1042    pub async fn execution_engine_execute(
1043        &mut self,
1044        tx_digest: &TransactionDigest,
1045        expensive_safety_check_config: ExpensiveSafetyCheckConfig,
1046    ) -> Result<ExecutionSandboxState, ReplayEngineError> {
1047        let sandbox_state = self
1048            .execution_engine_execute_impl(tx_digest, expensive_safety_check_config)
1049            .await?;
1050
1051        Ok(sandbox_state)
1052    }
1053
1054    #[allow(clippy::result_large_err)]
1055    pub async fn execute_state_dump(
1056        &mut self,
1057        expensive_safety_check_config: ExpensiveSafetyCheckConfig,
1058    ) -> Result<(ExecutionSandboxState, NodeStateDump), ReplayEngineError> {
1059        assert!(!self.is_remote_replay());
1060
1061        let d = match self.fetcher.clone() {
1062            Fetchers::NodeStateDump(d) => d,
1063            _ => panic!("Invalid fetcher for state dump"),
1064        };
1065        let tx_digest = d.node_state_dump.clone().tx_digest;
1066        let sandbox_state = self
1067            .execution_engine_execute_impl(&tx_digest, expensive_safety_check_config)
1068            .await?;
1069
1070        Ok((sandbox_state, d.node_state_dump))
1071    }
1072
1073    #[allow(clippy::result_large_err)]
1074    pub async fn execute_transaction(
1075        &mut self,
1076        tx_digest: &TransactionDigest,
1077        expensive_safety_check_config: ExpensiveSafetyCheckConfig,
1078        use_authority: bool,
1079        executor_version: Option<i64>,
1080        protocol_version: Option<i64>,
1081        config_and_versions: Option<Vec<(ObjectID, SequenceNumber)>>,
1082    ) -> Result<ExecutionSandboxState, ReplayEngineError> {
1083        self.executor_version = executor_version;
1084        self.protocol_version = protocol_version;
1085        self.config_and_versions = config_and_versions;
1086        if use_authority {
1087            self.certificate_execute(tx_digest, expensive_safety_check_config.clone())
1088                .await
1089        } else {
1090            self.execution_engine_execute(tx_digest, expensive_safety_check_config)
1091                .await
1092        }
1093    }
1094    fn system_package_ids(protocol_version: u64) -> Vec<ObjectID> {
1095        let mut ids = BuiltInFramework::all_package_ids();
1096
1097        if protocol_version < 5 {
1098            ids.retain(|id| *id != DEEPBOOK_PACKAGE_ID)
1099        }
1100        ids
1101    }
1102
1103    /// This is the only function which accesses the network during execution
1104    #[allow(clippy::result_large_err)]
1105    pub fn get_or_download_object(
1106        &self,
1107        obj_id: &ObjectID,
1108        package_expected: bool,
1109    ) -> Result<Option<Object>, ReplayEngineError> {
1110        if package_expected {
1111            if let Some(obj) = self
1112                .storage
1113                .package_cache
1114                .lock()
1115                .expect("Cannot lock")
1116                .get(obj_id)
1117            {
1118                return Ok(Some(obj.clone()));
1119            };
1120            // Check if its a system package because we must've downloaded all
1121            // TODO: Will return this check once we can download completely for other networks
1122            // assert!(
1123            //     !self.system_package_ids().contains(obj_id),
1124            //     "All system packages should be downloaded already"
1125            // );
1126        } else if let Some(obj) = self
1127            .storage
1128            .live_objects_store
1129            .lock()
1130            .expect("Can't lock")
1131            .get(obj_id)
1132        {
1133            return Ok(Some(obj.clone()));
1134        }
1135
1136        let Some(o) = self.download_latest_object(obj_id)? else {
1137            return Ok(None);
1138        };
1139
1140        if o.is_package() {
1141            assert!(
1142                package_expected,
1143                "Did not expect package but downloaded object is a package: {obj_id}"
1144            );
1145
1146            self.storage
1147                .package_cache
1148                .lock()
1149                .expect("Cannot lock")
1150                .insert(*obj_id, o.clone());
1151        }
1152        let o_ref = o.compute_object_reference();
1153        self.storage
1154            .object_version_cache
1155            .lock()
1156            .expect("Cannot lock")
1157            .insert((o_ref.0, o_ref.1), o.clone());
1158        Ok(Some(o))
1159    }
1160
1161    pub fn is_remote_replay(&self) -> bool {
1162        matches!(self.fetcher, Fetchers::Remote(_))
1163    }
1164
1165    /// Must be called after `populate_protocol_version_tables`
1166    #[allow(clippy::result_large_err)]
1167    pub fn system_package_versions_for_protocol_version(
1168        &self,
1169        protocol_version: u64,
1170    ) -> Result<Vec<(ObjectID, SequenceNumber)>, ReplayEngineError> {
1171        match &self.fetcher {
1172            Fetchers::Remote(_) => Ok(self
1173                .protocol_version_system_package_table
1174                .get(&protocol_version)
1175                .ok_or(ReplayEngineError::FrameworkObjectVersionTableNotPopulated {
1176                    protocol_version,
1177                })?
1178                .clone()
1179                .into_iter()
1180                .collect()),
1181
1182            Fetchers::NodeStateDump(d) => Ok(d
1183                .node_state_dump
1184                .relevant_system_packages
1185                .iter()
1186                .map(|w| (w.id, w.version, w.digest))
1187                .map(|q| (q.0, q.1))
1188                .collect()),
1189        }
1190    }
1191
1192    pub async fn protocol_ver_to_epoch_map(
1193        &self,
1194    ) -> Result<BTreeMap<u64, ProtocolVersionSummary>, ReplayEngineError> {
1195        let mut range_map = BTreeMap::new();
1196        let epoch_change_events = self.fetcher.get_epoch_change_events(false).await?;
1197
1198        // Exception for Genesis: Protocol version 1 at epoch 0
1199        let mut tx_digest = *self
1200            .fetcher
1201            .get_checkpoint_txs(0)
1202            .await?
1203            .first()
1204            .expect("Genesis TX must be in first checkpoint");
1205        // Somehow the genesis TX did not emit any event, but we know it was the start of version 1
1206        // So we need to manually add this range
1207        let (mut start_epoch, mut start_protocol_version, mut start_checkpoint) =
1208            (0, 1, Some(0u64));
1209
1210        let (mut curr_epoch, mut curr_protocol_version, mut curr_checkpoint) =
1211            (start_epoch, start_protocol_version, start_checkpoint);
1212
1213        (start_epoch, start_protocol_version, start_checkpoint) =
1214            (curr_epoch, curr_protocol_version, curr_checkpoint);
1215
1216        // This is the final tx digest for the epoch change. We need this to track the final checkpoint
1217        let mut end_epoch_tx_digest = tx_digest;
1218
1219        for event in epoch_change_events {
1220            (curr_epoch, curr_protocol_version) = extract_epoch_and_version(event.clone())?;
1221            end_epoch_tx_digest = event.id.tx_digest;
1222
1223            if start_protocol_version == curr_protocol_version {
1224                // Same range
1225                continue;
1226            }
1227
1228            // Change in prot version
1229            // Find the last checkpoint
1230            curr_checkpoint = self
1231                .fetcher
1232                .get_transaction(&event.id.tx_digest)
1233                .await?
1234                .checkpoint;
1235            // Insert the last range
1236            range_map.insert(
1237                start_protocol_version,
1238                ProtocolVersionSummary {
1239                    protocol_version: start_protocol_version,
1240                    epoch_start: start_epoch,
1241                    epoch_end: curr_epoch - 1,
1242                    checkpoint_start: start_checkpoint,
1243                    checkpoint_end: curr_checkpoint.map(|x| x - 1),
1244                    epoch_change_tx: tx_digest,
1245                },
1246            );
1247
1248            start_epoch = curr_epoch;
1249            start_protocol_version = curr_protocol_version;
1250            tx_digest = event.id.tx_digest;
1251            start_checkpoint = curr_checkpoint;
1252        }
1253
1254        // Insert the last range
1255        range_map.insert(
1256            curr_protocol_version,
1257            ProtocolVersionSummary {
1258                protocol_version: curr_protocol_version,
1259                epoch_start: start_epoch,
1260                epoch_end: curr_epoch,
1261                checkpoint_start: curr_checkpoint,
1262                checkpoint_end: self
1263                    .fetcher
1264                    .get_transaction(&end_epoch_tx_digest)
1265                    .await?
1266                    .checkpoint,
1267                epoch_change_tx: tx_digest,
1268            },
1269        );
1270
1271        Ok(range_map)
1272    }
1273
1274    pub fn protocol_version_for_epoch(
1275        epoch: u64,
1276        mp: &BTreeMap<u64, (TransactionDigest, u64, u64)>,
1277    ) -> u64 {
1278        // Naive impl but works for now
1279        // Can improve with range algos & data structures
1280        let mut version = 1;
1281        for (k, v) in mp.iter().rev() {
1282            if v.1 <= epoch {
1283                version = *k;
1284                break;
1285            }
1286        }
1287        version
1288    }
1289
1290    pub async fn populate_protocol_version_tables(&mut self) -> Result<(), ReplayEngineError> {
1291        self.protocol_version_epoch_table = self.protocol_ver_to_epoch_map().await?;
1292
1293        let system_package_revisions = self.system_package_versions().await?;
1294
1295        // This can be more efficient but small footprint so okay for now
1296        //Table is sorted from earliest to latest
1297        for (
1298            prot_ver,
1299            ProtocolVersionSummary {
1300                epoch_change_tx: tx_digest,
1301                ..
1302            },
1303        ) in self.protocol_version_epoch_table.clone()
1304        {
1305            // Use the previous versions protocol version table
1306            let mut working = if prot_ver <= 1 {
1307                BTreeMap::new()
1308            } else {
1309                self.protocol_version_system_package_table
1310                    .iter()
1311                    .rev()
1312                    .find(|(ver, _)| **ver <= prot_ver)
1313                    .expect("Prev entry must exist")
1314                    .1
1315                    .clone()
1316            };
1317
1318            for (id, versions) in system_package_revisions.iter() {
1319                // Oldest appears first in list, so reverse
1320                for ver in versions.iter().rev() {
1321                    if ver.1 == tx_digest {
1322                        // Found the version for this protocol version
1323                        working.insert(*id, ver.0);
1324                        break;
1325                    }
1326                }
1327            }
1328            self.protocol_version_system_package_table
1329                .insert(prot_ver, working);
1330        }
1331        Ok(())
1332    }
1333
1334    pub async fn system_package_versions(
1335        &self,
1336    ) -> Result<BTreeMap<ObjectID, Vec<(SequenceNumber, TransactionDigest)>>, ReplayEngineError>
1337    {
1338        let system_package_ids = Self::system_package_ids(
1339            *self
1340                .protocol_version_epoch_table
1341                .keys()
1342                .peekable()
1343                .last()
1344                .expect("Protocol version epoch table not populated"),
1345        );
1346        let mut system_package_objs = self.multi_download_latest(&system_package_ids).await?;
1347
1348        let mut mapping = BTreeMap::new();
1349
1350        // Extract all the transactions which created or mutated this object
1351        while !system_package_objs.is_empty() {
1352            // For the given object and its version, record the transaction which upgraded or created it
1353            let previous_txs: Vec<_> = system_package_objs
1354                .iter()
1355                .map(|o| (o.compute_object_reference(), o.previous_transaction))
1356                .collect();
1357
1358            previous_txs.iter().for_each(|((id, ver, _), tx)| {
1359                mapping.entry(*id).or_insert(vec![]).push((*ver, *tx));
1360            });
1361
1362            // Next round
1363            // Get the previous version of each object if exists
1364            let previous_ver_refs: Vec<_> = previous_txs
1365                .iter()
1366                .filter_map(|(q, _)| {
1367                    let prev_ver = u64::from(q.1) - 1;
1368                    if prev_ver == 0 {
1369                        None
1370                    } else {
1371                        Some((q.0, SequenceNumber::from(prev_ver)))
1372                    }
1373                })
1374                .collect();
1375            system_package_objs = match self.multi_download(&previous_ver_refs).await {
1376                Ok(packages) => packages,
1377                Err(ReplayEngineError::ObjectNotExist { id }) => {
1378                    // This happens when the RPC server prunes older object
1379                    // Replays in the current protocol version will work but old ones might not
1380                    // as we cannot fetch the package
1381                    warn!(
1382                        "Object {} does not exist on RPC server. This might be due to pruning. Historical replays might not work",
1383                        id
1384                    );
1385                    break;
1386                }
1387                Err(ReplayEngineError::ObjectVersionNotFound { id, version }) => {
1388                    // This happens when the RPC server prunes older object
1389                    // Replays in the current protocol version will work but old ones might not
1390                    // as we cannot fetch the package
1391                    warn!(
1392                        "Object {} at version {} does not exist on RPC server. This might be due to pruning. Historical replays might not work",
1393                        id, version
1394                    );
1395                    break;
1396                }
1397                Err(ReplayEngineError::ObjectVersionTooHigh {
1398                    id,
1399                    asked_version,
1400                    latest_version,
1401                }) => {
1402                    warn!(
1403                        "Object {} at version {} does not exist on RPC server. Latest version is {}. This might be due to pruning. Historical replays might not work",
1404                        id, asked_version, latest_version
1405                    );
1406                    break;
1407                }
1408                Err(ReplayEngineError::ObjectDeleted {
1409                    id,
1410                    version,
1411                    digest,
1412                }) => {
1413                    // This happens when the RPC server prunes older object
1414                    // Replays in the current protocol version will work but old ones might not
1415                    // as we cannot fetch the package
1416                    warn!(
1417                        "Object {} at version {} digest {} deleted from RPC server. This might be due to pruning. Historical replays might not work",
1418                        id, version, digest
1419                    );
1420                    break;
1421                }
1422                Err(e) => return Err(e),
1423            };
1424        }
1425        Ok(mapping)
1426    }
1427
1428    pub async fn get_protocol_config(
1429        &self,
1430        epoch_id: EpochId,
1431        chain: Chain,
1432    ) -> Result<ProtocolConfig, ReplayEngineError> {
1433        match self.protocol_version {
1434            Some(x) if x < 0 => Ok(ProtocolConfig::get_for_max_version_UNSAFE()),
1435            Some(v) => Ok(ProtocolConfig::get_for_version((v as u64).into(), chain)),
1436            None => self
1437                .protocol_version_epoch_table
1438                .iter()
1439                .rev()
1440                .find(|(_, rg)| epoch_id >= rg.epoch_start)
1441                .map(|(p, _rg)| Ok(ProtocolConfig::get_for_version((*p).into(), chain)))
1442                .unwrap_or_else(|| {
1443                    Err(ReplayEngineError::ProtocolVersionNotFound { epoch: epoch_id })
1444                }),
1445        }
1446    }
1447
1448    pub async fn checkpoints_for_epoch(
1449        &self,
1450        epoch_id: u64,
1451    ) -> Result<(u64, u64), ReplayEngineError> {
1452        let epoch_change_events = self
1453            .fetcher
1454            .get_epoch_change_events(true)
1455            .await?
1456            .into_iter()
1457            .collect::<Vec<_>>();
1458        let (start_checkpoint, start_epoch_idx) = if epoch_id == 0 {
1459            (0, 1)
1460        } else {
1461            let idx = epoch_change_events
1462                .iter()
1463                .position(|ev| match extract_epoch_and_version(ev.clone()) {
1464                    Ok((epoch, _)) => epoch == epoch_id,
1465                    Err(_) => false,
1466                })
1467                .ok_or(ReplayEngineError::EventNotFound { epoch: epoch_id })?;
1468            let epoch_change_tx = epoch_change_events[idx].id.tx_digest;
1469            (
1470                self.fetcher
1471                    .get_transaction(&epoch_change_tx)
1472                    .await?
1473                    .checkpoint
1474                    .unwrap_or_else(|| {
1475                        panic!(
1476                            "Checkpoint for transaction {} not present. Could be due to pruning",
1477                            epoch_change_tx
1478                        )
1479                    }),
1480                idx,
1481            )
1482        };
1483
1484        let next_epoch_change_tx = epoch_change_events
1485            .get(start_epoch_idx + 1)
1486            .map(|v| v.id.tx_digest)
1487            .ok_or(ReplayEngineError::UnableToDetermineCheckpoint { epoch: epoch_id })?;
1488
1489        let next_epoch_checkpoint = self
1490            .fetcher
1491            .get_transaction(&next_epoch_change_tx)
1492            .await?
1493            .checkpoint
1494            .unwrap_or_else(|| {
1495                panic!(
1496                    "Checkpoint for transaction {} not present. Could be due to pruning",
1497                    next_epoch_change_tx
1498                )
1499            });
1500
1501        Ok((start_checkpoint, next_epoch_checkpoint - 1))
1502    }
1503
1504    pub async fn get_epoch_start_timestamp_and_rgp(
1505        &self,
1506        epoch_id: u64,
1507        tx_digest: &TransactionDigest,
1508    ) -> Result<(u64, u64), ReplayEngineError> {
1509        if epoch_id == 0 {
1510            return Err(ReplayEngineError::TransactionNotSupported {
1511                digest: *tx_digest,
1512                reason: "Transactions from epoch 0 not supported".to_string(),
1513            });
1514        }
1515        self.fetcher
1516            .get_epoch_start_timestamp_and_rgp(epoch_id)
1517            .await
1518    }
1519
1520    fn add_config_objects_if_needed(
1521        &self,
1522        status: &SuiExecutionStatus,
1523    ) -> Vec<(ObjectID, SequenceNumber)> {
1524        match parse_effect_error_for_denied_coins(status) {
1525            Some(coin_type) => {
1526                let Some(mut config_id_and_version) = self.config_and_versions.clone() else {
1527                    panic!(
1528                        "Need to specify the config object ID and version for '{coin_type}' in order to replay this transaction"
1529                    );
1530                };
1531                // NB: the version of the deny list object doesn't matter
1532                if !config_id_and_version
1533                    .iter()
1534                    .any(|(id, _)| id == &SUI_DENY_LIST_OBJECT_ID)
1535                {
1536                    let deny_list_oid_version = self.download_latest_object(&SUI_DENY_LIST_OBJECT_ID)
1537                        .ok()
1538                        .flatten()
1539                        .expect("Unable to download the deny list object for a transaction that requires it")
1540                        .version();
1541                    config_id_and_version.push((SUI_DENY_LIST_OBJECT_ID, deny_list_oid_version));
1542                }
1543                config_id_and_version
1544            }
1545            None => vec![],
1546        }
1547    }
1548
1549    async fn resolve_tx_components(
1550        &self,
1551        tx_digest: &TransactionDigest,
1552    ) -> Result<OnChainTransactionInfo, ReplayEngineError> {
1553        assert!(self.is_remote_replay());
1554        // Fetch full transaction content
1555        let tx_info = self.fetcher.get_transaction(tx_digest).await?;
1556        let sender = match tx_info.clone().transaction.unwrap().data {
1557            sui_json_rpc_types::SuiTransactionBlockData::V1(tx) => tx.sender,
1558        };
1559        let SuiTransactionBlockEffects::V1(effects) = tx_info.clone().effects.unwrap();
1560
1561        let config_objects = self.add_config_objects_if_needed(effects.status());
1562
1563        let raw_tx_bytes = tx_info.clone().raw_transaction;
1564        let orig_tx: SenderSignedData = bcs::from_bytes(&raw_tx_bytes).unwrap();
1565        let input_objs = orig_tx
1566            .transaction_data()
1567            .input_objects()
1568            .map_err(|e| ReplayEngineError::UserInputError { err: e })?;
1569        let tx_kind_orig = orig_tx.transaction_data().kind();
1570
1571        // Download the objects at the version right before the execution of this TX
1572        let modified_at_versions: Vec<(ObjectID, SequenceNumber)> = effects.modified_at_versions();
1573
1574        let shared_object_refs: Vec<ObjectRef> = effects
1575            .shared_objects()
1576            .iter()
1577            .map(|so_ref| {
1578                if so_ref.digest == ObjectDigest::OBJECT_DIGEST_DELETED {
1579                    unimplemented!(
1580                        "Replay of deleted shared object transactions is not supported yet"
1581                    );
1582                } else {
1583                    so_ref.to_object_ref()
1584                }
1585            })
1586            .collect();
1587        let gas_data = match tx_info.clone().transaction.unwrap().data {
1588            sui_json_rpc_types::SuiTransactionBlockData::V1(tx) => tx.gas_data,
1589        };
1590        let gas_object_refs: Vec<_> = gas_data
1591            .payment
1592            .iter()
1593            .map(|obj_ref| obj_ref.to_object_ref())
1594            .collect();
1595        let receiving_objs = orig_tx
1596            .transaction_data()
1597            .receiving_objects()
1598            .into_iter()
1599            .map(|(obj_id, version, _)| (obj_id, version))
1600            .collect();
1601
1602        let epoch_id = effects.executed_epoch;
1603        let chain = chain_from_chain_id(self.fetcher.get_chain_id().await?.as_str());
1604
1605        // Extract the epoch start timestamp
1606        let (epoch_start_timestamp, reference_gas_price) = self
1607            .get_epoch_start_timestamp_and_rgp(epoch_id, tx_digest)
1608            .await?;
1609
1610        Ok(OnChainTransactionInfo {
1611            kind: tx_kind_orig.clone(),
1612            sender,
1613            modified_at_versions,
1614            input_objects: input_objs,
1615            shared_object_refs,
1616            gas: gas_object_refs,
1617            gas_owner: (gas_data.owner != sender).then_some(gas_data.owner),
1618            gas_price: gas_data.price,
1619            gas_budget: gas_data.budget,
1620            executed_epoch: epoch_id,
1621            dependencies: effects.dependencies().to_vec(),
1622            effects: SuiTransactionBlockEffects::V1(effects),
1623            receiving_objs,
1624            config_objects,
1625            // Find the protocol version for this epoch
1626            // This assumes we already initialized the protocol version table `protocol_version_epoch_table`
1627            protocol_version: self.get_protocol_config(epoch_id, chain).await?.version,
1628            tx_digest: *tx_digest,
1629            epoch_start_timestamp,
1630            sender_signed_data: orig_tx.clone(),
1631            reference_gas_price,
1632            chain,
1633        })
1634    }
1635
1636    async fn resolve_tx_components_from_dump(
1637        &self,
1638        tx_digest: &TransactionDigest,
1639    ) -> Result<OnChainTransactionInfo, ReplayEngineError> {
1640        assert!(!self.is_remote_replay());
1641
1642        let dp = self.fetcher.as_node_state_dump();
1643
1644        let sender = dp
1645            .node_state_dump
1646            .sender_signed_data
1647            .transaction_data()
1648            .sender();
1649        let orig_tx = dp.node_state_dump.sender_signed_data.clone();
1650        let effects = dp.node_state_dump.computed_effects.clone();
1651        let effects = SuiTransactionBlockEffects::try_from(effects).unwrap();
1652        // Config objects don't show up in the node state dump so they need to be provided.
1653        let config_objects = self.add_config_objects_if_needed(effects.status());
1654
1655        // Fetch full transaction content
1656        //let tx_info = self.fetcher.get_transaction(tx_digest).await?;
1657
1658        let input_objs = orig_tx
1659            .transaction_data()
1660            .input_objects()
1661            .map_err(|e| ReplayEngineError::UserInputError { err: e })?;
1662        let tx_kind_orig = orig_tx.transaction_data().kind();
1663
1664        // Download the objects at the version right before the execution of this TX
1665        let modified_at_versions: Vec<(ObjectID, SequenceNumber)> = effects.modified_at_versions();
1666
1667        let shared_object_refs: Vec<ObjectRef> = effects
1668            .shared_objects()
1669            .iter()
1670            .map(|so_ref| {
1671                if so_ref.digest == ObjectDigest::OBJECT_DIGEST_DELETED {
1672                    unimplemented!(
1673                        "Replay of deleted shared object transactions is not supported yet"
1674                    );
1675                } else {
1676                    so_ref.to_object_ref()
1677                }
1678            })
1679            .collect();
1680        let receiving_objs = orig_tx
1681            .transaction_data()
1682            .receiving_objects()
1683            .into_iter()
1684            .map(|(obj_id, version, _)| (obj_id, version))
1685            .collect();
1686
1687        let epoch_id = dp.node_state_dump.executed_epoch;
1688
1689        let chain = chain_from_chain_id(self.fetcher.get_chain_id().await?.as_str());
1690
1691        let protocol_config =
1692            ProtocolConfig::get_for_version(dp.node_state_dump.protocol_version.into(), chain);
1693        // Extract the epoch start timestamp
1694        let (epoch_start_timestamp, reference_gas_price) = self
1695            .get_epoch_start_timestamp_and_rgp(epoch_id, tx_digest)
1696            .await?;
1697        let gas_data = orig_tx.transaction_data().gas_data();
1698        let gas_object_refs: Vec<_> = gas_data.clone().payment.into_iter().collect();
1699
1700        Ok(OnChainTransactionInfo {
1701            kind: tx_kind_orig.clone(),
1702            sender,
1703            modified_at_versions,
1704            input_objects: input_objs,
1705            shared_object_refs,
1706            gas: gas_object_refs,
1707            gas_owner: (gas_data.owner != sender).then_some(gas_data.owner),
1708            gas_price: gas_data.price,
1709            gas_budget: gas_data.budget,
1710            executed_epoch: epoch_id,
1711            dependencies: effects.dependencies().to_vec(),
1712            effects,
1713            receiving_objs,
1714            config_objects,
1715            protocol_version: protocol_config.version,
1716            tx_digest: *tx_digest,
1717            epoch_start_timestamp,
1718            sender_signed_data: orig_tx.clone(),
1719            reference_gas_price,
1720            chain,
1721        })
1722    }
1723
1724    async fn resolve_download_input_objects(
1725        &mut self,
1726        tx_info: &OnChainTransactionInfo,
1727        deleted_shared_objects: Vec<ObjectRef>,
1728    ) -> Result<InputObjects, ReplayEngineError> {
1729        // Download the input objects
1730        let mut package_inputs = vec![];
1731        let mut imm_owned_inputs = vec![];
1732        let mut shared_inputs = vec![];
1733        let mut deleted_shared_info_map = BTreeMap::new();
1734
1735        // for deleted shared objects, we need to look at the transaction dependencies to find the
1736        // correct transaction dependency for a deleted shared object.
1737        if !deleted_shared_objects.is_empty() {
1738            for tx_digest in tx_info.dependencies.iter() {
1739                let tx_info = self.resolve_tx_components(tx_digest).await?;
1740                for (obj_id, version, _) in tx_info.shared_object_refs.iter() {
1741                    deleted_shared_info_map.insert(*obj_id, (tx_info.tx_digest, *version));
1742                }
1743            }
1744        }
1745
1746        tx_info
1747            .input_objects
1748            .iter()
1749            .map(|kind| match kind {
1750                InputObjectKind::MovePackage(i) => {
1751                    package_inputs.push(*i);
1752                    Ok(())
1753                }
1754                InputObjectKind::ImmOrOwnedMoveObject(o_ref) => {
1755                    imm_owned_inputs.push((o_ref.0, o_ref.1));
1756                    Ok(())
1757                }
1758                InputObjectKind::SharedMoveObject {
1759                    id,
1760                    initial_shared_version: _,
1761                    mutability: _,
1762                } if !deleted_shared_info_map.contains_key(id) => {
1763                    // We already downloaded
1764                    if let Some(o) = self
1765                        .storage
1766                        .live_objects_store
1767                        .lock()
1768                        .expect("Can't lock")
1769                        .get(id)
1770                    {
1771                        shared_inputs.push(o.clone());
1772                        Ok(())
1773                    } else {
1774                        Err(ReplayEngineError::InternalCacheInvariantViolation {
1775                            id: *id,
1776                            version: None,
1777                        })
1778                    }
1779                }
1780                _ => Ok(()),
1781            })
1782            .collect::<Result<Vec<_>, _>>()?;
1783
1784        // Download the imm and owned objects
1785        let mut in_objs = self.multi_download_and_store(&imm_owned_inputs).await?;
1786
1787        // For packages, download latest if non framework
1788        // If framework, download relevant for the current protocol version
1789        in_objs.extend(
1790            self.multi_download_relevant_packages_and_store(
1791                package_inputs,
1792                tx_info.protocol_version.as_u64(),
1793            )
1794            .await?,
1795        );
1796        // Add shared objects
1797        in_objs.extend(shared_inputs);
1798
1799        // TODO(Zhe): Account for cancelled transaction assigned version here, and tests.
1800        let resolved_input_objs = tx_info
1801            .input_objects
1802            .iter()
1803            .flat_map(|kind| match kind {
1804                InputObjectKind::MovePackage(i) => {
1805                    // Okay to unwrap since we downloaded it
1806                    Some(ObjectReadResult::new(
1807                        *kind,
1808                        self.storage
1809                            .package_cache
1810                            .lock()
1811                            .expect("Cannot lock")
1812                            .get(i)
1813                            .unwrap_or(
1814                                &self
1815                                    .download_latest_object(i)
1816                                    .expect("Object download failed")
1817                                    .expect("Object not found on chain"),
1818                            )
1819                            .clone()
1820                            .into(),
1821                    ))
1822                }
1823                InputObjectKind::ImmOrOwnedMoveObject(o_ref) => Some(ObjectReadResult::new(
1824                    *kind,
1825                    self.storage
1826                        .object_version_cache
1827                        .lock()
1828                        .expect("Cannot lock")
1829                        .get(&(o_ref.0, o_ref.1))
1830                        .unwrap()
1831                        .clone()
1832                        .into(),
1833                )),
1834                InputObjectKind::SharedMoveObject { id, .. }
1835                    if !deleted_shared_info_map.contains_key(id) =>
1836                {
1837                    // we already downloaded
1838                    Some(ObjectReadResult::new(
1839                        *kind,
1840                        self.storage
1841                            .live_objects_store
1842                            .lock()
1843                            .expect("Can't lock")
1844                            .get(id)
1845                            .unwrap()
1846                            .clone()
1847                            .into(),
1848                    ))
1849                }
1850                InputObjectKind::SharedMoveObject { id, .. } => {
1851                    let (digest, version) = deleted_shared_info_map.get(id).unwrap();
1852                    Some(ObjectReadResult::new(
1853                        *kind,
1854                        ObjectReadResultKind::ObjectConsensusStreamEnded(*version, *digest),
1855                    ))
1856                }
1857            })
1858            .collect();
1859
1860        Ok(InputObjects::new(resolved_input_objs))
1861    }
1862
1863    /// Given the OnChainTransactionInfo, download and store the input objects, and other info necessary
1864    /// for execution
1865    async fn initialize_execution_env_state(
1866        &mut self,
1867        tx_info: &OnChainTransactionInfo,
1868    ) -> Result<InputObjects, ReplayEngineError> {
1869        // We need this for other activities in this session
1870        self.current_protocol_version = tx_info.protocol_version.as_u64();
1871
1872        // Download the objects at the version right before the execution of this TX
1873        self.multi_download_and_store(&tx_info.modified_at_versions)
1874            .await?;
1875
1876        let (shared_refs, deleted_shared_refs): (Vec<ObjectRef>, Vec<ObjectRef>) = tx_info
1877            .shared_object_refs
1878            .iter()
1879            .partition(|r| r.2 != ObjectDigest::OBJECT_DIGEST_DELETED);
1880
1881        // Download shared objects at the version right before the execution of this TX
1882        let shared_refs: Vec<_> = shared_refs.iter().map(|r| (r.0, r.1)).collect();
1883        self.multi_download_and_store(&shared_refs).await?;
1884
1885        // Download gas (although this should already be in cache from modified at versions?)
1886        let gas_refs: Vec<_> = tx_info
1887            .gas
1888            .iter()
1889            .filter_map(|w| (w.0 != ObjectID::ZERO).then_some((w.0, w.1)))
1890            .collect();
1891        self.multi_download_and_store(&gas_refs).await?;
1892
1893        // Fetch the input objects we know from the raw transaction
1894        let input_objs = self
1895            .resolve_download_input_objects(tx_info, deleted_shared_refs)
1896            .await?;
1897
1898        // Fetch the receiving objects
1899        self.multi_download_and_store(&tx_info.receiving_objs)
1900            .await?;
1901
1902        // Fetch specified config objects if any
1903        self.multi_download_and_store(&tx_info.config_objects)
1904            .await?;
1905
1906        // Prep the object runtime for dynamic fields
1907        // Download the child objects accessed at the version right before the execution of this TX
1908        let loaded_child_refs = self.fetch_loaded_child_refs(&tx_info.tx_digest).await?;
1909        self.multi_download_and_store(&loaded_child_refs).await?;
1910        tokio::task::yield_now().await;
1911
1912        Ok(input_objs)
1913    }
1914}
1915
1916// <---------------------  Implement necessary traits for LocalExec to work with exec engine ----------------------->
1917
1918impl BackingPackageStore for LocalExec {
1919    /// In this case we might need to download a dependency package which was not present in the
1920    /// modified at versions list because packages are immutable
1921    fn get_package_object(&self, package_id: &ObjectID) -> SuiResult<Option<PackageObject>> {
1922        fn inner(self_: &LocalExec, package_id: &ObjectID) -> SuiResult<Option<Object>> {
1923            // If package not present fetch it from the network
1924            self_
1925                .get_or_download_object(package_id, true /* we expect a Move package*/)
1926                .map_err(|e| SuiErrorKind::Storage(e.to_string()).into())
1927        }
1928
1929        let res = inner(self, package_id);
1930        self.exec_store_events
1931            .lock()
1932            .expect("Unable to lock events list")
1933            .push(ExecutionStoreEvent::BackingPackageGetPackageObject {
1934                package_id: *package_id,
1935                result: res.clone(),
1936            });
1937        res.map(|o| o.map(PackageObject::new))
1938    }
1939}
1940
1941impl RuntimeObjectResolver for LocalExec {
1942    /// This uses `get_object`, which does not download from the network
1943    /// Hence all objects must be in store already
1944    fn read_child_object(
1945        &self,
1946        parent: &ObjectID,
1947        child: &ObjectID,
1948        child_version_upper_bound: SequenceNumber,
1949    ) -> SuiResult<Option<Object>> {
1950        fn inner(
1951            self_: &LocalExec,
1952            parent: &ObjectID,
1953            child: &ObjectID,
1954            child_version_upper_bound: SequenceNumber,
1955        ) -> SuiResult<Option<Object>> {
1956            let child_object =
1957                match self_.download_object_by_upper_bound(child, child_version_upper_bound)? {
1958                    None => return Ok(None),
1959                    Some(o) => o,
1960                };
1961            let child_version = child_object.version();
1962            if child_object.version() > child_version_upper_bound {
1963                return Err(SuiErrorKind::Unknown(format!(
1964                    "Invariant Violation. Replay loaded child_object {child} at version \
1965                    {child_version} but expected the version to be <= {child_version_upper_bound}"
1966                ))
1967                .into());
1968            }
1969            let parent = *parent;
1970            if child_object.owner != Owner::ObjectOwner(parent.into()) {
1971                return Err(SuiErrorKind::InvalidChildObjectAccess {
1972                    object: *child,
1973                    given_parent: parent,
1974                    actual_owner: child_object.owner.clone(),
1975                }
1976                .into());
1977            }
1978            Ok(Some(child_object))
1979        }
1980
1981        let res = inner(self, parent, child, child_version_upper_bound);
1982        self.exec_store_events
1983            .lock()
1984            .expect("Unable to lock events list")
1985            .push(
1986                ExecutionStoreEvent::RuntimeObjectResolverStoreReadChildObject {
1987                    parent: *parent,
1988                    child: *child,
1989                    result: res.clone(),
1990                },
1991            );
1992        res
1993    }
1994
1995    fn get_object_received_at_version(
1996        &self,
1997        owner: &ObjectID,
1998        receiving_object_id: &ObjectID,
1999        receive_object_at_version: SequenceNumber,
2000        _epoch_id: EpochId,
2001    ) -> SuiResult<Option<Object>> {
2002        fn inner(
2003            self_: &LocalExec,
2004            owner: &ObjectID,
2005            receiving_object_id: &ObjectID,
2006            receive_object_at_version: SequenceNumber,
2007        ) -> SuiResult<Option<Object>> {
2008            let recv_object = match self_.get_object(receiving_object_id) {
2009                None => return Ok(None),
2010                Some(o) => o,
2011            };
2012            if recv_object.version() != receive_object_at_version {
2013                return Err(SuiErrorKind::Unknown(format!(
2014                    "Invariant Violation. Replay loaded child_object {receiving_object_id} at version \
2015                    {receive_object_at_version} but expected the version to be == {receive_object_at_version}"
2016                )).into());
2017            }
2018            if recv_object.owner != Owner::AddressOwner((*owner).into()) {
2019                return Ok(None);
2020            }
2021            Ok(Some(recv_object))
2022        }
2023
2024        let res = inner(self, owner, receiving_object_id, receive_object_at_version);
2025        self.exec_store_events
2026            .lock()
2027            .expect("Unable to lock events list")
2028            .push(ExecutionStoreEvent::ReceiveObject {
2029                owner: *owner,
2030                receive: *receiving_object_id,
2031                receive_at_version: receive_object_at_version,
2032                result: res.clone(),
2033            });
2034        res
2035    }
2036}
2037
2038impl ParentSync for LocalExec {
2039    /// The objects here much already exist in the store because we downloaded them earlier
2040    /// No download from network
2041    fn get_latest_parent_entry_ref_deprecated(&self, object_id: ObjectID) -> Option<ObjectRef> {
2042        fn inner(self_: &LocalExec, object_id: ObjectID) -> Option<ObjectRef> {
2043            if let Some(v) = self_
2044                .storage
2045                .live_objects_store
2046                .lock()
2047                .expect("Can't lock")
2048                .get(&object_id)
2049            {
2050                return Some(v.compute_object_reference());
2051            }
2052            None
2053        }
2054        let res = inner(self, object_id);
2055        self.exec_store_events
2056            .lock()
2057            .expect("Unable to lock events list")
2058            .push(
2059                ExecutionStoreEvent::ParentSyncStoreGetLatestParentEntryRef {
2060                    object_id,
2061                    result: res,
2062                },
2063            );
2064        res
2065    }
2066}
2067
2068impl ModuleResolver for LocalExec {
2069    type Error = SuiError;
2070
2071    /// This fetches a module which must already be present in the store
2072    /// We do not download
2073    fn get_module(&self, module_id: &ModuleId) -> SuiResult<Option<Vec<u8>>> {
2074        fn inner(self_: &LocalExec, module_id: &ModuleId) -> SuiResult<Option<Vec<u8>>> {
2075            get_module(self_, module_id)
2076        }
2077
2078        let res = inner(self, module_id);
2079        self.exec_store_events
2080            .lock()
2081            .expect("Unable to lock events list")
2082            .push(ExecutionStoreEvent::ModuleResolverGetModule {
2083                module_id: module_id.clone(),
2084                result: res.clone(),
2085            });
2086        res
2087    }
2088
2089    fn get_packages_static<const N: usize>(
2090        &self,
2091        ids: [AccountAddress; N],
2092    ) -> Result<[Option<SerializedPackage>; N], Self::Error> {
2093        let mut res = [const { None }; N];
2094        for i in 0..N {
2095            res[i] = get_package(self, &ids[i].into())?;
2096        }
2097        Ok(res)
2098    }
2099
2100    fn get_packages<'a>(
2101        &self,
2102        ids: impl ExactSizeIterator<Item = &'a AccountAddress>,
2103    ) -> Result<Vec<Option<SerializedPackage>>, Self::Error> {
2104        ids.map(|id| get_package(self, &(*id).into())).collect()
2105    }
2106}
2107
2108impl ModuleResolver for &mut LocalExec {
2109    type Error = SuiError;
2110
2111    fn get_module(&self, module_id: &ModuleId) -> SuiResult<Option<Vec<u8>>> {
2112        // Recording event here will be double-counting since its already recorded in the get_module fn
2113        (**self).get_module(module_id)
2114    }
2115
2116    fn get_packages<'a>(
2117        &self,
2118        ids: impl ExactSizeIterator<Item = &'a AccountAddress>,
2119    ) -> Result<Vec<Option<SerializedPackage>>, Self::Error> {
2120        (**self).get_packages(ids)
2121    }
2122
2123    fn get_packages_static<const N: usize>(
2124        &self,
2125        ids: [AccountAddress; N],
2126    ) -> Result<[Option<SerializedPackage>; N], Self::Error> {
2127        (**self).get_packages_static(ids)
2128    }
2129}
2130
2131impl ObjectStore for LocalExec {
2132    /// The object must be present in store by normal process we used to backfill store in init
2133    /// We dont download if not present
2134    fn get_object(&self, object_id: &ObjectID) -> Option<Object> {
2135        let res = self
2136            .storage
2137            .live_objects_store
2138            .lock()
2139            .expect("Can't lock")
2140            .get(object_id)
2141            .cloned();
2142        self.exec_store_events
2143            .lock()
2144            .expect("Unable to lock events list")
2145            .push(ExecutionStoreEvent::ObjectStoreGetObject {
2146                object_id: *object_id,
2147                result: Ok(res.clone()),
2148            });
2149        res
2150    }
2151
2152    /// The object must be present in store by normal process we used to backfill store in init
2153    /// We dont download if not present
2154    fn get_object_by_key(&self, object_id: &ObjectID, version: VersionNumber) -> Option<Object> {
2155        let res = self
2156            .storage
2157            .live_objects_store
2158            .lock()
2159            .expect("Can't lock")
2160            .get(object_id)
2161            .and_then(|obj| {
2162                if obj.version() == version {
2163                    Some(obj.clone())
2164                } else {
2165                    None
2166                }
2167            });
2168
2169        self.exec_store_events
2170            .lock()
2171            .expect("Unable to lock events list")
2172            .push(ExecutionStoreEvent::ObjectStoreGetObjectByKey {
2173                object_id: *object_id,
2174                version,
2175                result: Ok(res.clone()),
2176            });
2177
2178        res
2179    }
2180}
2181
2182impl ObjectStore for &mut LocalExec {
2183    fn get_object(&self, object_id: &ObjectID) -> Option<Object> {
2184        // Recording event here will be double-counting since its already recorded in the get_module fn
2185        (**self).get_object(object_id)
2186    }
2187
2188    fn get_object_by_key(&self, object_id: &ObjectID, version: VersionNumber) -> Option<Object> {
2189        // Recording event here will be double-counting since its already recorded in the get_module fn
2190        (**self).get_object_by_key(object_id, version)
2191    }
2192}
2193
2194impl GetModule for LocalExec {
2195    type Error = SuiError;
2196    type Item = CompiledModule;
2197
2198    fn get_module_by_id(&self, id: &ModuleId) -> SuiResult<Option<Self::Item>> {
2199        let res = get_module_by_id(self, id);
2200
2201        self.exec_store_events
2202            .lock()
2203            .expect("Unable to lock events list")
2204            .push(ExecutionStoreEvent::GetModuleGetModuleByModuleId {
2205                id: id.clone(),
2206                result: res.clone(),
2207            });
2208        res
2209    }
2210}
2211
2212// <--------------------- Util functions ----------------------->
2213
2214pub fn get_executor(
2215    executor_version_override: Option<i64>,
2216    protocol_config: &ProtocolConfig,
2217    _expensive_safety_check_config: ExpensiveSafetyCheckConfig,
2218) -> Arc<dyn Executor + Send + Sync> {
2219    let protocol_config = executor_version_override
2220        .map(|q| {
2221            let ver = if q < 0 {
2222                ProtocolConfig::get_for_max_version_UNSAFE().execution_version()
2223            } else {
2224                q as u64
2225            };
2226
2227            let mut c = protocol_config.clone();
2228            c.set_execution_version_for_testing(ver);
2229            c
2230        })
2231        .unwrap_or(protocol_config.clone());
2232
2233    let silent = true;
2234    sui_execution::executor(&protocol_config, silent)
2235        .expect("Creating an executor should not fail here")
2236}
2237
2238fn parse_effect_error_for_denied_coins(status: &SuiExecutionStatus) -> Option<String> {
2239    let SuiExecutionStatus::Failure { error } = status else {
2240        return None;
2241    };
2242    parse_denied_error_string(error)
2243}
2244
2245fn parse_denied_error_string(error: &str) -> Option<String> {
2246    let regulated_regex = regex::Regex::new(
2247        r#"CoinTypeGlobalPause.*?"(.*?)"|AddressDeniedForCoin.*coin_type:.*?"(.*?)""#,
2248    )
2249    .unwrap();
2250
2251    let caps = regulated_regex.captures(error)?;
2252    Some(caps.get(1).or(caps.get(2))?.as_str().to_string())
2253}
2254
2255#[cfg(test)]
2256mod tests {
2257    use super::parse_denied_error_string;
2258    #[test]
2259    fn test_regex_regulated_coin_errors() {
2260        let test_bank = vec![
2261            "CoinTypeGlobalPause { coin_type: \"39a572c071784c280ee8ee8c683477e059d1381abc4366f9a58ffac3f350a254::rcoin::RCOIN\" }",
2262            "AddressDeniedForCoin { address: B, coin_type: \"39a572c071784c280ee8ee8c683477e059d1381abc4366f9a58ffac3f350a254::rcoin::RCOIN\" }",
2263        ];
2264        let expected_string =
2265            "39a572c071784c280ee8ee8c683477e059d1381abc4366f9a58ffac3f350a254::rcoin::RCOIN";
2266
2267        for test in &test_bank {
2268            assert!(parse_denied_error_string(test).unwrap() == expected_string);
2269        }
2270    }
2271}