Skip to main content

laminar_core/checkpoint/checkpoint_store/
subscription_segments.rs

1//! Immutable subscription segment persistence and validation.
2
3use 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/// Result of one bounded subscription-segment orphan scan.
15#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
16pub struct SubscriptionOrphanCleanup {
17    /// Objects examined under the subscription-output namespace.
18    pub objects_scanned: u64,
19    /// Grace-expired, unreachable objects deleted by exact key.
20    pub objects_deleted: u64,
21    /// Encoded bytes represented by deleted object metadata.
22    pub bytes_deleted: u64,
23    /// Unreachable bytes retained because their orphan grace period has not elapsed.
24    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}