Skip to main content

sui_rpc/client/
transaction_execution.rs

1use super::Client;
2use crate::field::FieldMaskUtil;
3use crate::proto::TryFromProtoError;
4use crate::proto::sui::rpc::v2::ExecuteTransactionRequest;
5use crate::proto::sui::rpc::v2::ExecuteTransactionResponse;
6use crate::proto::sui::rpc::v2::ExecutionError;
7use crate::proto::sui::rpc::v2::GetEpochRequest;
8use crate::proto::sui::rpc::v2::GetTransactionRequest;
9use crate::proto::sui::rpc::v2::GetTransactionResponse;
10use crate::proto::sui::rpc::v2::SubscribeCheckpointsRequest;
11use futures::TryStreamExt;
12use prost_types::FieldMask;
13use std::fmt;
14use std::time::Duration;
15use tonic::Response;
16
17/// Error types that can occur when executing a transaction and waiting for checkpoint
18#[derive(Debug)]
19#[non_exhaustive]
20pub enum ExecuteAndWaitError {
21    /// RPC Error (actual tonic::Status from the client/server)
22    RpcError(tonic::Status),
23    /// Request is missing the required transaction field
24    MissingTransaction,
25    /// Failed to parse/convert the transaction for digest calculation
26    ProtoConversionError(TryFromProtoError),
27    /// Transaction executed but checkpoint wait timed out
28    CheckpointTimeout(Response<ExecuteTransactionResponse>),
29    /// Transaction executed but checkpoint stream had an error
30    CheckpointStreamError {
31        response: Response<ExecuteTransactionResponse>,
32        error: tonic::Status,
33    },
34}
35
36impl std::fmt::Display for ExecuteAndWaitError {
37    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
38        match self {
39            Self::RpcError(status) => write!(f, "RPC error: {status}"),
40            Self::MissingTransaction => {
41                write!(f, "Request is missing the required transaction field")
42            }
43            Self::ProtoConversionError(e) => write!(f, "Failed to convert transaction: {e}"),
44            Self::CheckpointTimeout(_) => {
45                write!(f, "Transaction executed but checkpoint wait timed out")
46            }
47            Self::CheckpointStreamError { error, .. } => {
48                write!(
49                    f,
50                    "Transaction executed but checkpoint stream had an error: {error}"
51                )
52            }
53        }
54    }
55}
56
57impl std::error::Error for ExecuteAndWaitError {
58    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
59        match self {
60            Self::RpcError(status) => Some(status),
61            Self::ProtoConversionError(e) => Some(e),
62            Self::CheckpointStreamError { error, .. } => Some(error),
63            Self::MissingTransaction => None,
64            Self::CheckpointTimeout(_) => None,
65        }
66    }
67}
68
69impl Client {
70    /// Executes a transaction and waits for it to be included in a checkpoint.
71    ///
72    /// This method provides "read your writes" consistency by executing the transaction
73    /// and waiting for it to appear in a checkpoint, which gauruntees indexes have been updated on
74    /// this node.
75    ///
76    /// # Arguments
77    /// * `request` - The transaction execution request (ExecuteTransactionRequest)
78    /// * `timeout` - Maximum time to wait for indexing confirmation
79    ///
80    /// # Returns
81    /// A `Result` containing the response if the transaction was executed and checkpoint confirmed,
82    /// or an error that may still include the response if execution succeeded but checkpoint
83    /// confirmation failed.
84    ///
85    /// # Duplicate submissions
86    /// Submitting a transaction that has already been executed is handled
87    /// gracefully. While the execution RPC is in flight the ledger is probed
88    /// for the transaction, and a transaction that is already in a checkpoint
89    /// is returned without waiting for execution to finish. Likewise, if
90    /// execution fails but the ledger shows the transaction in a checkpoint
91    /// (for example when a resubmission races the original submission), the
92    /// execution error is discarded and the committed transaction is
93    /// returned. In both cases the response is assembled from
94    /// `GetTransaction` using the request's read mask, so it carries the same
95    /// fields an execution response would, with `digest`, `checkpoint`, and
96    /// `timestamp` always populated.
97    #[allow(clippy::result_large_err)]
98    pub async fn execute_transaction_and_wait_for_checkpoint(
99        &mut self,
100        request: impl tonic::IntoRequest<ExecuteTransactionRequest>,
101        timeout: Duration,
102    ) -> Result<Response<ExecuteTransactionResponse>, ExecuteAndWaitError> {
103        // Calculate digest from the input transaction to avoid relying on response read mask
104        let request = request.into_request();
105        let transaction = match request.get_ref().transaction_opt() {
106            Some(tx) => tx,
107            None => return Err(ExecuteAndWaitError::MissingTransaction),
108        };
109
110        let executed_txn_digest = match sui_sdk_types::Transaction::try_from(transaction) {
111            Ok(tx) => tx.digest().to_string(),
112            Err(e) => return Err(ExecuteAndWaitError::ProtoConversionError(e)),
113        };
114
115        // Read mask for answering from GetTransaction when execution cannot
116        // provide the response: a duplicate submission that already
117        // committed, or an execution error after the transaction landed.
118        // Both RPCs' masks select fields of `ExecutedTransaction`, so the
119        // caller's mask passes through unchanged; when the caller didn't set
120        // one, mirror ExecuteTransaction's documented default. `digest`,
121        // `checkpoint`, and `timestamp` are always included since this
122        // method's contract populates them.
123        let lookup_mask = {
124            let caller_paths = match &request.get_ref().read_mask {
125                Some(mask) => mask.paths.clone(),
126                None => vec!["effects.status".to_owned()],
127            };
128            FieldMask::from_paths(caller_paths.iter().map(String::as_str).chain([
129                "digest",
130                "checkpoint",
131                "timestamp",
132            ]))
133            .normalize()
134        };
135
136        // Subscribe to checkpoint stream before execution to avoid missing the transaction.
137        // Uses minimal read mask for efficiency since we only nee digest confirmation.
138        // Once server-side filtering is available, we should filter by transaction digest to
139        // further reduce bandwidth.
140        let mut checkpoint_stream = match self
141            .subscription_client()
142            .subscribe_checkpoints(SubscribeCheckpointsRequest::default().with_read_mask(
143                FieldMask::from_str("transactions.digest,sequence_number,summary.timestamp"),
144            ))
145            .await
146        {
147            Ok(stream) => stream.into_inner(),
148            Err(e) => return Err(ExecuteAndWaitError::RpcError(e)),
149        };
150
151        // Scan the subscription for the transaction's digest. Every RPC on
152        // this client shares one HTTP/2 connection, so this future must be
153        // polled concurrently with the execution phase below: a subscription
154        // parked while another call is awaited pins its flow-control window
155        // (checkpoints keep arriving whether or not anyone reads them) and,
156        // past the idle timeout, gets reset by the client's body watchdog.
157        //
158        // Both this future and the execution future below are boxed: their
159        // combined state (two full tonic call chains alive at once) would
160        // otherwise be inlined into this method's future, making it large
161        // enough to threaten a stack overflow in callers that hold it in
162        // deeply nested or spawned futures.
163        let mut scan = Box::pin(async {
164            while let Some(response) = checkpoint_stream.try_next().await? {
165                let checkpoint = response.checkpoint();
166
167                for tx in checkpoint.transactions() {
168                    if tx.digest() == executed_txn_digest {
169                        return Ok((checkpoint.sequence_number(), checkpoint.summary().timestamp));
170                    }
171                }
172            }
173            Err(tonic::Status::aborted(
174                "checkpoint stream ended unexpectedly",
175            ))
176        });
177
178        // The concurrent futures below each need a service client, and a
179        // single `&mut self` cannot back all of them at once, so give each
180        // its own client over the shared channel.
181        let mut execution_client = self.execution_client();
182        let mut post_exec_lookup_client = self.ledger_client();
183        let mut probe_client = self.ledger_client();
184
185        // Execute, then query the fullnode directly to see if it already has
186        // the txn in a checkpoint. This is to handle the case where an
187        // already executed transaction is sent multiple times.
188        let mut exec_and_check = Box::pin(async {
189            let response = execution_client.execute_transaction(request).await?;
190
191            let already_checkpointed = match post_exec_lookup_client
192                .get_transaction(
193                    GetTransactionRequest::default()
194                        .with_digest(&executed_txn_digest)
195                        .with_read_mask(FieldMask::from_str("digest,checkpoint,timestamp")),
196                )
197                .await
198            {
199                Ok(resp) if resp.get_ref().transaction().checkpoint_opt().is_some() => Some((
200                    resp.get_ref().transaction().checkpoint(),
201                    resp.get_ref().transaction().timestamp,
202                )),
203                _ => None,
204            };
205
206            Ok::<_, tonic::Status>((response, already_checkpointed))
207        });
208
209        // Probe the ledger while execution is in flight: a resubmission of a
210        // transaction that already committed can be answered from the ledger
211        // without waiting for, or succeeding at, execution.
212        let mut probe = Box::pin(async {
213            probe_client
214                .get_transaction(
215                    GetTransactionRequest::default()
216                        .with_digest(&executed_txn_digest)
217                        .with_read_mask(lookup_mask.clone()),
218                )
219                .await
220        });
221
222        // Drive execution, the scan, and the probe together. The scan can
223        // complete first (for example, when a duplicate of an already
224        // executed transaction lands in a checkpoint mid-execution), so
225        // remember its outcome; the guards keep completed futures from being
226        // polled again. A probe that finds the transaction in a checkpoint
227        // resolves the call on the spot; any other probe outcome (not found,
228        // not yet checkpointed, or an RPC error) means execution has to
229        // provide the answer.
230        let mut scan_result = None;
231        let mut probe_done = false;
232        let exec_result = loop {
233            tokio::select! {
234                exec = &mut exec_and_check => break exec,
235                result = &mut scan, if scan_result.is_none() => {
236                    scan_result = Some(result);
237                }
238                result = &mut probe, if !probe_done => {
239                    probe_done = true;
240                    if let Ok(lookup) = result
241                        && lookup.get_ref().transaction().checkpoint_opt().is_some()
242                    {
243                        return Ok(lookup_into_execute_response(lookup));
244                    }
245                }
246            }
247        };
248
249        let (mut response, already_checkpointed) = match exec_result {
250            Ok(ok) => ok,
251            Err(error) => {
252                // Execution can fail for a transaction that nonetheless
253                // committed, for example when a resubmission races the
254                // original submission. Consult the ledger before surfacing
255                // the error.
256                drop(probe);
257                if let Ok(lookup) = probe_client
258                    .get_transaction(
259                        GetTransactionRequest::default()
260                            .with_digest(&executed_txn_digest)
261                            .with_read_mask(lookup_mask),
262                    )
263                    .await
264                    && lookup.get_ref().transaction().checkpoint_opt().is_some()
265                {
266                    return Ok(lookup_into_execute_response(lookup));
267                }
268                return Err(ExecuteAndWaitError::RpcError(error));
269            }
270        };
271
272        // Wait for the transaction to appear in a checkpoint, at which point
273        // indexes will have been updated. The direct lookup takes precedence:
274        // when it already places the transaction in a checkpoint there is
275        // nothing to wait for, even if the scan failed in the meantime.
276        let (checkpoint, timestamp) = if let Some(found) = already_checkpointed {
277            found
278        } else {
279            let result = match scan_result {
280                Some(result) => result,
281                None => {
282                    tokio::select! {
283                        result = &mut scan => result,
284                        _ = tokio::time::sleep(timeout) => {
285                            return Err(ExecuteAndWaitError::CheckpointTimeout(response));
286                        }
287                    }
288                }
289            };
290            match result {
291                Ok(found) => found,
292                Err(e) => {
293                    return Err(ExecuteAndWaitError::CheckpointStreamError { response, error: e });
294                }
295            }
296        };
297
298        response
299            .get_mut()
300            .transaction_mut()
301            .set_checkpoint(checkpoint);
302        response.get_mut().transaction_mut().timestamp = timestamp;
303        Ok(response)
304    }
305
306    /// Retrieves the current reference gas price from the latest epoch information.
307    ///
308    /// # Returns
309    /// The reference gas price as a `u64`
310    ///
311    /// # Errors
312    /// Returns an error if there is an RPC error when fetching the epoch information
313    pub async fn get_reference_gas_price(&mut self) -> Result<u64, tonic::Status> {
314        let request = GetEpochRequest::latest()
315            .with_read_mask(FieldMask::from_paths(["reference_gas_price"]));
316        let response = self.ledger_client().get_epoch(request).await?.into_inner();
317        Ok(response.epoch().reference_gas_price())
318    }
319}
320
321/// Builds the response for a transaction answered from the ledger instead of
322/// from execution.
323fn lookup_into_execute_response(
324    response: Response<GetTransactionResponse>,
325) -> Response<ExecuteTransactionResponse> {
326    Response::new(ExecuteTransactionResponse {
327        transaction: response.into_inner().transaction,
328    })
329}
330
331impl fmt::Display for ExecutionError {
332    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
333        let description = self.description.as_deref().unwrap_or("No description");
334        write!(
335            f,
336            "ExecutionError: Kind: {}, Description: {}",
337            self.kind().as_str_name(),
338            description
339        )
340    }
341}