Skip to main content

sui_rpc_api/client/
mod.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4use bytes::Bytes;
5use fastcrypto::traits::ToFromBytes;
6use futures::stream::Stream;
7use futures::stream::TryStreamExt;
8use prost_types::FieldMask;
9use prost_types::value::Kind as ProtoValueKind;
10use std::time::Duration;
11use sui_rpc::field::FieldMaskUtil;
12use sui_rpc::proto::TryFromProtoError;
13use sui_rpc::proto::sui::rpc::v2::{self as proto, GetServiceInfoRequest};
14use sui_types::base_types::{ObjectID, SequenceNumber, SuiAddress};
15use sui_types::digests::ChainIdentifier;
16use sui_types::digests::TransactionDigest;
17use sui_types::effects::{TransactionEffects, TransactionEvents};
18use sui_types::full_checkpoint_content::Checkpoint;
19use sui_types::messages_checkpoint::{CertifiedCheckpointSummary, CheckpointSequenceNumber};
20use sui_types::object::Object;
21use sui_types::signature::GenericSignature;
22use sui_types::transaction::Transaction;
23use sui_types::transaction::TransactionData;
24use tap::Pipe;
25use tonic::Status;
26use tonic::metadata::MetadataMap;
27
28use sui_protocol_config::ProtocolVersion;
29pub use sui_rpc::client::HeadersInterceptor;
30pub use sui_rpc::client::ResponseExt;
31
32pub type Result<T, E = tonic::Status> = std::result::Result<T, E>;
33pub type BoxError = Box<dyn std::error::Error + Send + Sync + 'static>;
34
35pub struct Page<T> {
36    pub items: Vec<T>,
37    pub next_page_token: Option<Bytes>,
38}
39
40#[derive(Clone)]
41pub struct Client(sui_rpc::Client);
42
43impl Client {
44    pub fn new<T>(uri: T) -> Result<Self>
45    where
46        T: TryInto<http::Uri>,
47        T::Error: Into<BoxError>,
48    {
49        sui_rpc::Client::new(uri)
50            .map(Self)
51            .map(|client| client.with_headers(HeadersInterceptor::new()))
52    }
53
54    /// Sets the headers sent with every request. `x-sui-client-protocol-version` is always
55    /// included, since this client can decode everything up to this binary's max protocol version.
56    pub fn with_headers(self, mut headers: HeadersInterceptor) -> Self {
57        headers
58            .headers_mut()
59            .entry(crate::X_SUI_CLIENT_PROTOCOL_VERSION)
60            .expect("valid header name")
61            .or_insert(ProtocolVersion::MAX_ALLOWED.as_u64().into());
62        Self(self.0.with_headers(headers))
63    }
64
65    pub fn inner_mut(&mut self) -> &mut sui_rpc::Client {
66        &mut self.0
67    }
68
69    pub fn into_inner(self) -> sui_rpc::Client {
70        self.0
71    }
72
73    pub async fn get_latest_checkpoint(&mut self) -> Result<CertifiedCheckpointSummary> {
74        self.get_checkpoint_internal(None).await
75    }
76
77    pub async fn get_checkpoint_summary(
78        &mut self,
79        sequence_number: CheckpointSequenceNumber,
80    ) -> Result<CertifiedCheckpointSummary> {
81        self.get_checkpoint_internal(Some(sequence_number)).await
82    }
83
84    async fn get_checkpoint_internal(
85        &mut self,
86        sequence_number: Option<CheckpointSequenceNumber>,
87    ) -> Result<CertifiedCheckpointSummary> {
88        let mut request = proto::GetCheckpointRequest::default()
89            .with_read_mask(FieldMask::from_paths(["summary.bcs", "signature"]));
90        request.checkpoint_id = sequence_number.map(|sequence_number| {
91            proto::get_checkpoint_request::CheckpointId::SequenceNumber(sequence_number)
92        });
93
94        let (metadata, checkpoint, _extentions) = self
95            .0
96            .ledger_client()
97            .get_checkpoint(request)
98            .await?
99            .into_parts();
100
101        let checkpoint = checkpoint
102            .checkpoint
103            .ok_or_else(|| tonic::Status::not_found("no checkpoint returned"))?;
104        certified_checkpoint_summary_try_from_proto(&checkpoint)
105            .map_err(|e| status_from_error_with_metadata(e, metadata))
106    }
107
108    pub async fn get_full_checkpoint(
109        &mut self,
110        sequence_number: CheckpointSequenceNumber,
111    ) -> Result<Checkpoint> {
112        let request = proto::GetCheckpointRequest::by_sequence_number(sequence_number)
113            .with_read_mask(Checkpoint::proto_field_mask());
114
115        let (metadata, response, _extentions) = self
116            .0
117            .ledger_client()
118            .max_decoding_message_size(128 * 1024 * 1024)
119            .get_checkpoint(request)
120            .await?
121            .into_parts();
122
123        let checkpoint = response
124            .checkpoint
125            .ok_or_else(|| tonic::Status::not_found("no checkpoint returned"))?;
126        sui_types::full_checkpoint_content::Checkpoint::try_from(&checkpoint)
127            .map_err(|e| status_from_error_with_metadata(e, metadata))
128    }
129
130    pub async fn get_object(&mut self, object_id: ObjectID) -> Result<Object> {
131        self.get_object_internal(object_id, None).await
132    }
133
134    pub async fn get_object_with_version(
135        &mut self,
136        object_id: ObjectID,
137        version: SequenceNumber,
138    ) -> Result<Object> {
139        self.get_object_internal(object_id, Some(version.value()))
140            .await
141    }
142
143    async fn get_object_internal(
144        &mut self,
145        object_id: ObjectID,
146        version: Option<u64>,
147    ) -> Result<Object> {
148        let (object, _) = self
149            .get_object_with_content(object_id, version, false)
150            .await?;
151        Ok(object)
152    }
153
154    pub async fn get_object_with_json(
155        &mut self,
156        object_id: ObjectID,
157    ) -> Result<(Object, Option<serde_json::Value>)> {
158        self.get_object_with_content(object_id, None, true).await
159    }
160
161    async fn get_object_with_content(
162        &mut self,
163        object_id: ObjectID,
164        version: Option<u64>,
165        include_json: bool,
166    ) -> Result<(Object, Option<serde_json::Value>)> {
167        let paths: Vec<&str> = if include_json {
168            vec!["bcs", "json"]
169        } else {
170            vec!["bcs"]
171        };
172        let mut request = proto::GetObjectRequest::new(&object_id.into())
173            .with_read_mask(FieldMask::from_paths(paths));
174        request.version = version;
175
176        let (metadata, response, _extentions) = self
177            .0
178            .ledger_client()
179            .get_object(request)
180            .await?
181            .into_parts();
182
183        let proto_object = response
184            .object
185            .ok_or_else(|| tonic::Status::not_found("no object returned"))?;
186        let json_content = proto_object
187            .json
188            .as_ref()
189            .map(|v| proto_value_to_json_value(v));
190        let object = object_try_from_proto(&proto_object)
191            .map_err(|e| status_from_error_with_metadata(e, metadata))?;
192        Ok((object, json_content))
193    }
194
195    pub async fn batch_get_objects(&self, ids: &[ObjectID]) -> Result<Vec<Object>> {
196        let request = proto::BatchGetObjectsRequest::default()
197            .with_requests(
198                ids.iter()
199                    .map(|id| proto::GetObjectRequest::new(&(*id).into()))
200                    .collect(),
201            )
202            .with_read_mask(FieldMask::from_paths(["bcs"]));
203
204        let (metadata, response, _extentions) = self
205            .0
206            .clone()
207            .ledger_client()
208            .batch_get_objects(request)
209            .await?
210            .into_parts();
211
212        let objects = response
213            .objects
214            .into_iter()
215            .map(|o| o.to_result())
216            .collect::<Result<Vec<_>, _>>()
217            .map_err(|e| Status::not_found(e.message))?;
218
219        let objects = objects
220            .iter()
221            .map(object_try_from_proto)
222            .collect::<Result<_, _>>()
223            .map_err(|e| status_from_error_with_metadata(e, metadata))?;
224        Ok(objects)
225    }
226
227    pub async fn execute_transaction(
228        &mut self,
229        transaction: &Transaction,
230    ) -> Result<ExecutedTransaction> {
231        let request = Self::create_executed_transaction_request(transaction)?;
232
233        let (metadata, response, _extentions) = self
234            .0
235            .execution_client()
236            .execute_transaction(request)
237            .await?
238            .into_parts();
239
240        execute_transaction_response_try_from_proto(&response)
241            .map_err(|e| status_from_error_with_metadata(e, metadata))
242    }
243
244    pub async fn execute_transaction_and_wait_for_checkpoint(
245        &self,
246        transaction: &Transaction,
247    ) -> Result<ExecutedTransaction> {
248        const WAIT_FOR_CHECKPOINT_TIMEOUT: Duration = Duration::from_secs(30);
249
250        let request = Self::create_executed_transaction_request(transaction)?;
251
252        let (metadata, response, _extentions) = self
253            .0
254            .clone()
255            .execute_transaction_and_wait_for_checkpoint(request, WAIT_FOR_CHECKPOINT_TIMEOUT)
256            .await
257            .map_err(|e| Status::from_error(e.into()))?
258            .into_parts();
259
260        execute_transaction_response_try_from_proto(&response)
261            .map_err(|e| status_from_error_with_metadata(e, metadata))
262    }
263
264    fn create_executed_transaction_request(
265        transaction: &Transaction,
266    ) -> Result<proto::ExecuteTransactionRequest> {
267        let signatures = transaction
268            .inner()
269            .tx_signatures
270            .iter()
271            .map(|signature| {
272                let mut message = proto::UserSignature::default();
273                message.bcs = Some(signature.as_ref().to_vec().into());
274                message
275            })
276            .collect();
277
278        let request = proto::ExecuteTransactionRequest::new({
279            let mut tx = proto::Transaction::default();
280            tx.bcs = Some(
281                proto::Bcs::serialize(&transaction.inner().intent_message.value)
282                    .map_err(|e| Status::from_error(e.into()))?,
283            );
284            tx
285        })
286        .with_signatures(signatures)
287        .with_read_mask(ExecutedTransaction::proto_read_mask());
288
289        Ok(request)
290    }
291
292    pub async fn simulate_transaction(
293        &self,
294        tx: &TransactionData,
295        checks: bool,
296        do_gas_selection: bool,
297    ) -> Result<SimulateTransactionResponse> {
298        let mut request = proto::SimulateTransactionRequest::default();
299        request.set_checks(if checks {
300            proto::simulate_transaction_request::TransactionChecks::Enabled
301        } else {
302            proto::simulate_transaction_request::TransactionChecks::Disabled
303        });
304        request.set_do_gas_selection(do_gas_selection);
305        request.set_transaction(
306            proto::Transaction::default()
307                .with_bcs(proto::Bcs::serialize(&tx).map_err(|e| Status::from_error(e.into()))?),
308        );
309
310        let (metadata, response, _extentions) = self
311            .0
312            .clone()
313            .execution_client()
314            .simulate_transaction(request)
315            .await?
316            .into_parts();
317
318        let transaction = executed_transaction_try_from_proto(response.transaction())
319            .map_err(|e| status_from_error_with_metadata(e, metadata))?;
320
321        Ok(SimulateTransactionResponse {
322            transaction,
323            command_outputs: response.command_outputs,
324            suggested_gas_price: response.suggested_gas_price,
325        })
326    }
327
328    pub async fn get_transaction(
329        &mut self,
330        digest: &TransactionDigest,
331    ) -> Result<ExecutedTransaction> {
332        let request = proto::GetTransactionRequest::new(&(*digest).into())
333            .with_read_mask(ExecutedTransaction::proto_read_mask());
334
335        let (metadata, resp, _extentions) = self
336            .0
337            .ledger_client()
338            .get_transaction(request)
339            .await?
340            .into_parts();
341
342        let transaction = resp
343            .transaction
344            .ok_or_else(|| tonic::Status::not_found("no transaction returned"))?;
345        executed_transaction_try_from_proto(&transaction)
346            .map_err(|e| status_from_error_with_metadata(e, metadata))
347    }
348
349    pub async fn get_chain_identifier(&self) -> Result<ChainIdentifier> {
350        let response = self
351            .0
352            .clone()
353            .ledger_client()
354            .get_service_info(GetServiceInfoRequest::default())
355            .await?
356            .into_inner();
357        let chain_id = response
358            .chain_id()
359            .parse::<sui_sdk_types::Digest>()
360            .map_err(|e| TryFromProtoError::invalid("chain_id", e))
361            .map_err(|e| Status::from_error(e.into()))?;
362
363        Ok(ChainIdentifier::from(
364            sui_types::digests::CheckpointDigest::from(chain_id),
365        ))
366    }
367
368    pub async fn get_owned_objects(
369        &self,
370        owner: SuiAddress,
371        object_type: Option<move_core_types::language_storage::StructTag>,
372        page_size: Option<u32>,
373        page_token: Option<Bytes>,
374    ) -> Result<Page<Object>> {
375        let mut request = proto::ListOwnedObjectsRequest::default()
376            .with_owner(owner.to_string())
377            .with_read_mask(FieldMask::from_paths(["bcs"]));
378        if let Some(object_type) = object_type {
379            request.set_object_type(object_type.to_canonical_string(true));
380        }
381
382        if let Some(page_size) = page_size {
383            request.set_page_size(page_size);
384        }
385
386        if let Some(page_token) = page_token {
387            request.set_page_token(page_token);
388        }
389
390        let (metadata, response, _extentions) = self
391            .0
392            .clone()
393            .state_client()
394            .list_owned_objects(request)
395            .await?
396            .into_parts();
397
398        let objects = response
399            .objects()
400            .iter()
401            .map(object_try_from_proto)
402            .collect::<Result<_, _>>()
403            .map_err(|e| status_from_error_with_metadata(e, metadata))?;
404
405        Ok(Page {
406            items: objects,
407            next_page_token: response.next_page_token,
408        })
409    }
410
411    pub fn list_owned_objects(
412        &self,
413        owner: SuiAddress,
414        object_type: Option<move_core_types::language_storage::StructTag>,
415    ) -> impl Stream<Item = Result<Object>> + 'static {
416        let mut request = proto::ListOwnedObjectsRequest::default()
417            .with_owner(owner.to_string())
418            .with_read_mask(FieldMask::from_paths(["bcs"]));
419
420        if let Some(object_type) = object_type {
421            request.set_object_type(object_type.to_canonical_string(true));
422        }
423
424        self.0
425            .list_owned_objects(request)
426            .and_then(|object| async move {
427                object_try_from_proto(&object).map_err(|e| Status::from_error(e.into()))
428            })
429    }
430
431    pub async fn get_dynamic_fields(
432        &self,
433        parent: ObjectID,
434        page_size: Option<u32>,
435        page_token: Option<Bytes>,
436    ) -> Result<proto::ListDynamicFieldsResponse> {
437        let mut request = proto::ListDynamicFieldsRequest::default()
438            .with_parent(parent.to_string())
439            .with_read_mask(FieldMask::from_paths(["*"]));
440
441        if let Some(page_size) = page_size {
442            request.set_page_size(page_size);
443        }
444
445        if let Some(page_token) = page_token {
446            request.set_page_token(page_token);
447        }
448
449        let response = self
450            .0
451            .clone()
452            .state_client()
453            .list_dynamic_fields(request)
454            .await?
455            .into_inner();
456
457        Ok(response)
458    }
459
460    pub async fn get_reference_gas_price(&self) -> Result<u64> {
461        let request = proto::GetEpochRequest::default()
462            .with_read_mask(FieldMask::from_paths(["epoch", "reference_gas_price"]));
463
464        let response = self
465            .0
466            .clone()
467            .ledger_client()
468            .get_epoch(request)
469            .await?
470            .into_inner();
471
472        Ok(response.epoch().reference_gas_price())
473    }
474
475    pub async fn get_current_epoch(&self) -> Result<u64> {
476        let request =
477            proto::GetEpochRequest::default().with_read_mask(FieldMask::from_paths(["epoch"]));
478
479        let response = self
480            .0
481            .clone()
482            .ledger_client()
483            .get_epoch(request)
484            .await?
485            .into_inner();
486
487        Ok(response.epoch().epoch())
488    }
489
490    /// Wait for a transaction to be available in the ledger AND indexed (equivalent to WaitForLocalExecution)
491    pub async fn wait_for_transaction(
492        &self,
493        digest: &sui_types::digests::TransactionDigest,
494    ) -> Result<(), anyhow::Error> {
495        const WAIT_FOR_LOCAL_EXECUTION_TIMEOUT: Duration = Duration::from_secs(30);
496        const WAIT_FOR_LOCAL_EXECUTION_DELAY: Duration = Duration::from_millis(200);
497        const WAIT_FOR_LOCAL_EXECUTION_INTERVAL: Duration = Duration::from_millis(500);
498
499        let mut client = self.0.clone();
500        let mut client = client.ledger_client();
501
502        tokio::time::timeout(WAIT_FOR_LOCAL_EXECUTION_TIMEOUT, async {
503            // Apply a short delay to give the full node a chance to catch up.
504            tokio::time::sleep(WAIT_FOR_LOCAL_EXECUTION_DELAY).await;
505
506            let mut interval = tokio::time::interval(WAIT_FOR_LOCAL_EXECUTION_INTERVAL);
507            loop {
508                interval.tick().await;
509
510                let request = proto::GetTransactionRequest::default()
511                    .with_digest(digest.to_string())
512                    .with_read_mask(prost_types::FieldMask::from_paths(["digest", "checkpoint"]));
513
514                if let Ok(response) = client.get_transaction(request).await {
515                    let tx = response.into_inner().transaction;
516                    if let Some(executed_tx) = tx {
517                        // Check that transaction is indexed (checkpoint field is populated)
518                        if executed_tx.checkpoint.is_some() {
519                            break;
520                        }
521                    }
522                }
523            }
524        })
525        .await
526        .map_err(|_| anyhow::anyhow!("Timeout waiting for transaction indexing: {}", digest))?;
527
528        Ok(())
529    }
530
531    pub async fn get_protocol_config(&self, epoch: Option<u64>) -> Result<proto::ProtocolConfig> {
532        let mut request = proto::GetEpochRequest::default();
533        if let Some(epoch) = epoch {
534            request.set_epoch(epoch);
535        }
536        request.set_read_mask(FieldMask::from_paths([
537            proto::Epoch::path_builder().epoch(),
538            proto::Epoch::path_builder().protocol_config().finish(),
539        ]));
540        let mut response = self
541            .0
542            .clone()
543            .ledger_client()
544            .get_epoch(request)
545            .await?
546            .into_inner();
547
548        Ok(response
549            .epoch_mut()
550            .protocol_config
551            .take()
552            .unwrap_or_default())
553    }
554
555    pub async fn get_system_state(&self, epoch: Option<u64>) -> Result<Box<proto::SystemState>> {
556        let mut request = proto::GetEpochRequest::default();
557        if let Some(epoch) = epoch {
558            request.set_epoch(epoch);
559        }
560        request.set_read_mask(FieldMask::from_paths([
561            proto::Epoch::path_builder().epoch(),
562            proto::Epoch::path_builder().system_state().finish(),
563        ]));
564        let mut response = self
565            .0
566            .clone()
567            .ledger_client()
568            .get_epoch(request)
569            .await?
570            .into_inner();
571
572        Ok(response.epoch_mut().system_state.take().unwrap_or_default())
573    }
574
575    pub async fn get_system_state_summary(
576        &self,
577        epoch: Option<u64>,
578    ) -> Result<sui_types::sui_system_state::sui_system_state_summary::SuiSystemStateSummary> {
579        let system_state = self.get_system_state(epoch).await?;
580        system_state
581            .as_ref()
582            .try_into()
583            .map_err(|e: TryFromProtoError| tonic::Status::from_error(e.into()))
584    }
585
586    pub async fn get_committee(
587        &self,
588        epoch: Option<u64>,
589    ) -> Result<sui_types::committee::Committee> {
590        let mut request = proto::GetEpochRequest::default();
591        if let Some(epoch) = epoch {
592            request.set_epoch(epoch);
593        }
594        request.set_read_mask(FieldMask::from_paths([
595            proto::Epoch::path_builder().epoch(),
596            proto::Epoch::path_builder().committee().finish(),
597        ]));
598        let response = self
599            .0
600            .clone()
601            .ledger_client()
602            .get_epoch(request)
603            .await?
604            .into_inner();
605
606        response
607            .epoch()
608            .committee()
609            .try_into()
610            .map_err(|e: TryFromProtoError| tonic::Status::from_error(e.into()))
611    }
612
613    pub async fn get_coin_info(
614        &self,
615        coin_type: &move_core_types::language_storage::StructTag,
616    ) -> Result<proto::GetCoinInfoResponse> {
617        let resp = self
618            .0
619            .clone()
620            .state_client()
621            .get_coin_info(
622                proto::GetCoinInfoRequest::default()
623                    .with_coin_type(coin_type.to_canonical_string(true)),
624            )
625            .await?
626            .into_inner();
627        Ok(resp)
628    }
629
630    pub async fn get_balance(
631        &self,
632        owner: SuiAddress,
633        coin_type: &move_core_types::language_storage::StructTag,
634    ) -> Result<proto::Balance> {
635        let resp = self
636            .0
637            .clone()
638            .state_client()
639            .get_balance(
640                proto::GetBalanceRequest::default()
641                    .with_owner(owner.to_string())
642                    .with_coin_type(coin_type.to_canonical_string(true)),
643            )
644            .await?
645            .into_inner();
646
647        Ok(resp.balance.unwrap_or_default())
648    }
649
650    pub fn list_balances(
651        &self,
652        owner: SuiAddress,
653    ) -> impl Stream<Item = Result<proto::Balance>> + 'static {
654        self.0
655            .list_balances(proto::ListBalancesRequest::default().with_owner(owner.to_string()))
656    }
657
658    pub async fn list_delegated_stake(
659        &self,
660        owner: SuiAddress,
661    ) -> Result<Vec<sui_rpc::client::DelegatedStake>> {
662        self.0.clone().list_delegated_stake(&owner.into()).await
663    }
664
665    pub fn transaction_builder(&self) -> sui_transaction_builder::TransactionBuilder {
666        sui_transaction_builder::TransactionBuilder::new(std::sync::Arc::new(self.clone()) as _)
667    }
668}
669
670#[derive(Clone, Debug, serde::Serialize)]
671pub struct ExecutedTransaction {
672    pub transaction: TransactionData,
673    pub signatures: Vec<GenericSignature>,
674    pub effects: TransactionEffects,
675    pub clever_error: Option<proto::CleverError>,
676    pub events: Option<TransactionEvents>,
677    pub event_json: Vec<Option<serde_json::Value>>,
678    pub changed_objects: Vec<proto::ChangedObject>,
679    #[allow(unused)]
680    unchanged_loaded_runtime_objects: Vec<proto::ObjectReference>,
681    pub balance_changes: Vec<sui_sdk_types::BalanceChange>,
682    pub checkpoint: Option<u64>,
683    #[allow(unused)]
684    #[serde(skip)]
685    timestamp: Option<prost_types::Timestamp>,
686}
687
688impl ExecutedTransaction {
689    fn proto_read_mask() -> FieldMask {
690        use proto::ExecutedTransaction;
691        FieldMask::from_paths([
692            ExecutedTransaction::path_builder()
693                .transaction()
694                .bcs()
695                .finish(),
696            ExecutedTransaction::path_builder()
697                .signatures()
698                .bcs()
699                .finish(),
700            ExecutedTransaction::path_builder().effects().bcs().finish(),
701            ExecutedTransaction::path_builder()
702                .effects()
703                .status()
704                .error()
705                .abort()
706                .clever_error()
707                .finish(),
708            ExecutedTransaction::path_builder()
709                .effects()
710                .unchanged_loaded_runtime_objects()
711                .finish(),
712            ExecutedTransaction::path_builder()
713                .effects()
714                .changed_objects()
715                .finish(),
716            ExecutedTransaction::path_builder().events().bcs().finish(),
717            ExecutedTransaction::path_builder().events().events().json(),
718            ExecutedTransaction::path_builder()
719                .balance_changes()
720                .finish(),
721            ExecutedTransaction::path_builder().checkpoint(),
722            ExecutedTransaction::path_builder().timestamp(),
723        ])
724    }
725
726    pub fn get_new_package_obj(&self) -> Option<sui_types::base_types::ObjectRef> {
727        use sui_rpc::proto::sui::rpc::v2::changed_object::OutputObjectState;
728
729        self.changed_objects
730            .iter()
731            .find(|o| matches!(o.output_state(), OutputObjectState::PackageWrite))
732            .and_then(|o| {
733                let id = o.object_id().parse().ok()?;
734                let version = o.output_version().into();
735                let digest = o.output_digest().parse().ok()?;
736                Some((id, version, digest))
737            })
738    }
739
740    pub fn get_new_package_upgrade_cap(&self) -> Option<sui_types::base_types::ObjectRef> {
741        use sui_rpc::proto::sui::rpc::v2::changed_object::OutputObjectState;
742        use sui_rpc::proto::sui::rpc::v2::owner::OwnerKind;
743
744        const UPGRADE_CAP: &str = "0x0000000000000000000000000000000000000000000000000000000000000002::package::UpgradeCap";
745
746        self.changed_objects
747            .iter()
748            .find(|o| {
749                matches!(o.output_state(), OutputObjectState::ObjectWrite)
750                    && matches!(
751                        o.output_owner().kind(),
752                        OwnerKind::Address | OwnerKind::ConsensusAddress
753                    )
754                    && o.object_type() == UPGRADE_CAP
755            })
756            .and_then(|o| {
757                let id = o.object_id().parse().ok()?;
758                let version = o.output_version().into();
759                let digest = o.output_digest().parse().ok()?;
760                Some((id, version, digest))
761            })
762    }
763
764    pub fn timestamp_ms(&self) -> Option<u64> {
765        self.timestamp
766            .and_then(|timestamp| sui_rpc::proto::proto_to_timestamp_ms(timestamp).ok())
767    }
768}
769
770#[derive(Clone, Debug, serde::Serialize)]
771pub struct SimulateTransactionResponse {
772    pub transaction: ExecutedTransaction,
773    pub command_outputs: Vec<proto::CommandResult>,
774    pub suggested_gas_price: Option<u64>,
775}
776
777/// Attempts to parse `CertifiedCheckpointSummary` from a proto::Checkpoint
778#[allow(clippy::result_large_err)]
779fn certified_checkpoint_summary_try_from_proto(
780    checkpoint: &proto::Checkpoint,
781) -> Result<CertifiedCheckpointSummary, TryFromProtoError> {
782    let summary = checkpoint
783        .summary
784        .as_ref()
785        .and_then(|summary| summary.bcs.as_ref())
786        .ok_or_else(|| TryFromProtoError::missing("summary.bcs"))?
787        .deserialize()
788        .map_err(|e| TryFromProtoError::invalid("summary.bcs", e))?;
789
790    let signature = sui_types::crypto::AuthorityStrongQuorumSignInfo::from(
791        sui_sdk_types::ValidatorAggregatedSignature::try_from(
792            checkpoint
793                .signature
794                .as_ref()
795                .ok_or_else(|| TryFromProtoError::missing("signature"))?,
796        )
797        .map_err(|e| TryFromProtoError::invalid("signature", e))?,
798    );
799
800    Ok(CertifiedCheckpointSummary::new_from_data_and_sig(
801        summary, signature,
802    ))
803}
804
805/// Attempts to parse `Object` from the bcs fields in `GetObjectResponse`
806#[allow(clippy::result_large_err)]
807fn object_try_from_proto(object: &proto::Object) -> Result<Object, TryFromProtoError> {
808    object
809        .bcs
810        .as_ref()
811        .ok_or_else(|| TryFromProtoError::missing("bcs"))?
812        .deserialize()
813        .map_err(|e| TryFromProtoError::invalid("bcs", e))
814}
815
816/// Attempts to parse `ExecutedTransaction` from the fields in `proto::ExecuteTransactionResponse`
817#[allow(clippy::result_large_err)]
818fn execute_transaction_response_try_from_proto(
819    response: &proto::ExecuteTransactionResponse,
820) -> Result<ExecutedTransaction, TryFromProtoError> {
821    let executed_transaction = response
822        .transaction
823        .as_ref()
824        .ok_or_else(|| TryFromProtoError::missing("transaction"))?;
825
826    executed_transaction_try_from_proto(executed_transaction)
827}
828
829#[allow(clippy::result_large_err)]
830fn executed_transaction_try_from_proto(
831    executed_transaction: &proto::ExecutedTransaction,
832) -> Result<ExecutedTransaction, TryFromProtoError> {
833    let transaction = executed_transaction
834        .transaction()
835        .bcs()
836        .deserialize()
837        .map_err(|e| TryFromProtoError::invalid("transaction.bcs", e))?;
838
839    let effects = executed_transaction
840        .effects()
841        .bcs()
842        .deserialize()
843        .map_err(|e| TryFromProtoError::invalid("effects.bcs", e))?;
844    let signatures = executed_transaction
845        .signatures()
846        .iter()
847        .map(|sig| {
848            GenericSignature::from_bytes(sig.bcs().value())
849                .map_err(|e| TryFromProtoError::invalid("signatures.bcs", e))
850        })
851        .collect::<Result<_, _>>()?;
852    let clever_error = executed_transaction
853        .effects()
854        .status()
855        .error()
856        .abort()
857        .clever_error_opt()
858        .cloned();
859    let events = executed_transaction
860        .events
861        .as_ref()
862        .and_then(|events| events.bcs.as_ref())
863        .map(|bcs| bcs.deserialize())
864        .transpose()
865        .map_err(|e| TryFromProtoError::invalid("events.bcs", e))?;
866    let event_json = executed_transaction
867        .events_opt()
868        .map(|events| {
869            events
870                .events()
871                .iter()
872                .map(|event| event.json_opt().map(proto_value_to_json_value))
873                .collect::<Vec<_>>()
874        })
875        .unwrap_or_default();
876
877    let balance_changes = executed_transaction
878        .balance_changes
879        .iter()
880        .map(TryInto::try_into)
881        .collect::<Result<_, _>>()?;
882
883    ExecutedTransaction {
884        transaction,
885        signatures,
886        effects,
887        clever_error,
888        events,
889        event_json,
890        balance_changes,
891        checkpoint: executed_transaction.checkpoint,
892        changed_objects: executed_transaction.effects().changed_objects().to_owned(),
893        unchanged_loaded_runtime_objects: executed_transaction
894            .effects()
895            .unchanged_loaded_runtime_objects()
896            .to_owned(),
897        timestamp: executed_transaction.timestamp,
898    }
899    .pipe(Ok)
900}
901
902fn proto_value_to_json_value(proto: &prost_types::Value) -> serde_json::Value {
903    match proto.kind.as_ref() {
904        Some(ProtoValueKind::NullValue(_)) | None => serde_json::Value::Null,
905        Some(ProtoValueKind::NumberValue(n)) => serde_json::Value::from(*n),
906        Some(ProtoValueKind::StringValue(s)) => serde_json::Value::from(s.clone()),
907        Some(ProtoValueKind::BoolValue(b)) => serde_json::Value::from(*b),
908        Some(ProtoValueKind::StructValue(map)) => serde_json::Value::Object(
909            map.fields
910                .iter()
911                .map(|(k, v)| (k.clone(), proto_value_to_json_value(v)))
912                .collect(),
913        ),
914        Some(ProtoValueKind::ListValue(list_value)) => serde_json::Value::Array(
915            list_value
916                .values
917                .iter()
918                .map(proto_value_to_json_value)
919                .collect(),
920        ),
921    }
922}
923
924fn status_from_error_with_metadata<T: Into<BoxError>>(err: T, metadata: MetadataMap) -> Status {
925    let mut status = Status::from_error(err.into());
926    *status.metadata_mut() = metadata;
927    status
928}
929
930#[async_trait::async_trait]
931impl sui_transaction_builder::DataReader for Client {
932    async fn get_owned_objects(
933        &self,
934        address: SuiAddress,
935        object_type: move_core_types::language_storage::StructTag,
936    ) -> Result<Vec<sui_types::base_types::ObjectInfo>, anyhow::Error> {
937        self.list_owned_objects(address, Some(object_type))
938            .map_ok(|o| sui_types::base_types::ObjectInfo::from_object(&o))
939            .try_collect()
940            .await
941            .map_err(Into::into)
942    }
943
944    async fn get_object(&self, object_id: ObjectID) -> Result<Object, anyhow::Error> {
945        let mut client = self.clone();
946        Self::get_object(&mut client, object_id)
947            .await
948            .map_err(Into::into)
949    }
950
951    async fn get_reference_gas_price(&self) -> Result<u64, anyhow::Error> {
952        self.get_reference_gas_price().await.map_err(Into::into)
953    }
954}