Skip to main content

sui_rpc/client/ledger_streams/adapter/
transaction.rs

1use prost::bytes::Bytes;
2use tonic::Request;
3use tonic::Status;
4use tonic::codegen::BoxStream;
5
6use super::super::super::Client;
7use super::super::super::Result;
8use super::super::observability::LedgerStreamFamily;
9use super::super::types::TransactionStreamFrame;
10use super::ListResponseParts;
11use super::ListScanDirection;
12use super::LiveFrame;
13use super::PositionedItem;
14use super::Progress;
15use super::RpcFuture;
16use super::SubscriptionAdapter;
17use super::parse_opaque_live_frame;
18use crate::proto::sui::rpc::v2::ExecutedTransaction;
19use crate::proto::sui::rpc::v2::ListTransactionsRequest;
20use crate::proto::sui::rpc::v2::ListTransactionsResponse;
21use crate::proto::sui::rpc::v2::QueryEnd;
22use crate::proto::sui::rpc::v2::QueryOptions;
23use crate::proto::sui::rpc::v2::SubscribeTransactionsRequest;
24use crate::proto::sui::rpc::v2::SubscribeTransactionsResponse;
25use crate::proto::sui::rpc::v2::Watermark;
26
27pub(in crate::client::ledger_streams) struct TransactionAdapter;
28
29impl SubscriptionAdapter for TransactionAdapter {
30    const FAMILY: LedgerStreamFamily = LedgerStreamFamily::Transaction;
31    const REQUIRED_READ_MASK_FIELDS: &'static [&'static str] = &["checkpoint", "transaction_index"];
32    const READ_MASK_REQUIREMENT: &'static str =
33        "read_mask must include \"checkpoint\" and \"transaction_index\" or \"*\"";
34
35    type Item = PositionedItem<ExecutedTransaction, (u64, u64)>;
36    type ItemPosition = (u64, u64);
37    type Cursor = Bytes;
38    type Output = TransactionStreamFrame;
39    type ListRequest = ListTransactionsRequest;
40    type ListResponse = ListTransactionsResponse;
41    type SubscribeRequest = SubscribeTransactionsRequest;
42    type SubscribeResponse = SubscribeTransactionsResponse;
43
44    fn list_read_mask(request: &Self::ListRequest) -> Option<&prost_types::FieldMask> {
45        request.read_mask.as_ref()
46    }
47
48    fn list_request_from_subscribe(request: &Self::SubscribeRequest) -> Self::ListRequest {
49        ListTransactionsRequest {
50            read_mask: request.read_mask.clone(),
51            filter: request.filter.clone(),
52            start_checkpoint: None,
53            end_checkpoint: None,
54            options: None,
55        }
56    }
57
58    fn options(request: &Self::ListRequest) -> Option<&QueryOptions> {
59        request.options.as_ref()
60    }
61
62    fn options_mut(request: &mut Self::ListRequest) -> &mut QueryOptions {
63        request.options.get_or_insert_with(QueryOptions::default)
64    }
65
66    fn start_checkpoint(request: &Self::ListRequest) -> Option<u64> {
67        request.start_checkpoint
68    }
69
70    fn set_start_checkpoint(request: &mut Self::ListRequest, checkpoint: Option<u64>) {
71        request.start_checkpoint = checkpoint;
72    }
73
74    fn end_checkpoint(request: &Self::ListRequest) -> Option<u64> {
75        request.end_checkpoint
76    }
77
78    fn set_end_checkpoint(request: &mut Self::ListRequest, checkpoint: Option<u64>) {
79        request.end_checkpoint = checkpoint;
80    }
81
82    fn set_ascending_resume(
83        request: &mut Self::ListRequest,
84        progress: &Progress<Self::Cursor>,
85    ) -> Result<()> {
86        Self::options_mut(request).after = Some(progress.cursor.clone());
87        Ok(())
88    }
89
90    fn request_resume_position(request: &Self::ListRequest) -> Option<Progress<Self::Cursor>> {
91        request
92            .options
93            .as_ref()
94            .and_then(|options| options.after.clone())
95            .map(|cursor| Progress {
96                cursor,
97                checkpoint: None,
98            })
99    }
100    fn validate_checkpoint_bound(
101        _request: &Self::ListRequest,
102        _direction: ListScanDirection,
103        _checkpoint: Option<u64>,
104    ) -> Result<()> {
105        Ok(())
106    }
107    fn item_position(item: &Self::Item) -> &Self::ItemPosition {
108        item.position()
109    }
110
111    fn extract_metadata(
112        response: &Self::ListResponse,
113    ) -> (bool, Option<&Watermark>, Option<&QueryEnd>) {
114        (
115            response.transaction.is_some(),
116            response.watermark.as_ref(),
117            response.end.as_ref(),
118        )
119    }
120
121    fn split_list(response: Self::ListResponse) -> Result<ListResponseParts<Self>> {
122        let item = response.transaction.map(position_transaction).transpose()?;
123        Ok(ListResponseParts {
124            item,
125            watermark: response.watermark,
126        })
127    }
128
129    fn parse_live(
130        response: Self::SubscribeResponse,
131        _item_required: bool,
132    ) -> Result<LiveFrame<Self::Item, Progress<Self::Cursor>>> {
133        let item = response.transaction.map(position_transaction).transpose()?;
134        parse_opaque_live_frame(item, response.watermark)
135    }
136
137    fn into_output(item: Option<Self::Item>, progress: Progress<Self::Cursor>) -> Self::Output {
138        let transaction = item.map(PositionedItem::into_payload);
139        TransactionStreamFrame {
140            transaction,
141            cursor: progress.cursor,
142            covered_checkpoint: progress.checkpoint,
143        }
144    }
145
146    fn dispatch_list(
147        mut client: Client,
148        request: Request<Self::ListRequest>,
149    ) -> RpcFuture<Self::ListResponse> {
150        Box::pin(async move {
151            let stream = client
152                .ledger_client()
153                .list_transactions(request)
154                .await?
155                .into_inner();
156            Ok(Box::pin(stream) as BoxStream<Self::ListResponse>)
157        })
158    }
159
160    fn dispatch_subscribe(
161        mut client: Client,
162        request: Request<Self::SubscribeRequest>,
163    ) -> RpcFuture<Self::SubscribeResponse> {
164        Box::pin(async move {
165            let stream = client
166                .subscription_client()
167                .subscribe_transactions(request)
168                .await?
169                .into_inner();
170            Ok(Box::pin(stream) as BoxStream<Self::SubscribeResponse>)
171        })
172    }
173}
174
175fn position_transaction(
176    transaction: ExecutedTransaction,
177) -> Result<PositionedItem<ExecutedTransaction, (u64, u64)>> {
178    let checkpoint = transaction
179        .checkpoint
180        .ok_or_else(|| Status::data_loss("transaction item is missing its checkpoint"))?;
181    let transaction_index = transaction
182        .transaction_index
183        .ok_or_else(|| Status::data_loss("transaction item is missing its transaction index"))?;
184    Ok(PositionedItem::new(
185        transaction,
186        (checkpoint, transaction_index),
187    ))
188}