1use 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#[derive(Debug, Serialize, Deserialize)]
73pub struct ExecutionSandboxState {
74 pub transaction_info: OnChainTransactionInfo,
76 pub required_objects: Vec<Object>,
78 #[serde(skip)]
80 pub local_exec_temporary_store: Option<InnerTemporaryStore>,
81 pub local_exec_effects: SuiTransactionBlockEffects,
83 #[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 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 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 pub protocol_version: u64,
144 pub epoch_start: u64,
146 pub epoch_end: u64,
148 pub checkpoint_start: Option<u64>,
150 pub checkpoint_end: Option<u64>,
152 pub epoch_change_tx: TransactionDigest,
154}
155
156#[derive(Clone)]
157pub struct Storage {
158 pub live_objects_store: Arc<Mutex<BTreeMap<ObjectID, Object>>>,
163
164 pub package_cache: Arc<Mutex<BTreeMap<ObjectID, Object>>>,
167 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 pub protocol_version_epoch_table: BTreeMap<u64, ProtocolVersionSummary>,
240 pub protocol_version_system_package_table: BTreeMap<u64, BTreeMap<ObjectID, SequenceNumber>>,
242 pub current_protocol_version: u64,
244 pub storage: Storage,
246 pub exec_store_events: Arc<Mutex<Vec<ExecutionStoreEvent>>>,
248 pub metrics: Arc<ExecutionMetrics>,
250 pub fetcher: Fetchers,
252
253 pub executor_version: Option<i64>,
256 pub protocol_version: Option<i64>,
260 pub config_and_versions: Option<Vec<(ObjectID, SequenceNumber)>>,
261 pub num_retries_for_timeout: u32,
263 pub sleep_period_for_timeout: std::time::Duration,
264}
265
266impl LocalExec {
267 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 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 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 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 let registry = prometheus::Registry::new();
387 let metrics = Arc::new(ExecutionMetrics::new(®istry));
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 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 let registry = prometheus::Registry::new();
415 let metrics = Arc::new(ExecutionMetrics::new(®istry));
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 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 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 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 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 #[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 #[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 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 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 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 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 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 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 let executor = get_executor(ov, protocol_config, expensive_safety_check_config);
762
763 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 &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, 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 &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, 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 #[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 #[allow(clippy::result_large_err)]
946 pub async fn certificate_execute_with_sandbox_state(
947 pre_run_sandbox: &ExecutionSandboxState,
948 ) -> Result<ExecutionSandboxState, ReplayEngineError> {
949 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 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 &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, 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, local_exec_effects: effects,
1020 local_exec_status: Some(exec_res),
1021 })
1022 }
1023
1024 #[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 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 #[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 #[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 } 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 #[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 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 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 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 continue;
1226 }
1227
1228 curr_checkpoint = self
1231 .fetcher
1232 .get_transaction(&event.id.tx_digest)
1233 .await?
1234 .checkpoint;
1235 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 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 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 for (
1298 prot_ver,
1299 ProtocolVersionSummary {
1300 epoch_change_tx: tx_digest,
1301 ..
1302 },
1303 ) in self.protocol_version_epoch_table.clone()
1304 {
1305 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 for ver in versions.iter().rev() {
1321 if ver.1 == tx_digest {
1322 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 while !system_package_objs.is_empty() {
1352 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 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 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 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 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 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 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 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 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 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 let config_objects = self.add_config_objects_if_needed(effects.status());
1654
1655 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 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 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 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 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 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 let mut in_objs = self.multi_download_and_store(&imm_owned_inputs).await?;
1786
1787 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 in_objs.extend(shared_inputs);
1798
1799 let resolved_input_objs = tx_info
1801 .input_objects
1802 .iter()
1803 .flat_map(|kind| match kind {
1804 InputObjectKind::MovePackage(i) => {
1805 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 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 async fn initialize_execution_env_state(
1866 &mut self,
1867 tx_info: &OnChainTransactionInfo,
1868 ) -> Result<InputObjects, ReplayEngineError> {
1869 self.current_protocol_version = tx_info.protocol_version.as_u64();
1871
1872 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 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 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 let input_objs = self
1895 .resolve_download_input_objects(tx_info, deleted_shared_refs)
1896 .await?;
1897
1898 self.multi_download_and_store(&tx_info.receiving_objs)
1900 .await?;
1901
1902 self.multi_download_and_store(&tx_info.config_objects)
1904 .await?;
1905
1906 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
1916impl BackingPackageStore for LocalExec {
1919 fn get_package_object(&self, package_id: &ObjectID) -> SuiResult<Option<PackageObject>> {
1922 fn inner(self_: &LocalExec, package_id: &ObjectID) -> SuiResult<Option<Object>> {
1923 self_
1925 .get_or_download_object(package_id, true )
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 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 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 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 (**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 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 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 (**self).get_object(object_id)
2186 }
2187
2188 fn get_object_by_key(&self, object_id: &ObjectID, version: VersionNumber) -> Option<Object> {
2189 (**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
2212pub 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}