Skip to main content

laminar_core/state/
object_store.rs

1//! [`ObjectStoreBackend`] — durable partial-state storage backed by any
2//! `object_store` implementation (S3, GCS, Azure, `LocalFileSystem`).
3//!
4//! `seal_checkpoint` performs a CAS seal: if every vnode's `partial.bin`
5//! and every required commit descriptor is present, `put(_SEAL, Create)`
6//! seals the exact checkpoint attempt. The `_SEAL` marker is the durability boundary the
7//! checkpoint coordinator consults before releasing sinks. Retention advances one durable,
8//! monotonic prune floor before deleting artifacts, preventing concurrent or restarted writers
9//! from republishing a retired attempt without accumulating one tombstone per checkpoint.
10
11use std::sync::atomic::{AtomicU64, Ordering};
12use std::sync::Arc;
13
14use async_trait::async_trait;
15use bytes::Bytes;
16use object_store::path::Path as OsPath;
17use object_store::{
18    GetOptions, GetRange, ObjectStore, ObjectStoreExt, PutMode, PutOptions, PutPayload,
19    UpdateVersion,
20};
21
22use crate::checkpoint::{
23    CheckpointAssignmentFence, LeaderProof, LeaderProofOwner, PipelineIdentity,
24};
25
26use super::backend::{
27    digest_hex, sha256, CheckpointAttempt, CheckpointSeal, CheckpointSealInventory,
28    SealedCommitDescriptor, SealedCommitDescriptorWriter, SealedVnodePartial, SealedVnodeWriter,
29    StateBackend, StateBackendDurability, StateBackendError, StateNamespaceBinding,
30    CHECKPOINT_SEAL_VERSION, STATE_NAMESPACE_RESOURCE,
31};
32
33const VNODE_PARTIAL_MAGIC: &[u8; 8] = b"LDBVP2\0\0";
34const VNODE_PARTIAL_VERSION: u32 = 2;
35const VNODE_PARTIAL_HEADER_LEN: usize = 136;
36const PARTIAL_ATTESTATION_READ_CONCURRENCY: usize = 32;
37const COMMIT_DESCRIPTOR_MAGIC: &[u8; 8] = b"LDBCD2\0\0";
38const COMMIT_DESCRIPTOR_VERSION: u32 = 2;
39const COMMIT_DESCRIPTOR_HEADER_LEN: usize = 204;
40const DESCRIPTOR_ATTESTATION_READ_CONCURRENCY: usize = 32;
41const STATE_PRUNE_FLOOR_VERSION: u32 = 1;
42const STATE_PRUNE_FLOOR_MAX_BYTES: u64 = 512;
43const STATE_PRUNE_DELETE_BATCH_SIZE: usize = 256;
44const STATE_NAMESPACE_VERSION: u32 = 1;
45const STATE_NAMESPACE_MAX_BYTES: u64 = 512;
46// MAX_KEY_GROUP_COUNT (65,535) at a conservative 768 encoded bytes of provenance per
47// vnode is under 48 MiB; 64 MiB leaves over 16 MiB for the assignment and descriptors.
48const MAX_CHECKPOINT_SEAL_BYTES: u64 = 64 * 1024 * 1024;
49
50#[derive(Debug, serde::Serialize, serde::Deserialize)]
51#[serde(deny_unknown_fields)]
52struct StateNamespaceMarker {
53    version: u32,
54    deployment_id: String,
55    pipeline_identity: PipelineIdentity,
56}
57
58/// Monotonic publication fence for every state attempt below `before_epoch`.
59///
60/// `swept_before_epoch` is only a repair cursor. Readers and writers fence on
61/// `before_epoch`, so a crash between publishing the floor and deleting old objects cannot make a
62/// pruned checkpoint visible again.
63#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
64struct StatePruneFloor {
65    version: u32,
66    before_epoch: u64,
67    swept_before_epoch: u64,
68}
69
70#[derive(Debug)]
71struct VersionedStatePruneFloor {
72    floor: StatePruneFloor,
73    update_version: UpdateVersion,
74}
75
76/// Object-store-backed [`StateBackend`].
77pub struct ObjectStoreBackend {
78    store: Arc<dyn ObjectStore>,
79    empty_prefix_cleanup: Option<Arc<dyn crate::durable_local_store::EmptyPrefixCleanup>>,
80    durability_scope: StateBackendDurability,
81    instance_id: String,
82    /// Fresh for each backend construction, even when `instance_id` is stable across restarts.
83    execution_id: uuid::Uuid,
84    vnode_capacity: u32,
85    /// Split-brain fence: writes must match this exact `assignment_version`.
86    /// `0` disables the fence, accepting unconfigured single-instance callers.
87    authoritative_version: Arc<AtomicU64>,
88    /// Serializes the node-local fallback for stores that cannot perform a conditional update.
89    prune_floor_update_lock: tokio::sync::Mutex<()>,
90}
91
92impl std::fmt::Debug for ObjectStoreBackend {
93    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
94        f.debug_struct("ObjectStoreBackend")
95            .field("durability_scope", &self.durability_scope)
96            .field("instance_id", &self.instance_id)
97            .field("execution_id", &self.execution_id)
98            .field("vnode_capacity", &self.vnode_capacity)
99            .finish_non_exhaustive()
100    }
101}
102
103impl ObjectStoreBackend {
104    /// Wrap an existing [`ObjectStore`] without certifying persistence.
105    ///
106    /// The opaque trait object does not reveal whether it is an in-memory,
107    /// node-local, or shared implementation, so this conservative constructor
108    /// reports [`StateBackendDurability::Volatile`]. Production hosts should use
109    /// [`Self::node_durable`] or [`Self::cluster_shared`] after establishing the
110    /// storage topology.
111    #[must_use]
112    pub fn new(
113        store: Arc<dyn ObjectStore>,
114        instance_id: impl Into<String>,
115        vnode_capacity: u32,
116    ) -> Self {
117        Self::with_durability_scope(
118            store,
119            instance_id,
120            vnode_capacity,
121            StateBackendDurability::Volatile,
122        )
123    }
124
125    /// Wrap storage that survives restart on this node but is not guaranteed
126    /// to be reachable by cluster peers.
127    #[must_use]
128    pub fn node_durable(
129        store: Arc<dyn ObjectStore>,
130        instance_id: impl Into<String>,
131        vnode_capacity: u32,
132    ) -> Self {
133        Self::with_durability_scope(
134            store,
135            instance_id,
136            vnode_capacity,
137            StateBackendDurability::NodeDurable,
138        )
139    }
140
141    pub(crate) fn node_durable_with_empty_prefix_cleanup<T>(
142        store: Arc<T>,
143        instance_id: impl Into<String>,
144        vnode_capacity: u32,
145    ) -> Self
146    where
147        T: ObjectStore + crate::durable_local_store::EmptyPrefixCleanup + 'static,
148    {
149        let object_store: Arc<dyn ObjectStore> = store.clone();
150        let cleanup: Arc<dyn crate::durable_local_store::EmptyPrefixCleanup> = store;
151        let mut backend = Self::node_durable(object_store, instance_id, vnode_capacity);
152        backend.empty_prefix_cleanup = Some(cleanup);
153        backend
154    }
155
156    /// Wrap durable storage whose namespace is reachable by every cluster node.
157    #[must_use]
158    pub fn cluster_shared(
159        store: Arc<dyn ObjectStore>,
160        instance_id: impl Into<String>,
161        vnode_capacity: u32,
162    ) -> Self {
163        Self::with_durability_scope(
164            store,
165            instance_id,
166            vnode_capacity,
167            StateBackendDurability::ClusterShared,
168        )
169    }
170
171    fn with_durability_scope(
172        store: Arc<dyn ObjectStore>,
173        instance_id: impl Into<String>,
174        vnode_capacity: u32,
175        durability_scope: StateBackendDurability,
176    ) -> Self {
177        let instance_id = instance_id.into();
178        Self {
179            store,
180            empty_prefix_cleanup: None,
181            durability_scope,
182            instance_id,
183            execution_id: uuid::Uuid::new_v4(),
184            vnode_capacity,
185            authoritative_version: Arc::new(AtomicU64::new(0)),
186            prune_floor_update_lock: tokio::sync::Mutex::new(()),
187        }
188    }
189
190    /// Shared handle to the authoritative version counter, cloneable by a
191    /// single owner that drives it without relaying through the trait method.
192    #[must_use]
193    pub fn authoritative_version_handle(&self) -> Arc<AtomicU64> {
194        Arc::clone(&self.authoritative_version)
195    }
196
197    #[cfg(test)]
198    fn execution_id(&self) -> uuid::Uuid {
199        self.execution_id
200    }
201
202    fn check_vnode(&self, v: u32) -> Result<(), StateBackendError> {
203        if v >= self.vnode_capacity {
204            Err(StateBackendError::Io(format!(
205                "vnode {v} out of range (capacity {})",
206                self.vnode_capacity
207            )))
208        } else {
209            Ok(())
210        }
211    }
212
213    fn attempt_prefix(attempt: CheckpointAttempt) -> String {
214        format!(
215            "state-v2/epoch={}/checkpoint={}/",
216            attempt.epoch, attempt.checkpoint_id
217        )
218    }
219
220    fn ensure_canonical_attempt(attempt: CheckpointAttempt) -> Result<(), StateBackendError> {
221        if attempt.is_canonical() {
222            Ok(())
223        } else {
224            Err(StateBackendError::Conflict {
225                resource: Self::attempt_prefix(attempt),
226                message: "state attempt must use one nonzero canonical checkpoint ID".into(),
227            })
228        }
229    }
230
231    fn partial_path(attempt: CheckpointAttempt, vnode: u32) -> OsPath {
232        OsPath::from(format!(
233            "{}vnode={vnode}/partial.bin",
234            Self::attempt_prefix(attempt)
235        ))
236    }
237
238    fn seal_path(attempt: CheckpointAttempt) -> OsPath {
239        OsPath::from(format!("{}_SEAL", Self::attempt_prefix(attempt)))
240    }
241
242    fn descriptor_path(attempt: CheckpointAttempt, key: &str) -> OsPath {
243        OsPath::from(format!("{}commit/{key}", Self::attempt_prefix(attempt)))
244    }
245
246    fn prune_floor_path() -> OsPath {
247        OsPath::from("state-v2/_PRUNE_FLOOR")
248    }
249
250    fn namespace_path() -> OsPath {
251        OsPath::from(STATE_NAMESPACE_RESOURCE)
252    }
253
254    fn check_namespace_marker_size(path: &OsPath, size: u64) -> Result<(), StateBackendError> {
255        if size == 0 || size > STATE_NAMESPACE_MAX_BYTES {
256            return Err(StateBackendError::Conflict {
257                resource: path.to_string(),
258                message: format!(
259                    "state namespace marker is {size} bytes; expected 1..={STATE_NAMESPACE_MAX_BYTES}"
260                ),
261            });
262        }
263        Ok(())
264    }
265
266    fn encode_namespace_marker(
267        binding: &StateNamespaceBinding,
268    ) -> Result<Bytes, StateBackendError> {
269        let marker = StateNamespaceMarker {
270            version: STATE_NAMESPACE_VERSION,
271            deployment_id: binding.deployment_id.clone(),
272            pipeline_identity: binding.pipeline_identity.clone(),
273        };
274        let bytes = serde_json::to_vec(&marker)
275            .map(Bytes::from)
276            .map_err(|error| StateBackendError::Serialization(error.to_string()))?;
277        Self::check_namespace_marker_size(&Self::namespace_path(), bytes.len() as u64)?;
278        Ok(bytes)
279    }
280
281    fn decode_namespace_marker(
282        path: &OsPath,
283        bytes: &[u8],
284    ) -> Result<StateNamespaceBinding, StateBackendError> {
285        Self::check_namespace_marker_size(path, bytes.len() as u64)?;
286        let marker: StateNamespaceMarker =
287            serde_json::from_slice(bytes).map_err(|error| StateBackendError::Conflict {
288                resource: path.to_string(),
289                message: format!("state namespace marker is malformed: {error}"),
290            })?;
291        if marker.version != STATE_NAMESPACE_VERSION {
292            return Err(StateBackendError::Conflict {
293                resource: path.to_string(),
294                message: format!(
295                    "state namespace marker version {} is unsupported; expected {STATE_NAMESPACE_VERSION}",
296                    marker.version
297                ),
298            });
299        }
300        let canonical = serde_json::to_vec(&marker)
301            .map_err(|error| StateBackendError::Serialization(error.to_string()))?;
302        if canonical.as_slice() != bytes {
303            return Err(StateBackendError::Conflict {
304                resource: path.to_string(),
305                message: "state namespace marker is not canonical".into(),
306            });
307        }
308        StateNamespaceBinding::try_new(&marker.deployment_id, &marker.pipeline_identity)
309    }
310
311    async fn read_namespace_binding(
312        &self,
313        path: &OsPath,
314    ) -> Result<StateNamespaceBinding, StateBackendError> {
315        let result = self
316            .store
317            .get(path)
318            .await
319            .map_err(|error| StateBackendError::Io(error.to_string()))?;
320        Self::check_namespace_marker_size(path, result.meta.size)?;
321        let bytes = result
322            .bytes()
323            .await
324            .map_err(|error| StateBackendError::Io(error.to_string()))?;
325        Self::decode_namespace_marker(path, &bytes)
326    }
327
328    fn verify_namespace_binding(
329        path: &OsPath,
330        existing: &StateNamespaceBinding,
331        requested: &StateNamespaceBinding,
332    ) -> Result<(), StateBackendError> {
333        if existing.deployment_id != requested.deployment_id {
334            return Err(StateBackendError::Conflict {
335                resource: path.to_string(),
336                message: format!(
337                    "state root belongs to deployment {}; requested {}",
338                    existing.deployment_id, requested.deployment_id
339                ),
340            });
341        }
342        if existing.pipeline_identity != requested.pipeline_identity {
343            return Err(StateBackendError::Conflict {
344                resource: path.to_string(),
345                message: format!(
346                    "state root pipeline identity {} does not match requested {}",
347                    existing.pipeline_identity.sha256, requested.pipeline_identity.sha256
348                ),
349            });
350        }
351        Ok(())
352    }
353
354    async fn preflight_unbound_state_root(
355        &self,
356        namespace_path: &OsPath,
357        requested: &StateNamespaceBinding,
358    ) -> Result<bool, StateBackendError> {
359        use futures::StreamExt as _;
360
361        let prefix = OsPath::from("state-v2");
362        let mut objects = self.store.list(Some(&prefix));
363        if let Some(result) = objects.next().await {
364            let object = result.map_err(|error| StateBackendError::Io(error.to_string()))?;
365            if object.location == *namespace_path {
366                let existing = self.read_namespace_binding(namespace_path).await?;
367                Self::verify_namespace_binding(namespace_path, &existing, requested)?;
368                return Ok(true);
369            }
370
371            // A concurrent first binder may have published the marker and begun state I/O after
372            // this caller's initial marker lookup. Prefer its immutable binding over a false
373            // legacy-artifact rejection.
374            match self.store.get(namespace_path).await {
375                Ok(marker) => {
376                    Self::check_namespace_marker_size(namespace_path, marker.meta.size)?;
377                    let bytes = marker
378                        .bytes()
379                        .await
380                        .map_err(|error| StateBackendError::Io(error.to_string()))?;
381                    let existing = Self::decode_namespace_marker(namespace_path, &bytes)?;
382                    Self::verify_namespace_binding(namespace_path, &existing, requested)?;
383                    return Ok(true);
384                }
385                Err(object_store::Error::NotFound { .. }) => {
386                    return Err(StateBackendError::Conflict {
387                        resource: namespace_path.to_string(),
388                        message: format!(
389                            "state root contains unbound artifact {}; remove the old state root before reuse",
390                            object.location
391                        ),
392                    });
393                }
394                Err(error) => return Err(StateBackendError::Io(error.to_string())),
395            }
396        }
397        Ok(false)
398    }
399
400    async fn verify_object_size_from_metadata(
401        &self,
402        path: &OsPath,
403        listed_size: Option<u64>,
404        expected_size: u64,
405    ) -> Result<(), StateBackendError> {
406        if listed_size == Some(expected_size) {
407            return Ok(());
408        }
409        match self.store.head(path).await {
410            Ok(metadata) if metadata.size == expected_size => Ok(()),
411            Ok(metadata) => Err(StateBackendError::Conflict {
412                resource: path.to_string(),
413                message: format!(
414                    "sealed artifact is {} bytes in storage metadata; expected {expected_size}",
415                    metadata.size
416                ),
417            }),
418            Err(object_store::Error::NotFound { .. }) => Err(StateBackendError::Conflict {
419                resource: path.to_string(),
420                message: "sealed artifact is missing from storage metadata".into(),
421            }),
422            Err(error) => Err(StateBackendError::Io(error.to_string())),
423        }
424    }
425
426    /// Parse one immediate `state-v2/epoch=N` delimiter prefix.
427    fn epoch_from_prefix(prefix: &OsPath) -> Option<u64> {
428        let encoded = prefix.as_ref().strip_prefix("state-v2/epoch=")?;
429        if encoded.is_empty() || encoded.contains('/') {
430            return None;
431        }
432        let epoch = encoded.parse::<u64>().ok()?;
433        (epoch != 0 && epoch.to_string() == encoded).then_some(epoch)
434    }
435
436    /// Wrap raw operator state in a fixed-width provenance header. The fixed width lets the
437    /// durability gate validate hundreds of vnode generations with small concurrent range GETs
438    /// instead of downloading every state blob again.
439    fn encode_partial(
440        attempt: CheckpointAttempt,
441        vnode: u32,
442        assignment_version: u64,
443        writer: Option<&SealedVnodeWriter>,
444        payload: &Bytes,
445    ) -> Bytes {
446        let payload_digest = sha256(payload);
447        let mut encoded = Vec::with_capacity(VNODE_PARTIAL_HEADER_LEN + payload.len());
448        encoded.extend_from_slice(VNODE_PARTIAL_MAGIC);
449        encoded.extend_from_slice(&VNODE_PARTIAL_VERSION.to_be_bytes());
450        encoded.extend_from_slice(&attempt.epoch.to_be_bytes());
451        encoded.extend_from_slice(&attempt.checkpoint_id.to_be_bytes());
452        encoded.extend_from_slice(&vnode.to_be_bytes());
453        encoded.extend_from_slice(&assignment_version.to_be_bytes());
454        if let Some(writer) = writer {
455            encoded.extend_from_slice(&writer.node_id.to_be_bytes());
456            encoded.extend_from_slice(writer.boot_incarnation.as_bytes());
457            encoded.extend_from_slice(&writer.assignment_certificate_digest);
458        } else {
459            encoded.extend_from_slice(&0_u64.to_be_bytes());
460            encoded.extend_from_slice(uuid::Uuid::nil().as_bytes());
461            encoded.extend_from_slice(&[0; 32]);
462        }
463        encoded.extend_from_slice(&(payload.len() as u64).to_be_bytes());
464        encoded.extend_from_slice(&payload_digest);
465        debug_assert_eq!(encoded.len(), VNODE_PARTIAL_HEADER_LEN);
466        encoded.extend_from_slice(payload);
467        Bytes::from(encoded)
468    }
469
470    fn parse_partial_header(
471        header: &[u8],
472        expected_attempt: CheckpointAttempt,
473        expected_vnode: u32,
474    ) -> Result<SealedVnodePartial, StateBackendError> {
475        fn field<const N: usize>(
476            header: &[u8],
477            start: usize,
478        ) -> Result<[u8; N], StateBackendError> {
479            header
480                .get(start..start + N)
481                .and_then(|bytes| bytes.try_into().ok())
482                .ok_or_else(|| {
483                    StateBackendError::Serialization(
484                        "truncated vnode partial provenance header".into(),
485                    )
486                })
487        }
488
489        if header.len() < VNODE_PARTIAL_HEADER_LEN
490            || &header[..VNODE_PARTIAL_MAGIC.len()] != VNODE_PARTIAL_MAGIC
491        {
492            return Err(StateBackendError::Serialization(
493                "invalid vnode partial provenance header".into(),
494            ));
495        }
496        let version = u32::from_be_bytes(field(header, 8)?);
497        if version != VNODE_PARTIAL_VERSION {
498            return Err(StateBackendError::Serialization(format!(
499                "unsupported vnode partial version {version}; expected {VNODE_PARTIAL_VERSION}"
500            )));
501        }
502        let attempt = CheckpointAttempt::new(
503            u64::from_be_bytes(field(header, 12)?),
504            u64::from_be_bytes(field(header, 20)?),
505        );
506        let vnode = u32::from_be_bytes(field(header, 28)?);
507        if attempt != expected_attempt || vnode != expected_vnode {
508            return Err(StateBackendError::Conflict {
509                resource: Self::partial_path(expected_attempt, expected_vnode).to_string(),
510                message: format!(
511                    "partial header names attempt {attempt:?} vnode {vnode}, expected attempt \
512                     {expected_attempt:?} vnode {expected_vnode}"
513                ),
514            });
515        }
516        let assignment_version = u64::from_be_bytes(field(header, 32)?);
517        let writer_node_id = u64::from_be_bytes(field(header, 40)?);
518        let writer_boot_incarnation = uuid::Uuid::from_bytes(field(header, 48)?);
519        let assignment_certificate_digest = field::<32>(header, 64)?;
520        let writer = if writer_node_id == 0
521            && writer_boot_incarnation.is_nil()
522            && assignment_certificate_digest == [0; 32]
523        {
524            None
525        } else if writer_node_id != 0
526            && !writer_boot_incarnation.is_nil()
527            && assignment_certificate_digest != [0; 32]
528        {
529            Some(SealedVnodeWriter {
530                node_id: writer_node_id,
531                boot_incarnation: writer_boot_incarnation,
532                assignment_certificate_digest,
533            })
534        } else {
535            return Err(StateBackendError::Serialization(
536                "incomplete vnode partial writer certificate".into(),
537            ));
538        };
539        let payload_len = u64::from_be_bytes(field(header, 96)?);
540        let payload_digest = field::<32>(header, 104)?;
541        Ok(SealedVnodePartial {
542            vnode,
543            assignment_version,
544            writer,
545            payload_len,
546            payload_sha256: digest_hex(&payload_digest),
547        })
548    }
549
550    fn decode_partial(
551        bytes: &Bytes,
552        expected_attempt: CheckpointAttempt,
553        expected_vnode: u32,
554    ) -> Result<Bytes, StateBackendError> {
555        const ARCHIVE_ALIGNMENT: usize = rkyv::util::AlignedVec::<16>::ALIGNMENT;
556
557        let metadata = Self::parse_partial_header(bytes, expected_attempt, expected_vnode)?;
558        let payload_len = usize::try_from(metadata.payload_len).map_err(|_| {
559            StateBackendError::Serialization("vnode partial payload length overflows usize".into())
560        })?;
561        if bytes.len() != VNODE_PARTIAL_HEADER_LEN.saturating_add(payload_len) {
562            return Err(StateBackendError::Serialization(format!(
563                "vnode partial payload length mismatch: header={} actual={}",
564                metadata.payload_len,
565                bytes.len().saturating_sub(VNODE_PARTIAL_HEADER_LEN)
566            )));
567        }
568        let payload = bytes.slice(VNODE_PARTIAL_HEADER_LEN..);
569        if metadata.payload_sha256 != digest_hex(&sha256(&payload)) {
570            return Err(StateBackendError::Serialization(
571                "vnode partial payload checksum mismatch".into(),
572            ));
573        }
574        if payload.is_empty() || payload.as_ptr().align_offset(ARCHIVE_ALIGNMENT) == 0 {
575            return Ok(payload);
576        }
577
578        // Object-store clients may expose a view into an arbitrarily aligned network buffer.
579        // Normalize it once here so recovery-chain consumers can validate and decode the rkyv
580        // payload repeatedly without making a fresh aligned copy on every pass.
581        let mut aligned = rkyv::util::AlignedVec::<16>::with_capacity(payload.len());
582        aligned.extend_from_slice(&payload);
583        Ok(Bytes::from_owner(aligned))
584    }
585
586    async fn read_partial_attestation(
587        &self,
588        attempt: CheckpointAttempt,
589        vnode: u32,
590    ) -> Result<Option<SealedVnodePartial>, StateBackendError> {
591        let path = Self::partial_path(attempt, vnode);
592        match self
593            .store
594            .get_range(&path, 0..VNODE_PARTIAL_HEADER_LEN as u64)
595            .await
596        {
597            Ok(header) => Self::parse_partial_header(&header, attempt, vnode).map(Some),
598            Err(object_store::Error::NotFound { .. }) => Ok(None),
599            Err(error) => Err(StateBackendError::Io(error.to_string())),
600        }
601    }
602
603    /// Wrap a coordinated-commit descriptor in a fixed-width identity and content header. Seal
604    /// publication reads only this header; recovery validates the complete payload before use.
605    fn encode_commit_descriptor(
606        attempt: CheckpointAttempt,
607        key: &str,
608        assignment_version: u64,
609        writer: Option<&SealedCommitDescriptorWriter>,
610        payload: &Bytes,
611    ) -> Bytes {
612        let key_digest = sha256(key.as_bytes());
613        let payload_digest = sha256(payload);
614        let mut encoded = Vec::with_capacity(COMMIT_DESCRIPTOR_HEADER_LEN + payload.len());
615        encoded.extend_from_slice(COMMIT_DESCRIPTOR_MAGIC);
616        encoded.extend_from_slice(&COMMIT_DESCRIPTOR_VERSION.to_be_bytes());
617        encoded.extend_from_slice(&attempt.epoch.to_be_bytes());
618        encoded.extend_from_slice(&attempt.checkpoint_id.to_be_bytes());
619        encoded.extend_from_slice(&key_digest);
620        encoded.extend_from_slice(&assignment_version.to_be_bytes());
621        if let Some(writer) = writer {
622            encoded.extend_from_slice(&writer.assignment_certificate_digest);
623            encoded.extend_from_slice(&writer.participant.node_id.to_be_bytes());
624            encoded.extend_from_slice(writer.participant.boot_incarnation.as_bytes());
625            encoded.extend_from_slice(&writer.leader_proof.owner.node_id.to_be_bytes());
626            encoded.extend_from_slice(writer.leader_proof.owner.boot_id.as_bytes());
627            encoded.extend_from_slice(&writer.leader_proof.owner.process_term.to_be_bytes());
628            encoded.extend_from_slice(&writer.leader_proof.fencing_token.to_be_bytes());
629        } else {
630            encoded.extend_from_slice(&[0; 32]);
631            encoded.extend_from_slice(&0_u64.to_be_bytes());
632            encoded.extend_from_slice(uuid::Uuid::nil().as_bytes());
633            encoded.extend_from_slice(&0_u64.to_be_bytes());
634            encoded.extend_from_slice(uuid::Uuid::nil().as_bytes());
635            encoded.extend_from_slice(&0_u64.to_be_bytes());
636            encoded.extend_from_slice(&0_u64.to_be_bytes());
637        }
638        encoded.extend_from_slice(&(payload.len() as u64).to_be_bytes());
639        encoded.extend_from_slice(&payload_digest);
640        debug_assert_eq!(encoded.len(), COMMIT_DESCRIPTOR_HEADER_LEN);
641        encoded.extend_from_slice(payload);
642        Bytes::from(encoded)
643    }
644
645    fn parse_commit_descriptor_header(
646        header: &[u8],
647        expected_attempt: CheckpointAttempt,
648        expected_key: &str,
649    ) -> Result<SealedCommitDescriptor, StateBackendError> {
650        fn field<const N: usize>(
651            header: &[u8],
652            start: usize,
653        ) -> Result<[u8; N], StateBackendError> {
654            header
655                .get(start..start + N)
656                .and_then(|bytes| bytes.try_into().ok())
657                .ok_or_else(|| {
658                    StateBackendError::Serialization(
659                        "truncated commit descriptor provenance header".into(),
660                    )
661                })
662        }
663
664        if header.len() < COMMIT_DESCRIPTOR_HEADER_LEN
665            || &header[..COMMIT_DESCRIPTOR_MAGIC.len()] != COMMIT_DESCRIPTOR_MAGIC
666        {
667            return Err(StateBackendError::Serialization(
668                "invalid commit descriptor provenance header".into(),
669            ));
670        }
671        let version = u32::from_be_bytes(field(header, 8)?);
672        if version != COMMIT_DESCRIPTOR_VERSION {
673            return Err(StateBackendError::Serialization(format!(
674                "unsupported commit descriptor version {version}; expected \
675                 {COMMIT_DESCRIPTOR_VERSION}"
676            )));
677        }
678        let attempt = CheckpointAttempt::new(
679            u64::from_be_bytes(field(header, 12)?),
680            u64::from_be_bytes(field(header, 20)?),
681        );
682        let key_digest = field::<32>(header, 28)?;
683        if attempt != expected_attempt || key_digest != sha256(expected_key.as_bytes()) {
684            return Err(StateBackendError::Conflict {
685                resource: Self::descriptor_path(expected_attempt, expected_key).to_string(),
686                message: format!(
687                    "descriptor header names attempt {attempt:?} key digest {}, expected attempt \
688                     {expected_attempt:?} key digest {}",
689                    digest_hex(&key_digest),
690                    digest_hex(&sha256(expected_key.as_bytes()))
691                ),
692            });
693        }
694
695        let assignment_version = u64::from_be_bytes(field(header, 60)?);
696        let assignment_certificate_digest = field::<32>(header, 68)?;
697        let writer_node_id = u64::from_be_bytes(field(header, 100)?);
698        let writer_boot_incarnation = uuid::Uuid::from_bytes(field(header, 108)?);
699        let leader_node_id = u64::from_be_bytes(field(header, 124)?);
700        let leader_boot_id = uuid::Uuid::from_bytes(field(header, 132)?);
701        let leader_process_term = u64::from_be_bytes(field(header, 148)?);
702        let leader_fencing_token = u64::from_be_bytes(field(header, 156)?);
703        let local_provenance = assignment_version == 0
704            && assignment_certificate_digest == [0; 32]
705            && writer_node_id == 0
706            && writer_boot_incarnation.is_nil()
707            && leader_node_id == 0
708            && leader_boot_id.is_nil()
709            && leader_process_term == 0
710            && leader_fencing_token == 0;
711        let writer = if local_provenance {
712            None
713        } else {
714            let leader_proof = LeaderProof {
715                owner: LeaderProofOwner {
716                    node_id: leader_node_id,
717                    boot_id: leader_boot_id,
718                    process_term: leader_process_term,
719                },
720                fencing_token: leader_fencing_token,
721            };
722            if assignment_version == 0
723                || assignment_certificate_digest == [0; 32]
724                || writer_node_id == 0
725                || writer_boot_incarnation.is_nil()
726                || !leader_proof.is_canonical()
727            {
728                return Err(StateBackendError::Serialization(
729                    "incomplete commit descriptor writer certificate".into(),
730                ));
731            }
732            Some(SealedCommitDescriptorWriter {
733                participant: crate::checkpoint::CheckpointParticipant {
734                    node_id: writer_node_id,
735                    boot_incarnation: writer_boot_incarnation,
736                },
737                assignment_certificate_digest,
738                leader_proof,
739            })
740        };
741        let payload_len = u64::from_be_bytes(field(header, 164)?);
742        let payload_digest = field::<32>(header, 172)?;
743        Ok(SealedCommitDescriptor {
744            key: expected_key.to_owned(),
745            assignment_version,
746            writer,
747            payload_len,
748            payload_sha256: digest_hex(&payload_digest),
749        })
750    }
751
752    fn decode_commit_descriptor(
753        bytes: &Bytes,
754        expected_attempt: CheckpointAttempt,
755        expected_key: &str,
756    ) -> Result<Bytes, StateBackendError> {
757        Self::decode_commit_descriptor_with_attestation(bytes, expected_attempt, expected_key)
758            .map(|(_, payload)| payload)
759    }
760
761    fn decode_commit_descriptor_with_attestation(
762        bytes: &Bytes,
763        expected_attempt: CheckpointAttempt,
764        expected_key: &str,
765    ) -> Result<(SealedCommitDescriptor, Bytes), StateBackendError> {
766        let metadata = Self::parse_commit_descriptor_header(bytes, expected_attempt, expected_key)?;
767        let payload_len = usize::try_from(metadata.payload_len).map_err(|_| {
768            StateBackendError::Serialization(
769                "commit descriptor payload length overflows usize".into(),
770            )
771        })?;
772        if bytes.len() != COMMIT_DESCRIPTOR_HEADER_LEN.saturating_add(payload_len) {
773            return Err(StateBackendError::Serialization(format!(
774                "commit descriptor payload length mismatch: header={} actual={}",
775                metadata.payload_len,
776                bytes.len().saturating_sub(COMMIT_DESCRIPTOR_HEADER_LEN)
777            )));
778        }
779        let payload = bytes.slice(COMMIT_DESCRIPTOR_HEADER_LEN..);
780        if metadata.payload_sha256 != digest_hex(&sha256(&payload)) {
781            return Err(StateBackendError::Serialization(
782                "commit descriptor payload checksum mismatch".into(),
783            ));
784        }
785        Ok((metadata, payload))
786    }
787
788    async fn read_commit_descriptor_attestation(
789        &self,
790        attempt: CheckpointAttempt,
791        key: &str,
792    ) -> Result<Option<SealedCommitDescriptor>, StateBackendError> {
793        let path = Self::descriptor_path(attempt, key);
794        let options = GetOptions {
795            range: Some(GetRange::Bounded(0..COMMIT_DESCRIPTOR_HEADER_LEN as u64)),
796            ..GetOptions::default()
797        };
798        match self.store.get_opts(&path, options).await {
799            Ok(result) => {
800                let object_size = result.meta.size;
801                let header = result
802                    .bytes()
803                    .await
804                    .map_err(|error| StateBackendError::Io(error.to_string()))?;
805                let attestation = Self::parse_commit_descriptor_header(&header, attempt, key)?;
806                let expected_size = (COMMIT_DESCRIPTOR_HEADER_LEN as u64)
807                    .checked_add(attestation.payload_len)
808                    .ok_or_else(|| StateBackendError::Conflict {
809                        resource: path.to_string(),
810                        message: "commit descriptor declared length overflows object size".into(),
811                    })?;
812                if object_size != expected_size {
813                    return Err(StateBackendError::Conflict {
814                        resource: path.to_string(),
815                        message: format!(
816                            "commit descriptor declared {} payload bytes but its stored object is \
817                             {object_size} bytes",
818                            attestation.payload_len
819                        ),
820                    });
821                }
822                Ok(Some(attestation))
823            }
824            Err(object_store::Error::NotFound { .. }) => Ok(None),
825            Err(error) => Err(StateBackendError::Io(error.to_string())),
826        }
827    }
828
829    async fn read_prune_floor(
830        &self,
831    ) -> Result<Option<VersionedStatePruneFloor>, StateBackendError> {
832        let path = Self::prune_floor_path();
833        let result = match self.store.get(&path).await {
834            Ok(result) => result,
835            Err(object_store::Error::NotFound { .. }) => return Ok(None),
836            Err(error) => return Err(StateBackendError::Io(error.to_string())),
837        };
838        if result.meta.size == 0 || result.meta.size > STATE_PRUNE_FLOOR_MAX_BYTES {
839            return Err(StateBackendError::Conflict {
840                resource: path.to_string(),
841                message: format!(
842                    "state prune floor is {} bytes; expected 1..={STATE_PRUNE_FLOOR_MAX_BYTES}",
843                    result.meta.size
844                ),
845            });
846        }
847        let update_version = UpdateVersion {
848            e_tag: result.meta.e_tag.clone(),
849            version: result.meta.version.clone(),
850        };
851        let bytes = result
852            .bytes()
853            .await
854            .map_err(|error| StateBackendError::Io(error.to_string()))?;
855        let floor: StatePruneFloor =
856            serde_json::from_slice(&bytes).map_err(|error| StateBackendError::Conflict {
857                resource: path.to_string(),
858                message: format!("invalid state prune floor: {error}"),
859            })?;
860        if floor.version != STATE_PRUNE_FLOOR_VERSION
861            || floor.before_epoch == 0
862            || floor.swept_before_epoch > floor.before_epoch
863        {
864            return Err(StateBackendError::Conflict {
865                resource: path.to_string(),
866                message: "state prune floor has a non-canonical version or horizon".into(),
867            });
868        }
869        let canonical = serde_json::to_vec(&floor)
870            .map_err(|error| StateBackendError::Serialization(error.to_string()))?;
871        if canonical.as_slice() != bytes.as_ref() {
872            return Err(StateBackendError::Conflict {
873                resource: path.to_string(),
874                message: "state prune floor does not use its canonical body".into(),
875            });
876        }
877        Ok(Some(VersionedStatePruneFloor {
878            floor,
879            update_version,
880        }))
881    }
882
883    async fn attempt_is_pruned(
884        &self,
885        attempt: CheckpointAttempt,
886    ) -> Result<bool, StateBackendError> {
887        Self::ensure_canonical_attempt(attempt)?;
888        Ok(self
889            .read_prune_floor()
890            .await?
891            .is_some_and(|versioned| attempt.epoch < versioned.floor.before_epoch))
892    }
893
894    async fn ensure_attempt_live(
895        &self,
896        attempt: CheckpointAttempt,
897    ) -> Result<(), StateBackendError> {
898        Self::ensure_canonical_attempt(attempt)?;
899        if let Some(versioned) = self.read_prune_floor().await? {
900            if attempt.epoch < versioned.floor.before_epoch {
901                return Err(StateBackendError::Conflict {
902                    resource: Self::attempt_prefix(attempt),
903                    message: format!(
904                        "checkpoint epoch {} is below durable state prune floor {}",
905                        attempt.epoch, versioned.floor.before_epoch
906                    ),
907                });
908            }
909        }
910        Ok(())
911    }
912
913    async fn put_live_immutable(
914        &self,
915        attempt: CheckpointAttempt,
916        path: &OsPath,
917        bytes: Bytes,
918    ) -> Result<(), StateBackendError> {
919        self.ensure_attempt_live(attempt).await?;
920        self.put_immutable(path, bytes).await?;
921        let Some(floor) = self.read_prune_floor().await? else {
922            return Ok(());
923        };
924        if attempt.epoch >= floor.floor.before_epoch {
925            return Ok(());
926        }
927        match self.store.delete(path).await {
928            Ok(()) | Err(object_store::Error::NotFound { .. }) => {}
929            Err(delete_error) => tracing::warn!(
930                %delete_error,
931                path = %path,
932                "state prune: failed to remove a late immutable artifact"
933            ),
934        }
935        Err(StateBackendError::Conflict {
936            resource: Self::attempt_prefix(attempt),
937            message: format!(
938                "checkpoint epoch {} fell below durable state prune floor {} during publication",
939                attempt.epoch, floor.floor.before_epoch
940            ),
941        })
942    }
943
944    async fn compare_and_swap_prune_floor(
945        &self,
946        floor: &StatePruneFloor,
947        expected: Option<UpdateVersion>,
948    ) -> Result<bool, StateBackendError> {
949        let path = Self::prune_floor_path();
950        let bytes = serde_json::to_vec(floor)
951            .map(Bytes::from)
952            .map_err(|error| StateBackendError::Serialization(error.to_string()))?;
953        let options = PutOptions {
954            mode: expected.clone().map_or(PutMode::Create, PutMode::Update),
955            ..PutOptions::default()
956        };
957        match self
958            .store
959            .put_opts(&path, PutPayload::from(bytes.clone()), options)
960            .await
961        {
962            Ok(_) => Ok(true),
963            Err(
964                object_store::Error::Precondition { .. }
965                | object_store::Error::AlreadyExists { .. }
966                | object_store::Error::NotFound { .. },
967            ) => Ok(false),
968            Err(object_store::Error::NotImplemented { .. })
969                if expected.is_some()
970                    && self.durability_scope != StateBackendDurability::ClusterShared =>
971            {
972                // `LocalFileSystem` provides atomic overwrite but no conditional update. Local
973                // runtimes have one process owner for a state namespace, so serialize the
974                // read/compare/overwrite within that owner. Cluster-shared storage must provide
975                // native compare-and-swap and never takes this weaker path.
976                let _guard = self.prune_floor_update_lock.lock().await;
977                let current = self.read_prune_floor().await?;
978                if current.as_ref().map(|value| &value.update_version) != expected.as_ref() {
979                    return Ok(false);
980                }
981                let overwrite = PutOptions {
982                    mode: PutMode::Overwrite,
983                    ..PutOptions::default()
984                };
985                self.store
986                    .put_opts(&path, PutPayload::from(bytes), overwrite)
987                    .await
988                    .map(|_| true)
989                    .map_err(|error| StateBackendError::Io(error.to_string()))
990            }
991            Err(error) => match self.read_prune_floor().await? {
992                Some(current)
993                    if current.floor.before_epoch >= floor.before_epoch
994                        && current.floor.swept_before_epoch >= floor.swept_before_epoch =>
995                {
996                    Ok(true)
997                }
998                _ => Err(StateBackendError::Io(error.to_string())),
999            },
1000        }
1001    }
1002
1003    async fn delete_retired_prefix(&self, prefix: &OsPath) -> Result<(), StateBackendError> {
1004        use futures::StreamExt;
1005
1006        loop {
1007            // Materialize one bounded batch before mutating the prefix. LocalFileSystem resumes
1008            // WalkDir in chunks, and deleting from that live iterator can skip entries. Re-listing
1009            // from the prefix root also bounds memory for the maximum key-group topology.
1010            let mut entries = self.store.list(Some(prefix));
1011            let mut locations = Vec::with_capacity(STATE_PRUNE_DELETE_BATCH_SIZE);
1012            while locations.len() < STATE_PRUNE_DELETE_BATCH_SIZE {
1013                let Some(entry) = entries.next().await else {
1014                    break;
1015                };
1016                locations.push(
1017                    entry
1018                        .map_err(|error| StateBackendError::Io(error.to_string()))?
1019                        .location,
1020                );
1021            }
1022            drop(entries);
1023            if locations.is_empty() {
1024                if let Some(cleanup) = &self.empty_prefix_cleanup {
1025                    cleanup
1026                        .cleanup_empty_prefix(prefix)
1027                        .await
1028                        .map_err(|error| StateBackendError::Io(error.to_string()))?;
1029                }
1030                return Ok(());
1031            }
1032
1033            let expected = locations.len();
1034            let input =
1035                futures::stream::iter(locations.into_iter().map(Ok::<_, object_store::Error>))
1036                    .boxed();
1037            let mut deletes = self.store.delete_stream(input);
1038            let mut completed = 0_usize;
1039            while let Some(result) = deletes.next().await {
1040                match result {
1041                    Ok(_) | Err(object_store::Error::NotFound { .. }) => {
1042                        completed += 1;
1043                    }
1044                    Err(error) => {
1045                        return Err(StateBackendError::Io(format!(
1046                            "state backend prune failed to delete an artifact: {error}"
1047                        )));
1048                    }
1049                }
1050            }
1051            if completed != expected {
1052                return Err(StateBackendError::Io(format!(
1053                    "state backend prune delete stream ended after {completed} of {expected} artifacts"
1054                )));
1055            }
1056            tokio::task::yield_now().await;
1057        }
1058    }
1059}
1060
1061#[async_trait]
1062impl StateBackend for ObjectStoreBackend {
1063    fn key_group_capacity(&self) -> u32 {
1064        self.vnode_capacity
1065    }
1066
1067    async fn bind_state_namespace(
1068        &self,
1069        deployment_id: &str,
1070        pipeline_identity: &PipelineIdentity,
1071    ) -> Result<(), StateBackendError> {
1072        let requested = StateNamespaceBinding::try_new(deployment_id, pipeline_identity)?;
1073        let path = Self::namespace_path();
1074        match self.store.get(&path).await {
1075            Ok(result) => {
1076                Self::check_namespace_marker_size(&path, result.meta.size)?;
1077                let bytes = result
1078                    .bytes()
1079                    .await
1080                    .map_err(|error| StateBackendError::Io(error.to_string()))?;
1081                let existing = Self::decode_namespace_marker(&path, &bytes)?;
1082                return Self::verify_namespace_binding(&path, &existing, &requested);
1083            }
1084            Err(object_store::Error::NotFound { .. }) => {}
1085            Err(error) => return Err(StateBackendError::Io(error.to_string())),
1086        }
1087
1088        // No marker may claim an already-populated root. All conforming writers bind before state
1089        // I/O, so concurrent first binders can both observe an empty root; create-only publication
1090        // below chooses the winner and the loser compares that immutable winner.
1091        if self.preflight_unbound_state_root(&path, &requested).await? {
1092            return Ok(());
1093        }
1094        let bytes = Self::encode_namespace_marker(&requested)?;
1095        let options = PutOptions {
1096            mode: PutMode::Create,
1097            ..PutOptions::default()
1098        };
1099        match self
1100            .store
1101            .put_opts(&path, PutPayload::from(bytes), options)
1102            .await
1103        {
1104            Ok(_) => Ok(()),
1105            Err(object_store::Error::AlreadyExists { .. }) => {
1106                let existing = self.read_namespace_binding(&path).await?;
1107                Self::verify_namespace_binding(&path, &existing, &requested)
1108            }
1109            Err(error) => Err(StateBackendError::Io(error.to_string())),
1110        }
1111    }
1112
1113    fn durability_scope(&self) -> StateBackendDurability {
1114        self.durability_scope
1115    }
1116
1117    fn uses_exact_object_store(&self, expected: &Arc<dyn ObjectStore>) -> bool {
1118        Arc::ptr_eq(&self.store, expected)
1119    }
1120
1121    async fn write_partial(
1122        &self,
1123        attempt: CheckpointAttempt,
1124        vnode: u32,
1125        assignment_version: u64,
1126        bytes: Bytes,
1127    ) -> Result<(), StateBackendError> {
1128        self.check_vnode(vnode)?;
1129        self.check_assignment_version(assignment_version)?;
1130        let path = Self::partial_path(attempt, vnode);
1131        let bytes = Self::encode_partial(attempt, vnode, assignment_version, None, &bytes);
1132        self.put_live_immutable(attempt, &path, bytes).await
1133    }
1134
1135    async fn write_certified_partial(
1136        &self,
1137        attempt: CheckpointAttempt,
1138        vnode: u32,
1139        assignment_fence: &CheckpointAssignmentFence,
1140        writer_node_id: u64,
1141        bytes: Bytes,
1142    ) -> Result<(), StateBackendError> {
1143        self.check_vnode(vnode)?;
1144        if !assignment_fence.is_canonical() {
1145            return Err(StateBackendError::Conflict {
1146                resource: Self::partial_path(attempt, vnode).to_string(),
1147                message: "assignment certificate is not canonical".into(),
1148            });
1149        }
1150        self.check_assignment_version(assignment_fence.assignment_version)?;
1151        let writer =
1152            SealedVnodeWriter::from_fence(assignment_fence, writer_node_id).ok_or_else(|| {
1153                StateBackendError::Conflict {
1154                    resource: Self::partial_path(attempt, vnode).to_string(),
1155                    message: "partial writer is absent from the canonical assignment certificate"
1156                        .into(),
1157                }
1158            })?;
1159        let encoded = Self::encode_partial(
1160            attempt,
1161            vnode,
1162            assignment_fence.assignment_version,
1163            Some(&writer),
1164            &bytes,
1165        );
1166        let path = Self::partial_path(attempt, vnode);
1167        self.put_live_immutable(attempt, &path, encoded).await
1168    }
1169
1170    async fn read_partial(
1171        &self,
1172        attempt: CheckpointAttempt,
1173        vnode: u32,
1174    ) -> Result<Option<Bytes>, StateBackendError> {
1175        self.check_vnode(vnode)?;
1176        if self.attempt_is_pruned(attempt).await? {
1177            return Ok(None);
1178        }
1179        let path = Self::partial_path(attempt, vnode);
1180        match self.store.get(&path).await {
1181            Ok(res) => {
1182                let b = res
1183                    .bytes()
1184                    .await
1185                    .map_err(|e| StateBackendError::Io(e.to_string()))?;
1186                if self.attempt_is_pruned(attempt).await? {
1187                    Ok(None)
1188                } else {
1189                    Self::decode_partial(&b, attempt, vnode).map(Some)
1190                }
1191            }
1192            Err(object_store::Error::NotFound { .. }) => Ok(None),
1193            Err(e) => Err(StateBackendError::Io(e.to_string())),
1194        }
1195    }
1196
1197    async fn write_commit_descriptor(
1198        &self,
1199        attempt: CheckpointAttempt,
1200        key: &str,
1201        bytes: Bytes,
1202    ) -> Result<(), StateBackendError> {
1203        let path = Self::descriptor_path(attempt, key);
1204        let authoritative = self.authoritative_version();
1205        if authoritative != 0 {
1206            return Err(StateBackendError::Conflict {
1207                resource: path.to_string(),
1208                message: format!(
1209                    "uncertified commit descriptor write is disabled while assignment version \
1210                     {authoritative} is authoritative"
1211                ),
1212            });
1213        }
1214        let encoded = Self::encode_commit_descriptor(attempt, key, 0, None, &bytes);
1215        self.put_live_immutable(attempt, &path, encoded).await
1216    }
1217
1218    async fn write_certified_commit_descriptor(
1219        &self,
1220        attempt: CheckpointAttempt,
1221        key: &str,
1222        assignment_fence: &CheckpointAssignmentFence,
1223        writer_node_id: u64,
1224        leader_proof: &LeaderProof,
1225        bytes: Bytes,
1226    ) -> Result<(), StateBackendError> {
1227        let path = Self::descriptor_path(attempt, key);
1228        if !assignment_fence.is_canonical() {
1229            return Err(StateBackendError::Conflict {
1230                resource: path.to_string(),
1231                message: "assignment certificate is not canonical".into(),
1232            });
1233        }
1234        self.check_assignment_version(assignment_fence.assignment_version)?;
1235        let writer = SealedCommitDescriptorWriter::from_fence(
1236            assignment_fence,
1237            writer_node_id,
1238            leader_proof,
1239        )
1240        .ok_or_else(|| StateBackendError::Conflict {
1241            resource: path.to_string(),
1242            message: "descriptor writer or leader is absent from the canonical assignment \
1243                      certificate"
1244                .into(),
1245        })?;
1246        let encoded = Self::encode_commit_descriptor(
1247            attempt,
1248            key,
1249            assignment_fence.assignment_version,
1250            Some(&writer),
1251            &bytes,
1252        );
1253        self.put_live_immutable(attempt, &path, encoded).await
1254    }
1255
1256    async fn read_commit_descriptor(
1257        &self,
1258        attempt: CheckpointAttempt,
1259        key: &str,
1260    ) -> Result<Option<Bytes>, StateBackendError> {
1261        self.read_commit_descriptor_bounded(attempt, key, u64::MAX)
1262            .await
1263    }
1264
1265    async fn read_commit_descriptor_bounded(
1266        &self,
1267        attempt: CheckpointAttempt,
1268        key: &str,
1269        max_bytes: u64,
1270    ) -> Result<Option<Bytes>, StateBackendError> {
1271        if self.attempt_is_pruned(attempt).await? {
1272            return Ok(None);
1273        }
1274        let path = Self::descriptor_path(attempt, key);
1275        match self.store.get(&path).await {
1276            Ok(result) => {
1277                let max_object_bytes =
1278                    (COMMIT_DESCRIPTOR_HEADER_LEN as u64).saturating_add(max_bytes);
1279                if result.meta.size > max_object_bytes {
1280                    return Err(StateBackendError::Conflict {
1281                        resource: path.to_string(),
1282                        message: format!(
1283                            "commit descriptor payload exceeds its read bound; read bound is \
1284                             {max_bytes} bytes (stored object is {} bytes including its {}-byte \
1285                             header)",
1286                            result.meta.size, COMMIT_DESCRIPTOR_HEADER_LEN
1287                        ),
1288                    });
1289                }
1290                let bytes = result
1291                    .bytes()
1292                    .await
1293                    .map_err(|error| StateBackendError::Io(error.to_string()))?;
1294                if self.attempt_is_pruned(attempt).await? {
1295                    Ok(None)
1296                } else {
1297                    let payload = Self::decode_commit_descriptor(&bytes, attempt, key)?;
1298                    if payload.len() as u64 > max_bytes {
1299                        return Err(StateBackendError::Conflict {
1300                            resource: path.to_string(),
1301                            message: format!(
1302                                "commit descriptor payload is {} bytes; read bound is {max_bytes}",
1303                                payload.len()
1304                            ),
1305                        });
1306                    }
1307                    Ok(Some(payload))
1308                }
1309            }
1310            Err(object_store::Error::NotFound { .. }) => Ok(None),
1311            Err(error) => Err(StateBackendError::Io(error.to_string())),
1312        }
1313    }
1314
1315    async fn read_sealed_commit_descriptor_bounded(
1316        &self,
1317        attempt: CheckpointAttempt,
1318        sealed: &SealedCommitDescriptor,
1319        max_bytes: u64,
1320    ) -> Result<Option<Bytes>, StateBackendError> {
1321        let path = Self::descriptor_path(attempt, &sealed.key);
1322        if sealed.payload_len > max_bytes {
1323            return Err(StateBackendError::Conflict {
1324                resource: path.to_string(),
1325                message: format!(
1326                    "sealed commit descriptor declares {} bytes; read bound is {max_bytes}",
1327                    sealed.payload_len
1328                ),
1329            });
1330        }
1331        if self.attempt_is_pruned(attempt).await? {
1332            return Ok(None);
1333        }
1334        match self.store.get(&path).await {
1335            Ok(result) => {
1336                let expected_object_bytes = (COMMIT_DESCRIPTOR_HEADER_LEN as u64)
1337                    .checked_add(sealed.payload_len)
1338                    .ok_or_else(|| StateBackendError::Conflict {
1339                        resource: path.to_string(),
1340                        message: "sealed commit descriptor length overflows object size".into(),
1341                    })?;
1342                if result.meta.size != expected_object_bytes {
1343                    return Err(StateBackendError::Conflict {
1344                        resource: path.to_string(),
1345                        message: format!(
1346                            "stored commit descriptor is {} bytes; checkpoint seal requires \
1347                             {expected_object_bytes} bytes including its header",
1348                            result.meta.size
1349                        ),
1350                    });
1351                }
1352                let bytes = result
1353                    .bytes()
1354                    .await
1355                    .map_err(|error| StateBackendError::Io(error.to_string()))?;
1356                if self.attempt_is_pruned(attempt).await? {
1357                    return Ok(None);
1358                }
1359                let (current, payload) =
1360                    Self::decode_commit_descriptor_with_attestation(&bytes, attempt, &sealed.key)?;
1361                if &current != sealed {
1362                    return Err(StateBackendError::Conflict {
1363                        resource: path.to_string(),
1364                        message: "stored commit descriptor attestation does not match the checkpoint seal"
1365                            .into(),
1366                    });
1367                }
1368                Ok(Some(payload))
1369            }
1370            Err(object_store::Error::NotFound { .. }) => Ok(None),
1371            Err(error) => Err(StateBackendError::Io(error.to_string())),
1372        }
1373    }
1374
1375    async fn seal_checkpoint(
1376        &self,
1377        attempt: CheckpointAttempt,
1378        assignment_fence: Option<&CheckpointAssignmentFence>,
1379        vnodes: &[u32],
1380        required_descriptors: &[String],
1381    ) -> Result<bool, StateBackendError> {
1382        use rustc_hash::FxHashSet;
1383        use tokio_stream::StreamExt;
1384
1385        Self::ensure_canonical_attempt(attempt)?;
1386        let assignment_version = self.seal_assignment_version(attempt, assignment_fence)?;
1387        self.check_assignment_version(assignment_version)?;
1388        self.ensure_attempt_live(attempt).await?;
1389        let mut required_vnodes = vnodes.to_vec();
1390        required_vnodes.sort_unstable();
1391        required_vnodes.dedup();
1392        let mut required_descriptors = required_descriptors.to_vec();
1393        required_descriptors.sort_unstable();
1394        required_descriptors.dedup();
1395        if required_descriptors.iter().any(String::is_empty) {
1396            return Err(StateBackendError::Conflict {
1397                resource: Self::seal_path(attempt).to_string(),
1398                message: "checkpoint seal descriptor key cannot be empty".into(),
1399            });
1400        }
1401        let seal_path = Self::seal_path(attempt);
1402        match self.store.head(&seal_path).await {
1403            Ok(_) => {
1404                let existing = self.read_seal(&seal_path).await?;
1405                let expected = CheckpointSeal::new(
1406                    self.instance_id.clone(),
1407                    self.execution_id,
1408                    CheckpointSealInventory {
1409                        attempt,
1410                        assignment_fence: assignment_fence.cloned(),
1411                        assignment_version,
1412                        required_vnodes,
1413                        sealed_partials: existing.sealed_partials.clone(),
1414                        required_descriptors,
1415                        sealed_descriptors: existing.sealed_descriptors.clone(),
1416                    },
1417                );
1418                let result = if existing == expected {
1419                    Ok(true)
1420                } else {
1421                    Err(StateBackendError::Conflict {
1422                        resource: seal_path.to_string(),
1423                        message: "existing seal does not match this execution, assignment, or artifact inventory".into(),
1424                    })
1425                };
1426                self.ensure_attempt_live(attempt).await?;
1427                return result;
1428            }
1429            Err(object_store::Error::NotFound { .. }) => {}
1430            Err(e) => return Err(StateBackendError::Io(e.to_string())),
1431        }
1432
1433        for &v in &required_vnodes {
1434            self.check_vnode(v)?;
1435        }
1436
1437        // List once for presence only. Some providers do not populate object size in LIST;
1438        // descriptor length is checked from ranged-GET metadata below.
1439        let prefix = OsPath::from(Self::attempt_prefix(attempt));
1440        let mut entries = self.store.list(Some(&prefix));
1441        let mut found_objects: FxHashSet<OsPath> = FxHashSet::default();
1442        while let Some(entry) = entries.next().await {
1443            let entry = entry.map_err(|e| StateBackendError::Io(e.to_string()))?;
1444            found_objects.insert(entry.location);
1445        }
1446
1447        for &v in &required_vnodes {
1448            let path = Self::partial_path(attempt, v);
1449            if !found_objects.contains(&path) {
1450                return Ok(false);
1451            }
1452        }
1453        // Commit descriptors live under this attempt's `commit/` prefix.
1454        for key in &required_descriptors {
1455            if !found_objects.contains(&Self::descriptor_path(attempt, key)) {
1456                return Ok(false);
1457            }
1458        }
1459
1460        let Some(sealed_partials) = self
1461            .read_sealed_partials(
1462                attempt,
1463                &required_vnodes,
1464                assignment_version,
1465                assignment_fence,
1466            )
1467            .await?
1468        else {
1469            return Ok(false);
1470        };
1471        let Some(sealed_descriptors) = self
1472            .read_sealed_descriptors(attempt, &required_descriptors, assignment_fence)
1473            .await?
1474        else {
1475            return Ok(false);
1476        };
1477
1478        let expected_seal = CheckpointSeal::new(
1479            self.instance_id.clone(),
1480            self.execution_id,
1481            CheckpointSealInventory {
1482                attempt,
1483                assignment_fence: assignment_fence.cloned(),
1484                assignment_version,
1485                required_vnodes,
1486                sealed_partials,
1487                required_descriptors,
1488                sealed_descriptors,
1489            },
1490        );
1491        expected_seal
1492            .validate()
1493            .map_err(|message| StateBackendError::Conflict {
1494                resource: seal_path.to_string(),
1495                message,
1496            })?;
1497
1498        let encoded = serde_json::to_vec(&expected_seal)
1499            .map_err(|e| StateBackendError::Serialization(e.to_string()))?;
1500        Self::check_seal_encoded_size(&seal_path, encoded.len() as u64)?;
1501        let bytes = Bytes::from(encoded);
1502        self.put_live_immutable(attempt, &seal_path, bytes).await?;
1503        Ok(true)
1504    }
1505
1506    async fn checkpoint_seal_inventory(
1507        &self,
1508        attempt: CheckpointAttempt,
1509    ) -> Result<Option<CheckpointSealInventory>, StateBackendError> {
1510        if self.attempt_is_pruned(attempt).await? {
1511            return Ok(None);
1512        }
1513        let path = Self::seal_path(attempt);
1514        match self.store.get(&path).await {
1515            Ok(result) => {
1516                Self::check_seal_encoded_size(&path, result.meta.size)?;
1517                let bytes = match result.bytes().await {
1518                    Ok(bytes) => bytes,
1519                    Err(object_store::Error::NotFound { .. }) => return Ok(None),
1520                    Err(error) => return Err(StateBackendError::Io(error.to_string())),
1521                };
1522                if self.attempt_is_pruned(attempt).await? {
1523                    return Ok(None);
1524                }
1525                let seal = Self::decode_seal(&bytes)?;
1526                if seal.attempt != attempt {
1527                    return Err(StateBackendError::Conflict {
1528                        resource: path.to_string(),
1529                        message: format!(
1530                            "seal body names {:?}, requested {attempt:?}",
1531                            seal.attempt
1532                        ),
1533                    });
1534                }
1535                Ok(Some(seal.inventory()))
1536            }
1537            Err(object_store::Error::NotFound { .. }) => Ok(None),
1538            Err(error) => Err(StateBackendError::Io(error.to_string())),
1539        }
1540    }
1541
1542    async fn verify_checkpoint_artifact_metadata(
1543        &self,
1544        inventory: &CheckpointSealInventory,
1545    ) -> Result<(), StateBackendError> {
1546        use futures::StreamExt as _;
1547
1548        let attempt = inventory.attempt;
1549        if self.attempt_is_pruned(attempt).await? {
1550            return Err(StateBackendError::Conflict {
1551                resource: Self::attempt_prefix(attempt),
1552                message: "sealed attempt is below the durable state prune floor".into(),
1553            });
1554        }
1555
1556        let prefix = OsPath::from(Self::attempt_prefix(attempt));
1557        let mut objects = self.store.list(Some(&prefix));
1558        let mut listed_sizes = rustc_hash::FxHashMap::default();
1559        while let Some(entry) = objects.next().await {
1560            let entry = entry.map_err(|error| StateBackendError::Io(error.to_string()))?;
1561            listed_sizes.insert(entry.location, entry.size);
1562        }
1563
1564        for partial in &inventory.sealed_partials {
1565            let path = Self::partial_path(attempt, partial.vnode);
1566            let header_len = u64::try_from(VNODE_PARTIAL_HEADER_LEN).map_err(|_| {
1567                StateBackendError::Conflict {
1568                    resource: path.to_string(),
1569                    message: "vnode storage header length is not representable".into(),
1570                }
1571            })?;
1572            let expected_size = header_len.checked_add(partial.payload_len).ok_or_else(|| {
1573                StateBackendError::Conflict {
1574                    resource: path.to_string(),
1575                    message: "sealed vnode partial length overflows storage size".into(),
1576                }
1577            })?;
1578            self.verify_object_size_from_metadata(
1579                &path,
1580                listed_sizes.get(&path).copied(),
1581                expected_size,
1582            )
1583            .await?;
1584        }
1585        for descriptor in &inventory.sealed_descriptors {
1586            let path = Self::descriptor_path(attempt, &descriptor.key);
1587            let header_len = u64::try_from(COMMIT_DESCRIPTOR_HEADER_LEN).map_err(|_| {
1588                StateBackendError::Conflict {
1589                    resource: path.to_string(),
1590                    message: "descriptor storage header length is not representable".into(),
1591                }
1592            })?;
1593            let expected_size =
1594                header_len
1595                    .checked_add(descriptor.payload_len)
1596                    .ok_or_else(|| StateBackendError::Conflict {
1597                        resource: path.to_string(),
1598                        message: "sealed commit descriptor length overflows storage size".into(),
1599                    })?;
1600            self.verify_object_size_from_metadata(
1601                &path,
1602                listed_sizes.get(&path).copied(),
1603                expected_size,
1604            )
1605            .await?;
1606        }
1607
1608        if self.attempt_is_pruned(attempt).await? {
1609            return Err(StateBackendError::Conflict {
1610                resource: Self::attempt_prefix(attempt),
1611                message: "sealed attempt was pruned during metadata verification".into(),
1612            });
1613        }
1614        Ok(())
1615    }
1616
1617    async fn prune_before(&self, before: u64) -> Result<(), StateBackendError> {
1618        if before == 0 {
1619            return Ok(());
1620        }
1621
1622        // Publish the correctness boundary first. A failed or interrupted sweep leaves garbage,
1623        // never a readable checkpoint, and the durable sweep cursor makes the next caller repair
1624        // the exact unfinished range.
1625        loop {
1626            let observed = self.read_prune_floor().await?;
1627            if observed
1628                .as_ref()
1629                .is_some_and(|current| current.floor.before_epoch >= before)
1630            {
1631                break;
1632            }
1633            let floor = StatePruneFloor {
1634                version: STATE_PRUNE_FLOOR_VERSION,
1635                before_epoch: before,
1636                swept_before_epoch: observed
1637                    .as_ref()
1638                    .map_or(0, |current| current.floor.swept_before_epoch),
1639            };
1640            let expected = observed.map(|current| current.update_version);
1641            if self.compare_and_swap_prune_floor(&floor, expected).await? {
1642                break;
1643            }
1644            tokio::task::yield_now().await;
1645        }
1646
1647        'sweep: loop {
1648            let mut current =
1649                self.read_prune_floor()
1650                    .await?
1651                    .ok_or_else(|| StateBackendError::Conflict {
1652                        resource: Self::prune_floor_path().to_string(),
1653                        message: "state prune floor disappeared after publication".into(),
1654                    })?;
1655            let target = current.floor.before_epoch;
1656
1657            // Discover materialized epochs rather than issuing one LIST for every numeric ID in
1658            // the retired range. Sparse checkpoint IDs are normal after allocation failures and
1659            // can otherwise turn one retention pass into tens of thousands of remote requests.
1660            // Revisit prefixes below the durable cursor too: a writer whose publication raced the
1661            // floor can leave garbage after an ambiguously failed cleanup, even though readers
1662            // already reject it.
1663            let state_root = OsPath::from("state-v2");
1664            let discovered = self
1665                .store
1666                .list_with_delimiter(Some(&state_root))
1667                .await
1668                .map_err(|error| StateBackendError::Io(error.to_string()))?;
1669            let mut retired_prefixes = discovered
1670                .common_prefixes
1671                .into_iter()
1672                .filter_map(|prefix| {
1673                    let epoch = Self::epoch_from_prefix(&prefix)?;
1674                    (epoch < target).then_some((epoch, prefix))
1675                })
1676                .collect::<Vec<_>>();
1677            retired_prefixes.sort_unstable_by_key(|(epoch, _)| *epoch);
1678
1679            for (epoch, prefix) in retired_prefixes {
1680                self.delete_retired_prefix(&prefix).await?;
1681
1682                // The numeric cursor can represent completion of a sparse materialized prefix:
1683                // every lower prefix was absent from the same delimiter snapshot or was already
1684                // swept. Publishing after deletion makes a crash resume at this prefix unless the
1685                // entire prefix was removed.
1686                let next = epoch.saturating_add(1).min(target);
1687                if current.floor.swept_before_epoch < next {
1688                    let swept = StatePruneFloor {
1689                        swept_before_epoch: next,
1690                        ..current.floor.clone()
1691                    };
1692                    if !self
1693                        .compare_and_swap_prune_floor(&swept, Some(current.update_version.clone()))
1694                        .await?
1695                    {
1696                        tokio::task::yield_now().await;
1697                        continue 'sweep;
1698                    }
1699                    current = self.read_prune_floor().await?.ok_or_else(|| {
1700                        StateBackendError::Conflict {
1701                            resource: Self::prune_floor_path().to_string(),
1702                            message: "state prune floor disappeared after sweep progress".into(),
1703                        }
1704                    })?;
1705                    if current.floor.before_epoch != target {
1706                        continue 'sweep;
1707                    }
1708                }
1709            }
1710
1711            if current.floor.swept_before_epoch >= target {
1712                return Ok(());
1713            }
1714            let swept = StatePruneFloor {
1715                swept_before_epoch: target,
1716                ..current.floor.clone()
1717            };
1718            if self
1719                .compare_and_swap_prune_floor(&swept, Some(current.update_version.clone()))
1720                .await?
1721            {
1722                return Ok(());
1723            }
1724            tokio::task::yield_now().await;
1725        }
1726    }
1727
1728    fn set_authoritative_version(&self, version: u64) {
1729        // CAS loop avoids lowering the version on a late call.
1730        let mut cur = self.authoritative_version.load(Ordering::Acquire);
1731        while version > cur {
1732            match self.authoritative_version.compare_exchange(
1733                cur,
1734                version,
1735                Ordering::AcqRel,
1736                Ordering::Acquire,
1737            ) {
1738                Ok(_) => return,
1739                Err(observed) => cur = observed,
1740            }
1741        }
1742    }
1743
1744    fn authoritative_version(&self) -> u64 {
1745        self.authoritative_version.load(Ordering::Acquire)
1746    }
1747}
1748
1749impl ObjectStoreBackend {
1750    async fn read_sealed_partials(
1751        &self,
1752        attempt: CheckpointAttempt,
1753        required_vnodes: &[u32],
1754        assignment_version: u64,
1755        assignment_fence: Option<&CheckpointAssignmentFence>,
1756    ) -> Result<Option<Vec<SealedVnodePartial>>, StateBackendError> {
1757        let mut sealed_partials = Vec::with_capacity(required_vnodes.len());
1758        for chunk in required_vnodes.chunks(PARTIAL_ATTESTATION_READ_CONCURRENCY) {
1759            let attestations = futures::future::try_join_all(
1760                chunk
1761                    .iter()
1762                    .map(|&vnode| self.read_partial_attestation(attempt, vnode)),
1763            )
1764            .await?;
1765            for attestation in attestations {
1766                let Some(attestation) = attestation else {
1767                    return Ok(None);
1768                };
1769                if attestation.assignment_version != assignment_version {
1770                    return Err(StateBackendError::Conflict {
1771                        resource: Self::partial_path(attempt, attestation.vnode).to_string(),
1772                        message: format!(
1773                            "partial assignment version {} cannot satisfy seal version {assignment_version}",
1774                            attestation.assignment_version
1775                        ),
1776                    });
1777                }
1778                match (assignment_fence, &attestation.writer) {
1779                    (Some(fence), Some(writer)) if writer.matches_fence(fence) => {}
1780                    (None, None) => {}
1781                    _ => {
1782                        return Err(StateBackendError::Conflict {
1783                            resource: Self::partial_path(attempt, attestation.vnode).to_string(),
1784                            message: "partial writer certificate does not match the exact seal assignment"
1785                                .into(),
1786                        });
1787                    }
1788                }
1789                sealed_partials.push(attestation);
1790            }
1791        }
1792        Ok(Some(sealed_partials))
1793    }
1794
1795    async fn read_sealed_descriptors(
1796        &self,
1797        attempt: CheckpointAttempt,
1798        required_descriptors: &[String],
1799        assignment_fence: Option<&CheckpointAssignmentFence>,
1800    ) -> Result<Option<Vec<SealedCommitDescriptor>>, StateBackendError> {
1801        let mut sealed_descriptors = Vec::with_capacity(required_descriptors.len());
1802        for chunk in required_descriptors.chunks(DESCRIPTOR_ATTESTATION_READ_CONCURRENCY) {
1803            let attestations = futures::future::try_join_all(
1804                chunk
1805                    .iter()
1806                    .map(|key| self.read_commit_descriptor_attestation(attempt, key)),
1807            )
1808            .await?;
1809            for (key, attestation) in chunk.iter().zip(attestations) {
1810                let Some(attestation) = attestation else {
1811                    return Ok(None);
1812                };
1813                let path = Self::descriptor_path(attempt, key);
1814                match (assignment_fence, &attestation.writer) {
1815                    (Some(fence), Some(writer))
1816                        if attestation.assignment_version == fence.assignment_version
1817                            && writer.matches_fence(fence) => {}
1818                    (None, None) if attestation.assignment_version == 0 => {}
1819                    _ => {
1820                        return Err(StateBackendError::Conflict {
1821                            resource: path.to_string(),
1822                            message: "descriptor writer certificate does not match the exact seal \
1823                                      assignment"
1824                                .into(),
1825                        });
1826                    }
1827                }
1828                sealed_descriptors.push(attestation);
1829            }
1830        }
1831
1832        Ok(Some(sealed_descriptors))
1833    }
1834
1835    fn seal_assignment_version(
1836        &self,
1837        attempt: CheckpointAttempt,
1838        assignment_fence: Option<&CheckpointAssignmentFence>,
1839    ) -> Result<u64, StateBackendError> {
1840        if assignment_fence.is_some_and(|fence| !fence.is_canonical()) {
1841            return Err(StateBackendError::Conflict {
1842                resource: Self::seal_path(attempt).to_string(),
1843                message: "assignment certificate is not canonical".into(),
1844            });
1845        }
1846        Ok(assignment_fence.map_or_else(
1847            || self.authoritative_version(),
1848            |fence| fence.assignment_version,
1849        ))
1850    }
1851
1852    fn check_assignment_version(&self, caller: u64) -> Result<(), StateBackendError> {
1853        let authoritative = self.authoritative_version.load(Ordering::Acquire);
1854        if authoritative == 0 || caller == authoritative {
1855            return Ok(());
1856        }
1857        if caller < authoritative {
1858            return Err(StateBackendError::StaleVersion {
1859                caller,
1860                authoritative,
1861            });
1862        }
1863        Err(StateBackendError::FutureVersion {
1864            caller,
1865            authoritative,
1866        })
1867    }
1868
1869    fn check_seal_encoded_size(path: &OsPath, size: u64) -> Result<(), StateBackendError> {
1870        if size > MAX_CHECKPOINT_SEAL_BYTES {
1871            return Err(StateBackendError::Conflict {
1872                resource: path.to_string(),
1873                message: format!(
1874                    "checkpoint seal is {size} bytes; maximum is {MAX_CHECKPOINT_SEAL_BYTES}"
1875                ),
1876            });
1877        }
1878        Ok(())
1879    }
1880
1881    /// CAS-create immutable bytes. A retry of the exact bytes succeeds; a
1882    /// different payload at the same key is a hard conflict.
1883    async fn put_immutable(&self, path: &OsPath, bytes: Bytes) -> Result<(), StateBackendError> {
1884        let intended_size = bytes.len() as u64;
1885        let opts = PutOptions {
1886            mode: PutMode::Create,
1887            ..PutOptions::default()
1888        };
1889        match self
1890            .store
1891            .put_opts(path, PutPayload::from(bytes.clone()), opts)
1892            .await
1893        {
1894            Ok(_) => Ok(()),
1895            Err(object_store::Error::AlreadyExists { .. }) => {
1896                let result = self
1897                    .store
1898                    .get(path)
1899                    .await
1900                    .map_err(|e| StateBackendError::Io(e.to_string()))?;
1901                if result.meta.size != intended_size {
1902                    return Err(StateBackendError::Conflict {
1903                        resource: path.to_string(),
1904                        message: format!(
1905                            "existing immutable artifact is {} bytes; retry is {intended_size} bytes",
1906                            result.meta.size
1907                        ),
1908                    });
1909                }
1910                let existing = result
1911                    .bytes()
1912                    .await
1913                    .map_err(|e| StateBackendError::Io(e.to_string()))?;
1914                if existing == bytes {
1915                    Ok(())
1916                } else {
1917                    Err(StateBackendError::Conflict {
1918                        resource: path.to_string(),
1919                        message: "existing immutable artifact has different bytes".into(),
1920                    })
1921                }
1922            }
1923            Err(e) => Err(StateBackendError::Io(e.to_string())),
1924        }
1925    }
1926
1927    async fn read_seal_if_present(
1928        &self,
1929        path: &OsPath,
1930    ) -> Result<Option<CheckpointSeal>, StateBackendError> {
1931        let result = match self.store.get(path).await {
1932            Ok(result) => result,
1933            Err(object_store::Error::NotFound { .. }) => return Ok(None),
1934            Err(error) => return Err(StateBackendError::Io(error.to_string())),
1935        };
1936        Self::check_seal_encoded_size(path, result.meta.size)?;
1937        let bytes = match result.bytes().await {
1938            Ok(bytes) => bytes,
1939            Err(object_store::Error::NotFound { .. }) => return Ok(None),
1940            Err(error) => return Err(StateBackendError::Io(error.to_string())),
1941        };
1942        Self::decode_seal(&bytes).map(Some)
1943    }
1944
1945    async fn read_seal(&self, path: &OsPath) -> Result<CheckpointSeal, StateBackendError> {
1946        self.read_seal_if_present(path).await?.ok_or_else(|| {
1947            StateBackendError::Io(format!("checkpoint seal '{}' is absent", path.as_ref()))
1948        })
1949    }
1950
1951    fn decode_seal(bytes: &[u8]) -> Result<CheckpointSeal, StateBackendError> {
1952        let seal: CheckpointSeal = serde_json::from_slice(bytes).map_err(|e| {
1953            StateBackendError::Serialization(format!("invalid checkpoint seal: {e}"))
1954        })?;
1955        if seal.version != CHECKPOINT_SEAL_VERSION {
1956            return Err(StateBackendError::Serialization(format!(
1957                "unsupported checkpoint seal version {}; expected {CHECKPOINT_SEAL_VERSION}",
1958                seal.version
1959            )));
1960        }
1961        seal.validate().map_err(|error| {
1962            StateBackendError::Serialization(format!("invalid checkpoint seal: {error}"))
1963        })?;
1964        Ok(seal)
1965    }
1966}
1967
1968#[cfg(test)]
1969mod tests;