Skip to main content

laminar_core/cluster/control/snapshot/
mod.rs

1//! Durable vnode→instance assignment snapshots. One object per
2//! version at `control/assignment-snapshots/v{N:020}.json`. Chitchat
3//! carries the ephemeral copy; these files survive full-cluster
4//! restart.
5//!
6//! Rotation and drain finalization use `PutMode::Create` on separate per-version paths. The
7//! append-only winner works on every backend, including `LocalFileSystem`, without relying on
8//! conditional overwrite support.
9
10use std::collections::BTreeSet;
11use std::sync::Arc;
12
13use bytes::Bytes;
14use object_store::path::Path as OsPath;
15use object_store::{ObjectStore, ObjectStoreExt, PutMode, PutOptions, PutPayload};
16use tokio_stream::StreamExt;
17
18use crate::checkpoint::AssignmentDrainTransition;
19
20mod model;
21
22pub use model::{AssignmentSnapshot, AssignmentSnapshotRef};
23use model::{DrainFinalization, RecoveryMaterialization};
24
25const SNAPSHOT_PREFIX: &str = "control/assignment-snapshots/";
26const RECOVERY_PROPOSAL_PREFIX: &str = "control/assignment-recovery-proposals/v1/";
27const RECOVERY_MATERIALIZATION_RELATIVE_PREFIX: &str = "recovery-materializations/v1/";
28const RECOVERY_MATERIALIZATION_PREFIX: &str =
29    "control/assignment-snapshots/recovery-materializations/v1/";
30const DRAIN_FINALIZATION_PREFIX: &str = "control/assignment-drain-finalizations/";
31const SNAPSHOT_VERSION_WIDTH: usize = 20;
32const DRAIN_FINALIZATION_VERSION: u16 = 1;
33const RECOVERY_MATERIALIZATION_VERSION: u16 = 1;
34const MAX_RECOVERY_PROPOSAL_BYTES: usize = 8 * 1024 * 1024;
35const MAX_RECOVERY_MATERIALIZATION_BYTES: u64 = 8 * 1024 * 1024 + 1024;
36const RECOVERY_PROPOSAL_GC_BATCH: usize = 64;
37const RECOVERY_PROPOSAL_GC_MAX_BATCHES: usize = 4;
38
39fn snapshot_path(version: u64) -> OsPath {
40    // Fixed-width so lexicographic list order matches numeric order.
41    OsPath::from(format!(
42        "{SNAPSHOT_PREFIX}v{version:0SNAPSHOT_VERSION_WIDTH$}.json"
43    ))
44}
45
46fn drain_finalization_path(version: u64) -> OsPath {
47    OsPath::from(format!(
48        "{DRAIN_FINALIZATION_PREFIX}v{version:0SNAPSHOT_VERSION_WIDTH$}.json"
49    ))
50}
51
52fn recovery_proposal_path(reference: &AssignmentSnapshotRef) -> OsPath {
53    OsPath::from(format!(
54        "{RECOVERY_PROPOSAL_PREFIX}v{:0width$}/sha256={}.json",
55        reference.version,
56        reference.sha256,
57        width = SNAPSHOT_VERSION_WIDTH
58    ))
59}
60
61fn recovery_proposal_version_prefix(version: u64) -> OsPath {
62    OsPath::from(format!(
63        "{RECOVERY_PROPOSAL_PREFIX}v{version:0SNAPSHOT_VERSION_WIDTH$}/"
64    ))
65}
66
67fn recovery_materialization_path(version: u64) -> OsPath {
68    OsPath::from(format!(
69        "{RECOVERY_MATERIALIZATION_PREFIX}v{version:0SNAPSHOT_VERSION_WIDTH$}.json"
70    ))
71}
72
73fn current_time_millis() -> i64 {
74    std::time::SystemTime::now()
75        .duration_since(std::time::UNIX_EPOCH)
76        .map_or(0, |duration| {
77            i64::try_from(duration.as_millis()).unwrap_or(i64::MAX)
78        })
79}
80
81fn version_from_file(name: &str, kind: &str, minimum: u64) -> Result<u64, SnapshotError> {
82    let Some(number) = name
83        .strip_prefix('v')
84        .and_then(|name| name.strip_suffix(".json"))
85    else {
86        return Err(SnapshotError::Invalid(format!(
87            "non-canonical {kind} filename {name}"
88        )));
89    };
90    if number.len() != SNAPSHOT_VERSION_WIDTH || !number.bytes().all(|byte| byte.is_ascii_digit()) {
91        return Err(SnapshotError::Invalid(format!(
92            "non-canonical {kind} filename {name}"
93        )));
94    }
95    let version = number.parse::<u64>().map_err(|error| {
96        SnapshotError::Invalid(format!("invalid {kind} filename {name}: {error}"))
97    })?;
98    if version < minimum {
99        return Err(SnapshotError::Invalid(format!(
100            "{kind} version must be at least {minimum}"
101        )));
102    }
103    Ok(version)
104}
105
106/// I/O wrapper for [`AssignmentSnapshot`] on an object store.
107pub struct AssignmentSnapshotStore {
108    store: Arc<dyn ObjectStore>,
109    /// Exact kind returned by the last successful head inventory/load. Snapshot watchers audit
110    /// that same version immediately, so they can reuse this immutable provenance instead of
111    /// issuing a speculative overlay GET. Unknown versions still take the verified fallback.
112    last_loaded_head: parking_lot::Mutex<Option<(u64, SnapshotHeadKind)>>,
113}
114
115struct AssignmentVersionInventory {
116    versions: Vec<u64>,
117    recovery_materializations: BTreeSet<u64>,
118}
119
120#[derive(Clone, Copy, PartialEq, Eq)]
121enum SnapshotHeadKind {
122    Raw,
123    Recovery,
124}
125
126impl std::fmt::Debug for AssignmentSnapshotStore {
127    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
128        f.debug_struct("AssignmentSnapshotStore")
129            .finish_non_exhaustive()
130    }
131}
132
133/// Errors loading or saving an [`AssignmentSnapshot`].
134#[derive(Debug, thiserror::Error)]
135pub enum SnapshotError {
136    /// Underlying object store I/O failure.
137    #[error("object store I/O: {0}")]
138    Io(String),
139    /// JSON de/serialization failure.
140    #[error("JSON: {0}")]
141    Json(#[from] serde_json::Error),
142    /// Snapshot metadata, owner map, or process roster is non-canonical.
143    #[error("invalid snapshot: {0}")]
144    Invalid(String),
145}
146
147impl AssignmentSnapshotStore {
148    /// Wrap a pre-constructed object store.
149    #[must_use]
150    pub fn new(store: Arc<dyn ObjectStore>) -> Self {
151        Self {
152            store,
153            last_loaded_head: parking_lot::Mutex::new(None),
154        }
155    }
156
157    /// Stage one committed successor under its canonical content address.
158    ///
159    /// Identical retries converge on the same immutable object. This does not change the durable
160    /// assignment head; callers publish the returned reference through their fencing authority
161    /// before materialization.
162    ///
163    /// # Errors
164    /// Rejects an invalid/non-committed successor or a write that cannot be reconciled exactly.
165    pub async fn stage_recovery_proposal(
166        &self,
167        proposal: &AssignmentSnapshot,
168    ) -> Result<AssignmentSnapshotRef, SnapshotError> {
169        let (encoded, reference) = proposal.encode_recovery_proposal()?;
170        let path = recovery_proposal_path(&reference);
171        let options = PutOptions {
172            mode: PutMode::Create,
173            ..PutOptions::default()
174        };
175        let put_error = self
176            .store
177            .put_opts(
178                &path,
179                PutPayload::from(Bytes::copy_from_slice(&encoded)),
180                options,
181            )
182            .await
183            .err();
184
185        match self.load_recovery_proposal(&reference).await {
186            Ok(stored) if stored == *proposal => Ok(reference),
187            Ok(_) => Err(SnapshotError::Invalid(format!(
188                "recovery proposal '{}' differs from the proposed snapshot",
189                reference.sha256
190            ))),
191            Err(reconcile_error) => {
192                if let Some(put_error) = put_error {
193                    Err(SnapshotError::Io(format!(
194                        "recovery proposal write failed ({put_error}); reconciliation failed ({reconcile_error})"
195                    )))
196                } else {
197                    Err(reconcile_error)
198                }
199            }
200        }
201    }
202
203    /// Load and verify one exact immutable recovery proposal.
204    ///
205    /// # Errors
206    /// Rejects a missing, malformed, non-canonical, or reference-mismatched object.
207    pub async fn load_recovery_proposal(
208        &self,
209        reference: &AssignmentSnapshotRef,
210    ) -> Result<AssignmentSnapshot, SnapshotError> {
211        reference.validate()?;
212        let result = match self.store.get(&recovery_proposal_path(reference)).await {
213            Ok(result) => result,
214            Err(object_store::Error::NotFound { .. }) => {
215                return Err(SnapshotError::Invalid(format!(
216                    "recovery proposal '{}' is missing",
217                    reference.sha256
218                )));
219            }
220            Err(error) => return Err(SnapshotError::Io(error.to_string())),
221        };
222        if result.meta.size != reference.encoded_len {
223            return Err(SnapshotError::Invalid(format!(
224                "recovery proposal '{}' is {} bytes, expected {}",
225                reference.sha256, result.meta.size, reference.encoded_len
226            )));
227        }
228        let bytes = result
229            .bytes()
230            .await
231            .map_err(|error| SnapshotError::Io(error.to_string()))?;
232        if u64::try_from(bytes.len()).ok() != Some(reference.encoded_len) {
233            return Err(SnapshotError::Invalid(format!(
234                "recovery proposal '{}' payload length changed while reading",
235                reference.sha256
236            )));
237        }
238        let proposal: AssignmentSnapshot = serde_json::from_slice(&bytes).map_err(|error| {
239            SnapshotError::Invalid(format!("recovery proposal '{}': {error}", reference.sha256))
240        })?;
241        let (canonical, actual_reference) = proposal.encode_recovery_proposal()?;
242        if actual_reference != *reference || canonical.as_slice() != bytes.as_ref() {
243            return Err(SnapshotError::Invalid(format!(
244                "recovery proposal '{}' does not match its content-addressed reference",
245                reference.sha256
246            )));
247        }
248        Ok(proposal)
249    }
250
251    async fn load_recovery_materialization(
252        &self,
253        version: u64,
254    ) -> Result<Option<RecoveryMaterialization>, SnapshotError> {
255        let result = match self
256            .store
257            .get(&recovery_materialization_path(version))
258            .await
259        {
260            Ok(result) => result,
261            Err(object_store::Error::NotFound { .. }) => return Ok(None),
262            Err(error) => return Err(SnapshotError::Io(error.to_string())),
263        };
264        if result.meta.size == 0 || result.meta.size > MAX_RECOVERY_MATERIALIZATION_BYTES {
265            return Err(SnapshotError::Invalid(format!(
266                "recovery materialization is {} bytes; expected 1..={MAX_RECOVERY_MATERIALIZATION_BYTES}",
267                result.meta.size
268            )));
269        }
270        let bytes = result
271            .bytes()
272            .await
273            .map_err(|error| SnapshotError::Io(error.to_string()))?;
274        let materialization: RecoveryMaterialization = serde_json::from_slice(&bytes)?;
275        materialization.validate()?;
276        if materialization.proposal.version != version {
277            return Err(SnapshotError::Invalid(format!(
278                "recovery materialization path version {version} references proposal version {}",
279                materialization.proposal.version
280            )));
281        }
282        let canonical = serde_json::to_vec(&materialization)?;
283        if canonical.as_slice() != bytes.as_ref() {
284            return Err(SnapshotError::Invalid(format!(
285                "recovery materialization {version} does not use its canonical body"
286            )));
287        }
288        Ok(Some(materialization))
289    }
290
291    /// Verify and publish a staged successor as the monotonic durable assignment head.
292    ///
293    /// The create-only materialization is the version reservation. It lives outside the raw
294    /// graceful-drain snapshot namespace, so a drain write already in flight under a superseded
295    /// leader cannot occupy or replace the recovery winner. Readers always prefer this record for
296    /// its exact version.
297    ///
298    /// # Errors
299    /// Rejects an invalid proposal reference or a durable head other than its predecessor.
300    pub(super) async fn materialize_recovery(
301        &self,
302        reference: &AssignmentSnapshotRef,
303    ) -> Result<RotateOutcome, SnapshotError> {
304        let proposal = self.load_recovery_proposal(reference).await?;
305        let predecessor_version = reference.version.checked_sub(1).ok_or_else(|| {
306            SnapshotError::Invalid("recovery proposal has no predecessor generation".into())
307        })?;
308        let head = self.list_versions().await?.last().copied();
309        if head != Some(predecessor_version) && head != Some(reference.version) {
310            return Err(SnapshotError::Invalid(format!(
311                "recovery materialization requires durable head {predecessor_version} or {}, observed {head:?}",
312                reference.version
313            )));
314        }
315
316        let materialization = RecoveryMaterialization::new(reference.clone(), proposal.clone())?;
317        let options = PutOptions {
318            mode: PutMode::Create,
319            ..PutOptions::default()
320        };
321        let result = self
322            .store
323            .put_opts(
324                &recovery_materialization_path(reference.version),
325                PutPayload::from(Bytes::from(serde_json::to_vec(&materialization)?)),
326                options,
327            )
328            .await;
329        let winner = self
330            .load_recovery_materialization(reference.version)
331            .await?
332            .ok_or_else(|| {
333                SnapshotError::Io(format!(
334                    "recovery materialization {} was not durably visible",
335                    reference.version
336                ))
337            })?;
338        let winner_snapshot = winner.snapshot.clone();
339        if result.is_ok() {
340            if winner != materialization || winner_snapshot != proposal {
341                return Err(SnapshotError::Invalid(format!(
342                    "recovery materialization {} changed after its create succeeded",
343                    reference.version
344                )));
345            }
346            return Ok(RotateOutcome::Rotated);
347        }
348        Ok(RotateOutcome::Conflict(Box::new(winner_snapshot)))
349    }
350
351    /// Enumerate raw snapshots and recovery materializations with one object-store LIST.
352    async fn list_version_inventory(&self) -> Result<AssignmentVersionInventory, SnapshotError> {
353        let prefix = OsPath::from(SNAPSHOT_PREFIX);
354        let mut entries = self.store.list(Some(&prefix));
355        let mut versions = Vec::new();
356        let mut recovery_materializations = BTreeSet::new();
357        while let Some(entry) = entries.next().await {
358            let entry = entry.map_err(|e| SnapshotError::Io(e.to_string()))?;
359            let loc = entry.location.as_ref();
360            // Accept only canonical fixed-width snapshot names. Unrelated siblings are ignored,
361            // but a snapshot-like name with another shape is a split-history risk and fails load.
362            let Some(rest) = loc.strip_prefix(SNAPSHOT_PREFIX) else {
363                continue;
364            };
365            if let Some(name) = rest.strip_prefix(RECOVERY_MATERIALIZATION_RELATIVE_PREFIX) {
366                let version = version_from_file(name, "recovery materialization", 2)?;
367                versions.push(version);
368                recovery_materializations.insert(version);
369                continue;
370            }
371            if !rest.starts_with('v') {
372                continue;
373            }
374            versions.push(version_from_file(rest, "assignment snapshot", 1)?);
375        }
376        versions.sort_unstable();
377        versions.dedup();
378        if versions.windows(2).any(|pair| {
379            pair[0]
380                .checked_add(1)
381                .is_none_or(|expected| expected != pair[1])
382        }) {
383            return Err(SnapshotError::Invalid(
384                "assignment snapshot versions are not contiguous".into(),
385            ));
386        }
387        Ok(AssignmentVersionInventory {
388            versions,
389            recovery_materializations,
390        })
391    }
392
393    /// Scan the shared history prefix and return every logical version in ascending order.
394    async fn list_versions(&self) -> Result<Vec<u64>, SnapshotError> {
395        Ok(self.list_version_inventory().await?.versions)
396    }
397
398    async fn list_drain_finalization_versions(&self) -> Result<Vec<u64>, SnapshotError> {
399        let prefix = OsPath::from(DRAIN_FINALIZATION_PREFIX);
400        let mut entries = self.store.list(Some(&prefix));
401        let mut versions = Vec::new();
402        while let Some(entry) = entries.next().await {
403            let entry = entry.map_err(|error| SnapshotError::Io(error.to_string()))?;
404            let location = entry.location.as_ref();
405            let Some(rest) = location.strip_prefix(DRAIN_FINALIZATION_PREFIX) else {
406                continue;
407            };
408            let Some(number) = rest
409                .strip_prefix('v')
410                .and_then(|name| name.strip_suffix(".json"))
411            else {
412                return Err(SnapshotError::Invalid(format!(
413                    "non-canonical drain finalization filename {rest}"
414                )));
415            };
416            if number.len() != SNAPSHOT_VERSION_WIDTH
417                || !number.bytes().all(|byte| byte.is_ascii_digit())
418            {
419                return Err(SnapshotError::Invalid(format!(
420                    "non-canonical drain finalization filename {rest}"
421                )));
422            }
423            let version = number.parse::<u64>().map_err(|error| {
424                SnapshotError::Invalid(format!(
425                    "invalid drain finalization filename {rest}: {error}"
426                ))
427            })?;
428            if version == 0 {
429                return Err(SnapshotError::Invalid(
430                    "assignment drain finalization version zero is not durable".into(),
431                ));
432            }
433            versions.push(version);
434        }
435        versions.sort_unstable();
436        versions.dedup();
437        Ok(versions)
438    }
439
440    /// Load the current (highest-versioned) snapshot; `Ok(None)` on
441    /// fresh cluster.
442    ///
443    /// # Errors
444    /// Object-store I/O or JSON decode failure.
445    pub async fn load(&self) -> Result<Option<AssignmentSnapshot>, SnapshotError> {
446        let inventory = self.list_version_inventory().await?;
447        let Some(&latest) = inventory.versions.last() else {
448            return Ok(None);
449        };
450        if inventory.recovery_materializations.contains(&latest) {
451            let materialization = self
452                .load_recovery_materialization(latest)
453                .await?
454                .ok_or_else(|| {
455                    SnapshotError::Io(format!(
456                        "listed recovery materialization {latest} disappeared before load"
457                    ))
458                })?;
459            self.last_loaded_head
460                .lock()
461                .replace((latest, SnapshotHeadKind::Recovery));
462            return Ok(Some(materialization.snapshot));
463        }
464        let loaded = self.load_base_version(latest).await?;
465        if loaded.is_some() {
466            let mut last_loaded_head = self.last_loaded_head.lock();
467            if *last_loaded_head != Some((latest, SnapshotHeadKind::Recovery)) {
468                last_loaded_head.replace((latest, SnapshotHeadKind::Raw));
469            }
470        }
471        Ok(loaded)
472    }
473
474    /// Load a specific version's snapshot. `Ok(None)` if that version
475    /// was never written or has been pruned.
476    ///
477    /// # Errors
478    /// Object-store I/O or JSON decode failure.
479    pub async fn load_version(
480        &self,
481        version: u64,
482    ) -> Result<Option<AssignmentSnapshot>, SnapshotError> {
483        if let Some(materialization) = self.load_recovery_materialization(version).await? {
484            return Ok(Some(materialization.snapshot));
485        }
486        self.load_base_version(version).await
487    }
488
489    async fn load_base_version(
490        &self,
491        version: u64,
492    ) -> Result<Option<AssignmentSnapshot>, SnapshotError> {
493        let Some(snapshot) = self.load_snapshot_object(version).await? else {
494            return Ok(None);
495        };
496        if !snapshot.draining {
497            return Ok(Some(snapshot));
498        }
499        match self.load_drain_finalization(version).await? {
500            Some(finalization) => {
501                finalization.validate_against(&snapshot)?;
502                Ok(Some(finalization.proposal))
503            }
504            None => Ok(Some(snapshot)),
505        }
506    }
507
508    /// Load the immutable drain transition underlying a materialized assignment version.
509    ///
510    /// A terminal `load_version` result intentionally contains only the installed assignment.
511    /// Cluster readers use this accessor to bind that materialized result back to the shared
512    /// authority decision before adoption. Ordinary assignment versions return `None`.
513    ///
514    /// # Errors
515    /// Object-store I/O, JSON decode failure, or a malformed base snapshot.
516    pub async fn load_drain_transition(
517        &self,
518        version: u64,
519    ) -> Result<Option<AssignmentDrainTransition>, SnapshotError> {
520        let last_loaded_head = *self.last_loaded_head.lock();
521        match last_loaded_head {
522            Some((loaded, SnapshotHeadKind::Recovery)) if loaded == version => return Ok(None),
523            Some((loaded, SnapshotHeadKind::Raw)) if loaded == version => {
524                return Ok(self
525                    .load_snapshot_object(version)
526                    .await?
527                    .and_then(|snapshot| snapshot.drain_transition));
528            }
529            _ => {}
530        }
531        if self.load_recovery_materialization(version).await?.is_some() {
532            return Ok(None);
533        }
534        Ok(self
535            .load_snapshot_object(version)
536            .await?
537            .and_then(|snapshot| snapshot.drain_transition))
538    }
539
540    async fn load_snapshot_object(
541        &self,
542        version: u64,
543    ) -> Result<Option<AssignmentSnapshot>, SnapshotError> {
544        let path = snapshot_path(version);
545        match self.store.get(&path).await {
546            Ok(res) => {
547                let bytes = res
548                    .bytes()
549                    .await
550                    .map_err(|e| SnapshotError::Io(e.to_string()))?;
551                let snap: AssignmentSnapshot = serde_json::from_slice(&bytes)?;
552                if snap.version != version {
553                    return Err(SnapshotError::Invalid(format!(
554                        "snapshot path version {version} contains payload version {}",
555                        snap.version
556                    )));
557                }
558                snap.validate()?;
559                Ok(Some(snap))
560            }
561            Err(object_store::Error::NotFound { .. }) => Ok(None),
562            Err(e) => Err(SnapshotError::Io(e.to_string())),
563        }
564    }
565
566    async fn load_drain_finalization(
567        &self,
568        version: u64,
569    ) -> Result<Option<DrainFinalization>, SnapshotError> {
570        let path = drain_finalization_path(version);
571        match self.store.get(&path).await {
572            Ok(result) => {
573                let bytes = result
574                    .bytes()
575                    .await
576                    .map_err(|error| SnapshotError::Io(error.to_string()))?;
577                let finalization: DrainFinalization = serde_json::from_slice(&bytes)?;
578                if finalization.proposal.version != version {
579                    return Err(SnapshotError::Invalid(format!(
580                        "drain finalization path version {version} contains payload version {}",
581                        finalization.proposal.version
582                    )));
583                }
584                Ok(Some(finalization))
585            }
586            Err(object_store::Error::NotFound { .. }) => Ok(None),
587            Err(error) => Err(SnapshotError::Io(error.to_string())),
588        }
589    }
590
591    async fn create_if_absent(
592        &self,
593        snapshot: &AssignmentSnapshot,
594    ) -> Result<Option<AssignmentSnapshot>, SnapshotError> {
595        snapshot.validate()?;
596        let path = snapshot_path(snapshot.version);
597        let bytes = serde_json::to_vec_pretty(snapshot)?;
598        let opts = PutOptions {
599            mode: PutMode::Create,
600            ..PutOptions::default()
601        };
602        match self
603            .store
604            .put_opts(&path, PutPayload::from(Bytes::from(bytes)), opts)
605            .await
606        {
607            Ok(_) => Ok(Some(snapshot.clone())),
608            Err(object_store::Error::AlreadyExists { .. }) => Ok(None),
609            Err(e) => Err(SnapshotError::Io(e.to_string())),
610        }
611    }
612
613    async fn prune_recovery_proposals_for_version(
614        &self,
615        version: u64,
616    ) -> Result<(), SnapshotError> {
617        let prefix = recovery_proposal_version_prefix(version);
618        for _ in 0..RECOVERY_PROPOSAL_GC_MAX_BATCHES {
619            let mut entries = self.store.list(Some(&prefix));
620            let mut candidates = Vec::with_capacity(RECOVERY_PROPOSAL_GC_BATCH);
621            while candidates.len() < RECOVERY_PROPOSAL_GC_BATCH {
622                let Some(entry) = entries.next().await else {
623                    break;
624                };
625                candidates.push(
626                    entry
627                        .map_err(|error| SnapshotError::Io(error.to_string()))?
628                        .location,
629                );
630            }
631            if candidates.is_empty() {
632                return Ok(());
633            }
634            let deletions =
635                futures::stream::iter(candidates.into_iter().map(Ok::<_, object_store::Error>));
636            let mut results = self.store.delete_stream(Box::pin(deletions));
637            while let Some(result) = results.next().await {
638                if let Err(error) = result {
639                    if !matches!(error, object_store::Error::NotFound { .. }) {
640                        return Err(SnapshotError::Io(error.to_string()));
641                    }
642                }
643            }
644            tokio::task::yield_now().await;
645        }
646
647        let mut remaining = self.store.list(Some(&prefix));
648        match remaining.next().await {
649            None => Ok(()),
650            Some(Ok(_)) => Err(SnapshotError::Io(format!(
651                "recovery proposal garbage for assignment {version} exceeds the bounded cleanup budget"
652            ))),
653            Some(Err(error)) => Err(SnapshotError::Io(error.to_string())),
654        }
655    }
656
657    /// CAS-create the version-one seed. `Ok(None)` means another initial writer won.
658    ///
659    /// # Errors
660    /// Object-store I/O or JSON encode failure.
661    pub async fn save_if_absent(
662        &self,
663        snapshot: &AssignmentSnapshot,
664    ) -> Result<Option<AssignmentSnapshot>, SnapshotError> {
665        if snapshot.version != 1 {
666            return Err(SnapshotError::Invalid(format!(
667                "save_if_absent only accepts the version-one seed, got {}",
668                snapshot.version
669            )));
670        }
671        if let Some(head) = self
672            .list_versions()
673            .await?
674            .last()
675            .copied()
676            .filter(|head| *head != 1)
677        {
678            return Err(SnapshotError::Invalid(format!(
679                "cannot seed assignment history with durable head {head}"
680            )));
681        }
682        self.create_if_absent(snapshot).await
683    }
684
685    /// Rotate to `snapshot` assuming the current durable version is
686    /// `prior_version`. Returns [`RotateOutcome::Conflict`] carrying
687    /// the winner's snapshot if a racer produced `prior_version + 1`
688    /// first.
689    ///
690    /// # Errors
691    /// Object-store I/O, JSON encode, or a non-monotonic version bump
692    /// (caller bug).
693    pub async fn save_if_version(
694        &self,
695        snapshot: &AssignmentSnapshot,
696        prior_version: u64,
697    ) -> Result<RotateOutcome, SnapshotError> {
698        snapshot.validate()?;
699        let expected = prior_version
700            .checked_add(1)
701            .ok_or_else(|| SnapshotError::Invalid("assignment snapshot version overflow".into()))?;
702        if snapshot.version != expected {
703            return Err(SnapshotError::Invalid(format!(
704                "save_if_version requires monotonic +1 bump: prior={prior_version}, \
705                 proposed={}",
706                snapshot.version,
707            )));
708        }
709        let head = self.list_versions().await?.last().copied();
710        if head == Some(expected) {
711            let winner = self.load_version(expected).await?.ok_or_else(|| {
712                SnapshotError::Io("durable head disappeared while loading CAS winner".into())
713            })?;
714            return Ok(RotateOutcome::Conflict(Box::new(winner)));
715        }
716        if head != Some(prior_version) {
717            return Err(SnapshotError::Invalid(format!(
718                "save_if_version requires durable head {prior_version}, observed {head:?}"
719            )));
720        }
721        if self.create_if_absent(snapshot).await?.is_some() {
722            return Ok(RotateOutcome::Rotated);
723        }
724        let winner = self.load_version(snapshot.version).await?.ok_or_else(|| {
725            SnapshotError::Io("CAS conflict but load of winner returned None".into())
726        })?;
727        Ok(RotateOutcome::Conflict(Box::new(winner)))
728    }
729
730    /// Append exactly one immutable winner for a draining object: its target or a rollback.
731    ///
732    /// The object version is intentionally unchanged: source receipts certify the target
733    /// assignment version, so committing the map under another version would discard the very
734    /// identity they proved. `PutMode::Create` makes commit versus abort a store-level race with
735    /// one winner on local and cloud backends; the original transition remains auditable.
736    /// Cluster callers must first serialize the verdict through `LeaderLeaseStore`; this method
737    /// only materializes that already-authoritative verdict.
738    ///
739    /// # Errors
740    /// Rejects a stale/non-draining expected value, an unrelated proposal, or a non-head object.
741    pub async fn finalize_drain(
742        &self,
743        draining: &AssignmentSnapshot,
744        proposal: &AssignmentSnapshot,
745    ) -> Result<RotateOutcome, SnapshotError> {
746        let finalization = DrainFinalization::new(draining, proposal.clone())?;
747        if self.list_versions().await?.last().copied() != Some(draining.version) {
748            return Err(SnapshotError::Invalid(format!(
749                "draining assignment {} is no longer the durable head",
750                draining.version
751            )));
752        }
753
754        let current = self
755            .load_snapshot_object(draining.version)
756            .await?
757            .ok_or_else(|| SnapshotError::Io("draining assignment disappeared".into()))?;
758        if current != *draining {
759            let winner = self
760                .load_version(draining.version)
761                .await?
762                .ok_or_else(|| SnapshotError::Io("drain conflict winner disappeared".into()))?;
763            return Ok(RotateOutcome::Conflict(Box::new(winner)));
764        }
765        if let Some(winner) = self.load_drain_finalization(draining.version).await? {
766            winner.validate_against(draining)?;
767            return Ok(RotateOutcome::Conflict(Box::new(winner.proposal)));
768        }
769
770        let path = drain_finalization_path(draining.version);
771        let payload = PutPayload::from(Bytes::from(serde_json::to_vec_pretty(&finalization)?));
772        let options = PutOptions {
773            mode: PutMode::Create,
774            ..PutOptions::default()
775        };
776        match self.store.put_opts(&path, payload, options).await {
777            Ok(_) => Ok(RotateOutcome::Rotated),
778            Err(error) => match self.load_drain_finalization(draining.version).await {
779                Ok(Some(winner)) => {
780                    winner.validate_against(draining)?;
781                    Ok(RotateOutcome::Conflict(Box::new(winner.proposal)))
782                }
783                Ok(None) | Err(_) => Err(SnapshotError::Io(error.to_string())),
784            },
785        }
786    }
787
788    /// Delete every snapshot object with `version < before`.
789    /// Idempotent — missing objects are tolerated.
790    ///
791    /// # Errors
792    /// Object-store I/O.
793    pub async fn prune_before(&self, before: u64) -> Result<(), SnapshotError> {
794        if before == 0 {
795            return Ok(());
796        }
797        let inventory = self.list_version_inventory().await?;
798        for version in inventory.versions {
799            if version >= before {
800                break;
801            }
802            // Remove every winning and losing staged body while the version marker still exists.
803            // A crash leaves that marker discoverable, so the next retention pass resumes GC
804            // instead of leaking an orphaned body of up to 8 MiB.
805            self.prune_recovery_proposals_for_version(version).await?;
806            match self.store.delete(&snapshot_path(version)).await {
807                Ok(()) | Err(object_store::Error::NotFound { .. }) => {}
808                Err(e) => return Err(SnapshotError::Io(e.to_string())),
809            }
810            if inventory.recovery_materializations.contains(&version) {
811                match self
812                    .store
813                    .delete(&recovery_materialization_path(version))
814                    .await
815                {
816                    Ok(()) | Err(object_store::Error::NotFound { .. }) => {}
817                    Err(error) => return Err(SnapshotError::Io(error.to_string())),
818                }
819            }
820        }
821        // Finalization records are in a separate append-only namespace. Scan it independently so
822        // a prior failure after deleting the snapshot can be repaired without leaking orphans.
823        for version in self.list_drain_finalization_versions().await? {
824            if version >= before {
825                break;
826            }
827            let path = drain_finalization_path(version);
828            match self.store.delete(&path).await {
829                Ok(()) | Err(object_store::Error::NotFound { .. }) => {}
830                Err(error) => return Err(SnapshotError::Io(error.to_string())),
831            }
832        }
833        Ok(())
834    }
835}
836
837/// Outcome of [`AssignmentSnapshotStore::save_if_version`].
838#[derive(Debug, Clone)]
839pub enum RotateOutcome {
840    /// Our write landed. The snapshot we passed in is now canonical.
841    Rotated,
842    /// Another writer (a racing leader) won the CAS. The attached
843    /// snapshot is what's currently durable; the caller must adopt it
844    /// rather than retry with a stale view.
845    Conflict(Box<AssignmentSnapshot>),
846}
847
848#[cfg(test)]
849mod tests;