sui_rpc/client/ledger_streams/adapter/
checkpoint.rs1use tonic::Request;
2use tonic::Status;
3use tonic::codegen::BoxStream;
4
5use super::super::super::Client;
6use super::super::super::Result;
7use super::super::observability::LedgerStreamFamily;
8use super::super::types::CheckpointStreamFrame;
9use super::CHECKPOINT_CURSOR_OVERFLOW;
10use super::ListResponseParts;
11use super::ListScanDirection;
12use super::LiveFrame;
13use super::PositionedItem;
14use super::Progress;
15use super::RpcFuture;
16use super::SubscriptionAdapter;
17use crate::proto::sui::rpc::v2::Checkpoint;
18use crate::proto::sui::rpc::v2::ListCheckpointsRequest;
19use crate::proto::sui::rpc::v2::ListCheckpointsResponse;
20use crate::proto::sui::rpc::v2::QueryEnd;
21use crate::proto::sui::rpc::v2::QueryOptions;
22use crate::proto::sui::rpc::v2::SubscribeCheckpointsRequest;
23use crate::proto::sui::rpc::v2::SubscribeCheckpointsResponse;
24use crate::proto::sui::rpc::v2::Watermark;
25
26pub(in crate::client::ledger_streams) struct CheckpointAdapter;
27
28impl SubscriptionAdapter for CheckpointAdapter {
29 const FAMILY: LedgerStreamFamily = LedgerStreamFamily::Checkpoint;
30 const REQUIRED_READ_MASK_FIELDS: &'static [&'static str] = &["sequence_number"];
31 const READ_MASK_REQUIREMENT: &'static str =
32 "read_mask must include \"sequence_number\" or \"*\"";
33
34 type Item = PositionedItem<Checkpoint, u64>;
35 type ItemPosition = u64;
36 type Cursor = u64;
37 type Output = CheckpointStreamFrame;
38 type ListRequest = ListCheckpointsRequest;
39 type ListResponse = ListCheckpointsResponse;
40 type SubscribeRequest = SubscribeCheckpointsRequest;
41 type SubscribeResponse = SubscribeCheckpointsResponse;
42
43 fn list_read_mask(request: &Self::ListRequest) -> Option<&prost_types::FieldMask> {
44 request.read_mask.as_ref()
45 }
46
47 fn list_request_from_subscribe(request: &Self::SubscribeRequest) -> Self::ListRequest {
48 ListCheckpointsRequest {
49 read_mask: request.read_mask.clone(),
50 filter: request.filter.clone(),
51 start_checkpoint: None,
52 end_checkpoint: None,
53 options: None,
54 }
55 }
56
57 fn options(request: &Self::ListRequest) -> Option<&QueryOptions> {
58 request.options.as_ref()
59 }
60
61 fn options_mut(request: &mut Self::ListRequest) -> &mut QueryOptions {
62 request.options.get_or_insert_with(QueryOptions::default)
63 }
64
65 fn start_checkpoint(request: &Self::ListRequest) -> Option<u64> {
66 request.start_checkpoint
67 }
68
69 fn set_start_checkpoint(request: &mut Self::ListRequest, checkpoint: Option<u64>) {
70 request.start_checkpoint = checkpoint;
71 }
72
73 fn end_checkpoint(request: &Self::ListRequest) -> Option<u64> {
74 request.end_checkpoint
75 }
76
77 fn set_end_checkpoint(request: &mut Self::ListRequest, checkpoint: Option<u64>) {
78 request.end_checkpoint = checkpoint;
79 }
80
81 fn set_ascending_resume(
82 request: &mut Self::ListRequest,
83 progress: &Progress<Self::Cursor>,
84 ) -> Result<()> {
85 request.start_checkpoint = Some(
86 progress
87 .cursor
88 .checked_add(1)
89 .ok_or_else(|| Status::out_of_range(CHECKPOINT_CURSOR_OVERFLOW))?,
90 );
91 Self::options_mut(request).after = None;
92 Ok(())
93 }
94
95 fn request_resume_position(request: &Self::ListRequest) -> Option<Progress<Self::Cursor>> {
96 request
97 .start_checkpoint
98 .and_then(|start| start.checked_sub(1))
99 .map(|cursor| Progress {
100 cursor,
101 checkpoint: Some(cursor),
102 })
103 }
104
105 fn validate_checkpoint_bound(
106 request: &Self::ListRequest,
107 direction: ListScanDirection,
108 checkpoint: Option<u64>,
109 ) -> Result<()> {
110 super::validate_typed_checkpoint_bound(
111 request.start_checkpoint,
112 request.end_checkpoint,
113 direction,
114 checkpoint,
115 )
116 }
117 fn item_position(item: &Self::Item) -> &Self::ItemPosition {
118 item.position()
119 }
120
121 fn extract_metadata(
122 response: &Self::ListResponse,
123 ) -> (bool, Option<&Watermark>, Option<&QueryEnd>) {
124 (
125 response.checkpoint.is_some(),
126 response.watermark.as_ref(),
127 response.end.as_ref(),
128 )
129 }
130
131 fn split_list(response: Self::ListResponse) -> Result<ListResponseParts<Self>> {
132 let item = response.checkpoint.map(|checkpoint| -> Result<Self::Item> {
133 let position = checkpoint.sequence_number.ok_or_else(|| {
134 Status::data_loss("List checkpoint item is missing its sequence number")
135 })?;
136 Ok(PositionedItem::new(checkpoint, position))
137 });
138 Ok(ListResponseParts {
139 item: item.transpose()?,
140 watermark: response.watermark,
141 })
142 }
143
144 fn item_required(request: &Self::SubscribeRequest) -> bool {
145 request.filter.is_none()
146 }
147
148 fn parse_live(
149 response: Self::SubscribeResponse,
150 item_required: bool,
151 ) -> Result<LiveFrame<Self::Item, Progress<Self::Cursor>>> {
152 let cursor = response.cursor.ok_or_else(|| {
153 Status::data_loss("checkpoint subscription frame is missing its cursor")
154 })?;
155 if response.checkpoint.is_none() && item_required {
156 return Err(Status::data_loss(
157 "unfiltered checkpoint subscription frame is missing its checkpoint",
158 ));
159 }
160 let item = response.checkpoint.map(|checkpoint| -> Result<Self::Item> {
161 let position = checkpoint.sequence_number.ok_or_else(|| {
162 Status::data_loss("checkpoint subscription item is missing its sequence number")
163 })?;
164 if position != cursor {
165 return Err(Status::data_loss(
166 "checkpoint subscription item sequence does not match its cursor",
167 ));
168 }
169 Ok(PositionedItem::new(checkpoint, position))
170 });
171 Ok(LiveFrame {
172 item: item.transpose()?,
173 progress: Progress {
174 cursor,
175 checkpoint: Some(cursor),
176 },
177 })
178 }
179
180 fn into_output(item: Option<Self::Item>, progress: Progress<Self::Cursor>) -> Self::Output {
181 let checkpoint = item.map(PositionedItem::into_payload);
182 CheckpointStreamFrame {
183 checkpoint,
184 cursor: progress.cursor,
185 }
186 }
187
188 fn dispatch_list(
189 mut client: Client,
190 request: Request<Self::ListRequest>,
191 ) -> RpcFuture<Self::ListResponse> {
192 Box::pin(async move {
193 let stream = client
194 .ledger_client()
195 .list_checkpoints(request)
196 .await?
197 .into_inner();
198 Ok(Box::pin(stream) as BoxStream<Self::ListResponse>)
199 })
200 }
201
202 fn dispatch_subscribe(
203 mut client: Client,
204 request: Request<Self::SubscribeRequest>,
205 ) -> RpcFuture<Self::SubscribeResponse> {
206 Box::pin(async move {
207 let stream = client
208 .subscription_client()
209 .subscribe_checkpoints(request)
210 .await?
211 .into_inner();
212 Ok(Box::pin(stream) as BoxStream<Self::SubscribeResponse>)
213 })
214 }
215}