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}