Skip to main content

coven_storage/
provider_probe.rs

1//! Cross-principal provider probe execution: reserving, creating, and
2//! settling exact probe slots on the primary and peer provider storage, over
3//! the probe transcript model in [`coven_protocol::provider`].
4
5use 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(&current.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(&current.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}