1use 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;
46const 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#[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
76pub 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 execution_id: uuid::Uuid,
84 vnode_capacity: u32,
85 authoritative_version: Arc<AtomicU64>,
88 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 #[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 #[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 #[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 #[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 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 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 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 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 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 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 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 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 ¤t != 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 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 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 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 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 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 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 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;