sui_rpc/client/ledger_streams/adapter/
transaction.rs1use 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}