1use std::sync::Arc;
6
7use crate::cloud::{
8 CloudHomeError, ConditionalWriteOutcome, ExactCloudHome, ExactCreateOutcome, ExactSlotStorage,
9 ExactUpload,
10};
11use coven_keys::keys::UserKeypair;
12use coven_protocol::objects::{ExactObjectRef, ObjectSlot, StorageError};
13use coven_protocol::provider::*;
14use coven_protocol::provider::{
15 advance_cross_completion, advance_exact, cross_challenge_hash, cross_response_hash, invalid,
16 validate_cross_provider_evidence, validate_cross_provider_evidence_context,
17};
18use coven_protocol::store_commit::ObjectHash;
19use coven_protocol::StoreProviderBinding;
20
21pub struct ProviderProbeStorage {
22 storage: Arc<dyn ExactCloudHome>,
23}
24
25impl ProviderProbeStorage {
26 pub fn new(storage: Arc<dyn ExactCloudHome>) -> Self {
27 Self { storage }
28 }
29
30 pub async fn reserve_cross_principal_response_slot(
31 &self,
32 probe_id: ProviderProbeId,
33 ) -> Result<ObjectSlot, ProviderProbeError> {
34 let logical = cross_peer_logical_key(probe_id);
35 let slot = self
36 .storage
37 .allocate_slot(&logical)
38 .await
39 .map_err(StorageError::from)?;
40 if slot.logical_key() != logical {
41 return invalid("cross-principal response slot changed its logical key");
42 }
43 Ok(slot)
44 }
45
46 pub async fn prepare_cross_principal_challenge(
47 &self,
48 publication_journal: &dyn DeviceJoinChallengePublicationJournal,
49 probe_id: ProviderProbeId,
50 store: &StoreProviderBinding,
51 context: &CrossPrincipalChallengeContext,
52 administrator_signer: &dyn coven_keys::keys::DeviceSigningAuthority,
53 ) -> Result<CrossPrincipalProbeChallenge, ProviderProbeError> {
54 let administrator_live = self
55 .storage
56 .provider_binding()
57 .await
58 .map_err(StorageError::from)?;
59 if administrator_live.store != *store
60 || administrator_live.device != context.administrator_binding
61 {
62 return invalid("cross-principal administrator does not match the challenge context");
63 }
64 validate_cross_provider_evidence_context(store, context)?;
65 let suffix = hex::encode(probe_id.as_bytes());
66 let administrator_key = format!("__coven_probe__/cross/{suffix}/administrator");
67 let administrator_slot = self
68 .storage
69 .allocate_slot(&administrator_key)
70 .await
71 .map_err(StorageError::from)?;
72 if administrator_slot.logical_key() != administrator_key {
73 return invalid("cross-principal administrator slot changed its logical key");
74 }
75 let administrator_payload = probe_payload(&probe_id, ProbePayloadLabel::CrossAdministrator);
76 let administrator_object = ProbeExactObjectReceipt {
77 slot: administrator_slot.clone(),
78 payload_hash: ObjectHash::digest(&administrator_payload),
79 object: ExactObjectRef::new(
80 administrator_slot,
81 administrator_payload.len() as u64,
82 ObjectHash::digest(&administrator_payload),
83 ),
84 };
85 let unsigned = CrossPrincipalProbeChallenge {
86 probe_id,
87 administrator_object,
88 challenge_hash: ObjectHash::digest(&[]),
89 administrator_signature: String::new(),
90 };
91 let challenge_hash = cross_challenge_hash(store, context, &unsigned);
92 let challenge = CrossPrincipalProbeChallenge {
93 challenge_hash,
94 administrator_signature: hex::encode(
95 administrator_signer.sign(challenge_hash.as_bytes()),
96 ),
97 ..unsigned
98 };
99 challenge.verify(context, store, &administrator_signer.public_key_hex())?;
100 let durable = publication_journal.prepare(&challenge).await?;
101 if durable.challenge != challenge {
102 return invalid("durable cross-principal challenge differs from its prepared bytes");
103 }
104 Ok(challenge)
105 }
106
107 pub async fn settle_cross_principal_challenge(
108 &self,
109 publication_journal: &dyn DeviceJoinChallengePublicationJournal,
110 authorization: &DeviceJoinChallengePublicationAuthorization,
111 challenge: &CrossPrincipalProbeChallenge,
112 context: &CrossPrincipalChallengeContext,
113 store: &StoreProviderBinding,
114 ) -> Result<CrossPrincipalProbeChallenge, ProviderProbeError> {
115 let live = self
116 .storage
117 .provider_binding()
118 .await
119 .map_err(StorageError::from)?;
120 if live.store != *store || live.device != context.administrator_binding {
121 return invalid("cross-principal administrator does not match the published challenge");
122 }
123 publication_journal
124 .claim_published(authorization, challenge)
125 .await?;
126 let payload = probe_payload(&challenge.probe_id, ProbePayloadLabel::CrossAdministrator);
127 self.settle_exact_create(&challenge.administrator_object.slot, &payload)
128 .await?;
129 let observed = self
130 .storage
131 .read_at(&challenge.administrator_object.slot)
132 .await
133 .map_err(StorageError::from)?;
134 if observed != payload {
135 return invalid("published cross-principal challenge differs from its signed bytes");
136 }
137 Ok(challenge.clone())
138 }
139
140 pub async fn create_cross_principal_response(
141 &self,
142 challenge: &CrossPrincipalProbeChallenge,
143 context: &CrossPrincipalResponseContext,
144 store: &StoreProviderBinding,
145 administrator_signing_pubkey: &str,
146 peer_signer: &UserKeypair,
147 ) -> Result<CrossPrincipalProbeResponse, ProviderProbeError> {
148 challenge.verify(&context.challenge, store, administrator_signing_pubkey)?;
149 let peer_pubkey = coven_keys::keys::public_key_hex(peer_signer);
150 if context.challenge.member_pubkey != peer_pubkey {
151 return invalid("cross-principal peer signer is not the joining member");
152 }
153 let live = self
154 .storage
155 .provider_binding()
156 .await
157 .map_err(StorageError::from)?;
158 if live.store != *store || live.device != context.challenge.peer_binding {
159 return invalid("cross-principal peer does not match the response context");
160 }
161 let evidence = self
162 .storage
163 .cross_principal_evidence()
164 .await
165 .map_err(StorageError::from)?;
166 validate_cross_provider_evidence(
167 store,
168 &context.challenge.administrator_binding,
169 &context.challenge.peer_binding,
170 &evidence,
171 )?;
172 let expected_peer_key = cross_peer_logical_key(challenge.probe_id);
173 if context.response_slot.logical_key() != expected_peer_key {
174 return invalid("cross-principal response slot uses the wrong logical key");
175 }
176 let administrator_payload =
177 probe_payload(&challenge.probe_id, ProbePayloadLabel::CrossAdministrator);
178 let administrator_read = self
179 .storage
180 .read_at(&challenge.administrator_object.slot)
181 .await
182 .map_err(StorageError::from)?;
183 if administrator_read != administrator_payload {
184 return invalid("peer read bytes differ from the signed cross-principal challenge");
185 }
186 let peer_payload = probe_payload(&challenge.probe_id, ProbePayloadLabel::CrossPeer);
187 self.settle_exact_create(&context.response_slot, &peer_payload)
188 .await?;
189 let peer_read = self
190 .storage
191 .read_at(&context.response_slot)
192 .await
193 .map_err(StorageError::from)?;
194 if peer_read != peer_payload {
195 return invalid("peer response readback differs from its deterministic bytes");
196 }
197 let peer_object = ProbeExactObjectReceipt {
198 slot: context.response_slot.clone(),
199 payload_hash: ObjectHash::digest(&peer_payload),
200 object: ExactObjectRef::new(
201 context.response_slot.clone(),
202 peer_payload.len() as u64,
203 ObjectHash::digest(&peer_payload),
204 ),
205 };
206 let unsigned = CrossPrincipalProbeResponse {
207 challenge_hash: challenge.challenge_hash,
208 provider_evidence: evidence,
209 peer_object,
210 peer_read_administrator_hash: ObjectHash::digest(&administrator_read),
211 response_hash: ObjectHash::digest(&[]),
212 peer_signature: String::new(),
213 };
214 let response_hash = cross_response_hash(store, context, challenge, &unsigned);
215 let response = CrossPrincipalProbeResponse {
216 response_hash,
217 peer_signature: hex::encode(peer_signer.sign(response_hash.as_bytes())),
218 ..unsigned
219 };
220 response.verify(
221 challenge,
222 context,
223 store,
224 administrator_signing_pubkey,
225 &peer_pubkey,
226 )?;
227 Ok(response)
228 }
229
230 pub async fn complete_cross_principal_probe(
231 &self,
232 journal: &dyn ProviderProbeJournal,
233 challenge: &CrossPrincipalProbeChallenge,
234 response: &CrossPrincipalProbeResponse,
235 context: &CrossPrincipalResponseContext,
236 store: &StoreProviderBinding,
237 administrator_signer: &dyn coven_keys::keys::DeviceSigningAuthority,
238 peer_signing_pubkey: &str,
239 ) -> Result<CrossPrincipalProbeReceipt, ProviderProbeError> {
240 let administrator_pubkey = administrator_signer.public_key_hex();
241 challenge.verify(&context.challenge, store, &administrator_pubkey)?;
242 response.verify(
243 challenge,
244 context,
245 store,
246 &administrator_pubkey,
247 peer_signing_pubkey,
248 )?;
249 let live = self
250 .storage
251 .provider_binding()
252 .await
253 .map_err(StorageError::from)?;
254 if live.store != *store || live.device != context.challenge.administrator_binding {
255 return invalid("cross-principal administrator does not match the completion context");
256 }
257 let prepared =
258 ProviderProbeJournalRecord::CrossPrincipal(CrossPrincipalCompletionJournal {
259 probe_id: challenge.probe_id,
260 store: store.clone(),
261 context: context.clone(),
262 challenge: challenge.clone(),
263 response: response.clone(),
264 progress: CrossPrincipalCompletionProgress::Prepared,
265 });
266 let mut durable = match journal.load(challenge.probe_id).await? {
267 Some(existing) => existing,
268 None => journal.begin(prepared).await?,
269 };
270 let ProviderProbeJournalRecord::CrossPrincipal(mut record) = durable.clone() else {
271 return invalid("cross-principal probe id belongs to another durable probe kind");
272 };
273 if record.probe_id != challenge.probe_id
274 || record.store != *store
275 || record.context != *context
276 || record.challenge != *challenge
277 || record.response != *response
278 {
279 return invalid("durable cross-principal completion differs from the requested proof");
280 }
281 if let CrossPrincipalCompletionProgress::ReceiptReady { receipt } = &record.progress {
282 receipt.verify(context, store, &administrator_pubkey, peer_signing_pubkey)?;
283 return Ok(receipt.clone());
284 }
285 let peer_payload = probe_payload(&challenge.probe_id, ProbePayloadLabel::CrossPeer);
286 if matches!(record.progress, CrossPrincipalCompletionProgress::Prepared) {
287 let observed = self
288 .storage
289 .read_at(&response.peer_object.slot)
290 .await
291 .map_err(StorageError::from)?;
292 if observed != peer_payload {
293 return invalid("administrator read differs from the signed peer response");
294 }
295 advance_cross_completion(
296 journal,
297 &mut durable,
298 &mut record,
299 CrossPrincipalCompletionProgress::ReadsVerified {
300 administrator_read_peer_hash: ObjectHash::digest(&observed),
301 },
302 )
303 .await?;
304 }
305 let administrator_read_peer_hash = cross_completion_read_hash(&record.progress)?;
306 if matches!(
307 record.progress,
308 CrossPrincipalCompletionProgress::ReadsVerified { .. }
309 ) {
310 self.storage
311 .delete_and_verify_absent(&response.peer_object.slot)
312 .await
313 .map_err(StorageError::from)?;
314 advance_cross_completion(
315 journal,
316 &mut durable,
317 &mut record,
318 CrossPrincipalCompletionProgress::PeerAbsent {
319 administrator_read_peer_hash,
320 },
321 )
322 .await?;
323 }
324 if matches!(
325 record.progress,
326 CrossPrincipalCompletionProgress::PeerAbsent { .. }
327 ) {
328 self.storage
329 .delete_and_verify_absent(&challenge.administrator_object.slot)
330 .await
331 .map_err(StorageError::from)?;
332 advance_cross_completion(
333 journal,
334 &mut durable,
335 &mut record,
336 CrossPrincipalCompletionProgress::Absent {
337 administrator_read_peer_hash,
338 },
339 )
340 .await?;
341 }
342 let transcript = CrossPrincipalProbeTranscript {
343 challenge: challenge.clone(),
344 response: response.clone(),
345 administrator_read_peer_hash,
346 };
347 let receipt =
348 CrossPrincipalProbeReceipt::signed(transcript, context, store, administrator_signer)?;
349 advance_cross_completion(
350 journal,
351 &mut durable,
352 &mut record,
353 CrossPrincipalCompletionProgress::ReceiptReady {
354 receipt: receipt.clone(),
355 },
356 )
357 .await?;
358 Ok(receipt)
359 }
360
361 async fn settle_exact_create(
362 &self,
363 slot: &ObjectSlot,
364 payload: &[u8],
365 ) -> Result<(), ProviderProbeError> {
366 match self.storage.read_at(slot).await {
367 Ok(bytes) if bytes == payload => Ok(()),
368 Ok(_) => invalid("durable provider probe slot contains different bytes"),
369 Err(CloudHomeError::NotFound(_)) => {
370 create_exact_bytes(self.storage.as_ref(), slot, payload)
371 .await
372 .map(drop)
373 .map_err(StorageError::from)
374 .map_err(ProviderProbeError::Storage)
375 }
376 Err(error) => Err(ProviderProbeError::Storage(StorageError::from(error))),
377 }
378 }
379
380 pub async fn probe_exact_slots(
381 &self,
382 journal: &dyn ProviderProbeJournal,
383 probe_id: ProviderProbeId,
384 binding: &coven_protocol::objects::ResolvedProviderBinding,
385 ) -> Result<ExactSlotProbeReceipt, ProviderProbeError> {
386 let first = self.storage.as_ref();
387 let second = self.storage.as_ref();
388 binding.validate().map_err(ProviderProbeError::Storage)?;
389 let first_binding = first.provider_binding().await.map_err(StorageError::from)?;
390 let second_binding = second
391 .provider_binding()
392 .await
393 .map_err(StorageError::from)?;
394 if first_binding != *binding || second_binding != *binding {
395 return invalid("exact-slot probe clients do not match the receipt binding");
396 }
397 let id = hex::encode(probe_id.as_bytes());
398 let logical_key = format!("__coven_probe__/exact/{id}");
399 let conditional_logical_key = format!("__coven_probe__/conditional/{id}");
400 let lost_logical_key = format!("__coven_probe__/lost-response/{id}");
401 let mut durable = match journal.load(probe_id).await? {
402 Some(existing) => existing,
403 None => {
404 let allocated_slot = first
405 .allocate_slot(&logical_key)
406 .await
407 .map_err(StorageError::from)?;
408 let allocated_lost_slot = first
409 .allocate_slot(&lost_logical_key)
410 .await
411 .map_err(StorageError::from)?;
412 let allocated_conditional_slot = first
413 .allocate_slot(&conditional_logical_key)
414 .await
415 .map_err(StorageError::from)?;
416 journal
417 .begin(ProviderProbeJournalRecord::Exact(ExactProbeJournal {
418 probe_id,
419 binding: binding.clone(),
420 slot: allocated_slot,
421 conditional_slot: allocated_conditional_slot,
422 lost_response_slot: allocated_lost_slot,
423 progress: ExactProbeProgress::Prepared,
424 }))
425 .await?
426 }
427 };
428 let ProviderProbeJournalRecord::Exact(mut record) = durable.clone() else {
429 return invalid("exact probe id belongs to a different durable probe kind");
430 };
431 if record.probe_id != probe_id || record.binding != *binding {
432 return invalid("durable exact probe differs from its requested binding or id");
433 }
434 let slot = record.slot.clone();
435 let conditional_slot = record.conditional_slot.clone();
436 let lost_slot = record.lost_response_slot.clone();
437 if slot.logical_key() != logical_key
438 || conditional_slot.logical_key() != conditional_logical_key
439 || lost_slot.logical_key() != lost_logical_key
440 {
441 return invalid("exact-slot allocator changed the probe logical key");
442 }
443 let payloads = [
444 probe_payload(&probe_id, ProbePayloadLabel::ExactCreateFirst),
445 probe_payload(&probe_id, ProbePayloadLabel::ExactCreateSecond),
446 ];
447 if matches!(record.progress, ExactProbeProgress::Prepared) {
448 let (outcomes, _winner) = match first.read_at(&slot).await {
449 Err(CloudHomeError::NotFound(_)) => {
450 let (left, right) = tokio::join!(
451 create_exact_bytes(first, &slot, &payloads[0]),
452 create_exact_bytes(second, &slot, &payloads[1]),
453 );
454 classify_exact_create_race(left, right)?
455 }
456 Ok(bytes) if bytes == payloads[0] => {
457 require_occupied_rejection(
458 create_exact_bytes(second, &slot, &payloads[1]).await,
459 )?;
460 (
461 [
462 ProbeCreateOutcome::Created,
463 ProbeCreateOutcome::RejectedOccupied,
464 ],
465 0,
466 )
467 }
468 Ok(bytes) if bytes == payloads[1] => {
469 require_occupied_rejection(
470 create_exact_bytes(first, &slot, &payloads[0]).await,
471 )?;
472 (
473 [
474 ProbeCreateOutcome::RejectedOccupied,
475 ProbeCreateOutcome::Created,
476 ],
477 1,
478 )
479 }
480 Ok(_) => return invalid("durable exact probe slot contains unknown bytes"),
481 Err(error) => return Err(ProviderProbeError::Storage(StorageError::from(error))),
482 };
483 advance_exact(
484 journal,
485 &mut durable,
486 &mut record,
487 ExactProbeProgress::Created { outcomes },
488 )
489 .await?;
490 }
491 let (outcomes, winner) = exact_race_state(&record.progress)?;
492 let (full, range) = if matches!(record.progress, ExactProbeProgress::Created { .. }) {
493 let full = first.read_at(&slot).await.map_err(StorageError::from)?;
494 if full != payloads[winner] {
495 return invalid("authoritative exact read does not match the create winner");
496 }
497 let range = first
498 .read_range_at(&slot, PROBE_RANGE_START, PROBE_RANGE_END)
499 .await
500 .map_err(StorageError::from)?;
501 if range != full[PROBE_RANGE_START as usize..PROBE_RANGE_END as usize] {
502 return invalid("exact range read does not match the authoritative full read");
503 }
504 (full, range)
505 } else {
506 (
507 payloads[winner].clone(),
508 payloads[winner][PROBE_RANGE_START as usize..PROBE_RANGE_END as usize].to_vec(),
509 )
510 };
511 if matches!(record.progress, ExactProbeProgress::Created { .. }) {
512 advance_exact(
513 journal,
514 &mut durable,
515 &mut record,
516 ExactProbeProgress::ReadsVerified { outcomes },
517 )
518 .await?;
519 }
520 let accepted =
521 ExactObjectRef::new(slot.clone(), full.len() as u64, ObjectHash::digest(&full));
522 if matches!(record.progress, ExactProbeProgress::ReadsVerified { .. }) {
523 let initial = probe_payload(&probe_id, ProbePayloadLabel::ConditionalInitial);
524 let conditional_payloads = [
525 probe_payload(&probe_id, ProbePayloadLabel::ConditionalFirst),
526 probe_payload(&probe_id, ProbePayloadLabel::ConditionalSecond),
527 ];
528 let mut current = match first.read_versioned_at(&conditional_slot).await {
529 Ok(current) => current,
530 Err(CloudHomeError::NotFound(_)) => {
531 create_versioned_bytes(first, &conditional_slot, &initial).await?;
532 first
533 .read_versioned_at(&conditional_slot)
534 .await
535 .map_err(StorageError::from)?
536 }
537 Err(error) => return Err(ProviderProbeError::Storage(StorageError::from(error))),
538 };
539 if current.bytes != initial
540 && current.bytes != conditional_payloads[0]
541 && current.bytes != conditional_payloads[1]
542 {
543 return invalid("conditional-update probe slot contains unknown bytes");
544 }
545 let starting_payload_hash = ObjectHash::digest(¤t.bytes);
546 let expected = current.version.clone();
547 let (left, right) = tokio::join!(
548 first.replace_at_if_version(
549 &conditional_slot,
550 &expected,
551 conditional_payloads[0].clone(),
552 ),
553 second.replace_at_if_version(
554 &conditional_slot,
555 &expected,
556 conditional_payloads[1].clone(),
557 ),
558 );
559 let (conditional_outcomes, conditional_winner) =
560 classify_conditional_update_race(left, right)?;
561 current = first
562 .read_versioned_at(&conditional_slot)
563 .await
564 .map_err(StorageError::from)?;
565 if current.bytes != conditional_payloads[conditional_winner] {
566 return invalid("conditional-update readback does not match its winning write");
567 }
568 let conditional = ConditionalUpdateProbeReceipt {
569 logical_key: conditional_logical_key.clone(),
570 slot: conditional_slot.clone(),
571 starting_payload_hash,
572 contenders: [
573 ProbeConditionalAttempt {
574 payload_hash: ObjectHash::digest(&conditional_payloads[0]),
575 outcome: conditional_outcomes[0],
576 },
577 ProbeConditionalAttempt {
578 payload_hash: ObjectHash::digest(&conditional_payloads[1]),
579 outcome: conditional_outcomes[1],
580 },
581 ],
582 accepted_payload_hash: ObjectHash::digest(¤t.bytes),
583 };
584 advance_exact(
585 journal,
586 &mut durable,
587 &mut record,
588 ExactProbeProgress::ConditionalVerified {
589 outcomes,
590 conditional,
591 },
592 )
593 .await?;
594 }
595 let conditional = exact_conditional_evidence(&record.progress)?.clone();
596 if matches!(
597 record.progress,
598 ExactProbeProgress::ConditionalVerified { .. }
599 ) {
600 first
601 .delete_and_verify_absent(&slot)
602 .await
603 .map_err(StorageError::from)?;
604 first
605 .delete_versioned_at(&conditional_slot)
606 .await
607 .map_err(StorageError::from)?;
608 match first.read_versioned_at(&conditional_slot).await {
609 Err(CloudHomeError::NotFound(_)) => {}
610 Ok(_) => return invalid("conditional-update probe record remains after deletion"),
611 Err(error) => return Err(ProviderProbeError::Storage(StorageError::from(error))),
612 }
613 advance_exact(
614 journal,
615 &mut durable,
616 &mut record,
617 ExactProbeProgress::PrimaryAbsent {
618 outcomes,
619 conditional: conditional.clone(),
620 },
621 )
622 .await?;
623 }
624 let lost_payload = probe_payload(&probe_id, ProbePayloadLabel::LostResponse);
625 if matches!(record.progress, ExactProbeProgress::PrimaryAbsent { .. }) {
626 match first.read_at(&lost_slot).await {
627 Ok(bytes) if bytes == lost_payload => {}
628 Ok(_) => return invalid("lost-response slot contains unknown bytes"),
629 Err(CloudHomeError::NotFound(_)) => {
630 create_exact_bytes(first, &lost_slot, &lost_payload)
631 .await
632 .map_err(StorageError::from)?;
633 }
634 Err(error) => return Err(ProviderProbeError::Storage(StorageError::from(error))),
635 }
636 advance_exact(
637 journal,
638 &mut durable,
639 &mut record,
640 ExactProbeProgress::LostResponseCreated {
641 outcomes,
642 conditional: conditional.clone(),
643 },
644 )
645 .await?;
646 }
647 let lost_readback = if matches!(
648 record.progress,
649 ExactProbeProgress::LostResponseCreated { .. }
650 ) {
651 let readback = first
652 .read_at(&lost_slot)
653 .await
654 .map_err(StorageError::from)?;
655 if readback != lost_payload {
656 return invalid(
657 "lost-response authoritative readback differs from committed bytes",
658 );
659 }
660 readback
661 } else {
662 lost_payload.clone()
663 };
664 let settled = ExactObjectRef::new(
665 lost_slot.clone(),
666 lost_readback.len() as u64,
667 ObjectHash::digest(&lost_readback),
668 );
669 if matches!(
670 record.progress,
671 ExactProbeProgress::LostResponseCreated { .. }
672 ) {
673 advance_exact(
674 journal,
675 &mut durable,
676 &mut record,
677 ExactProbeProgress::LostResponseReadVerified {
678 outcomes,
679 conditional: conditional.clone(),
680 },
681 )
682 .await?;
683 }
684 if matches!(
685 record.progress,
686 ExactProbeProgress::LostResponseReadVerified { .. }
687 ) {
688 first
689 .delete_and_verify_absent(&lost_slot)
690 .await
691 .map_err(StorageError::from)?;
692 advance_exact(
693 journal,
694 &mut durable,
695 &mut record,
696 ExactProbeProgress::Absent {
697 outcomes,
698 conditional: conditional.clone(),
699 },
700 )
701 .await?;
702 }
703
704 if let ExactProbeProgress::ReceiptReady { receipt } = &record.progress {
705 receipt.verify(&binding.store, &binding.device)?;
706 return Ok(receipt.clone());
707 }
708 let transcript = ExactSlotProbeTranscript {
709 probe_id,
710 logical_key,
711 slot,
712 contenders: [
713 ProbeCreateAttempt {
714 payload_hash: ObjectHash::digest(&payloads[0]),
715 outcome: outcomes[0],
716 },
717 ProbeCreateAttempt {
718 payload_hash: ObjectHash::digest(&payloads[1]),
719 outcome: outcomes[1],
720 },
721 ],
722 accepted,
723 full_read_hash: ObjectHash::digest(&full),
724 range: ProbeRangeReceipt {
725 start: PROBE_RANGE_START,
726 end: PROBE_RANGE_END,
727 bytes_hash: ObjectHash::digest(&range),
728 },
729 conditional,
730 lost_response: LostResponseProbeReceipt {
731 logical_key: lost_logical_key,
732 slot: lost_slot,
733 payload_hash: ObjectHash::digest(&lost_payload),
734 settled,
735 readback_hash: ObjectHash::digest(&lost_readback),
736 },
737 };
738 let receipt =
739 ExactSlotProbeReceipt::from_transcript(transcript, &binding.store, &binding.device);
740 receipt.verify(&binding.store, &binding.device)?;
741 advance_exact(
742 journal,
743 &mut durable,
744 &mut record,
745 ExactProbeProgress::ReceiptReady {
746 receipt: receipt.clone(),
747 },
748 )
749 .await?;
750 Ok(receipt)
751 }
752}
753
754fn cross_completion_read_hash(
755 progress: &CrossPrincipalCompletionProgress,
756) -> Result<ObjectHash, ProviderProbeError> {
757 match progress {
758 CrossPrincipalCompletionProgress::ReadsVerified {
759 administrator_read_peer_hash,
760 }
761 | CrossPrincipalCompletionProgress::PeerAbsent {
762 administrator_read_peer_hash,
763 }
764 | CrossPrincipalCompletionProgress::Absent {
765 administrator_read_peer_hash,
766 } => Ok(*administrator_read_peer_hash),
767 CrossPrincipalCompletionProgress::Prepared
768 | CrossPrincipalCompletionProgress::ReceiptReady { .. } => {
769 invalid("cross-principal completion has no durable administrator read")
770 }
771 }
772}
773
774fn exact_race_state(
775 progress: &ExactProbeProgress,
776) -> Result<([ProbeCreateOutcome; 2], usize), ProviderProbeError> {
777 let (outcomes, winner) = match progress {
778 ExactProbeProgress::Prepared => return invalid("exact probe has no durable create result"),
779 ExactProbeProgress::Created { outcomes }
780 | ExactProbeProgress::ReadsVerified { outcomes }
781 | ExactProbeProgress::ConditionalVerified { outcomes, .. }
782 | ExactProbeProgress::PrimaryAbsent { outcomes, .. }
783 | ExactProbeProgress::LostResponseCreated { outcomes, .. }
784 | ExactProbeProgress::LostResponseReadVerified { outcomes, .. }
785 | ExactProbeProgress::Absent { outcomes, .. } => {
786 let winner = outcomes
787 .iter()
788 .position(|outcome| *outcome == ProbeCreateOutcome::Created)
789 .ok_or_else(|| {
790 ProviderProbeError::InvalidReceipt(
791 "durable exact probe has no create winner".to_string(),
792 )
793 })?;
794 (*outcomes, winner)
795 }
796 ExactProbeProgress::ReceiptReady { receipt } => {
797 let winner = receipt
798 .transcript
799 .contenders
800 .iter()
801 .position(|attempt| attempt.outcome == ProbeCreateOutcome::Created)
802 .ok_or_else(|| {
803 ProviderProbeError::InvalidReceipt(
804 "durable exact receipt has no create winner".to_string(),
805 )
806 })?;
807 (
808 [
809 receipt.transcript.contenders[0].outcome,
810 receipt.transcript.contenders[1].outcome,
811 ],
812 winner,
813 )
814 }
815 };
816 if winner > 1 || outcomes[winner] != ProbeCreateOutcome::Created {
817 return invalid("durable exact probe has an invalid winner");
818 }
819 Ok((outcomes, winner))
820}
821
822fn exact_conditional_evidence(
823 progress: &ExactProbeProgress,
824) -> Result<&ConditionalUpdateProbeReceipt, ProviderProbeError> {
825 match progress {
826 ExactProbeProgress::ConditionalVerified { conditional, .. }
827 | ExactProbeProgress::PrimaryAbsent { conditional, .. }
828 | ExactProbeProgress::LostResponseCreated { conditional, .. }
829 | ExactProbeProgress::LostResponseReadVerified { conditional, .. }
830 | ExactProbeProgress::Absent { conditional, .. } => Ok(conditional),
831 ExactProbeProgress::ReceiptReady { receipt } => Ok(&receipt.transcript.conditional),
832 ExactProbeProgress::Prepared
833 | ExactProbeProgress::Created { .. }
834 | ExactProbeProgress::ReadsVerified { .. } => {
835 invalid("exact probe has no conditional-update evidence")
836 }
837 }
838}
839
840fn classify_conditional_update_race(
841 left: Result<ConditionalWriteOutcome, CloudHomeError>,
842 right: Result<ConditionalWriteOutcome, CloudHomeError>,
843) -> Result<([ProbeConditionalOutcome; 2], usize), ProviderProbeError> {
844 match (left, right) {
845 (
846 Ok(ConditionalWriteOutcome::Replaced(_)),
847 Ok(ConditionalWriteOutcome::VersionChanged),
848 ) => Ok((
849 [
850 ProbeConditionalOutcome::Replaced,
851 ProbeConditionalOutcome::RejectedRevision,
852 ],
853 0,
854 )),
855 (
856 Ok(ConditionalWriteOutcome::VersionChanged),
857 Ok(ConditionalWriteOutcome::Replaced(_)),
858 ) => Ok((
859 [
860 ProbeConditionalOutcome::RejectedRevision,
861 ProbeConditionalOutcome::Replaced,
862 ],
863 1,
864 )),
865 (left, right) => invalid(&format!(
866 "conditional-update race did not produce one replacement and one revision rejection: left={left:?}, right={right:?}"
867 )),
868 }
869}
870
871fn classify_exact_create_race(
872 left: Result<ExactCreateOutcome, CloudHomeError>,
873 right: Result<ExactCreateOutcome, CloudHomeError>,
874) -> Result<([ProbeCreateOutcome; 2], usize), ProviderProbeError> {
875 match (left, right) {
876 (
877 Ok(ExactCreateOutcome::Created),
878 Err(CloudHomeError::SlotCollision(_) | CloudHomeError::AlreadyExists(_)),
879 ) => Ok((
880 [
881 ProbeCreateOutcome::Created,
882 ProbeCreateOutcome::RejectedOccupied,
883 ],
884 0,
885 )),
886 (
887 Err(CloudHomeError::SlotCollision(_) | CloudHomeError::AlreadyExists(_)),
888 Ok(ExactCreateOutcome::Created),
889 ) => Ok((
890 [
891 ProbeCreateOutcome::RejectedOccupied,
892 ProbeCreateOutcome::Created,
893 ],
894 1,
895 )),
896 (left, right) => invalid(&format!(
897 "exact-slot race did not produce one create and one occupied rejection: left={left:?}, right={right:?}"
898 )),
899 }
900}
901
902fn require_occupied_rejection(
903 result: Result<ExactCreateOutcome, CloudHomeError>,
904) -> Result<(), ProviderProbeError> {
905 match result {
906 Err(CloudHomeError::SlotCollision(_) | CloudHomeError::AlreadyExists(_)) => Ok(()),
907 Ok(ExactCreateOutcome::Created) => {
908 invalid("settled exact probe contender unexpectedly created a second object")
909 }
910 result => invalid(&format!(
911 "settled exact probe contender was not rejected as occupied: result={result:?}"
912 )),
913 }
914}
915
916async fn create_exact_bytes(
917 storage: &dyn ExactSlotStorage,
918 slot: &ObjectSlot,
919 bytes: &[u8],
920) -> Result<ExactCreateOutcome, CloudHomeError> {
921 let object = ExactObjectRef::new(slot.clone(), bytes.len() as u64, ObjectHash::digest(bytes));
922 let upload = ExactUpload::from_bytes(&object, bytes).map_err(CloudHomeError::from)?;
923 storage
924 .create_at(
925 &upload,
926 &crate::cloud::UploadControl::running(crate::cloud::no_progress()),
927 )
928 .await
929}
930
931async fn create_versioned_bytes(
932 storage: &dyn ExactSlotStorage,
933 slot: &ObjectSlot,
934 bytes: &[u8],
935) -> Result<(), ProviderProbeError> {
936 let object = ExactObjectRef::new(slot.clone(), bytes.len() as u64, ObjectHash::digest(bytes));
937 let upload = ExactUpload::from_bytes(&object, bytes)?;
938 storage
939 .create_versioned_at(
940 &upload,
941 &crate::cloud::UploadControl::running(crate::cloud::no_progress()),
942 )
943 .await
944 .map(drop)
945 .map_err(StorageError::from)
946 .map_err(ProviderProbeError::Storage)
947}