Skip to main content

sui_rpc/client/ledger_streams/adapter/
checkpoint.rs

1use 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}