1use 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 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
106pub struct AssignmentSnapshotStore {
108 store: Arc<dyn ObjectStore>,
109 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#[derive(Debug, thiserror::Error)]
135pub enum SnapshotError {
136 #[error("object store I/O: {0}")]
138 Io(String),
139 #[error("JSON: {0}")]
141 Json(#[from] serde_json::Error),
142 #[error("invalid snapshot: {0}")]
144 Invalid(String),
145}
146
147impl AssignmentSnapshotStore {
148 #[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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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#[derive(Debug, Clone)]
839pub enum RotateOutcome {
840 Rotated,
842 Conflict(Box<AssignmentSnapshot>),
846}
847
848#[cfg(test)]
849mod tests;