blokli_client/client/
transactions.rs1use std::time::Duration;
2
3use cynic::{MutationBuilder, SubscriptionBuilder};
4use futures::TryStreamExt;
5use futures_time::{future::FutureExt, time::Duration as FuturesTimeDuration};
6use hex::ToHex;
7
8use super::{BlokliClient, GraphQlQueries, response_to_data};
9use crate::{
10 api::{
11 BlokliTransactionClient, Result, TxId, TxReceipt,
12 internal::{
13 ConfirmTransactionVariables, MutateConfirmTransaction, MutateSendTransaction, MutateTrackTransaction,
14 SendTransactionVariables, SubscribeTransaction, TransactionsVariables,
15 },
16 types::{Transaction, TransactionStatus},
17 },
18 errors::{ErrorKind, TrackingErrorKind},
19};
20
21impl GraphQlQueries {
22 pub fn mutate_submit_transaction(
24 signed_tx: &[u8],
25 ) -> cynic::Operation<MutateSendTransaction, SendTransactionVariables> {
26 MutateSendTransaction::build(SendTransactionVariables {
27 raw_transaction: signed_tx.encode_hex(),
28 })
29 }
30
31 pub fn mutate_submit_and_track_transaction(
33 signed_tx: &[u8],
34 ) -> cynic::Operation<MutateTrackTransaction, SendTransactionVariables> {
35 MutateTrackTransaction::build(SendTransactionVariables {
36 raw_transaction: signed_tx.encode_hex(),
37 })
38 }
39
40 pub fn mutate_submit_and_confirm_transaction(
42 signed_tx: &[u8],
43 confirmations: usize,
44 ) -> cynic::Operation<MutateConfirmTransaction, ConfirmTransactionVariables> {
45 MutateConfirmTransaction::build(ConfirmTransactionVariables {
46 raw_transaction: signed_tx.encode_hex(),
47 confirmations: confirmations.min(128) as i32,
48 })
49 }
50
51 pub fn subscribe_track_transaction(
53 tx_id: TxId,
54 ) -> cynic::StreamingOperation<SubscribeTransaction, TransactionsVariables> {
55 SubscribeTransaction::build(TransactionsVariables { id: tx_id.into() })
56 }
57}
58
59#[async_trait::async_trait]
60impl BlokliTransactionClient for BlokliClient {
61 async fn submit_transaction(&self, signed_tx: &[u8]) -> Result<TxReceipt> {
62 let resp = self
63 .build_query(GraphQlQueries::mutate_submit_transaction(signed_tx))?
64 .await?;
65
66 response_to_data(resp)?.send_transaction.into()
67 }
68
69 async fn submit_and_track_transaction(&self, signed_tx: &[u8]) -> Result<TxId> {
70 let resp = self
71 .build_query(GraphQlQueries::mutate_submit_and_track_transaction(signed_tx))?
72 .await?;
73
74 let tx: Result<Transaction> = response_to_data(resp)?.send_transaction_async.into();
75
76 Ok(tx?.id.into_inner())
77 }
78
79 async fn submit_and_confirm_transaction(&self, signed_tx: &[u8], num_confirmations: usize) -> Result<TxReceipt> {
80 let resp = self
81 .build_query(GraphQlQueries::mutate_submit_and_confirm_transaction(
82 signed_tx,
83 num_confirmations,
84 ))?
85 .await?;
86
87 let tx: Result<Transaction> = response_to_data(resp)?.send_transaction_sync.into();
88
89 let hash = tx?.transaction_hash.0.to_lowercase();
90 Ok(hex::decode(hash.trim_start_matches("0x"))
91 .map_err(|_| ErrorKind::ParseError)
92 .and_then(|d| d.try_into().map_err(|_| ErrorKind::ParseError))?)
93 }
94
95 async fn track_transaction(&self, tx_id: TxId, client_timeout: Duration) -> Result<Transaction> {
96 self.build_subscription_stream(GraphQlQueries::subscribe_track_transaction(tx_id))?
97 .try_filter_map(|item| {
98 futures::future::ready(match &item.transaction_updated.status {
99 TransactionStatus::Confirmed => Ok(Some(item.transaction_updated)),
100 TransactionStatus::Pending | TransactionStatus::Submitted => Ok(None),
101 TransactionStatus::Reverted => Err(ErrorKind::TrackingError(TrackingErrorKind::Reverted).into()),
102 TransactionStatus::SubmissionFailed => {
103 Err(ErrorKind::TrackingError(TrackingErrorKind::SubmissionFailed).into())
104 }
105 TransactionStatus::Timeout => Err(ErrorKind::TrackingError(TrackingErrorKind::Timeout).into()),
106 TransactionStatus::ValidationFailed => {
107 Err(ErrorKind::TrackingError(TrackingErrorKind::ValidationFailed).into())
108 }
109 })
110 })
111 .try_next()
112 .timeout(FuturesTimeDuration::from(client_timeout))
113 .await
114 .map_err(|_| ErrorKind::Timeout)??
115 .ok_or(ErrorKind::NoData.into())
116 }
117}