laminar_core/checkpoint/checkpoint_store/
subscription_segments.rs1use bytes::Bytes;
4use futures::StreamExt;
5use object_store::{ObjectStoreExt, PutPayload};
6
7use super::{CheckpointStoreError, ObjectStoreCheckpointStore};
8use crate::checkpoint::{
9 OutputSegmentRef, SubscriptionDigest, MAX_OUTPUT_FRAMES_PER_SEGMENT, MAX_OUTPUT_SEGMENT_BYTES,
10};
11
12const MAX_ORPHAN_SCAN_OBJECTS: u64 = 1_000_000;
13
14#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
16pub struct SubscriptionOrphanCleanup {
17 pub objects_scanned: u64,
19 pub objects_deleted: u64,
21 pub bytes_deleted: u64,
23 pub bytes_remaining: u64,
25}
26
27pub(super) fn unsupported_store_error() -> CheckpointStoreError {
28 CheckpointStoreError::Invalid(
29 "checkpoint store does not support durable subscription segments".into(),
30 )
31}
32
33pub(super) async fn save(
34 store: &ObjectStoreCheckpointStore,
35 segment: &OutputSegmentRef,
36 payload: Bytes,
37) -> Result<(), CheckpointStoreError> {
38 validate_reference(segment)?;
39 validate_payload(segment, &payload)?;
40 let path = segment_path(store, &segment.object_key)?;
41 if store
42 .create_immutable(&path, PutPayload::from_bytes(payload.clone()))
43 .await?
44 {
45 return Ok(());
46 }
47 match load(store, segment).await? {
48 Some(existing) if existing == payload => Ok(()),
49 Some(_) => Err(CheckpointStoreError::Invalid(format!(
50 "subscription segment '{}' already exists with conflicting immutable content",
51 segment.object_key
52 ))),
53 None => Err(CheckpointStoreError::Invalid(format!(
54 "subscription segment '{}' create conflicted but no object exists",
55 segment.object_key
56 ))),
57 }
58}
59
60pub(super) async fn load(
61 store: &ObjectStoreCheckpointStore,
62 segment: &OutputSegmentRef,
63) -> Result<Option<Bytes>, CheckpointStoreError> {
64 validate_reference(segment)?;
65 let path = segment_path(store, &segment.object_key)?;
66 let result = match store.store.get(&path).await {
67 Ok(result) => result,
68 Err(object_store::Error::NotFound { .. }) => return Ok(None),
69 Err(error) => return Err(error.into()),
70 };
71 if result.meta.size != segment.encoded_length {
72 return Err(CheckpointStoreError::Invalid(format!(
73 "subscription segment '{}' is {} bytes; expected {}",
74 segment.object_key, result.meta.size, segment.encoded_length
75 )));
76 }
77 let bytes = result.bytes().await?;
78 validate_payload(segment, &bytes)?;
79 Ok(Some(bytes))
80}
81
82pub(super) async fn delete(
83 store: &ObjectStoreCheckpointStore,
84 object_key: &str,
85) -> Result<(), CheckpointStoreError> {
86 let path = segment_path(store, object_key)?;
87 store.delete_exact(&path).await
88}
89
90pub(super) async fn delete_orphans(
91 store: &ObjectStoreCheckpointStore,
92 reachable: &std::collections::BTreeSet<String>,
93 through_checkpoint_id: u64,
94 grace_before_ms: i64,
95) -> Result<SubscriptionOrphanCleanup, CheckpointStoreError> {
96 if through_checkpoint_id == 0 || grace_before_ms <= 0 {
97 return Err(CheckpointStoreError::Invalid(
98 "subscription orphan cleanup authority is not canonical".into(),
99 ));
100 }
101 let prefix = object_store::path::Path::from(format!("{}subscription-output", store.prefix));
102 let mut listed = store.store.list(Some(&prefix));
103 let mut report = SubscriptionOrphanCleanup::default();
104 while let Some(entry) = listed.next().await {
105 let entry = entry?;
106 report.objects_scanned = report.objects_scanned.checked_add(1).ok_or_else(|| {
107 CheckpointStoreError::Invalid("subscription orphan scan count overflow".into())
108 })?;
109 if report.objects_scanned > MAX_ORPHAN_SCAN_OBJECTS {
110 return Err(CheckpointStoreError::Invalid(format!(
111 "subscription orphan scan exceeds {MAX_ORPHAN_SCAN_OBJECTS} objects"
112 )));
113 }
114 let full_key = entry.location.to_string();
115 let object_key = full_key.strip_prefix(&store.prefix).ok_or_else(|| {
116 CheckpointStoreError::Invalid(
117 "subscription object lies outside its checkpoint namespace".into(),
118 )
119 })?;
120 let checkpoint_id = checkpoint_id_from_canonical_key(object_key)?;
121 if checkpoint_id > through_checkpoint_id || reachable.contains(object_key) {
122 continue;
123 }
124 if entry.last_modified.timestamp_millis() > grace_before_ms {
125 report.bytes_remaining =
126 report
127 .bytes_remaining
128 .checked_add(entry.size)
129 .ok_or_else(|| {
130 CheckpointStoreError::Invalid(
131 "subscription orphan byte count overflow".into(),
132 )
133 })?;
134 continue;
135 }
136 store.delete_exact(&entry.location).await?;
137 report.objects_deleted = report.objects_deleted.checked_add(1).ok_or_else(|| {
138 CheckpointStoreError::Invalid("subscription orphan delete count overflow".into())
139 })?;
140 report.bytes_deleted = report
141 .bytes_deleted
142 .checked_add(entry.size)
143 .ok_or_else(|| {
144 CheckpointStoreError::Invalid("subscription orphan byte count overflow".into())
145 })?;
146 }
147 Ok(report)
148}
149
150fn segment_path(
151 store: &ObjectStoreCheckpointStore,
152 object_key: &str,
153) -> Result<object_store::path::Path, CheckpointStoreError> {
154 validate_object_key(object_key)?;
155 Ok(object_store::path::Path::from(format!(
156 "{}{object_key}",
157 store.prefix
158 )))
159}
160
161fn validate_object_key(object_key: &str) -> Result<(), CheckpointStoreError> {
162 if object_key.is_empty()
163 || object_key.len() > 2_048
164 || !object_key.starts_with("subscription-output/")
165 || object_key.starts_with('/')
166 || object_key.contains('\\')
167 || object_key
168 .split('/')
169 .any(|component| component.is_empty() || matches!(component, "." | ".."))
170 {
171 return Err(CheckpointStoreError::Invalid(
172 "subscription segment object key is not canonical".into(),
173 ));
174 }
175 Ok(())
176}
177
178fn checkpoint_id_from_canonical_key(object_key: &str) -> Result<u64, CheckpointStoreError> {
179 let components = object_key.split('/').collect::<Vec<_>>();
180 let canonical = components.len() == 7
181 && components[0] == "subscription-output"
182 && uuid::Uuid::parse_str(components[1])
183 .is_ok_and(|deployment| deployment.to_string() == components[1])
184 && is_lower_sha256(components[2])
185 && is_lower_sha256(components[3])
186 && components[4]
187 .parse::<u16>()
188 .is_ok_and(|partition| partition.to_string() == components[4]);
189 let checkpoint_id = components
190 .get(5)
191 .and_then(|component| component.strip_prefix("checkpoint="))
192 .filter(|value| value.len() == 20)
193 .and_then(|value| value.parse::<u64>().ok())
194 .filter(|checkpoint_id| *checkpoint_id != 0);
195 if !canonical
196 || checkpoint_id.is_none()
197 || !components
198 .get(6)
199 .is_some_and(|filename| canonical_segment_filename(filename))
200 {
201 return Err(CheckpointStoreError::Invalid(
202 "listed subscription segment object key is not canonical".into(),
203 ));
204 }
205 checkpoint_id.ok_or_else(|| {
206 CheckpointStoreError::Invalid("subscription segment checkpoint ID is absent".into())
207 })
208}
209
210fn canonical_segment_filename(filename: &str) -> bool {
211 let Some(stem) = filename.strip_suffix(".arrow") else {
212 return false;
213 };
214 let parts = stem.split('-').collect::<Vec<_>>();
215 if parts.len() != 3
216 || parts[0].len() != 20
217 || parts[1].len() != 20
218 || !is_lower_sha256(parts[2])
219 {
220 return false;
221 }
222 let Some(first) = parts[0].parse::<u64>().ok() else {
223 return false;
224 };
225 parts[1].parse::<u64>().is_ok_and(|end| first < end)
226}
227
228fn is_lower_sha256(value: &str) -> bool {
229 value.len() == 64
230 && value
231 .bytes()
232 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
233}
234
235fn validate_reference(segment: &OutputSegmentRef) -> Result<(), CheckpointStoreError> {
236 validate_object_key(&segment.object_key)?;
237 if segment.encoded_length == 0
238 || segment.encoded_length > u64::try_from(MAX_OUTPUT_SEGMENT_BYTES).unwrap_or(u64::MAX)
239 || segment.frame_count == 0
240 || segment.frame_count > MAX_OUTPUT_FRAMES_PER_SEGMENT
241 || segment.row_count == 0
242 || segment.first_sequence >= segment.exclusive_end_sequence
243 || segment
244 .exclusive_end_sequence
245 .get()
246 .checked_sub(segment.first_sequence.get())
247 != Some(segment.frame_count)
248 {
249 return Err(CheckpointStoreError::Invalid(
250 "subscription segment reference is not canonical".into(),
251 ));
252 }
253 Ok(())
254}
255
256fn validate_payload(
257 segment: &OutputSegmentRef,
258 payload: &[u8],
259) -> Result<(), CheckpointStoreError> {
260 if u64::try_from(payload.len()).ok() != Some(segment.encoded_length) {
261 return Err(CheckpointStoreError::Invalid(format!(
262 "subscription segment '{}' payload length mismatch",
263 segment.object_key
264 )));
265 }
266 let digest = SubscriptionDigest::for_bytes(b"laminardb-subscription-segment-v1", payload);
267 if digest != segment.payload_digest {
268 return Err(CheckpointStoreError::Invalid(format!(
269 "subscription segment '{}' payload digest mismatch",
270 segment.object_key
271 )));
272 }
273 Ok(())
274}