Skip to main content

coven_core/sync/
provider.rs

1use std::collections::{BTreeMap, BTreeSet};
2use std::fmt;
3
4use async_trait::async_trait;
5use serde::{Deserialize, Deserializer, Serialize, Serializer};
6use sha2::{Digest, Sha256};
7
8use crate::keys::UserKeypair;
9use crate::storage::cloud::{BlobBody, CloudHomeError, ExactSlotStorage, ObjectSlot};
10use crate::sync::membership::{
11    MembershipCoord, MembershipEntry, MembershipGrantId, OwnerStreamBarrier,
12};
13use crate::sync::storage::{
14    CoordinationError, CoordinationStorage, CreateHeadError, ExactObjectRef, ProviderDeviceBinding,
15    ReplaceHeadError, StorageError, StoreProviderBinding, SyncStorage, VersionedObject,
16};
17use crate::sync::store_commit::{
18    DeviceJoinAttemptId, DeviceJoinAttemptRef, DeviceJoinOutcomeRef, ObjectHash,
19    StoreBatchCommitRef, StoreDeviceRegistration, StoreDeviceRegistrationRef, StoreRootRef,
20};
21
22const EXACT_TRANSCRIPT_DOMAIN: &[u8] = b"coven.provider-exact-slot-probe.v1\0";
23const SERIAL_TRANSCRIPT_DOMAIN: &[u8] = b"coven.provider-serial-probe.v1\0";
24const CROSS_TRANSCRIPT_DOMAIN: &[u8] = b"coven.provider-cross-principal-probe.v1\0";
25const CROSS_CHALLENGE_DOMAIN: &[u8] = b"coven.provider-cross-principal-challenge.v1\0";
26const CROSS_RESPONSE_DOMAIN: &[u8] = b"coven.provider-cross-principal-response.v1\0";
27const PAYLOAD_DOMAIN: &[u8] = b"coven.provider-probe-payload.v1\0";
28const MEMBER_ACCESS_GRANT_DOMAIN: &[u8] = b"coven.provider-member-access-grant.v1\0";
29const MEMBER_ACCESS_WITHDRAWAL_DOMAIN: &[u8] = b"coven.provider-member-access-withdrawal.v1\0";
30pub const PROBE_PAYLOAD_LEN: usize = 256;
31pub const PROBE_RANGE_START: u64 = 31;
32pub const PROBE_RANGE_END: u64 = 173;
33
34#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
35pub struct ProviderProbeId([u8; 32]);
36
37impl ProviderProbeId {
38    pub fn from_bytes(bytes: [u8; 32]) -> Self {
39        Self(bytes)
40    }
41
42    pub fn as_bytes(&self) -> &[u8; 32] {
43        &self.0
44    }
45}
46
47impl fmt::Debug for ProviderProbeId {
48    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
49        formatter.write_str(&hex::encode(self.0))
50    }
51}
52
53impl Serialize for ProviderProbeId {
54    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
55    where
56        S: Serializer,
57    {
58        serializer.serialize_str(&hex::encode(self.0))
59    }
60}
61
62impl<'de> Deserialize<'de> for ProviderProbeId {
63    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
64    where
65        D: Deserializer<'de>,
66    {
67        let value = String::deserialize(deserializer)?;
68        if value.len() != 64
69            || value
70                .bytes()
71                .any(|byte| !byte.is_ascii_digit() && !(b'a'..=b'f').contains(&byte))
72        {
73            return Err(serde::de::Error::custom(
74                "provider probe id must be 64 lowercase hexadecimal characters",
75            ));
76        }
77        let bytes: [u8; 32] = hex::decode(value)
78            .map_err(serde::de::Error::custom)?
79            .try_into()
80            .map_err(|_| serde::de::Error::custom("provider probe id has the wrong length"))?;
81        Ok(Self(bytes))
82    }
83}
84
85#[derive(Clone, Copy)]
86pub enum ProbePayloadLabel {
87    ExactCreateFirst,
88    ExactCreateSecond,
89    LostResponse,
90    SerialCreateFirst,
91    SerialCreateSecond,
92    SerialReplaceFirst,
93    SerialReplaceSecond,
94    CrossAdministrator,
95    CrossPeer,
96}
97
98impl ProbePayloadLabel {
99    fn bytes(self) -> &'static [u8] {
100        match self {
101            Self::ExactCreateFirst => b"exact-create-first",
102            Self::ExactCreateSecond => b"exact-create-second",
103            Self::LostResponse => b"lost-response",
104            Self::SerialCreateFirst => b"serial-create-first",
105            Self::SerialCreateSecond => b"serial-create-second",
106            Self::SerialReplaceFirst => b"serial-replace-first",
107            Self::SerialReplaceSecond => b"serial-replace-second",
108            Self::CrossAdministrator => b"cross-administrator",
109            Self::CrossPeer => b"cross-peer",
110        }
111    }
112}
113
114pub fn probe_payload(probe_id: &ProviderProbeId, label: ProbePayloadLabel) -> Vec<u8> {
115    let mut output = Vec::with_capacity(PROBE_PAYLOAD_LEN);
116    let mut counter = 0u32;
117    while output.len() < PROBE_PAYLOAD_LEN {
118        let mut digest = Sha256::new();
119        digest.update(PAYLOAD_DOMAIN);
120        digest.update(probe_id.as_bytes());
121        digest.update(label.bytes());
122        digest.update(counter.to_be_bytes());
123        output.extend_from_slice(&digest.finalize());
124        counter += 1;
125    }
126    output.truncate(PROBE_PAYLOAD_LEN);
127    output
128}
129
130pub fn canonical_custom_s3_origin(input: &str) -> Result<String, StorageError> {
131    if input.ends_with('/') {
132        return Err(StorageError::Configuration(
133            "custom S3 endpoint must not have a trailing slash".to_string(),
134        ));
135    }
136    let parsed = url::Url::parse(input).map_err(|error| {
137        StorageError::Configuration(format!("invalid custom S3 endpoint: {error}"))
138    })?;
139    if !matches!(parsed.scheme(), "http" | "https")
140        || !parsed.username().is_empty()
141        || parsed.password().is_some()
142        || parsed.query().is_some()
143        || parsed.fragment().is_some()
144        || parsed.path() != "/"
145    {
146        return Err(StorageError::Configuration(
147            "custom S3 endpoint must be an HTTP origin without user info, path, query, or fragment"
148                .to_string(),
149        ));
150    }
151    let host = parsed
152        .host_str()
153        .ok_or_else(|| StorageError::Configuration("custom S3 endpoint has no host".to_string()))?;
154    let port = parsed.port();
155    let default_port = matches!(
156        (parsed.scheme(), port),
157        ("http", Some(80)) | ("https", Some(443))
158    );
159    Ok(if let Some(port) = port.filter(|_| !default_port) {
160        format!("{}://{}:{port}", parsed.scheme(), host.to_ascii_lowercase())
161    } else {
162        format!("{}://{}", parsed.scheme(), host.to_ascii_lowercase())
163    })
164}
165
166#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
167#[serde(deny_unknown_fields)]
168pub struct ProviderCapabilityProof {
169    pub exact_slots: ExactSlotProbeReceipt,
170    pub serial_coordination: Option<SerialCoordinationProbeReceipt>,
171}
172
173#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
174#[serde(deny_unknown_fields)]
175pub struct FounderProviderAdminGrant {
176    pub grant_id: ProviderAdminGrantId,
177    pub provider: ProviderDeviceBinding,
178    pub access: ProviderAccessLocator,
179    pub capability: ProviderCapabilityProof,
180}
181
182impl ProviderCapabilityProof {
183    pub fn verify(
184        &self,
185        store: &StoreProviderBinding,
186        device: &ProviderDeviceBinding,
187        serial_required: bool,
188    ) -> Result<(), ProviderProbeError> {
189        self.exact_slots.verify(store, device)?;
190        match (serial_required, &self.serial_coordination) {
191            (true, Some(receipt)) => receipt.verify(store, device),
192            (false, None) => Ok(()),
193            (true, None) => Err(ProviderProbeError::InvalidReceipt(
194                "Serial Store has no serial coordination receipt".to_string(),
195            )),
196            (false, Some(_)) => Err(ProviderProbeError::InvalidReceipt(
197                "Merge-concurrent Store carries a serial coordination receipt".to_string(),
198            )),
199        }
200    }
201}
202
203#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
204#[serde(deny_unknown_fields)]
205pub struct ExactSlotProbeReceipt {
206    pub transcript: ExactSlotProbeTranscript,
207    pub transcript_hash: ObjectHash,
208}
209
210impl ExactSlotProbeReceipt {
211    pub fn from_transcript(
212        transcript: ExactSlotProbeTranscript,
213        store: &StoreProviderBinding,
214        device: &ProviderDeviceBinding,
215    ) -> Self {
216        let transcript_hash = exact_transcript_hash(store, device, &transcript);
217        Self {
218            transcript,
219            transcript_hash,
220        }
221    }
222
223    pub fn verify(
224        &self,
225        store: &StoreProviderBinding,
226        device: &ProviderDeviceBinding,
227    ) -> Result<(), ProviderProbeError> {
228        store.validate().map_err(ProviderProbeError::Storage)?;
229        device
230            .validate_for(store)
231            .map_err(ProviderProbeError::Storage)?;
232        let t = &self.transcript;
233        if self.transcript_hash != exact_transcript_hash(store, device, t) {
234            return invalid("exact-slot transcript hash does not match its context");
235        }
236        if t.logical_key != t.slot.logical_key() || t.accepted.slot() != &t.slot {
237            return invalid("exact-slot transcript disagrees with its allocated slot");
238        }
239        let payloads = [
240            probe_payload(&t.probe_id, ProbePayloadLabel::ExactCreateFirst),
241            probe_payload(&t.probe_id, ProbePayloadLabel::ExactCreateSecond),
242        ];
243        let expected_hashes = [
244            ObjectHash::digest(&payloads[0]),
245            ObjectHash::digest(&payloads[1]),
246        ];
247        if t.contenders[0].payload_hash != expected_hashes[0]
248            || t.contenders[1].payload_hash != expected_hashes[1]
249        {
250            return invalid("exact-slot contender payload hashes are not deterministic");
251        }
252        let winners: Vec<_> = t
253            .contenders
254            .iter()
255            .enumerate()
256            .filter_map(|(index, attempt)| {
257                (attempt.outcome == ProbeCreateOutcome::Created).then_some(index)
258            })
259            .collect();
260        let rejected = t
261            .contenders
262            .iter()
263            .filter(|attempt| attempt.outcome == ProbeCreateOutcome::RejectedOccupied)
264            .count();
265        if winners.len() != 1 || rejected != 1 {
266            return invalid("exact-slot race must contain one create and one occupied rejection");
267        }
268        let winner = &payloads[winners[0]];
269        if t.accepted.stored_size() != winner.len() as u64
270            || t.accepted.stored_hash() != ObjectHash::digest(winner)
271            || t.full_read_hash != ObjectHash::digest(winner)
272            || t.range.start != PROBE_RANGE_START
273            || t.range.end != PROBE_RANGE_END
274            || t.range.bytes_hash
275                != ObjectHash::digest(&winner[PROBE_RANGE_START as usize..PROBE_RANGE_END as usize])
276            || !t.delete_verified_absent
277        {
278            return invalid("exact-slot read, range, reference, or deletion evidence is invalid");
279        }
280        let lost = probe_payload(&t.probe_id, ProbePayloadLabel::LostResponse);
281        let lost_hash = ObjectHash::digest(&lost);
282        if t.lost_response.logical_key != t.lost_response.slot.logical_key()
283            || t.lost_response.settled.slot() != &t.lost_response.slot
284            || t.lost_response.payload_hash != lost_hash
285            || t.lost_response.settled.stored_size() != lost.len() as u64
286            || t.lost_response.settled.stored_hash() != lost_hash
287            || t.lost_response.readback_hash != lost_hash
288            || !t.lost_response.delete_verified_absent
289        {
290            return invalid("lost-response exact-slot evidence is invalid");
291        }
292        Ok(())
293    }
294}
295
296#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
297#[serde(deny_unknown_fields)]
298pub struct ExactSlotProbeTranscript {
299    pub probe_id: ProviderProbeId,
300    pub logical_key: String,
301    pub slot: ObjectSlot,
302    pub contenders: [ProbeCreateAttempt; 2],
303    pub accepted: ExactObjectRef,
304    pub full_read_hash: ObjectHash,
305    pub range: ProbeRangeReceipt,
306    pub delete_verified_absent: bool,
307    pub lost_response: LostResponseProbeReceipt,
308}
309
310#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
311#[serde(deny_unknown_fields)]
312pub struct ProbeCreateAttempt {
313    pub payload_hash: ObjectHash,
314    pub outcome: ProbeCreateOutcome,
315}
316
317#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
318#[serde(rename_all = "snake_case")]
319pub enum ProbeCreateOutcome {
320    Created,
321    RejectedOccupied,
322}
323
324#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
325#[serde(deny_unknown_fields)]
326pub struct ProbeReplaceAttempt {
327    pub payload_hash: ObjectHash,
328    pub outcome: ProbeReplaceOutcome,
329}
330
331#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
332#[serde(rename_all = "snake_case")]
333pub enum ProbeReplaceOutcome {
334    Replaced,
335    RejectedVersionMismatch,
336}
337
338#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
339#[serde(deny_unknown_fields)]
340pub struct LostResponseProbeReceipt {
341    pub logical_key: String,
342    pub slot: ObjectSlot,
343    pub payload_hash: ObjectHash,
344    pub settled: ExactObjectRef,
345    pub readback_hash: ObjectHash,
346    pub delete_verified_absent: bool,
347}
348
349#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
350#[serde(deny_unknown_fields)]
351pub struct ProbeRangeReceipt {
352    pub start: u64,
353    pub end: u64,
354    pub bytes_hash: ObjectHash,
355}
356
357#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
358#[serde(deny_unknown_fields)]
359pub struct SerialCoordinationProbeReceipt {
360    pub transcript: SerialCoordinationProbeTranscript,
361    pub transcript_hash: ObjectHash,
362}
363
364impl SerialCoordinationProbeReceipt {
365    pub fn from_transcript(
366        transcript: SerialCoordinationProbeTranscript,
367        store: &StoreProviderBinding,
368        device: &ProviderDeviceBinding,
369    ) -> Self {
370        Self {
371            transcript_hash: serial_transcript_hash(store, device, &transcript),
372            transcript,
373        }
374    }
375
376    pub fn verify(
377        &self,
378        store: &StoreProviderBinding,
379        device: &ProviderDeviceBinding,
380    ) -> Result<(), ProviderProbeError> {
381        if self.transcript_hash != serial_transcript_hash(store, device, &self.transcript) {
382            return invalid("serial transcript hash does not match its context");
383        }
384        let t = &self.transcript;
385        let create_payloads = [
386            probe_payload(&t.probe_id, ProbePayloadLabel::SerialCreateFirst),
387            probe_payload(&t.probe_id, ProbePayloadLabel::SerialCreateSecond),
388        ];
389        let replace_payloads = [
390            probe_payload(&t.probe_id, ProbePayloadLabel::SerialReplaceFirst),
391            probe_payload(&t.probe_id, ProbePayloadLabel::SerialReplaceSecond),
392        ];
393        verify_race_create(&t.create_attempts, &create_payloads)?;
394        verify_race_replace(&t.replace_attempts, &replace_payloads)?;
395        let created_winner = t
396            .create_attempts
397            .iter()
398            .position(|attempt| attempt.outcome == ProbeCreateOutcome::Created)
399            .expect("race verifier requires a winner");
400        let replaced_winner = t
401            .replace_attempts
402            .iter()
403            .position(|attempt| attempt.outcome == ProbeReplaceOutcome::Replaced)
404            .expect("race verifier requires a winner");
405        if t.created_bytes_hash != ObjectHash::digest(&create_payloads[created_winner])
406            || t.replaced_bytes_hash != ObjectHash::digest(&replace_payloads[replaced_winner])
407            || t.authoritative_read_bytes_hash != t.replaced_bytes_hash
408            || t.created_version_hash == ObjectHash::digest(&[])
409            || t.replaced_version_hash == ObjectHash::digest(&[])
410            || t.authoritative_read_version_hash != t.replaced_version_hash
411            || !t.delete_verified_absent
412        {
413            return invalid(
414                "serial winner, version, authoritative read, or deletion evidence is invalid",
415            );
416        }
417        Ok(())
418    }
419}
420
421#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
422#[serde(deny_unknown_fields)]
423pub struct SerialCoordinationProbeTranscript {
424    pub probe_id: ProviderProbeId,
425    pub logical_key: String,
426    pub create_attempts: [ProbeCreateAttempt; 2],
427    pub created_bytes_hash: ObjectHash,
428    pub created_version_hash: ObjectHash,
429    pub replace_attempts: [ProbeReplaceAttempt; 2],
430    pub replaced_bytes_hash: ObjectHash,
431    pub replaced_version_hash: ObjectHash,
432    pub authoritative_read_bytes_hash: ObjectHash,
433    pub authoritative_read_version_hash: ObjectHash,
434    pub delete_verified_absent: bool,
435}
436
437#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
438#[serde(deny_unknown_fields)]
439pub struct ProbeExactObjectReceipt {
440    pub slot: ObjectSlot,
441    pub payload_hash: ObjectHash,
442    pub object: ExactObjectRef,
443}
444
445#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
446#[serde(rename_all = "snake_case", deny_unknown_fields)]
447pub enum CrossPrincipalProviderEvidence {
448    GoogleSharedDrive,
449    DropboxSharedNamespace,
450    OneDriveSharedFolder,
451    CloudKit(CloudKitAcceptedShare),
452}
453
454#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
455#[serde(deny_unknown_fields)]
456pub struct CloudKitAcceptedShare {
457    pub share: ExactObjectRef,
458    pub share_record_name: String,
459    pub owner_name: String,
460    pub zone_name: String,
461    pub participant_record_name: String,
462}
463
464#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
465#[serde(deny_unknown_fields)]
466pub struct CrossPrincipalProbeTranscript {
467    pub challenge: CrossPrincipalProbeChallenge,
468    pub response: CrossPrincipalProbeResponse,
469    pub administrator_read_peer_hash: ObjectHash,
470    pub administrator_delete_peer_verified_absent: bool,
471    pub administrator_delete_own_verified_absent: bool,
472}
473
474#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
475#[serde(deny_unknown_fields)]
476pub struct CrossPrincipalProbeChallenge {
477    pub probe_id: ProviderProbeId,
478    pub administrator_object: ProbeExactObjectReceipt,
479    pub challenge_hash: ObjectHash,
480    pub administrator_signature: String,
481}
482
483#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
484#[serde(deny_unknown_fields)]
485pub struct CrossPrincipalProbeResponse {
486    pub challenge_hash: ObjectHash,
487    pub provider_evidence: CrossPrincipalProviderEvidence,
488    pub peer_object: ProbeExactObjectReceipt,
489    pub peer_read_administrator_hash: ObjectHash,
490    pub response_hash: ObjectHash,
491    pub peer_signature: String,
492}
493
494#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
495#[serde(deny_unknown_fields)]
496pub struct CrossPrincipalProbeReceipt {
497    pub transcript: CrossPrincipalProbeTranscript,
498    pub transcript_hash: ObjectHash,
499    pub administrator_completion_signature: String,
500}
501
502#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
503#[serde(deny_unknown_fields)]
504pub struct CrossPrincipalChallengeContext {
505    pub root: StoreRootRef,
506    pub attempt_id: DeviceJoinAttemptId,
507    pub access_request_hash: ObjectHash,
508    pub provider_admin_grant: ProviderAdminGrantId,
509    pub owner_registration: StoreDeviceRegistrationRef,
510    pub member_pubkey: String,
511    pub administrator_binding: ProviderDeviceBinding,
512    pub peer_binding: ProviderDeviceBinding,
513}
514
515#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
516#[serde(deny_unknown_fields)]
517pub struct CrossPrincipalResponseContext {
518    pub challenge: CrossPrincipalChallengeContext,
519    pub expected_registration_hash: ObjectHash,
520    pub response_slot: ObjectSlot,
521}
522
523#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
524#[serde(deny_unknown_fields)]
525pub struct DeviceJoinChallengePublicationAuthorization {
526    pub attempt: DeviceJoinAttemptRef,
527    pub attempt_activation: StoreBatchCommitRef,
528}
529
530#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
531#[serde(deny_unknown_fields)]
532pub struct DeviceJoinChallengePublicationRecord {
533    pub challenge: CrossPrincipalProbeChallenge,
534    pub progress: DeviceJoinChallengePublicationProgress,
535}
536
537#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
538#[serde(rename_all = "snake_case", deny_unknown_fields)]
539pub enum DeviceJoinChallengePublicationProgress {
540    Prepared,
541    Published {
542        authorization: DeviceJoinChallengePublicationAuthorization,
543    },
544    ProducerClosed {
545        authorization: DeviceJoinChallengePublicationAuthorization,
546    },
547    CancelledBeforeCreate {
548        authorization: DeviceJoinChallengePublicationAuthorization,
549        cancellation: DeviceJoinOutcomeRef,
550    },
551}
552
553#[async_trait]
554pub trait DeviceJoinChallengePublicationJournal: Send + Sync {
555    async fn prepare(
556        &self,
557        challenge: &CrossPrincipalProbeChallenge,
558    ) -> Result<DeviceJoinChallengePublicationRecord, StorageError>;
559
560    /// Atomically claims publication for these exact signed facts. An exact
561    /// replay of an existing `Published` claim succeeds; a producer closure or
562    /// cancellation claim rejects publication.
563    async fn claim_published(
564        &self,
565        authorization: &DeviceJoinChallengePublicationAuthorization,
566        challenge: &CrossPrincipalProbeChallenge,
567    ) -> Result<(), StorageError>;
568
569    async fn close_published(
570        &self,
571        authorization: &DeviceJoinChallengePublicationAuthorization,
572        challenge: &CrossPrincipalProbeChallenge,
573    ) -> Result<(), StorageError>;
574
575    async fn cancel_before_create(
576        &self,
577        authorization: &DeviceJoinChallengePublicationAuthorization,
578        challenge: &CrossPrincipalProbeChallenge,
579        cancellation: &DeviceJoinOutcomeRef,
580    ) -> Result<(), StorageError>;
581}
582
583#[async_trait]
584impl DeviceJoinChallengePublicationJournal for crate::database::Database {
585    async fn prepare(
586        &self,
587        challenge: &CrossPrincipalProbeChallenge,
588    ) -> Result<DeviceJoinChallengePublicationRecord, StorageError> {
589        self.prepare_device_join_challenge_publication(challenge.clone())
590            .await
591            .map_err(|error| StorageError::Storage(error.to_string()))
592    }
593
594    async fn claim_published(
595        &self,
596        authorization: &DeviceJoinChallengePublicationAuthorization,
597        challenge: &CrossPrincipalProbeChallenge,
598    ) -> Result<(), StorageError> {
599        self.publish_device_join_challenge(authorization.clone(), challenge.clone())
600            .await
601            .map_err(|error| StorageError::Storage(error.to_string()))
602    }
603
604    async fn close_published(
605        &self,
606        authorization: &DeviceJoinChallengePublicationAuthorization,
607        challenge: &CrossPrincipalProbeChallenge,
608    ) -> Result<(), StorageError> {
609        self.close_published_device_join_challenge(authorization.clone(), challenge.clone())
610            .await
611            .map_err(|error| StorageError::Storage(error.to_string()))
612    }
613
614    async fn cancel_before_create(
615        &self,
616        authorization: &DeviceJoinChallengePublicationAuthorization,
617        challenge: &CrossPrincipalProbeChallenge,
618        cancellation: &DeviceJoinOutcomeRef,
619    ) -> Result<(), StorageError> {
620        self.cancel_unpublished_device_join_challenge(
621            authorization.clone(),
622            challenge.clone(),
623            cancellation.clone(),
624        )
625        .await
626        .map_err(|error| StorageError::Storage(error.to_string()))
627    }
628}
629
630impl CrossPrincipalProbeReceipt {
631    fn signed(
632        transcript: CrossPrincipalProbeTranscript,
633        context: &CrossPrincipalResponseContext,
634        store: &StoreProviderBinding,
635        administrator_signer: &UserKeypair,
636    ) -> Result<Self, ProviderProbeError> {
637        validate_cross_transcript_payloads(&transcript, context)?;
638        let transcript_hash = cross_transcript_hash(store, context, &transcript);
639        Ok(Self {
640            transcript,
641            transcript_hash,
642            administrator_completion_signature: hex::encode(
643                administrator_signer.sign(transcript_hash.as_bytes()),
644            ),
645        })
646    }
647
648    pub fn verify(
649        &self,
650        context: &CrossPrincipalResponseContext,
651        store: &StoreProviderBinding,
652        administrator_signing_pubkey: &str,
653        peer_signing_pubkey: &str,
654    ) -> Result<(), ProviderProbeError> {
655        validate_cross_provider_evidence(
656            store,
657            &context.challenge.administrator_binding,
658            &context.challenge.peer_binding,
659            &self.transcript.response.provider_evidence,
660        )?;
661        self.transcript.challenge.verify(
662            &context.challenge,
663            store,
664            administrator_signing_pubkey,
665        )?;
666        self.transcript.response.verify(
667            &self.transcript.challenge,
668            context,
669            store,
670            administrator_signing_pubkey,
671            peer_signing_pubkey,
672        )?;
673        validate_cross_transcript_payloads(&self.transcript, context)?;
674        let expected_hash = cross_transcript_hash(store, context, &self.transcript);
675        if self.transcript_hash != expected_hash {
676            return invalid("cross-principal transcript hash does not match its join context");
677        }
678        if !crate::keys::verify_signature_hex(
679            administrator_signing_pubkey,
680            &self.administrator_completion_signature,
681            self.transcript_hash.as_bytes(),
682        ) {
683            return invalid("cross-principal completion signature is invalid");
684        }
685        Ok(())
686    }
687}
688
689impl CrossPrincipalProbeChallenge {
690    pub fn verify(
691        &self,
692        context: &CrossPrincipalChallengeContext,
693        store: &StoreProviderBinding,
694        administrator_signing_pubkey: &str,
695    ) -> Result<(), ProviderProbeError> {
696        validate_cross_challenge_payload(self)?;
697        validate_cross_provider_evidence_context(store, context)?;
698        let expected_hash = cross_challenge_hash(store, context, self);
699        if self.challenge_hash != expected_hash {
700            return invalid("cross-principal challenge hash does not match its join context");
701        }
702        if !crate::keys::verify_signature_hex(
703            administrator_signing_pubkey,
704            &self.administrator_signature,
705            self.challenge_hash.as_bytes(),
706        ) {
707            return invalid("cross-principal challenge signature is invalid");
708        }
709        Ok(())
710    }
711}
712
713impl CrossPrincipalProbeResponse {
714    pub fn verify(
715        &self,
716        challenge: &CrossPrincipalProbeChallenge,
717        context: &CrossPrincipalResponseContext,
718        store: &StoreProviderBinding,
719        administrator_signing_pubkey: &str,
720        peer_signing_pubkey: &str,
721    ) -> Result<(), ProviderProbeError> {
722        challenge.verify(&context.challenge, store, administrator_signing_pubkey)?;
723        if context.challenge.member_pubkey != peer_signing_pubkey {
724            return invalid("cross-principal response signer is not the joining member");
725        }
726        validate_cross_provider_evidence(
727            store,
728            &context.challenge.administrator_binding,
729            &context.challenge.peer_binding,
730            &self.provider_evidence,
731        )?;
732        validate_cross_response_payload(self, challenge, context)?;
733        let expected_hash = cross_response_hash(store, context, challenge, self);
734        if self.response_hash != expected_hash {
735            return invalid("cross-principal response hash does not match its join context");
736        }
737        if !crate::keys::verify_signature_hex(
738            peer_signing_pubkey,
739            &self.peer_signature,
740            self.response_hash.as_bytes(),
741        ) {
742            return invalid("cross-principal response signature is invalid");
743        }
744        Ok(())
745    }
746}
747
748#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
749#[serde(transparent)]
750pub struct ProviderAdminGrantId(pub ObjectHash);
751
752impl ProviderAdminGrantId {
753    pub fn from_random_bytes(bytes: [u8; 32]) -> Self {
754        Self(ObjectHash::from_digest(bytes))
755    }
756}
757
758#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
759#[serde(transparent)]
760pub struct ProviderAccessGrantId(pub ObjectHash);
761
762impl ProviderAccessGrantId {
763    pub fn from_random_bytes(bytes: [u8; 32]) -> Self {
764        Self(ObjectHash::from_digest(bytes))
765    }
766}
767
768/// Stable provider authority that can be withdrawn without rediscovering a
769/// member by mutable account metadata.
770#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
771#[serde(rename_all = "snake_case", deny_unknown_fields)]
772pub enum ProviderAccessLocator {
773    S3SharedCredentialGeneration {
774        generation: u64,
775        access_key_id_hash: ObjectHash,
776    },
777    GoogleDrivePermission {
778        drive_id: String,
779        permission_id: String,
780    },
781    DropboxSharedFolderMember {
782        namespace_id: String,
783        account_id: String,
784    },
785    OneDrivePermission {
786        drive_id: String,
787        item_id: String,
788        permission_id: String,
789    },
790    CloudKitPrivateZoneOwner {
791        owner_name: String,
792        zone_name: String,
793        owner_record_name: String,
794    },
795    CloudKitParticipant {
796        share_record_name: String,
797        owner_name: String,
798        zone_name: String,
799        participant_record_name: String,
800    },
801}
802
803impl ProviderAccessLocator {
804    pub fn for_current_administrator(
805        binding: &crate::sync::storage::ResolvedProviderBinding,
806    ) -> Result<Self, StorageError> {
807        binding.validate()?;
808        match (&binding.store, &binding.device.principal) {
809            (
810                StoreProviderBinding::S3 { .. },
811                crate::sync::storage::ProviderPrincipalId::CustomS3Credential {
812                    access_key_id_hash,
813                },
814            ) => Ok(Self::S3SharedCredentialGeneration {
815                generation: 1,
816                access_key_id_hash: *access_key_id_hash,
817            }),
818            (
819                StoreProviderBinding::GoogleDrive {
820                    corpus: crate::sync::storage::GoogleDriveCorpus::SharedDrive { drive_id, .. },
821                },
822                crate::sync::storage::ProviderPrincipalId::GoogleDrive { permission_id },
823            ) => Ok(Self::GoogleDrivePermission {
824                drive_id: drive_id.clone(),
825                permission_id: permission_id.clone(),
826            }),
827            (
828                StoreProviderBinding::Dropbox { namespace_id },
829                crate::sync::storage::ProviderPrincipalId::Dropbox { account_id },
830            ) => Ok(Self::DropboxSharedFolderMember {
831                namespace_id: namespace_id.clone(),
832                account_id: account_id.clone(),
833            }),
834            (
835                StoreProviderBinding::CloudKit {
836                    owner_name,
837                    zone_name,
838                    ..
839                },
840                crate::sync::storage::ProviderPrincipalId::CloudKitPrivateZoneOwner { record_name },
841            ) => Ok(Self::CloudKitPrivateZoneOwner {
842                owner_name: owner_name.clone(),
843                zone_name: zone_name.clone(),
844                owner_record_name: record_name.clone(),
845            }),
846            _ => Err(StorageError::Configuration(
847                "provider adapter did not expose the administrator's exact access locator"
848                    .to_string(),
849            )),
850        }
851    }
852
853    pub fn validate_for(
854        &self,
855        store: &StoreProviderBinding,
856        provider: &ProviderDeviceBinding,
857    ) -> Result<(), StorageError> {
858        provider.validate_for(store)?;
859        let valid = match (store, &provider.principal, self) {
860            (
861                StoreProviderBinding::S3 { .. },
862                crate::sync::storage::ProviderPrincipalId::CustomS3Credential {
863                    access_key_id_hash: provider_hash,
864                },
865                Self::S3SharedCredentialGeneration {
866                    generation,
867                    access_key_id_hash,
868                },
869            ) => *generation > 0 && provider_hash == access_key_id_hash,
870            (
871                StoreProviderBinding::S3 { .. },
872                crate::sync::storage::ProviderPrincipalId::Aws { .. },
873                Self::S3SharedCredentialGeneration { generation, .. },
874            ) => *generation > 0,
875            (
876                StoreProviderBinding::GoogleDrive {
877                    corpus: crate::sync::storage::GoogleDriveCorpus::SharedDrive { drive_id, .. },
878                },
879                crate::sync::storage::ProviderPrincipalId::GoogleDrive { permission_id },
880                Self::GoogleDrivePermission {
881                    drive_id: locator_drive,
882                    permission_id: locator_permission,
883                },
884            ) => drive_id == locator_drive && permission_id == locator_permission,
885            (
886                StoreProviderBinding::Dropbox { namespace_id },
887                crate::sync::storage::ProviderPrincipalId::Dropbox { account_id },
888                Self::DropboxSharedFolderMember {
889                    namespace_id: locator_namespace,
890                    account_id: locator_account,
891                },
892            ) => namespace_id == locator_namespace && account_id == locator_account,
893            (
894                StoreProviderBinding::OneDrive {
895                    drive_id,
896                    folder_id,
897                },
898                crate::sync::storage::ProviderPrincipalId::OneDrive { .. },
899                Self::OneDrivePermission {
900                    drive_id: locator_drive,
901                    item_id,
902                    permission_id,
903                },
904            ) => drive_id == locator_drive && folder_id == item_id && !permission_id.is_empty(),
905            (
906                StoreProviderBinding::CloudKit {
907                    owner_name,
908                    zone_name,
909                    ..
910                },
911                crate::sync::storage::ProviderPrincipalId::CloudKitPrivateZoneOwner { record_name },
912                Self::CloudKitPrivateZoneOwner {
913                    owner_name: locator_owner,
914                    zone_name: locator_zone,
915                    owner_record_name,
916                },
917            ) => {
918                owner_name == locator_owner
919                    && zone_name == locator_zone
920                    && record_name == owner_record_name
921            }
922            (
923                StoreProviderBinding::CloudKit {
924                    owner_name,
925                    zone_name,
926                    ..
927                },
928                crate::sync::storage::ProviderPrincipalId::CloudKitSharedZoneParticipant {
929                    record_name,
930                },
931                Self::CloudKitParticipant {
932                    share_record_name,
933                    owner_name: locator_owner,
934                    zone_name: locator_zone,
935                    participant_record_name,
936                },
937            ) => {
938                !share_record_name.is_empty()
939                    && owner_name == locator_owner
940                    && zone_name == locator_zone
941                    && record_name == participant_record_name
942            }
943            _ => false,
944        };
945        if valid {
946            Ok(())
947        } else {
948            Err(StorageError::Configuration(
949                "provider access locator differs from its Store and provider binding".to_string(),
950            ))
951        }
952    }
953}
954
955#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
956#[serde(deny_unknown_fields)]
957pub struct StoreMemberProviderAccessGrant {
958    pub grant_id: ProviderAccessGrantId,
959    pub member_pubkey: String,
960    pub provider: ProviderDeviceBinding,
961    pub locator: ProviderAccessLocator,
962    pub administrator_grant: ProviderAdminGrantId,
963    pub administrator: StoreDeviceRegistrationRef,
964    pub signature: String,
965}
966
967#[derive(Serialize)]
968struct StoreMemberProviderAccessGrantSignedFields<'a> {
969    grant_id: &'a ProviderAccessGrantId,
970    member_pubkey: &'a str,
971    provider: &'a ProviderDeviceBinding,
972    locator: &'a ProviderAccessLocator,
973    administrator_grant: &'a ProviderAdminGrantId,
974    administrator: &'a StoreDeviceRegistrationRef,
975}
976
977impl StoreMemberProviderAccessGrant {
978    #[allow(clippy::too_many_arguments)]
979    pub fn signed(
980        grant_id: ProviderAccessGrantId,
981        member_pubkey: String,
982        provider: ProviderDeviceBinding,
983        locator: ProviderAccessLocator,
984        administrator_grant: ProviderAdminGrantId,
985        administrator: StoreDeviceRegistrationRef,
986        store: &StoreProviderBinding,
987        administrator_registration: &StoreDeviceRegistration,
988        administrator_signer: &UserKeypair,
989    ) -> Result<Self, ProviderProbeError> {
990        administrator
991            .verify_registration(administrator_registration)
992            .map_err(|error| ProviderProbeError::InvalidReceipt(error.to_string()))?;
993        if crate::keys::public_key_hex(administrator_signer)
994            != administrator_registration.device_signing_pubkey
995        {
996            return invalid("provider access grant signer is not the administrator device");
997        }
998        locator.validate_for(store, &provider)?;
999        let mut grant = Self {
1000            grant_id,
1001            member_pubkey,
1002            provider,
1003            locator,
1004            administrator_grant,
1005            administrator,
1006            signature: String::new(),
1007        };
1008        grant.signature = hex::encode(
1009            administrator_signer
1010                .sign(ObjectHash::digest(&grant.canonical_signed_bytes()).as_bytes()),
1011        );
1012        Ok(grant)
1013    }
1014
1015    fn canonical_signed_bytes(&self) -> Vec<u8> {
1016        domain_json(
1017            MEMBER_ACCESS_GRANT_DOMAIN,
1018            &StoreMemberProviderAccessGrantSignedFields {
1019                grant_id: &self.grant_id,
1020                member_pubkey: &self.member_pubkey,
1021                provider: &self.provider,
1022                locator: &self.locator,
1023                administrator_grant: &self.administrator_grant,
1024                administrator: &self.administrator,
1025            },
1026        )
1027    }
1028
1029    pub fn grant_hash(&self) -> ObjectHash {
1030        ObjectHash::digest(&self.canonical_signed_bytes())
1031    }
1032
1033    pub fn to_bytes(&self) -> Vec<u8> {
1034        serde_json::to_vec(self).expect("provider member access grant serialization cannot fail")
1035    }
1036
1037    pub fn verify(
1038        &self,
1039        store: &StoreProviderBinding,
1040        administrator: &StoreDeviceRegistration,
1041    ) -> Result<(), ProviderProbeError> {
1042        self.administrator
1043            .verify_registration(administrator)
1044            .map_err(|error| ProviderProbeError::InvalidReceipt(error.to_string()))?;
1045        self.provider
1046            .validate_for(store)
1047            .map_err(ProviderProbeError::Storage)?;
1048        self.locator
1049            .validate_for(store, &self.provider)
1050            .map_err(ProviderProbeError::Storage)?;
1051        if !crate::keys::verify_signature_hex(
1052            &administrator.device_signing_pubkey,
1053            &self.signature,
1054            self.grant_hash().as_bytes(),
1055        ) {
1056            return invalid("provider access grant signature is invalid");
1057        }
1058        Ok(())
1059    }
1060}
1061
1062#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
1063#[serde(deny_unknown_fields)]
1064pub struct StoreMemberProviderAccessGrantRef {
1065    pub grant_id: ProviderAccessGrantId,
1066    pub grant_hash: ObjectHash,
1067    pub object: ExactObjectRef,
1068}
1069
1070impl StoreMemberProviderAccessGrantRef {
1071    pub fn from_grant(grant: &StoreMemberProviderAccessGrant, object: ExactObjectRef) -> Self {
1072        Self {
1073            grant_id: grant.grant_id.clone(),
1074            grant_hash: grant.grant_hash(),
1075            object,
1076        }
1077    }
1078
1079    pub fn verify(&self, grant: &StoreMemberProviderAccessGrant) -> Result<(), ProviderProbeError> {
1080        if self.grant_id != grant.grant_id || self.grant_hash != grant.grant_hash() {
1081            return invalid("provider access grant reference differs from its signed grant");
1082        }
1083        Ok(())
1084    }
1085}
1086
1087#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1088#[serde(deny_unknown_fields)]
1089pub struct ActivatedStoreMemberProviderAccessGrant {
1090    pub grant: StoreMemberProviderAccessGrant,
1091    pub grant_ref: StoreMemberProviderAccessGrantRef,
1092    pub activation: StoreBatchCommitRef,
1093}
1094
1095#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1096#[serde(rename_all = "snake_case", deny_unknown_fields)]
1097pub enum ProviderAccessWithdrawal {
1098    Direct {
1099        locator: ProviderAccessLocator,
1100        verified_absent: bool,
1101    },
1102    S3CredentialRotation {
1103        retired_generation: u64,
1104        active_generation: u64,
1105        retired_credential_verified_rejected: bool,
1106    },
1107}
1108
1109#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1110#[serde(deny_unknown_fields)]
1111pub struct StoreMemberProviderAccessWithdrawalReceipt {
1112    pub grant: StoreMemberProviderAccessGrantRef,
1113    pub withdrawal: ProviderAccessWithdrawal,
1114    pub administrator_grant: ProviderAdminGrantId,
1115    pub administrator: StoreDeviceRegistrationRef,
1116    pub signature: String,
1117}
1118
1119#[derive(Serialize)]
1120struct StoreMemberProviderAccessWithdrawalSignedFields<'a> {
1121    grant: &'a StoreMemberProviderAccessGrantRef,
1122    withdrawal: &'a ProviderAccessWithdrawal,
1123    administrator_grant: &'a ProviderAdminGrantId,
1124    administrator: &'a StoreDeviceRegistrationRef,
1125}
1126
1127impl StoreMemberProviderAccessWithdrawalReceipt {
1128    pub fn signed(
1129        grant: StoreMemberProviderAccessGrantRef,
1130        withdrawal: ProviderAccessWithdrawal,
1131        administrator_grant: ProviderAdminGrantId,
1132        administrator: StoreDeviceRegistrationRef,
1133        administrator_registration: &StoreDeviceRegistration,
1134        administrator_signer: &UserKeypair,
1135    ) -> Result<Self, ProviderProbeError> {
1136        administrator
1137            .verify_registration(administrator_registration)
1138            .map_err(|error| ProviderProbeError::InvalidReceipt(error.to_string()))?;
1139        if crate::keys::public_key_hex(administrator_signer)
1140            != administrator_registration.device_signing_pubkey
1141        {
1142            return invalid("provider access withdrawal signer is not the administrator device");
1143        }
1144        validate_withdrawal(&withdrawal)?;
1145        let mut receipt = Self {
1146            grant,
1147            withdrawal,
1148            administrator_grant,
1149            administrator,
1150            signature: String::new(),
1151        };
1152        receipt.signature = hex::encode(
1153            administrator_signer
1154                .sign(ObjectHash::digest(&receipt.canonical_signed_bytes()).as_bytes()),
1155        );
1156        Ok(receipt)
1157    }
1158
1159    fn canonical_signed_bytes(&self) -> Vec<u8> {
1160        domain_json(
1161            MEMBER_ACCESS_WITHDRAWAL_DOMAIN,
1162            &StoreMemberProviderAccessWithdrawalSignedFields {
1163                grant: &self.grant,
1164                withdrawal: &self.withdrawal,
1165                administrator_grant: &self.administrator_grant,
1166                administrator: &self.administrator,
1167            },
1168        )
1169    }
1170
1171    pub fn receipt_hash(&self) -> ObjectHash {
1172        ObjectHash::digest(&self.canonical_signed_bytes())
1173    }
1174
1175    pub fn to_bytes(&self) -> Vec<u8> {
1176        serde_json::to_vec(self)
1177            .expect("provider member access withdrawal serialization cannot fail")
1178    }
1179
1180    pub fn verify(
1181        &self,
1182        administrator: &StoreDeviceRegistration,
1183    ) -> Result<(), ProviderProbeError> {
1184        self.administrator
1185            .verify_registration(administrator)
1186            .map_err(|error| ProviderProbeError::InvalidReceipt(error.to_string()))?;
1187        validate_withdrawal(&self.withdrawal)?;
1188        if !crate::keys::verify_signature_hex(
1189            &administrator.device_signing_pubkey,
1190            &self.signature,
1191            self.receipt_hash().as_bytes(),
1192        ) {
1193            return invalid("provider access withdrawal signature is invalid");
1194        }
1195        Ok(())
1196    }
1197}
1198
1199#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
1200#[serde(deny_unknown_fields)]
1201pub struct StoreMemberProviderAccessWithdrawalReceiptRef {
1202    pub grant_id: ProviderAccessGrantId,
1203    pub receipt_hash: ObjectHash,
1204    pub object: ExactObjectRef,
1205}
1206
1207impl StoreMemberProviderAccessWithdrawalReceiptRef {
1208    pub fn from_receipt(
1209        receipt: &StoreMemberProviderAccessWithdrawalReceipt,
1210        object: ExactObjectRef,
1211    ) -> Self {
1212        Self {
1213            grant_id: receipt.grant.grant_id.clone(),
1214            receipt_hash: receipt.receipt_hash(),
1215            object,
1216        }
1217    }
1218
1219    pub fn verify(
1220        &self,
1221        receipt: &StoreMemberProviderAccessWithdrawalReceipt,
1222    ) -> Result<(), ProviderProbeError> {
1223        if self.grant_id != receipt.grant.grant_id || self.receipt_hash != receipt.receipt_hash() {
1224            return invalid("provider access withdrawal reference differs from its signed receipt");
1225        }
1226        Ok(())
1227    }
1228}
1229
1230fn validate_withdrawal(withdrawal: &ProviderAccessWithdrawal) -> Result<(), ProviderProbeError> {
1231    let valid = match withdrawal {
1232        ProviderAccessWithdrawal::Direct {
1233            verified_absent, ..
1234        } => *verified_absent,
1235        ProviderAccessWithdrawal::S3CredentialRotation {
1236            retired_generation,
1237            active_generation,
1238            retired_credential_verified_rejected,
1239        } => {
1240            *retired_generation > 0
1241                && *active_generation == retired_generation.saturating_add(1)
1242                && *retired_credential_verified_rejected
1243        }
1244    };
1245    if valid {
1246        Ok(())
1247    } else {
1248        invalid("provider access withdrawal does not prove the stored authority is unusable")
1249    }
1250}
1251
1252#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1253#[serde(deny_unknown_fields)]
1254pub struct ProviderAdminGrantRecord {
1255    pub grant_id: ProviderAdminGrantId,
1256    pub administrator: StoreDeviceRegistrationRef,
1257    pub provider: ProviderDeviceBinding,
1258    pub access: ProviderAccessLocator,
1259    pub capability: ProviderCapabilityProof,
1260    pub created_at: ProviderAdminGrantOrigin,
1261}
1262
1263#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1264#[serde(rename_all = "snake_case", deny_unknown_fields)]
1265pub enum ProviderAdminGrantOrigin {
1266    Founder { root: StoreRootRef },
1267    MergeMembership { coord: MembershipCoord },
1268    SerialCommit { commit: StoreBatchCommitRef },
1269}
1270
1271#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1272#[serde(rename_all = "snake_case", deny_unknown_fields)]
1273pub enum ProviderAdminMembershipChange {
1274    MergeConcurrent {
1275        change: ProviderAdminChange,
1276        #[serde(with = "ordered_owner_barriers")]
1277        owner_barriers: BTreeMap<MembershipGrantId, OwnerStreamBarrier>,
1278    },
1279    Serial {
1280        change: ProviderAdminChange,
1281    },
1282}
1283
1284#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1285#[serde(rename_all = "snake_case", deny_unknown_fields)]
1286pub enum ProviderAdminChange {
1287    Set {
1288        administrator: StoreDeviceRegistrationRef,
1289        provider: ProviderDeviceBinding,
1290        access: ProviderAccessLocator,
1291        capability: ProviderCapabilityProof,
1292        grant_id: ProviderAdminGrantId,
1293        replaces: BTreeSet<ProviderAdminGrantId>,
1294    },
1295    Remove {
1296        removes: BTreeSet<ProviderAdminGrantId>,
1297    },
1298}
1299
1300#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1301#[serde(deny_unknown_fields)]
1302pub struct ProviderAdminState {
1303    records: BTreeMap<ProviderAdminGrantId, ProviderAdminGrantRecord>,
1304    tombstones: BTreeSet<ProviderAdminGrantId>,
1305}
1306
1307#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1308#[serde(deny_unknown_fields)]
1309pub struct ProviderAdminBranch {
1310    pub heads: Vec<MembershipCoord>,
1311    pub state: ProviderAdminState,
1312}
1313
1314#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1315#[serde(deny_unknown_fields)]
1316pub struct ProviderAdminConflict {
1317    pub raw_heads: Vec<MembershipCoord>,
1318    pub cyclic_sources: Vec<MembershipCoord>,
1319    pub involved_grants: BTreeSet<ProviderAdminGrantId>,
1320    pub maximal_valid_branches: Vec<ProviderAdminBranch>,
1321    pub combined: ProviderAdminState,
1322}
1323
1324#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1325#[serde(rename_all = "snake_case", deny_unknown_fields)]
1326pub enum ProviderAdminResolution {
1327    Resolved(ProviderAdminState),
1328    RevocationConflict(ProviderAdminConflict),
1329}
1330
1331impl ProviderAdminResolution {
1332    pub fn combined_state(&self) -> &ProviderAdminState {
1333        match self {
1334            Self::Resolved(state) => state,
1335            Self::RevocationConflict(conflict) => &conflict.combined,
1336        }
1337    }
1338
1339    pub fn state_hash(&self) -> ObjectHash {
1340        ObjectHash::digest(&domain_json(b"coven.provider-admin-resolution.v1\0", self))
1341    }
1342}
1343
1344impl ProviderAdminState {
1345    pub fn founder(grant: ProviderAdminGrantRecord) -> Self {
1346        let grant_id = grant.grant_id.clone();
1347        Self {
1348            records: BTreeMap::from([(grant_id.clone(), grant)]),
1349            tombstones: BTreeSet::new(),
1350        }
1351    }
1352
1353    pub fn founder_from_root(
1354        root: StoreRootRef,
1355        administrator: StoreDeviceRegistrationRef,
1356        grant: &FounderProviderAdminGrant,
1357    ) -> Self {
1358        Self::founder(ProviderAdminGrantRecord {
1359            grant_id: grant.grant_id.clone(),
1360            administrator,
1361            provider: grant.provider.clone(),
1362            access: grant.access.clone(),
1363            capability: grant.capability.clone(),
1364            created_at: ProviderAdminGrantOrigin::Founder { root },
1365        })
1366    }
1367
1368    pub fn authorizes(
1369        &self,
1370        grant_id: &ProviderAdminGrantId,
1371        administrator: &StoreDeviceRegistrationRef,
1372    ) -> bool {
1373        !self.tombstones.contains(grant_id)
1374            && self
1375                .records
1376                .get(grant_id)
1377                .is_some_and(|record| &record.administrator == administrator)
1378    }
1379
1380    pub fn records(&self) -> &BTreeMap<ProviderAdminGrantId, ProviderAdminGrantRecord> {
1381        &self.records
1382    }
1383
1384    pub fn active(&self) -> BTreeSet<ProviderAdminGrantId> {
1385        self.records
1386            .keys()
1387            .filter(|grant_id| !self.tombstones.contains(*grant_id))
1388            .cloned()
1389            .collect()
1390    }
1391
1392    pub fn tombstones(&self) -> &BTreeSet<ProviderAdminGrantId> {
1393        &self.tombstones
1394    }
1395
1396    pub fn apply(
1397        &mut self,
1398        change: ProviderAdminChange,
1399        origin: ProviderAdminGrantOrigin,
1400    ) -> Result<(), ProviderAdminReducerError> {
1401        let mut next = self.clone();
1402        next.apply_unchecked(change, origin)?;
1403        if next.active().is_empty() {
1404            return Err(ProviderAdminReducerError::NoEffectiveAdministrator);
1405        }
1406        *self = next;
1407        Ok(())
1408    }
1409
1410    fn apply_unchecked(
1411        &mut self,
1412        change: ProviderAdminChange,
1413        origin: ProviderAdminGrantOrigin,
1414    ) -> Result<(), ProviderAdminReducerError> {
1415        match change {
1416            ProviderAdminChange::Set {
1417                administrator,
1418                provider,
1419                access,
1420                capability,
1421                grant_id,
1422                replaces,
1423            } => {
1424                let record = ProviderAdminGrantRecord {
1425                    grant_id: grant_id.clone(),
1426                    administrator,
1427                    provider,
1428                    access,
1429                    capability,
1430                    created_at: origin,
1431                };
1432                if let Some(existing) = self.records.get(&grant_id) {
1433                    if existing != &record {
1434                        return Err(ProviderAdminReducerError::GrantIdReuse);
1435                    }
1436                    if !replaces.iter().all(|id| self.tombstones.contains(id)) {
1437                        return Err(ProviderAdminReducerError::UnknownReplacement);
1438                    }
1439                    return Ok(());
1440                }
1441                if !replaces
1442                    .iter()
1443                    .all(|id| self.records.contains_key(id) && !self.tombstones.contains(id))
1444                {
1445                    return Err(ProviderAdminReducerError::UnknownReplacement);
1446                }
1447                for replaced in replaces {
1448                    self.tombstones.insert(replaced);
1449                }
1450                self.records.insert(grant_id, record);
1451            }
1452            ProviderAdminChange::Remove { removes } => {
1453                if removes.is_empty()
1454                    || !removes
1455                        .iter()
1456                        .all(|id| self.records.contains_key(id) || self.tombstones.contains(id))
1457                {
1458                    return Err(ProviderAdminReducerError::UnknownRemoval);
1459                }
1460                for removed in removes {
1461                    self.tombstones.insert(removed);
1462                }
1463            }
1464        }
1465        Ok(())
1466    }
1467
1468    pub(crate) fn apply_membership_change(
1469        &mut self,
1470        change: ProviderAdminMembershipChange,
1471        origin: ProviderAdminGrantOrigin,
1472    ) -> Result<(), ProviderAdminReducerError> {
1473        let change = match (change, &origin) {
1474            (
1475                ProviderAdminMembershipChange::MergeConcurrent {
1476                    change,
1477                    owner_barriers: _,
1478                },
1479                ProviderAdminGrantOrigin::MergeMembership { .. },
1480            ) => change,
1481            (
1482                ProviderAdminMembershipChange::Serial { change },
1483                ProviderAdminGrantOrigin::SerialCommit { .. },
1484            ) => change,
1485            _ => return Err(ProviderAdminReducerError::PolicyOriginMismatch),
1486        };
1487        self.apply(change, origin)
1488    }
1489
1490    pub fn state_hash(&self) -> ObjectHash {
1491        ObjectHash::digest(&domain_json(
1492            b"coven.provider-admin-state.v1\0",
1493            &(self.records(), self.tombstones()),
1494        ))
1495    }
1496
1497    pub fn merge(
1498        states: impl IntoIterator<Item = Self>,
1499    ) -> Result<Self, ProviderAdminReducerError> {
1500        let mut records = BTreeMap::new();
1501        let mut tombstones = BTreeSet::new();
1502        for state in states {
1503            for (grant_id, record) in state.records {
1504                if records
1505                    .insert(grant_id.clone(), record.clone())
1506                    .is_some_and(|current| current != record)
1507                {
1508                    return Err(ProviderAdminReducerError::GrantIdReuse);
1509                }
1510            }
1511            tombstones.extend(state.tombstones);
1512        }
1513        Ok(Self {
1514            records,
1515            tombstones,
1516        })
1517    }
1518
1519    pub(crate) fn reduce_merge(
1520        genesis: &Self,
1521        entries: &[MembershipEntry],
1522        included: &BTreeSet<MembershipCoord>,
1523    ) -> Result<ProviderAdminResolution, ProviderAdminReducerError> {
1524        let by_coord = entries
1525            .iter()
1526            .filter(|entry| included.contains(&entry.coord()))
1527            .map(|entry| (entry.coord(), entry))
1528            .collect::<BTreeMap<_, _>>();
1529        let mut states = BTreeMap::<MembershipCoord, Self>::new();
1530        let mut pending = by_coord.keys().cloned().collect::<BTreeSet<_>>();
1531        while !pending.is_empty() {
1532            let ready = pending.iter().find(|coord| {
1533                let entry = by_coord[*coord];
1534                let predecessor = (entry.seq > 1)
1535                    .then(|| {
1536                        by_coord.keys().find(|candidate| {
1537                            candidate.author_pubkey == entry.author_pubkey
1538                                && candidate.author_owner_grant == entry.author_owner_grant
1539                                && candidate.stream_id == entry.stream_id
1540                                && candidate.seq + 1 == entry.seq
1541                                && Some(candidate.entry_hash) == entry.previous_hash
1542                        })
1543                    })
1544                    .flatten();
1545                (entry.seq == 1 || predecessor.is_some_and(|value| states.contains_key(value)))
1546                    && entry
1547                        .dependencies
1548                        .iter()
1549                        .filter(|dependency| included.contains(*dependency))
1550                        .all(|dependency| states.contains_key(dependency))
1551            });
1552            let Some(coord) = ready.cloned() else {
1553                if pending.iter().any(|coord| {
1554                    let entry = by_coord[coord];
1555                    entry.seq > 1
1556                        && !by_coord.keys().any(|candidate| {
1557                            candidate.author_pubkey == entry.author_pubkey
1558                                && candidate.author_owner_grant == entry.author_owner_grant
1559                                && candidate.stream_id == entry.stream_id
1560                                && candidate.seq + 1 == entry.seq
1561                                && Some(candidate.entry_hash) == entry.previous_hash
1562                        })
1563                }) {
1564                    return Err(ProviderAdminReducerError::MissingPredecessor);
1565                }
1566                return Err(ProviderAdminReducerError::CausalCycle);
1567            };
1568            let entry = by_coord[&coord];
1569            let mut causal_states = entry
1570                .dependencies
1571                .iter()
1572                .filter_map(|dependency| states.get(dependency).cloned())
1573                .collect::<Vec<_>>();
1574            if entry.seq > 1 {
1575                if let Some(predecessor) = by_coord.keys().find(|candidate| {
1576                    candidate.author_pubkey == entry.author_pubkey
1577                        && candidate.author_owner_grant == entry.author_owner_grant
1578                        && candidate.stream_id == entry.stream_id
1579                        && candidate.seq + 1 == entry.seq
1580                        && Some(candidate.entry_hash) == entry.previous_hash
1581                }) {
1582                    if !entry.dependencies.contains(predecessor) {
1583                        causal_states.push(states[predecessor].clone());
1584                    }
1585                }
1586            }
1587            let mut state = if causal_states.is_empty() {
1588                genesis.clone()
1589            } else {
1590                Self::merge(causal_states)?
1591            };
1592            if let Some(change) = entry.provider_admin.clone() {
1593                state.apply_membership_change(
1594                    change,
1595                    ProviderAdminGrantOrigin::MergeMembership {
1596                        coord: coord.clone(),
1597                    },
1598                )?;
1599            }
1600            states.insert(coord.clone(), state);
1601            pending.remove(&coord);
1602        }
1603        let raw_heads = by_coord
1604            .keys()
1605            .filter(|coord| {
1606                !by_coord.values().any(|entry| {
1607                    entry.dependencies.contains(*coord)
1608                        || (entry.seq == coord.seq + 1
1609                            && entry.author_pubkey == coord.author_pubkey
1610                            && entry.author_owner_grant == coord.author_owner_grant
1611                            && entry.stream_id == coord.stream_id
1612                            && entry.previous_hash == Some(coord.entry_hash))
1613                })
1614            })
1615            .cloned()
1616            .collect::<Vec<_>>();
1617        let combined =
1618            Self::merge(std::iter::once(genesis.clone()).chain(states.values().cloned()))?;
1619        if !combined.active().is_empty() {
1620            return Ok(ProviderAdminResolution::Resolved(combined));
1621        }
1622        if raw_heads.len() > 12 {
1623            return Err(ProviderAdminReducerError::ConflictTooWide(raw_heads.len()));
1624        }
1625        let head_states = raw_heads
1626            .iter()
1627            .map(|head| (head.clone(), states[head].clone()))
1628            .collect::<Vec<_>>();
1629        let mut valid = Vec::<ProviderAdminBranch>::new();
1630        for mask in 1usize..(1usize << head_states.len()) {
1631            let heads = head_states
1632                .iter()
1633                .enumerate()
1634                .filter(|(index, _)| mask & (1usize << index) != 0)
1635                .map(|(_, (head, _))| head.clone())
1636                .collect::<Vec<_>>();
1637            let state = Self::merge(
1638                head_states
1639                    .iter()
1640                    .enumerate()
1641                    .filter(|(index, _)| mask & (1usize << index) != 0)
1642                    .map(|(_, (_, state))| state.clone()),
1643            )?;
1644            if !state.active().is_empty() {
1645                valid.push(ProviderAdminBranch { heads, state });
1646            }
1647        }
1648        let valid_head_sets = valid
1649            .iter()
1650            .map(|branch| branch.heads.iter().cloned().collect::<BTreeSet<_>>())
1651            .collect::<Vec<_>>();
1652        let maximal_valid_branches = valid
1653            .into_iter()
1654            .enumerate()
1655            .filter(|(index, _)| {
1656                !valid_head_sets.iter().enumerate().any(|(other, heads)| {
1657                    other != *index && valid_head_sets[*index].is_subset(heads)
1658                })
1659            })
1660            .map(|(_, branch)| branch)
1661            .collect();
1662        let mut cyclic_sources = Vec::new();
1663        let mut involved_grants = BTreeSet::new();
1664        for (coord, entry) in &by_coord {
1665            if let Some(ProviderAdminMembershipChange::MergeConcurrent {
1666                change: ProviderAdminChange::Remove { removes },
1667                ..
1668            }) = &entry.provider_admin
1669            {
1670                cyclic_sources.push(coord.clone());
1671                involved_grants.extend(removes.iter().cloned());
1672            }
1673        }
1674        cyclic_sources.sort();
1675        Ok(ProviderAdminResolution::RevocationConflict(
1676            ProviderAdminConflict {
1677                raw_heads,
1678                cyclic_sources,
1679                involved_grants,
1680                maximal_valid_branches,
1681                combined,
1682            },
1683        ))
1684    }
1685}
1686
1687#[derive(Debug, thiserror::Error, PartialEq, Eq)]
1688pub enum ProviderAdminReducerError {
1689    #[error("provider administrator grant id was reused with different facts")]
1690    GrantIdReuse,
1691    #[error("provider administrator replacement names an inactive grant")]
1692    UnknownReplacement,
1693    #[error("provider administrator removal names an inactive grant")]
1694    UnknownRemoval,
1695    #[error("provider administrator change leaves no effective administrator")]
1696    NoEffectiveAdministrator,
1697    #[error("provider administrator change policy does not match its derived origin")]
1698    PolicyOriginMismatch,
1699    #[error("provider administrator causal history is missing an exact stream predecessor")]
1700    MissingPredecessor,
1701    #[error("provider administrator causal history contains a cycle")]
1702    CausalCycle,
1703    #[error("provider administrator revocation conflict has {0} heads, exceeding 12")]
1704    ConflictTooWide(usize),
1705}
1706
1707#[derive(Debug, thiserror::Error)]
1708pub enum ProviderProbeError {
1709    #[error(transparent)]
1710    Storage(#[from] StorageError),
1711    #[error("provider capability receipt is invalid: {0}")]
1712    InvalidReceipt(String),
1713}
1714
1715#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1716#[serde(rename_all = "snake_case", deny_unknown_fields)]
1717pub enum ProviderProbeJournalRecord {
1718    Exact(ExactProbeJournal),
1719    Serial(SerialProbeJournal),
1720    CrossPrincipal(CrossPrincipalCompletionJournal),
1721}
1722
1723impl ProviderProbeJournalRecord {
1724    pub fn probe_id(&self) -> ProviderProbeId {
1725        match self {
1726            Self::Exact(record) => record.probe_id,
1727            Self::Serial(record) => record.probe_id,
1728            Self::CrossPrincipal(record) => record.probe_id,
1729        }
1730    }
1731
1732    pub fn validate_begin(&self) -> Result<(), ProviderProbeJournalError> {
1733        let prepared = match self {
1734            Self::Exact(record) => matches!(record.progress, ExactProbeProgress::Prepared),
1735            Self::Serial(record) => matches!(record.progress, SerialProbeProgress::Prepared),
1736            Self::CrossPrincipal(record) => {
1737                matches!(record.progress, CrossPrincipalCompletionProgress::Prepared)
1738            }
1739        };
1740        if !prepared {
1741            return Err(ProviderProbeJournalError::BeginNotPrepared);
1742        }
1743        Ok(())
1744    }
1745
1746    pub fn validate_transition(&self, next: &Self) -> Result<(), ProviderProbeJournalError> {
1747        match (self, next) {
1748            (Self::Exact(previous), Self::Exact(next)) => {
1749                if previous.probe_id != next.probe_id
1750                    || previous.binding != next.binding
1751                    || previous.slot != next.slot
1752                    || previous.lost_response_slot != next.lost_response_slot
1753                {
1754                    return Err(ProviderProbeJournalError::ImmutableFactsChanged);
1755                }
1756                validate_exact_progress_transition(&previous.progress, &next.progress)
1757            }
1758            (Self::Serial(previous), Self::Serial(next)) => {
1759                if previous.probe_id != next.probe_id
1760                    || previous.binding != next.binding
1761                    || previous.logical_key != next.logical_key
1762                {
1763                    return Err(ProviderProbeJournalError::ImmutableFactsChanged);
1764                }
1765                validate_serial_progress_transition(&previous.progress, &next.progress)
1766            }
1767            (Self::CrossPrincipal(previous), Self::CrossPrincipal(next)) => {
1768                if previous.probe_id != next.probe_id
1769                    || previous.store != next.store
1770                    || previous.context != next.context
1771                    || previous.challenge != next.challenge
1772                    || previous.response != next.response
1773                {
1774                    return Err(ProviderProbeJournalError::ImmutableFactsChanged);
1775                }
1776                let expected_read_hash = ObjectHash::digest(&probe_payload(
1777                    &previous.probe_id,
1778                    ProbePayloadLabel::CrossPeer,
1779                ));
1780                if cross_progress_evidence_hash(&next.progress)
1781                    .is_some_and(|hash| hash != expected_read_hash)
1782                {
1783                    return Err(ProviderProbeJournalError::EvidenceChanged);
1784                }
1785                validate_cross_progress_transition(&previous.progress, &next.progress)
1786            }
1787            _ => Err(ProviderProbeJournalError::ProbeKindChanged),
1788        }
1789    }
1790}
1791
1792#[derive(Debug, thiserror::Error, PartialEq, Eq)]
1793pub enum ProviderProbeJournalError {
1794    #[error("provider probe journal must begin at prepared")]
1795    BeginNotPrepared,
1796    #[error("provider probe journal advance changes immutable facts")]
1797    ImmutableFactsChanged,
1798    #[error("provider probe journal advance changes the probe kind")]
1799    ProbeKindChanged,
1800    #[error("provider probe journal advance skips or reverses progress")]
1801    NonAdjacentProgress,
1802    #[error("provider probe journal advance changes established evidence")]
1803    EvidenceChanged,
1804}
1805
1806#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1807#[serde(deny_unknown_fields)]
1808pub struct SerialProbeJournal {
1809    pub probe_id: ProviderProbeId,
1810    pub binding: crate::sync::storage::ResolvedProviderBinding,
1811    pub logical_key: String,
1812    pub progress: SerialProbeProgress,
1813}
1814
1815#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1816#[serde(rename_all = "snake_case", deny_unknown_fields)]
1817pub enum SerialProbeProgress {
1818    Prepared,
1819    Created {
1820        attempts: [ProbeCreateAttempt; 2],
1821        object: VersionedObject,
1822    },
1823    Replaced {
1824        create_attempts: [ProbeCreateAttempt; 2],
1825        created: VersionedObject,
1826        replace_attempts: [ProbeReplaceAttempt; 2],
1827        replaced: VersionedObject,
1828    },
1829    ReadVerified {
1830        create_attempts: [ProbeCreateAttempt; 2],
1831        created: VersionedObject,
1832        replace_attempts: [ProbeReplaceAttempt; 2],
1833        replaced: VersionedObject,
1834    },
1835    Absent {
1836        create_attempts: [ProbeCreateAttempt; 2],
1837        created: VersionedObject,
1838        replace_attempts: [ProbeReplaceAttempt; 2],
1839        replaced: VersionedObject,
1840    },
1841    ReceiptReady {
1842        receipt: SerialCoordinationProbeReceipt,
1843    },
1844}
1845
1846#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1847#[serde(deny_unknown_fields)]
1848pub struct ExactProbeJournal {
1849    pub probe_id: ProviderProbeId,
1850    pub binding: crate::sync::storage::ResolvedProviderBinding,
1851    pub slot: ObjectSlot,
1852    pub lost_response_slot: ObjectSlot,
1853    pub progress: ExactProbeProgress,
1854}
1855
1856#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1857#[serde(rename_all = "snake_case", deny_unknown_fields)]
1858pub enum ExactProbeProgress {
1859    Prepared,
1860    Created { outcomes: [ProbeCreateOutcome; 2] },
1861    ReadsVerified { outcomes: [ProbeCreateOutcome; 2] },
1862    PrimaryAbsent { outcomes: [ProbeCreateOutcome; 2] },
1863    LostResponseCreated { outcomes: [ProbeCreateOutcome; 2] },
1864    LostResponseReadVerified { outcomes: [ProbeCreateOutcome; 2] },
1865    Absent { outcomes: [ProbeCreateOutcome; 2] },
1866    ReceiptReady { receipt: ExactSlotProbeReceipt },
1867}
1868
1869#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1870#[serde(deny_unknown_fields)]
1871pub struct CrossPrincipalCompletionJournal {
1872    pub probe_id: ProviderProbeId,
1873    pub store: StoreProviderBinding,
1874    pub context: CrossPrincipalResponseContext,
1875    pub challenge: CrossPrincipalProbeChallenge,
1876    pub response: CrossPrincipalProbeResponse,
1877    pub progress: CrossPrincipalCompletionProgress,
1878}
1879
1880#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
1881#[serde(rename_all = "snake_case", deny_unknown_fields)]
1882pub enum CrossPrincipalCompletionProgress {
1883    Prepared,
1884    ReadsVerified {
1885        administrator_read_peer_hash: ObjectHash,
1886    },
1887    PeerAbsent {
1888        administrator_read_peer_hash: ObjectHash,
1889    },
1890    Absent {
1891        administrator_read_peer_hash: ObjectHash,
1892    },
1893    ReceiptReady {
1894        receipt: CrossPrincipalProbeReceipt,
1895    },
1896}
1897
1898fn validate_exact_progress_transition(
1899    previous: &ExactProbeProgress,
1900    next: &ExactProbeProgress,
1901) -> Result<(), ProviderProbeJournalError> {
1902    let evidence_matches = match (previous, next) {
1903        (ExactProbeProgress::Prepared, ExactProbeProgress::Created { .. }) => true,
1904        (
1905            ExactProbeProgress::Created { outcomes: previous },
1906            ExactProbeProgress::ReadsVerified { outcomes: next },
1907        )
1908        | (
1909            ExactProbeProgress::ReadsVerified { outcomes: previous },
1910            ExactProbeProgress::PrimaryAbsent { outcomes: next },
1911        )
1912        | (
1913            ExactProbeProgress::PrimaryAbsent { outcomes: previous },
1914            ExactProbeProgress::LostResponseCreated { outcomes: next },
1915        )
1916        | (
1917            ExactProbeProgress::LostResponseCreated { outcomes: previous },
1918            ExactProbeProgress::LostResponseReadVerified { outcomes: next },
1919        )
1920        | (
1921            ExactProbeProgress::LostResponseReadVerified { outcomes: previous },
1922            ExactProbeProgress::Absent { outcomes: next },
1923        ) => previous == next,
1924        (ExactProbeProgress::Absent { outcomes }, ExactProbeProgress::ReceiptReady { receipt }) => {
1925            receipt
1926                .transcript
1927                .contenders
1928                .iter()
1929                .map(|attempt| attempt.outcome)
1930                .eq(outcomes.iter().copied())
1931        }
1932        _ => return Err(ProviderProbeJournalError::NonAdjacentProgress),
1933    };
1934    if !evidence_matches {
1935        return Err(ProviderProbeJournalError::EvidenceChanged);
1936    }
1937    Ok(())
1938}
1939
1940fn validate_serial_progress_transition(
1941    previous: &SerialProbeProgress,
1942    next: &SerialProbeProgress,
1943) -> Result<(), ProviderProbeJournalError> {
1944    let evidence_matches = match (previous, next) {
1945        (SerialProbeProgress::Prepared, SerialProbeProgress::Created { .. }) => true,
1946        (
1947            SerialProbeProgress::Created { attempts, object },
1948            SerialProbeProgress::Replaced {
1949                create_attempts,
1950                created,
1951                ..
1952            },
1953        ) => attempts == create_attempts && object == created,
1954        (
1955            SerialProbeProgress::Replaced {
1956                create_attempts: previous_create,
1957                created: previous_created,
1958                replace_attempts: previous_replace,
1959                replaced: previous_replaced,
1960            },
1961            SerialProbeProgress::ReadVerified {
1962                create_attempts: next_create,
1963                created: next_created,
1964                replace_attempts: next_replace,
1965                replaced: next_replaced,
1966            },
1967        )
1968        | (
1969            SerialProbeProgress::ReadVerified {
1970                create_attempts: previous_create,
1971                created: previous_created,
1972                replace_attempts: previous_replace,
1973                replaced: previous_replaced,
1974            },
1975            SerialProbeProgress::Absent {
1976                create_attempts: next_create,
1977                created: next_created,
1978                replace_attempts: next_replace,
1979                replaced: next_replaced,
1980            },
1981        ) => {
1982            previous_create == next_create
1983                && previous_created == next_created
1984                && previous_replace == next_replace
1985                && previous_replaced == next_replaced
1986        }
1987        (
1988            SerialProbeProgress::Absent {
1989                create_attempts,
1990                created,
1991                replace_attempts,
1992                replaced,
1993            },
1994            SerialProbeProgress::ReceiptReady { receipt },
1995        ) => {
1996            receipt.transcript.create_attempts == *create_attempts
1997                && receipt.transcript.created_bytes_hash == ObjectHash::digest(&created.bytes)
1998                && receipt.transcript.replace_attempts == *replace_attempts
1999                && receipt.transcript.replaced_bytes_hash == ObjectHash::digest(&replaced.bytes)
2000        }
2001        _ => return Err(ProviderProbeJournalError::NonAdjacentProgress),
2002    };
2003    if !evidence_matches {
2004        return Err(ProviderProbeJournalError::EvidenceChanged);
2005    }
2006    Ok(())
2007}
2008
2009fn validate_cross_progress_transition(
2010    previous: &CrossPrincipalCompletionProgress,
2011    next: &CrossPrincipalCompletionProgress,
2012) -> Result<(), ProviderProbeJournalError> {
2013    let evidence_matches = match (previous, next) {
2014        (
2015            CrossPrincipalCompletionProgress::Prepared,
2016            CrossPrincipalCompletionProgress::ReadsVerified { .. },
2017        ) => true,
2018        (
2019            CrossPrincipalCompletionProgress::ReadsVerified {
2020                administrator_read_peer_hash: previous,
2021            },
2022            CrossPrincipalCompletionProgress::PeerAbsent {
2023                administrator_read_peer_hash: next,
2024            },
2025        )
2026        | (
2027            CrossPrincipalCompletionProgress::PeerAbsent {
2028                administrator_read_peer_hash: previous,
2029            },
2030            CrossPrincipalCompletionProgress::Absent {
2031                administrator_read_peer_hash: next,
2032            },
2033        ) => previous == next,
2034        (
2035            CrossPrincipalCompletionProgress::Absent {
2036                administrator_read_peer_hash,
2037            },
2038            CrossPrincipalCompletionProgress::ReceiptReady { receipt },
2039        ) => receipt.transcript.administrator_read_peer_hash == *administrator_read_peer_hash,
2040        _ => return Err(ProviderProbeJournalError::NonAdjacentProgress),
2041    };
2042    if !evidence_matches {
2043        return Err(ProviderProbeJournalError::EvidenceChanged);
2044    }
2045    Ok(())
2046}
2047
2048fn cross_progress_evidence_hash(progress: &CrossPrincipalCompletionProgress) -> Option<ObjectHash> {
2049    match progress {
2050        CrossPrincipalCompletionProgress::Prepared => None,
2051        CrossPrincipalCompletionProgress::ReadsVerified {
2052            administrator_read_peer_hash,
2053        }
2054        | CrossPrincipalCompletionProgress::PeerAbsent {
2055            administrator_read_peer_hash,
2056        }
2057        | CrossPrincipalCompletionProgress::Absent {
2058            administrator_read_peer_hash,
2059        } => Some(*administrator_read_peer_hash),
2060        CrossPrincipalCompletionProgress::ReceiptReady { receipt } => {
2061            Some(receipt.transcript.administrator_read_peer_hash)
2062        }
2063    }
2064}
2065
2066#[async_trait]
2067pub trait ProviderProbeJournal: Send + Sync {
2068    async fn load(
2069        &self,
2070        probe_id: ProviderProbeId,
2071    ) -> Result<Option<ProviderProbeJournalRecord>, StorageError>;
2072
2073    /// Atomically inserts `prepared` when absent or returns the exact existing
2074    /// record for this probe id. A different record under the id is corruption.
2075    async fn begin(
2076        &self,
2077        prepared: ProviderProbeJournalRecord,
2078    ) -> Result<ProviderProbeJournalRecord, StorageError>;
2079
2080    /// Atomically replaces the exact current record. Implementations reject a
2081    /// stale predecessor instead of merging progress.
2082    async fn advance(
2083        &self,
2084        previous: &ProviderProbeJournalRecord,
2085        next: ProviderProbeJournalRecord,
2086    ) -> Result<(), StorageError>;
2087}
2088
2089#[async_trait]
2090impl ProviderProbeJournal for crate::database::Database {
2091    async fn load(
2092        &self,
2093        probe_id: ProviderProbeId,
2094    ) -> Result<Option<ProviderProbeJournalRecord>, StorageError> {
2095        self.load_provider_probe_journal(probe_id)
2096            .await
2097            .map_err(|error| StorageError::Storage(error.to_string()))
2098    }
2099
2100    async fn begin(
2101        &self,
2102        prepared: ProviderProbeJournalRecord,
2103    ) -> Result<ProviderProbeJournalRecord, StorageError> {
2104        self.begin_provider_probe_journal(prepared)
2105            .await
2106            .map_err(|error| StorageError::Storage(error.to_string()))
2107    }
2108
2109    async fn advance(
2110        &self,
2111        previous: &ProviderProbeJournalRecord,
2112        next: ProviderProbeJournalRecord,
2113    ) -> Result<(), StorageError> {
2114        self.advance_provider_probe_journal(previous.clone(), next)
2115            .await
2116            .map_err(|error| StorageError::Storage(error.to_string()))
2117    }
2118}
2119
2120pub async fn prepare_cross_principal_challenge(
2121    administrator: &dyn ExactSlotStorage,
2122    publication_journal: &dyn DeviceJoinChallengePublicationJournal,
2123    probe_id: ProviderProbeId,
2124    store: &StoreProviderBinding,
2125    context: &CrossPrincipalChallengeContext,
2126    administrator_signer: &UserKeypair,
2127) -> Result<CrossPrincipalProbeChallenge, ProviderProbeError> {
2128    let administrator_live = administrator
2129        .provider_binding()
2130        .await
2131        .map_err(StorageError::from)?;
2132    if administrator_live.store != *store
2133        || administrator_live.device != context.administrator_binding
2134    {
2135        return invalid("cross-principal administrator does not match the challenge context");
2136    }
2137    validate_cross_provider_evidence_context(store, context)?;
2138    let suffix = hex::encode(probe_id.as_bytes());
2139    let administrator_key = format!("__coven_probe__/cross/{suffix}/administrator");
2140    let administrator_slot = administrator
2141        .allocate_slot(&administrator_key)
2142        .await
2143        .map_err(StorageError::from)?;
2144    if administrator_slot.logical_key() != administrator_key {
2145        return invalid("cross-principal administrator slot changed its logical key");
2146    }
2147    let administrator_payload = probe_payload(&probe_id, ProbePayloadLabel::CrossAdministrator);
2148    let administrator_object = ProbeExactObjectReceipt {
2149        slot: administrator_slot.clone(),
2150        payload_hash: ObjectHash::digest(&administrator_payload),
2151        object: ExactObjectRef::new(
2152            administrator_slot,
2153            administrator_payload.len() as u64,
2154            ObjectHash::digest(&administrator_payload),
2155        ),
2156    };
2157    let unsigned = CrossPrincipalProbeChallenge {
2158        probe_id,
2159        administrator_object,
2160        challenge_hash: ObjectHash::digest(&[]),
2161        administrator_signature: String::new(),
2162    };
2163    let challenge_hash = cross_challenge_hash(store, context, &unsigned);
2164    let challenge = CrossPrincipalProbeChallenge {
2165        challenge_hash,
2166        administrator_signature: hex::encode(administrator_signer.sign(challenge_hash.as_bytes())),
2167        ..unsigned
2168    };
2169    challenge.verify(
2170        context,
2171        store,
2172        &crate::keys::public_key_hex(administrator_signer),
2173    )?;
2174    let durable = publication_journal.prepare(&challenge).await?;
2175    if durable.challenge != challenge {
2176        return invalid("durable cross-principal challenge differs from its prepared bytes");
2177    }
2178    Ok(challenge)
2179}
2180
2181pub async fn publish_cross_principal_challenge(
2182    protocol_storage: &dyn SyncStorage,
2183    administrator: &dyn ExactSlotStorage,
2184    publication_journal: &dyn DeviceJoinChallengePublicationJournal,
2185    authorization: &DeviceJoinChallengePublicationAuthorization,
2186    challenge: &CrossPrincipalProbeChallenge,
2187    context: &CrossPrincipalChallengeContext,
2188    store: &StoreProviderBinding,
2189    attempt_owner: &StoreDeviceRegistration,
2190    activation_author: &StoreDeviceRegistration,
2191    administrator_signing_pubkey: &str,
2192) -> Result<CrossPrincipalProbeChallenge, ProviderProbeError> {
2193    challenge.verify(context, store, administrator_signing_pubkey)?;
2194    if authorization.attempt.attempt_id != context.attempt_id {
2195        return invalid("challenge publication authorization names another join attempt");
2196    }
2197    let attempt = crate::sync::store_objects::load_device_join_attempt_ref(
2198        protocol_storage,
2199        &context.root,
2200        &authorization.attempt,
2201        attempt_owner,
2202    )
2203    .await
2204    .map_err(|error| ProviderProbeError::Storage(StorageError::Storage(error.to_string())))?;
2205    if attempt.value.store_root != context.root
2206        || attempt.value.attempt_id != context.attempt_id
2207        || attempt.value.owner_registration != context.owner_registration
2208    {
2209        return invalid("activated join attempt differs from the challenge publication context");
2210    }
2211    let activation = crate::sync::store_objects::load_commit_ref(
2212        protocol_storage,
2213        context.root.store_root_hash,
2214        &authorization.attempt_activation,
2215        activation_author,
2216    )
2217    .await
2218    .map_err(|error| ProviderProbeError::Storage(StorageError::Storage(error.to_string())))?;
2219    if activation
2220        .value
2221        .device_join_attempts()
2222        .binary_search(&authorization.attempt)
2223        .is_err()
2224    {
2225        return invalid("activation commit does not activate the authorized join attempt");
2226    }
2227    let live = administrator
2228        .provider_binding()
2229        .await
2230        .map_err(StorageError::from)?;
2231    if live.store != *store || live.device != context.administrator_binding {
2232        return invalid("cross-principal administrator does not match the published challenge");
2233    }
2234    publication_journal
2235        .claim_published(authorization, challenge)
2236        .await?;
2237    let payload = probe_payload(&challenge.probe_id, ProbePayloadLabel::CrossAdministrator);
2238    settle_exact_create(
2239        administrator,
2240        &challenge.administrator_object.slot,
2241        &payload,
2242    )
2243    .await?;
2244    let observed = administrator
2245        .read_at(&challenge.administrator_object.slot)
2246        .await
2247        .map_err(StorageError::from)?;
2248    if observed != payload {
2249        return invalid("published cross-principal challenge differs from its signed bytes");
2250    }
2251    Ok(challenge.clone())
2252}
2253
2254pub async fn create_cross_principal_response(
2255    peer: &dyn ExactSlotStorage,
2256    challenge: &CrossPrincipalProbeChallenge,
2257    context: &CrossPrincipalResponseContext,
2258    store: &StoreProviderBinding,
2259    administrator_signing_pubkey: &str,
2260    peer_signer: &UserKeypair,
2261) -> Result<CrossPrincipalProbeResponse, ProviderProbeError> {
2262    challenge.verify(&context.challenge, store, administrator_signing_pubkey)?;
2263    let peer_pubkey = crate::keys::public_key_hex(peer_signer);
2264    if context.challenge.member_pubkey != peer_pubkey {
2265        return invalid("cross-principal peer signer is not the joining member");
2266    }
2267    let live = peer.provider_binding().await.map_err(StorageError::from)?;
2268    if live.store != *store || live.device != context.challenge.peer_binding {
2269        return invalid("cross-principal peer does not match the response context");
2270    }
2271    let evidence = peer
2272        .cross_principal_evidence()
2273        .await
2274        .map_err(StorageError::from)?;
2275    validate_cross_provider_evidence(
2276        store,
2277        &context.challenge.administrator_binding,
2278        &context.challenge.peer_binding,
2279        &evidence,
2280    )?;
2281    let expected_peer_key = cross_peer_logical_key(challenge.probe_id);
2282    if context.response_slot.logical_key() != expected_peer_key {
2283        return invalid("cross-principal response slot uses the wrong logical key");
2284    }
2285    let administrator_payload =
2286        probe_payload(&challenge.probe_id, ProbePayloadLabel::CrossAdministrator);
2287    let administrator_read = peer
2288        .read_at(&challenge.administrator_object.slot)
2289        .await
2290        .map_err(StorageError::from)?;
2291    if administrator_read != administrator_payload {
2292        return invalid("peer read bytes differ from the signed cross-principal challenge");
2293    }
2294    let peer_payload = probe_payload(&challenge.probe_id, ProbePayloadLabel::CrossPeer);
2295    settle_exact_create(peer, &context.response_slot, &peer_payload).await?;
2296    let peer_read = peer
2297        .read_at(&context.response_slot)
2298        .await
2299        .map_err(StorageError::from)?;
2300    if peer_read != peer_payload {
2301        return invalid("peer response readback differs from its deterministic bytes");
2302    }
2303    let peer_object = ProbeExactObjectReceipt {
2304        slot: context.response_slot.clone(),
2305        payload_hash: ObjectHash::digest(&peer_payload),
2306        object: ExactObjectRef::new(
2307            context.response_slot.clone(),
2308            peer_payload.len() as u64,
2309            ObjectHash::digest(&peer_payload),
2310        ),
2311    };
2312    let unsigned = CrossPrincipalProbeResponse {
2313        challenge_hash: challenge.challenge_hash,
2314        provider_evidence: evidence,
2315        peer_object,
2316        peer_read_administrator_hash: ObjectHash::digest(&administrator_read),
2317        response_hash: ObjectHash::digest(&[]),
2318        peer_signature: String::new(),
2319    };
2320    let response_hash = cross_response_hash(store, context, challenge, &unsigned);
2321    let response = CrossPrincipalProbeResponse {
2322        response_hash,
2323        peer_signature: hex::encode(peer_signer.sign(response_hash.as_bytes())),
2324        ..unsigned
2325    };
2326    response.verify(
2327        challenge,
2328        context,
2329        store,
2330        administrator_signing_pubkey,
2331        &peer_pubkey,
2332    )?;
2333    Ok(response)
2334}
2335
2336pub async fn complete_cross_principal_probe(
2337    administrator: &dyn ExactSlotStorage,
2338    journal: &dyn ProviderProbeJournal,
2339    challenge: &CrossPrincipalProbeChallenge,
2340    response: &CrossPrincipalProbeResponse,
2341    context: &CrossPrincipalResponseContext,
2342    store: &StoreProviderBinding,
2343    administrator_signer: &UserKeypair,
2344    peer_signing_pubkey: &str,
2345) -> Result<CrossPrincipalProbeReceipt, ProviderProbeError> {
2346    let administrator_pubkey = crate::keys::public_key_hex(administrator_signer);
2347    challenge.verify(&context.challenge, store, &administrator_pubkey)?;
2348    response.verify(
2349        challenge,
2350        context,
2351        store,
2352        &administrator_pubkey,
2353        peer_signing_pubkey,
2354    )?;
2355    let live = administrator
2356        .provider_binding()
2357        .await
2358        .map_err(StorageError::from)?;
2359    if live.store != *store || live.device != context.challenge.administrator_binding {
2360        return invalid("cross-principal administrator does not match the completion context");
2361    }
2362    let prepared = ProviderProbeJournalRecord::CrossPrincipal(CrossPrincipalCompletionJournal {
2363        probe_id: challenge.probe_id,
2364        store: store.clone(),
2365        context: context.clone(),
2366        challenge: challenge.clone(),
2367        response: response.clone(),
2368        progress: CrossPrincipalCompletionProgress::Prepared,
2369    });
2370    let mut durable = match journal.load(challenge.probe_id).await? {
2371        Some(existing) => existing,
2372        None => journal.begin(prepared).await?,
2373    };
2374    let ProviderProbeJournalRecord::CrossPrincipal(mut record) = durable.clone() else {
2375        return invalid("cross-principal probe id belongs to another durable probe kind");
2376    };
2377    if record.probe_id != challenge.probe_id
2378        || record.store != *store
2379        || record.context != *context
2380        || record.challenge != *challenge
2381        || record.response != *response
2382    {
2383        return invalid("durable cross-principal completion differs from the requested proof");
2384    }
2385    if let CrossPrincipalCompletionProgress::ReceiptReady { receipt } = &record.progress {
2386        receipt.verify(context, store, &administrator_pubkey, peer_signing_pubkey)?;
2387        return Ok(receipt.clone());
2388    }
2389    let peer_payload = probe_payload(&challenge.probe_id, ProbePayloadLabel::CrossPeer);
2390    if matches!(record.progress, CrossPrincipalCompletionProgress::Prepared) {
2391        let observed = administrator
2392            .read_at(&response.peer_object.slot)
2393            .await
2394            .map_err(StorageError::from)?;
2395        if observed != peer_payload {
2396            return invalid("administrator read differs from the signed peer response");
2397        }
2398        advance_cross_completion(
2399            journal,
2400            &mut durable,
2401            &mut record,
2402            CrossPrincipalCompletionProgress::ReadsVerified {
2403                administrator_read_peer_hash: ObjectHash::digest(&observed),
2404            },
2405        )
2406        .await?;
2407    }
2408    let administrator_read_peer_hash = cross_completion_read_hash(&record.progress)?;
2409    if matches!(
2410        record.progress,
2411        CrossPrincipalCompletionProgress::ReadsVerified { .. }
2412    ) {
2413        cleanup_exact_slot(administrator, &response.peer_object.slot).await?;
2414        advance_cross_completion(
2415            journal,
2416            &mut durable,
2417            &mut record,
2418            CrossPrincipalCompletionProgress::PeerAbsent {
2419                administrator_read_peer_hash,
2420            },
2421        )
2422        .await?;
2423    }
2424    if matches!(
2425        record.progress,
2426        CrossPrincipalCompletionProgress::PeerAbsent { .. }
2427    ) {
2428        cleanup_exact_slot(administrator, &challenge.administrator_object.slot).await?;
2429        advance_cross_completion(
2430            journal,
2431            &mut durable,
2432            &mut record,
2433            CrossPrincipalCompletionProgress::Absent {
2434                administrator_read_peer_hash,
2435            },
2436        )
2437        .await?;
2438    }
2439    let transcript = CrossPrincipalProbeTranscript {
2440        challenge: challenge.clone(),
2441        response: response.clone(),
2442        administrator_read_peer_hash,
2443        administrator_delete_peer_verified_absent: true,
2444        administrator_delete_own_verified_absent: true,
2445    };
2446    let receipt =
2447        CrossPrincipalProbeReceipt::signed(transcript, context, store, administrator_signer)?;
2448    advance_cross_completion(
2449        journal,
2450        &mut durable,
2451        &mut record,
2452        CrossPrincipalCompletionProgress::ReceiptReady {
2453            receipt: receipt.clone(),
2454        },
2455    )
2456    .await?;
2457    Ok(receipt)
2458}
2459
2460pub async fn cleanup_published_cross_principal_challenge(
2461    administrator: &dyn ExactSlotStorage,
2462    challenge: &CrossPrincipalProbeChallenge,
2463    context: &CrossPrincipalChallengeContext,
2464    store: &StoreProviderBinding,
2465    administrator_signing_pubkey: &str,
2466) -> Result<(), ProviderProbeError> {
2467    challenge.verify(context, store, administrator_signing_pubkey)?;
2468    let live = administrator
2469        .provider_binding()
2470        .await
2471        .map_err(StorageError::from)?;
2472    if live.store != *store || live.device != context.administrator_binding {
2473        return invalid("cross-principal administrator does not match challenge cleanup");
2474    }
2475    cleanup_exact_slot(administrator, &challenge.administrator_object.slot).await
2476}
2477
2478async fn advance_cross_completion(
2479    journal: &dyn ProviderProbeJournal,
2480    durable: &mut ProviderProbeJournalRecord,
2481    record: &mut CrossPrincipalCompletionJournal,
2482    progress: CrossPrincipalCompletionProgress,
2483) -> Result<(), ProviderProbeError> {
2484    record.progress = progress;
2485    let next = ProviderProbeJournalRecord::CrossPrincipal(record.clone());
2486    journal.advance(durable, next.clone()).await?;
2487    *durable = next;
2488    Ok(())
2489}
2490
2491fn cross_completion_read_hash(
2492    progress: &CrossPrincipalCompletionProgress,
2493) -> Result<ObjectHash, ProviderProbeError> {
2494    match progress {
2495        CrossPrincipalCompletionProgress::ReadsVerified {
2496            administrator_read_peer_hash,
2497        }
2498        | CrossPrincipalCompletionProgress::PeerAbsent {
2499            administrator_read_peer_hash,
2500        }
2501        | CrossPrincipalCompletionProgress::Absent {
2502            administrator_read_peer_hash,
2503        } => Ok(*administrator_read_peer_hash),
2504        CrossPrincipalCompletionProgress::Prepared
2505        | CrossPrincipalCompletionProgress::ReceiptReady { .. } => {
2506            invalid("cross-principal completion has no durable administrator read")
2507        }
2508    }
2509}
2510
2511async fn settle_exact_create(
2512    storage: &dyn ExactSlotStorage,
2513    slot: &ObjectSlot,
2514    payload: &[u8],
2515) -> Result<(), ProviderProbeError> {
2516    match storage.read_at(slot).await {
2517        Ok(bytes) if bytes == payload => Ok(()),
2518        Ok(_) => invalid("durable provider probe slot contains different bytes"),
2519        Err(CloudHomeError::NotFound(_)) => storage
2520            .create_at(slot, BlobBody::from_bytes(payload.to_vec()), &|_| {})
2521            .await
2522            .map_err(StorageError::from)
2523            .map_err(ProviderProbeError::Storage),
2524        Err(error) => Err(ProviderProbeError::Storage(StorageError::from(error))),
2525    }
2526}
2527
2528pub async fn probe_exact_slots(
2529    first: &dyn ExactSlotStorage,
2530    second: &dyn ExactSlotStorage,
2531    journal: &dyn ProviderProbeJournal,
2532    probe_id: ProviderProbeId,
2533    binding: &crate::sync::storage::ResolvedProviderBinding,
2534) -> Result<ExactSlotProbeReceipt, ProviderProbeError> {
2535    binding.validate().map_err(ProviderProbeError::Storage)?;
2536    let first_binding = first.provider_binding().await.map_err(StorageError::from)?;
2537    let second_binding = second
2538        .provider_binding()
2539        .await
2540        .map_err(StorageError::from)?;
2541    if first_binding != *binding || second_binding != *binding {
2542        return invalid("exact-slot probe clients do not match the receipt binding");
2543    }
2544    let id = hex::encode(probe_id.as_bytes());
2545    let logical_key = format!("__coven_probe__/exact/{id}");
2546    let lost_logical_key = format!("__coven_probe__/lost-response/{id}");
2547    let mut durable = match journal.load(probe_id).await? {
2548        Some(existing) => existing,
2549        None => {
2550            let allocated_slot = first
2551                .allocate_slot(&logical_key)
2552                .await
2553                .map_err(StorageError::from)?;
2554            let allocated_lost_slot = first
2555                .allocate_slot(&lost_logical_key)
2556                .await
2557                .map_err(StorageError::from)?;
2558            journal
2559                .begin(ProviderProbeJournalRecord::Exact(ExactProbeJournal {
2560                    probe_id,
2561                    binding: binding.clone(),
2562                    slot: allocated_slot,
2563                    lost_response_slot: allocated_lost_slot,
2564                    progress: ExactProbeProgress::Prepared,
2565                }))
2566                .await?
2567        }
2568    };
2569    let ProviderProbeJournalRecord::Exact(mut record) = durable.clone() else {
2570        return invalid("exact probe id belongs to a different durable probe kind");
2571    };
2572    if record.probe_id != probe_id || record.binding != *binding {
2573        return invalid("durable exact probe differs from its requested binding or id");
2574    }
2575    let slot = record.slot.clone();
2576    let lost_slot = record.lost_response_slot.clone();
2577    if slot.logical_key() != logical_key || lost_slot.logical_key() != lost_logical_key {
2578        return invalid("exact-slot allocator changed the probe logical key");
2579    }
2580    let payloads = [
2581        probe_payload(&probe_id, ProbePayloadLabel::ExactCreateFirst),
2582        probe_payload(&probe_id, ProbePayloadLabel::ExactCreateSecond),
2583    ];
2584    if matches!(record.progress, ExactProbeProgress::Prepared) {
2585        let (outcomes, _winner) = match first.read_at(&slot).await {
2586            Err(CloudHomeError::NotFound(_)) => {
2587                let (left, right) = tokio::join!(
2588                    first.create_at(&slot, BlobBody::from_bytes(payloads[0].clone()), &|_| {}),
2589                    second.create_at(&slot, BlobBody::from_bytes(payloads[1].clone()), &|_| {}),
2590                );
2591                classify_exact_create_race(left, right)?
2592            }
2593            Ok(bytes) if bytes == payloads[0] => {
2594                require_occupied_rejection(
2595                    second
2596                        .create_at(&slot, BlobBody::from_bytes(payloads[1].clone()), &|_| {})
2597                        .await,
2598                )?;
2599                (
2600                    [
2601                        ProbeCreateOutcome::Created,
2602                        ProbeCreateOutcome::RejectedOccupied,
2603                    ],
2604                    0,
2605                )
2606            }
2607            Ok(bytes) if bytes == payloads[1] => {
2608                require_occupied_rejection(
2609                    first
2610                        .create_at(&slot, BlobBody::from_bytes(payloads[0].clone()), &|_| {})
2611                        .await,
2612                )?;
2613                (
2614                    [
2615                        ProbeCreateOutcome::RejectedOccupied,
2616                        ProbeCreateOutcome::Created,
2617                    ],
2618                    1,
2619                )
2620            }
2621            Ok(_) => return invalid("durable exact probe slot contains unknown bytes"),
2622            Err(error) => return Err(ProviderProbeError::Storage(StorageError::from(error))),
2623        };
2624        advance_exact(
2625            journal,
2626            &mut durable,
2627            &mut record,
2628            ExactProbeProgress::Created { outcomes },
2629        )
2630        .await?;
2631    }
2632    let (outcomes, winner) = exact_race_state(&record.progress)?;
2633    let (full, range) = if matches!(record.progress, ExactProbeProgress::Created { .. }) {
2634        let full = first.read_at(&slot).await.map_err(StorageError::from)?;
2635        if full != payloads[winner] {
2636            return invalid("authoritative exact read does not match the create winner");
2637        }
2638        let range = first
2639            .read_range_at(&slot, PROBE_RANGE_START, PROBE_RANGE_END)
2640            .await
2641            .map_err(StorageError::from)?;
2642        if range != full[PROBE_RANGE_START as usize..PROBE_RANGE_END as usize] {
2643            return invalid("exact range read does not match the authoritative full read");
2644        }
2645        (full, range)
2646    } else {
2647        (
2648            payloads[winner].clone(),
2649            payloads[winner][PROBE_RANGE_START as usize..PROBE_RANGE_END as usize].to_vec(),
2650        )
2651    };
2652    if matches!(record.progress, ExactProbeProgress::Created { .. }) {
2653        advance_exact(
2654            journal,
2655            &mut durable,
2656            &mut record,
2657            ExactProbeProgress::ReadsVerified { outcomes },
2658        )
2659        .await?;
2660    }
2661    let accepted = ExactObjectRef::new(slot.clone(), full.len() as u64, ObjectHash::digest(&full));
2662    if matches!(record.progress, ExactProbeProgress::ReadsVerified { .. }) {
2663        cleanup_exact_slot(first, &slot).await?;
2664        advance_exact(
2665            journal,
2666            &mut durable,
2667            &mut record,
2668            ExactProbeProgress::PrimaryAbsent { outcomes },
2669        )
2670        .await?;
2671    }
2672    let lost_payload = probe_payload(&probe_id, ProbePayloadLabel::LostResponse);
2673    if matches!(record.progress, ExactProbeProgress::PrimaryAbsent { .. }) {
2674        match first.read_at(&lost_slot).await {
2675            Ok(bytes) if bytes == lost_payload => {}
2676            Ok(_) => return invalid("lost-response slot contains unknown bytes"),
2677            Err(CloudHomeError::NotFound(_)) => first
2678                .create_at(
2679                    &lost_slot,
2680                    BlobBody::from_bytes(lost_payload.clone()),
2681                    &|_| {},
2682                )
2683                .await
2684                .map_err(StorageError::from)?,
2685            Err(error) => return Err(ProviderProbeError::Storage(StorageError::from(error))),
2686        }
2687        advance_exact(
2688            journal,
2689            &mut durable,
2690            &mut record,
2691            ExactProbeProgress::LostResponseCreated { outcomes },
2692        )
2693        .await?;
2694    }
2695    let lost_readback = if matches!(
2696        record.progress,
2697        ExactProbeProgress::LostResponseCreated { .. }
2698    ) {
2699        let readback = first
2700            .read_at(&lost_slot)
2701            .await
2702            .map_err(StorageError::from)?;
2703        if readback != lost_payload {
2704            return invalid("lost-response authoritative readback differs from committed bytes");
2705        }
2706        readback
2707    } else {
2708        lost_payload.clone()
2709    };
2710    let settled = ExactObjectRef::new(
2711        lost_slot.clone(),
2712        lost_readback.len() as u64,
2713        ObjectHash::digest(&lost_readback),
2714    );
2715    if matches!(
2716        record.progress,
2717        ExactProbeProgress::LostResponseCreated { .. }
2718    ) {
2719        advance_exact(
2720            journal,
2721            &mut durable,
2722            &mut record,
2723            ExactProbeProgress::LostResponseReadVerified { outcomes },
2724        )
2725        .await?;
2726    }
2727    if matches!(
2728        record.progress,
2729        ExactProbeProgress::LostResponseReadVerified { .. }
2730    ) {
2731        cleanup_exact_slot(first, &lost_slot).await?;
2732        advance_exact(
2733            journal,
2734            &mut durable,
2735            &mut record,
2736            ExactProbeProgress::Absent { outcomes },
2737        )
2738        .await?;
2739    }
2740
2741    if let ExactProbeProgress::ReceiptReady { receipt } = &record.progress {
2742        receipt.verify(&binding.store, &binding.device)?;
2743        return Ok(receipt.clone());
2744    }
2745    let transcript = ExactSlotProbeTranscript {
2746        probe_id,
2747        logical_key,
2748        slot,
2749        contenders: [
2750            ProbeCreateAttempt {
2751                payload_hash: ObjectHash::digest(&payloads[0]),
2752                outcome: outcomes[0],
2753            },
2754            ProbeCreateAttempt {
2755                payload_hash: ObjectHash::digest(&payloads[1]),
2756                outcome: outcomes[1],
2757            },
2758        ],
2759        accepted,
2760        full_read_hash: ObjectHash::digest(&full),
2761        range: ProbeRangeReceipt {
2762            start: PROBE_RANGE_START,
2763            end: PROBE_RANGE_END,
2764            bytes_hash: ObjectHash::digest(&range),
2765        },
2766        delete_verified_absent: true,
2767        lost_response: LostResponseProbeReceipt {
2768            logical_key: lost_logical_key,
2769            slot: lost_slot,
2770            payload_hash: ObjectHash::digest(&lost_payload),
2771            settled,
2772            readback_hash: ObjectHash::digest(&lost_readback),
2773            delete_verified_absent: true,
2774        },
2775    };
2776    let receipt =
2777        ExactSlotProbeReceipt::from_transcript(transcript, &binding.store, &binding.device);
2778    receipt.verify(&binding.store, &binding.device)?;
2779    advance_exact(
2780        journal,
2781        &mut durable,
2782        &mut record,
2783        ExactProbeProgress::ReceiptReady {
2784            receipt: receipt.clone(),
2785        },
2786    )
2787    .await?;
2788    Ok(receipt)
2789}
2790
2791pub async fn probe_serial_coordination_receipt(
2792    first: &dyn CoordinationStorage,
2793    second: &dyn CoordinationStorage,
2794    journal: &dyn ProviderProbeJournal,
2795    probe_id: ProviderProbeId,
2796    binding: &crate::sync::storage::ResolvedProviderBinding,
2797) -> Result<SerialCoordinationProbeReceipt, ProviderProbeError> {
2798    binding.validate()?;
2799    let first_binding = first
2800        .provider_binding()
2801        .await
2802        .map_err(|error| ProviderProbeError::Storage(StorageError::Storage(error.to_string())))?;
2803    let second_binding = second
2804        .provider_binding()
2805        .await
2806        .map_err(|error| ProviderProbeError::Storage(StorageError::Storage(error.to_string())))?;
2807    if first_binding != *binding || second_binding != *binding {
2808        return invalid("serial coordination clients do not match the receipt binding");
2809    }
2810    let logical_key = format!(
2811        "__coven_probe__/serial/{}",
2812        hex::encode(probe_id.as_bytes())
2813    );
2814    let mut durable = match journal.load(probe_id).await? {
2815        Some(existing) => existing,
2816        None => {
2817            journal
2818                .begin(ProviderProbeJournalRecord::Serial(SerialProbeJournal {
2819                    probe_id,
2820                    binding: binding.clone(),
2821                    logical_key: logical_key.clone(),
2822                    progress: SerialProbeProgress::Prepared,
2823                }))
2824                .await?
2825        }
2826    };
2827    let ProviderProbeJournalRecord::Serial(mut record) = durable.clone() else {
2828        return invalid("serial probe id belongs to a different durable probe kind");
2829    };
2830    if record.binding != *binding
2831        || record.probe_id != probe_id
2832        || record.logical_key != logical_key
2833    {
2834        return invalid("durable serial probe differs from its requested context");
2835    }
2836    if let SerialProbeProgress::ReceiptReady { receipt } = &record.progress {
2837        receipt.verify(&binding.store, &binding.device)?;
2838        return Ok(receipt.clone());
2839    }
2840    let create_payloads = [
2841        probe_payload(&probe_id, ProbePayloadLabel::SerialCreateFirst),
2842        probe_payload(&probe_id, ProbePayloadLabel::SerialCreateSecond),
2843    ];
2844    if matches!(record.progress, SerialProbeProgress::Prepared) {
2845        let (attempts, object) = match first.read_head(&logical_key).await {
2846            Err(CoordinationError::NotFound(_)) => {
2847                let (left, right) = tokio::join!(
2848                    first.create_head(&logical_key, &create_payloads[0]),
2849                    second.create_head(&logical_key, &create_payloads[1]),
2850                );
2851                classify_serial_create(left, right, &create_payloads)?
2852            }
2853            Ok(object) if object.bytes == create_payloads[0] => {
2854                require_serial_create_rejection(
2855                    second.create_head(&logical_key, &create_payloads[1]).await,
2856                )?;
2857                (
2858                    [
2859                        ProbeCreateAttempt {
2860                            payload_hash: ObjectHash::digest(&create_payloads[0]),
2861                            outcome: ProbeCreateOutcome::Created,
2862                        },
2863                        ProbeCreateAttempt {
2864                            payload_hash: ObjectHash::digest(&create_payloads[1]),
2865                            outcome: ProbeCreateOutcome::RejectedOccupied,
2866                        },
2867                    ],
2868                    object,
2869                )
2870            }
2871            Ok(object) if object.bytes == create_payloads[1] => {
2872                require_serial_create_rejection(
2873                    first.create_head(&logical_key, &create_payloads[0]).await,
2874                )?;
2875                (
2876                    [
2877                        ProbeCreateAttempt {
2878                            payload_hash: ObjectHash::digest(&create_payloads[0]),
2879                            outcome: ProbeCreateOutcome::RejectedOccupied,
2880                        },
2881                        ProbeCreateAttempt {
2882                            payload_hash: ObjectHash::digest(&create_payloads[1]),
2883                            outcome: ProbeCreateOutcome::Created,
2884                        },
2885                    ],
2886                    object,
2887                )
2888            }
2889            Ok(_) => return invalid("durable serial probe key contains unknown create bytes"),
2890            Err(error) => {
2891                return Err(ProviderProbeError::Storage(StorageError::Storage(
2892                    error.to_string(),
2893                )))
2894            }
2895        };
2896        advance_serial(
2897            journal,
2898            &mut durable,
2899            &mut record,
2900            SerialProbeProgress::Created { attempts, object },
2901        )
2902        .await?;
2903    }
2904    let (create_attempts, created) = serial_created_state(&record.progress)?;
2905    let replace_payloads = [
2906        probe_payload(&probe_id, ProbePayloadLabel::SerialReplaceFirst),
2907        probe_payload(&probe_id, ProbePayloadLabel::SerialReplaceSecond),
2908    ];
2909    if matches!(record.progress, SerialProbeProgress::Created { .. }) {
2910        let observed = first.read_head(&logical_key).await.map_err(|error| {
2911            ProviderProbeError::Storage(StorageError::Storage(error.to_string()))
2912        })?;
2913        let (replace_attempts, replaced) = if observed == created {
2914            let (left, right) = tokio::join!(
2915                first.replace_head(&logical_key, &created.version, &replace_payloads[0]),
2916                second.replace_head(&logical_key, &created.version, &replace_payloads[1]),
2917            );
2918            classify_serial_replace(left, right, &replace_payloads)?
2919        } else if observed.bytes == replace_payloads[0] {
2920            require_serial_replace_rejection(
2921                second
2922                    .replace_head(&logical_key, &created.version, &replace_payloads[1])
2923                    .await,
2924            )?;
2925            (
2926                [
2927                    ProbeReplaceAttempt {
2928                        payload_hash: ObjectHash::digest(&replace_payloads[0]),
2929                        outcome: ProbeReplaceOutcome::Replaced,
2930                    },
2931                    ProbeReplaceAttempt {
2932                        payload_hash: ObjectHash::digest(&replace_payloads[1]),
2933                        outcome: ProbeReplaceOutcome::RejectedVersionMismatch,
2934                    },
2935                ],
2936                observed,
2937            )
2938        } else if observed.bytes == replace_payloads[1] {
2939            require_serial_replace_rejection(
2940                first
2941                    .replace_head(&logical_key, &created.version, &replace_payloads[0])
2942                    .await,
2943            )?;
2944            (
2945                [
2946                    ProbeReplaceAttempt {
2947                        payload_hash: ObjectHash::digest(&replace_payloads[0]),
2948                        outcome: ProbeReplaceOutcome::RejectedVersionMismatch,
2949                    },
2950                    ProbeReplaceAttempt {
2951                        payload_hash: ObjectHash::digest(&replace_payloads[1]),
2952                        outcome: ProbeReplaceOutcome::Replaced,
2953                    },
2954                ],
2955                observed,
2956            )
2957        } else {
2958            return invalid("durable serial probe key contains unknown replacement bytes");
2959        };
2960        advance_serial(
2961            journal,
2962            &mut durable,
2963            &mut record,
2964            SerialProbeProgress::Replaced {
2965                create_attempts,
2966                created: created.clone(),
2967                replace_attempts,
2968                replaced,
2969            },
2970        )
2971        .await?;
2972    }
2973    let (create_attempts, created, replace_attempts, replaced) =
2974        serial_replaced_state(&record.progress)?;
2975    if matches!(record.progress, SerialProbeProgress::Replaced { .. }) {
2976        let observed = second.read_head(&logical_key).await.map_err(|error| {
2977            ProviderProbeError::Storage(StorageError::Storage(error.to_string()))
2978        })?;
2979        if observed != replaced {
2980            return invalid("serial authoritative read differs from replacement winner");
2981        }
2982        advance_serial(
2983            journal,
2984            &mut durable,
2985            &mut record,
2986            SerialProbeProgress::ReadVerified {
2987                create_attempts: create_attempts.clone(),
2988                created: created.clone(),
2989                replace_attempts: replace_attempts.clone(),
2990                replaced: replaced.clone(),
2991            },
2992        )
2993        .await?;
2994    }
2995    if matches!(record.progress, SerialProbeProgress::ReadVerified { .. }) {
2996        match first.read_head(&logical_key).await {
2997            Err(CoordinationError::NotFound(_)) => {}
2998            Ok(object) if object == replaced => {
2999                first
3000                    .delete_probe_head(&logical_key)
3001                    .await
3002                    .map_err(|error| {
3003                        ProviderProbeError::Storage(StorageError::Storage(error.to_string()))
3004                    })?;
3005                match first.read_head(&logical_key).await {
3006                    Err(CoordinationError::NotFound(_)) => {}
3007                    Ok(_) => return invalid("serial coordination object remains after deletion"),
3008                    Err(error) => {
3009                        return Err(ProviderProbeError::Storage(StorageError::Storage(
3010                            error.to_string(),
3011                        )))
3012                    }
3013                }
3014            }
3015            Ok(_) => return invalid("serial cleanup found different authoritative bytes/version"),
3016            Err(error) => {
3017                return Err(ProviderProbeError::Storage(StorageError::Storage(
3018                    error.to_string(),
3019                )))
3020            }
3021        }
3022        advance_serial(
3023            journal,
3024            &mut durable,
3025            &mut record,
3026            SerialProbeProgress::Absent {
3027                create_attempts: create_attempts.clone(),
3028                created: created.clone(),
3029                replace_attempts: replace_attempts.clone(),
3030                replaced: replaced.clone(),
3031            },
3032        )
3033        .await?;
3034    }
3035    let transcript = SerialCoordinationProbeTranscript {
3036        probe_id,
3037        logical_key,
3038        create_attempts,
3039        created_bytes_hash: ObjectHash::digest(&created.bytes),
3040        created_version_hash: ObjectHash::digest(created.version.cloud().as_provider().as_bytes()),
3041        replace_attempts,
3042        replaced_bytes_hash: ObjectHash::digest(&replaced.bytes),
3043        replaced_version_hash: ObjectHash::digest(
3044            replaced.version.cloud().as_provider().as_bytes(),
3045        ),
3046        authoritative_read_bytes_hash: ObjectHash::digest(&replaced.bytes),
3047        authoritative_read_version_hash: ObjectHash::digest(
3048            replaced.version.cloud().as_provider().as_bytes(),
3049        ),
3050        delete_verified_absent: true,
3051    };
3052    let receipt = SerialCoordinationProbeReceipt::from_transcript(
3053        transcript,
3054        &binding.store,
3055        &binding.device,
3056    );
3057    receipt.verify(&binding.store, &binding.device)?;
3058    advance_serial(
3059        journal,
3060        &mut durable,
3061        &mut record,
3062        SerialProbeProgress::ReceiptReady {
3063            receipt: receipt.clone(),
3064        },
3065    )
3066    .await?;
3067    Ok(receipt)
3068}
3069
3070async fn advance_serial(
3071    journal: &dyn ProviderProbeJournal,
3072    durable: &mut ProviderProbeJournalRecord,
3073    record: &mut SerialProbeJournal,
3074    progress: SerialProbeProgress,
3075) -> Result<(), ProviderProbeError> {
3076    record.progress = progress;
3077    let next = ProviderProbeJournalRecord::Serial(record.clone());
3078    journal.advance(durable, next.clone()).await?;
3079    *durable = next;
3080    Ok(())
3081}
3082
3083fn classify_serial_create(
3084    left: Result<VersionedObject, CreateHeadError>,
3085    right: Result<VersionedObject, CreateHeadError>,
3086    payloads: &[Vec<u8>; 2],
3087) -> Result<([ProbeCreateAttempt; 2], VersionedObject), ProviderProbeError> {
3088    let (outcomes, object) = match (left, right) {
3089        (Ok(object), Err(CreateHeadError::AlreadyExists)) => (
3090            [
3091                ProbeCreateOutcome::Created,
3092                ProbeCreateOutcome::RejectedOccupied,
3093            ],
3094            object,
3095        ),
3096        (Err(CreateHeadError::AlreadyExists), Ok(object)) => (
3097            [
3098                ProbeCreateOutcome::RejectedOccupied,
3099                ProbeCreateOutcome::Created,
3100            ],
3101            object,
3102        ),
3103        (left, right) => {
3104            return invalid(&format!(
3105                "serial create race did not produce one create and one occupied rejection: left={left:?}, right={right:?}"
3106            ));
3107        }
3108    };
3109    Ok((
3110        [
3111            ProbeCreateAttempt {
3112                payload_hash: ObjectHash::digest(&payloads[0]),
3113                outcome: outcomes[0],
3114            },
3115            ProbeCreateAttempt {
3116                payload_hash: ObjectHash::digest(&payloads[1]),
3117                outcome: outcomes[1],
3118            },
3119        ],
3120        object,
3121    ))
3122}
3123
3124fn require_serial_create_rejection(
3125    result: Result<VersionedObject, CreateHeadError>,
3126) -> Result<(), ProviderProbeError> {
3127    match result {
3128        Err(CreateHeadError::AlreadyExists) => Ok(()),
3129        Ok(_) => invalid("settled serial contender unexpectedly created a second head"),
3130        Err(error) => Err(ProviderProbeError::Storage(StorageError::Storage(
3131            error.to_string(),
3132        ))),
3133    }
3134}
3135
3136fn classify_serial_replace(
3137    left: Result<VersionedObject, ReplaceHeadError>,
3138    right: Result<VersionedObject, ReplaceHeadError>,
3139    payloads: &[Vec<u8>; 2],
3140) -> Result<([ProbeReplaceAttempt; 2], VersionedObject), ProviderProbeError> {
3141    let (outcomes, object) = match (left, right) {
3142        (Ok(object), Err(ReplaceHeadError::VersionMismatch)) => (
3143            [
3144                ProbeReplaceOutcome::Replaced,
3145                ProbeReplaceOutcome::RejectedVersionMismatch,
3146            ],
3147            object,
3148        ),
3149        (Err(ReplaceHeadError::VersionMismatch), Ok(object)) => (
3150            [
3151                ProbeReplaceOutcome::RejectedVersionMismatch,
3152                ProbeReplaceOutcome::Replaced,
3153            ],
3154            object,
3155        ),
3156        (left, right) => {
3157            return invalid(&format!(
3158                "serial replacement race did not produce one replace and one version rejection: left={left:?}, right={right:?}"
3159            ));
3160        }
3161    };
3162    Ok((
3163        [
3164            ProbeReplaceAttempt {
3165                payload_hash: ObjectHash::digest(&payloads[0]),
3166                outcome: outcomes[0],
3167            },
3168            ProbeReplaceAttempt {
3169                payload_hash: ObjectHash::digest(&payloads[1]),
3170                outcome: outcomes[1],
3171            },
3172        ],
3173        object,
3174    ))
3175}
3176
3177fn require_serial_replace_rejection(
3178    result: Result<VersionedObject, ReplaceHeadError>,
3179) -> Result<(), ProviderProbeError> {
3180    match result {
3181        Err(ReplaceHeadError::VersionMismatch) => Ok(()),
3182        Ok(_) => invalid("settled serial contender unexpectedly replaced the head"),
3183        Err(error) => Err(ProviderProbeError::Storage(StorageError::Storage(
3184            error.to_string(),
3185        ))),
3186    }
3187}
3188
3189fn serial_created_state(
3190    progress: &SerialProbeProgress,
3191) -> Result<([ProbeCreateAttempt; 2], VersionedObject), ProviderProbeError> {
3192    match progress {
3193        SerialProbeProgress::Created { attempts, object } => Ok((attempts.clone(), object.clone())),
3194        SerialProbeProgress::Replaced {
3195            create_attempts,
3196            created,
3197            ..
3198        }
3199        | SerialProbeProgress::ReadVerified {
3200            create_attempts,
3201            created,
3202            ..
3203        }
3204        | SerialProbeProgress::Absent {
3205            create_attempts,
3206            created,
3207            ..
3208        } => Ok((create_attempts.clone(), created.clone())),
3209        SerialProbeProgress::Prepared | SerialProbeProgress::ReceiptReady { .. } => {
3210            invalid("serial journal has no durable create result")
3211        }
3212    }
3213}
3214
3215fn serial_replaced_state(
3216    progress: &SerialProbeProgress,
3217) -> Result<
3218    (
3219        [ProbeCreateAttempt; 2],
3220        VersionedObject,
3221        [ProbeReplaceAttempt; 2],
3222        VersionedObject,
3223    ),
3224    ProviderProbeError,
3225> {
3226    match progress {
3227        SerialProbeProgress::Replaced {
3228            create_attempts,
3229            created,
3230            replace_attempts,
3231            replaced,
3232        }
3233        | SerialProbeProgress::ReadVerified {
3234            create_attempts,
3235            created,
3236            replace_attempts,
3237            replaced,
3238        }
3239        | SerialProbeProgress::Absent {
3240            create_attempts,
3241            created,
3242            replace_attempts,
3243            replaced,
3244        } => Ok((
3245            create_attempts.clone(),
3246            created.clone(),
3247            replace_attempts.clone(),
3248            replaced.clone(),
3249        )),
3250        _ => invalid("serial journal has no durable replacement result"),
3251    }
3252}
3253
3254async fn advance_exact(
3255    journal: &dyn ProviderProbeJournal,
3256    durable: &mut ProviderProbeJournalRecord,
3257    record: &mut ExactProbeJournal,
3258    progress: ExactProbeProgress,
3259) -> Result<(), ProviderProbeError> {
3260    record.progress = progress;
3261    let next = ProviderProbeJournalRecord::Exact(record.clone());
3262    journal.advance(durable, next.clone()).await?;
3263    *durable = next;
3264    Ok(())
3265}
3266
3267fn exact_race_state(
3268    progress: &ExactProbeProgress,
3269) -> Result<([ProbeCreateOutcome; 2], usize), ProviderProbeError> {
3270    let (outcomes, winner) = match progress {
3271        ExactProbeProgress::Prepared => return invalid("exact probe has no durable create result"),
3272        ExactProbeProgress::Created { outcomes }
3273        | ExactProbeProgress::ReadsVerified { outcomes }
3274        | ExactProbeProgress::PrimaryAbsent { outcomes }
3275        | ExactProbeProgress::LostResponseCreated { outcomes }
3276        | ExactProbeProgress::LostResponseReadVerified { outcomes }
3277        | ExactProbeProgress::Absent { outcomes } => {
3278            let winner = outcomes
3279                .iter()
3280                .position(|outcome| *outcome == ProbeCreateOutcome::Created)
3281                .ok_or_else(|| {
3282                    ProviderProbeError::InvalidReceipt(
3283                        "durable exact probe has no create winner".to_string(),
3284                    )
3285                })?;
3286            (*outcomes, winner)
3287        }
3288        ExactProbeProgress::ReceiptReady { receipt } => {
3289            let winner = receipt
3290                .transcript
3291                .contenders
3292                .iter()
3293                .position(|attempt| attempt.outcome == ProbeCreateOutcome::Created)
3294                .ok_or_else(|| {
3295                    ProviderProbeError::InvalidReceipt(
3296                        "durable exact receipt has no create winner".to_string(),
3297                    )
3298                })?;
3299            (
3300                [
3301                    receipt.transcript.contenders[0].outcome,
3302                    receipt.transcript.contenders[1].outcome,
3303                ],
3304                winner,
3305            )
3306        }
3307    };
3308    if winner > 1 || outcomes[winner] != ProbeCreateOutcome::Created {
3309        return invalid("durable exact probe has an invalid winner");
3310    }
3311    Ok((outcomes, winner))
3312}
3313
3314fn classify_exact_create_race(
3315    left: Result<(), CloudHomeError>,
3316    right: Result<(), CloudHomeError>,
3317) -> Result<([ProbeCreateOutcome; 2], usize), ProviderProbeError> {
3318    match (left, right) {
3319        (Ok(()), Err(CloudHomeError::AlreadyExists(_))) => Ok((
3320            [
3321                ProbeCreateOutcome::Created,
3322                ProbeCreateOutcome::RejectedOccupied,
3323            ],
3324            0,
3325        )),
3326        (Err(CloudHomeError::AlreadyExists(_)), Ok(())) => Ok((
3327            [
3328                ProbeCreateOutcome::RejectedOccupied,
3329                ProbeCreateOutcome::Created,
3330            ],
3331            1,
3332        )),
3333        (left, right) => invalid(&format!(
3334            "exact-slot race did not produce one create and one occupied rejection: left={left:?}, right={right:?}"
3335        )),
3336    }
3337}
3338
3339fn require_occupied_rejection(
3340    result: Result<(), CloudHomeError>,
3341) -> Result<(), ProviderProbeError> {
3342    match result {
3343        Err(CloudHomeError::AlreadyExists(_)) => Ok(()),
3344        Ok(()) => invalid("settled exact probe contender unexpectedly created a second object"),
3345        Err(error) => Err(ProviderProbeError::Storage(StorageError::from(error))),
3346    }
3347}
3348
3349async fn cleanup_exact_slot(
3350    storage: &dyn ExactSlotStorage,
3351    slot: &ObjectSlot,
3352) -> Result<(), ProviderProbeError> {
3353    match storage.read_at(slot).await {
3354        Err(CloudHomeError::NotFound(_)) => Ok(()),
3355        Ok(_) => {
3356            storage.delete_at(slot).await.map_err(StorageError::from)?;
3357            match storage.read_at(slot).await {
3358                Err(CloudHomeError::NotFound(_)) => Ok(()),
3359                Ok(_) => invalid("exact-slot probe object remains after deletion"),
3360                Err(error) => Err(ProviderProbeError::Storage(StorageError::from(error))),
3361            }
3362        }
3363        Err(error) => Err(ProviderProbeError::Storage(StorageError::from(error))),
3364    }
3365}
3366
3367fn cross_transcript_hash(
3368    store: &StoreProviderBinding,
3369    context: &CrossPrincipalResponseContext,
3370    transcript: &CrossPrincipalProbeTranscript,
3371) -> ObjectHash {
3372    ObjectHash::digest(&domain_json(
3373        CROSS_TRANSCRIPT_DOMAIN,
3374        &(store, context, transcript),
3375    ))
3376}
3377
3378fn validate_cross_transcript_payloads(
3379    transcript: &CrossPrincipalProbeTranscript,
3380    context: &CrossPrincipalResponseContext,
3381) -> Result<(), ProviderProbeError> {
3382    validate_cross_challenge_payload(&transcript.challenge)?;
3383    validate_cross_response_payload(&transcript.response, &transcript.challenge, context)?;
3384    let peer = probe_payload(&transcript.challenge.probe_id, ProbePayloadLabel::CrossPeer);
3385    if transcript.administrator_read_peer_hash != ObjectHash::digest(&peer)
3386        || !transcript.administrator_delete_peer_verified_absent
3387        || !transcript.administrator_delete_own_verified_absent
3388    {
3389        return invalid("cross-principal object, read, or deletion evidence is invalid");
3390    }
3391    Ok(())
3392}
3393
3394fn cross_challenge_hash(
3395    store: &StoreProviderBinding,
3396    context: &CrossPrincipalChallengeContext,
3397    challenge: &CrossPrincipalProbeChallenge,
3398) -> ObjectHash {
3399    ObjectHash::digest(&domain_json(
3400        CROSS_CHALLENGE_DOMAIN,
3401        &(
3402            store,
3403            context,
3404            challenge.probe_id,
3405            &challenge.administrator_object,
3406        ),
3407    ))
3408}
3409
3410fn cross_response_hash(
3411    store: &StoreProviderBinding,
3412    context: &CrossPrincipalResponseContext,
3413    challenge: &CrossPrincipalProbeChallenge,
3414    response: &CrossPrincipalProbeResponse,
3415) -> ObjectHash {
3416    ObjectHash::digest(&domain_json(
3417        CROSS_RESPONSE_DOMAIN,
3418        &(
3419            store,
3420            context,
3421            challenge.challenge_hash,
3422            &response.provider_evidence,
3423            &response.peer_object,
3424            response.peer_read_administrator_hash,
3425        ),
3426    ))
3427}
3428
3429fn validate_cross_challenge_payload(
3430    challenge: &CrossPrincipalProbeChallenge,
3431) -> Result<(), ProviderProbeError> {
3432    let expected_key = cross_administrator_logical_key(challenge.probe_id);
3433    let payload = probe_payload(&challenge.probe_id, ProbePayloadLabel::CrossAdministrator);
3434    validate_probe_exact_object(
3435        &challenge.administrator_object,
3436        &expected_key,
3437        &payload,
3438        "cross-principal challenge",
3439    )
3440}
3441
3442fn validate_cross_response_payload(
3443    response: &CrossPrincipalProbeResponse,
3444    challenge: &CrossPrincipalProbeChallenge,
3445    context: &CrossPrincipalResponseContext,
3446) -> Result<(), ProviderProbeError> {
3447    let administrator = probe_payload(&challenge.probe_id, ProbePayloadLabel::CrossAdministrator);
3448    let peer = probe_payload(&challenge.probe_id, ProbePayloadLabel::CrossPeer);
3449    if response.challenge_hash != challenge.challenge_hash
3450        || response.peer_object.slot != context.response_slot
3451        || response.peer_read_administrator_hash != ObjectHash::digest(&administrator)
3452    {
3453        return invalid(
3454            "cross-principal response disagrees with its challenge or response context",
3455        );
3456    }
3457    validate_probe_exact_object(
3458        &response.peer_object,
3459        &cross_peer_logical_key(challenge.probe_id),
3460        &peer,
3461        "cross-principal response",
3462    )
3463}
3464
3465fn validate_probe_exact_object(
3466    receipt: &ProbeExactObjectReceipt,
3467    expected_logical_key: &str,
3468    payload: &[u8],
3469    label: &str,
3470) -> Result<(), ProviderProbeError> {
3471    let payload_hash = ObjectHash::digest(payload);
3472    if receipt.slot.logical_key() != expected_logical_key
3473        || receipt.slot != *receipt.object.slot()
3474        || receipt.payload_hash != payload_hash
3475        || receipt.object.stored_size() != payload.len() as u64
3476        || receipt.object.stored_hash() != payload_hash
3477    {
3478        return invalid(&format!(
3479            "{label} object reference or payload hash is invalid"
3480        ));
3481    }
3482    Ok(())
3483}
3484
3485fn cross_administrator_logical_key(probe_id: ProviderProbeId) -> String {
3486    format!(
3487        "__coven_probe__/cross/{}/administrator",
3488        hex::encode(probe_id.as_bytes())
3489    )
3490}
3491
3492pub(crate) fn cross_peer_logical_key(probe_id: ProviderProbeId) -> String {
3493    format!(
3494        "__coven_probe__/cross/{}/peer",
3495        hex::encode(probe_id.as_bytes())
3496    )
3497}
3498
3499fn validate_cross_provider_evidence_context(
3500    store: &StoreProviderBinding,
3501    context: &CrossPrincipalChallengeContext,
3502) -> Result<(), ProviderProbeError> {
3503    context
3504        .administrator_binding
3505        .validate_for(store)
3506        .map_err(ProviderProbeError::Storage)?;
3507    context
3508        .peer_binding
3509        .validate_for(store)
3510        .map_err(ProviderProbeError::Storage)?;
3511    if context.administrator_binding == context.peer_binding {
3512        return invalid("cross-principal context uses the same provider principal twice");
3513    }
3514    Ok(())
3515}
3516
3517fn validate_cross_provider_evidence(
3518    store: &StoreProviderBinding,
3519    administrator: &ProviderDeviceBinding,
3520    peer: &ProviderDeviceBinding,
3521    evidence: &CrossPrincipalProviderEvidence,
3522) -> Result<(), ProviderProbeError> {
3523    administrator
3524        .validate_for(store)
3525        .map_err(ProviderProbeError::Storage)?;
3526    peer.validate_for(store)
3527        .map_err(ProviderProbeError::Storage)?;
3528    if administrator == peer {
3529        return invalid("cross-principal receipt uses the same provider principal twice");
3530    }
3531    let compatible = matches!(
3532        (store, evidence),
3533        (
3534            StoreProviderBinding::GoogleDrive {
3535                corpus: crate::sync::storage::GoogleDriveCorpus::SharedDrive { .. }
3536            },
3537            CrossPrincipalProviderEvidence::GoogleSharedDrive
3538        ) | (
3539            StoreProviderBinding::Dropbox { .. },
3540            CrossPrincipalProviderEvidence::DropboxSharedNamespace
3541        ) | (
3542            StoreProviderBinding::OneDrive { .. },
3543            CrossPrincipalProviderEvidence::OneDriveSharedFolder
3544        ) | (
3545            StoreProviderBinding::CloudKit { .. },
3546            CrossPrincipalProviderEvidence::CloudKit(_)
3547        )
3548    );
3549    if !compatible {
3550        return invalid("provider binding does not permit the cross-principal evidence");
3551    }
3552    if let (
3553        StoreProviderBinding::CloudKit {
3554            owner_name,
3555            zone_name,
3556            ..
3557        },
3558        CrossPrincipalProviderEvidence::CloudKit(accepted),
3559    ) = (store, evidence)
3560    {
3561        let crate::sync::storage::ProviderPrincipalId::CloudKitSharedZoneParticipant {
3562            record_name,
3563        } = &peer.principal
3564        else {
3565            return invalid("CloudKit peer is not a shared-zone participant");
3566        };
3567        if accepted.owner_name != *owner_name
3568            || accepted.zone_name != *zone_name
3569            || accepted.participant_record_name != *record_name
3570            || accepted.share_record_name.is_empty()
3571        {
3572            return invalid("CloudKit accepted-share evidence differs from the Store binding");
3573        }
3574    }
3575    Ok(())
3576}
3577
3578fn invalid<T>(reason: &str) -> Result<T, ProviderProbeError> {
3579    Err(ProviderProbeError::InvalidReceipt(reason.to_string()))
3580}
3581
3582fn domain_json<T: Serialize>(domain: &[u8], value: &T) -> Vec<u8> {
3583    let mut bytes = domain.to_vec();
3584    bytes.extend(
3585        serde_json::to_vec(value).expect("closed provider transcript serialization cannot fail"),
3586    );
3587    bytes
3588}
3589
3590fn exact_transcript_hash(
3591    store: &StoreProviderBinding,
3592    device: &ProviderDeviceBinding,
3593    transcript: &ExactSlotProbeTranscript,
3594) -> ObjectHash {
3595    ObjectHash::digest(&domain_json(
3596        EXACT_TRANSCRIPT_DOMAIN,
3597        &(store, device, transcript),
3598    ))
3599}
3600
3601fn serial_transcript_hash(
3602    store: &StoreProviderBinding,
3603    device: &ProviderDeviceBinding,
3604    transcript: &SerialCoordinationProbeTranscript,
3605) -> ObjectHash {
3606    ObjectHash::digest(&domain_json(
3607        SERIAL_TRANSCRIPT_DOMAIN,
3608        &(store, device, transcript),
3609    ))
3610}
3611
3612fn verify_race_create(
3613    attempts: &[ProbeCreateAttempt; 2],
3614    payloads: &[Vec<u8>; 2],
3615) -> Result<(), ProviderProbeError> {
3616    for (attempt, payload) in attempts.iter().zip(payloads) {
3617        if attempt.payload_hash != ObjectHash::digest(payload) {
3618            return invalid("serial create payload hash is not deterministic");
3619        }
3620    }
3621    if attempts
3622        .iter()
3623        .filter(|attempt| attempt.outcome == ProbeCreateOutcome::Created)
3624        .count()
3625        != 1
3626        || attempts
3627            .iter()
3628            .filter(|attempt| attempt.outcome == ProbeCreateOutcome::RejectedOccupied)
3629            .count()
3630            != 1
3631    {
3632        return invalid("serial create race must contain one create and one occupied rejection");
3633    }
3634    Ok(())
3635}
3636
3637fn verify_race_replace(
3638    attempts: &[ProbeReplaceAttempt; 2],
3639    payloads: &[Vec<u8>; 2],
3640) -> Result<(), ProviderProbeError> {
3641    for (attempt, payload) in attempts.iter().zip(payloads) {
3642        if attempt.payload_hash != ObjectHash::digest(payload) {
3643            return invalid("serial replacement payload hash is not deterministic");
3644        }
3645    }
3646    if attempts
3647        .iter()
3648        .filter(|attempt| attempt.outcome == ProbeReplaceOutcome::Replaced)
3649        .count()
3650        != 1
3651        || attempts
3652            .iter()
3653            .filter(|attempt| attempt.outcome == ProbeReplaceOutcome::RejectedVersionMismatch)
3654            .count()
3655            != 1
3656    {
3657        return invalid(
3658            "serial replacement race must contain one replace and one version rejection",
3659        );
3660    }
3661    Ok(())
3662}
3663
3664mod ordered_owner_barriers {
3665    use super::*;
3666
3667    pub(super) fn serialize<S>(
3668        map: &BTreeMap<MembershipGrantId, OwnerStreamBarrier>,
3669        serializer: S,
3670    ) -> Result<S::Ok, S::Error>
3671    where
3672        S: Serializer,
3673    {
3674        map.iter().collect::<Vec<_>>().serialize(serializer)
3675    }
3676
3677    pub(super) fn deserialize<'de, D>(
3678        deserializer: D,
3679    ) -> Result<BTreeMap<MembershipGrantId, OwnerStreamBarrier>, D::Error>
3680    where
3681        D: Deserializer<'de>,
3682    {
3683        let entries = Vec::<(MembershipGrantId, OwnerStreamBarrier)>::deserialize(deserializer)?;
3684        let count = entries.len();
3685        let map = entries.into_iter().collect::<BTreeMap<_, _>>();
3686        if map.len() != count {
3687            return Err(serde::de::Error::custom(
3688                "provider administrator owner barriers contain a duplicate grant",
3689            ));
3690        }
3691        Ok(map)
3692    }
3693}
3694
3695#[cfg(test)]
3696mod tests {
3697    use super::*;
3698    use crate::sync::store_commit::ObjectHash;
3699
3700    #[test]
3701    fn exact_probe_verifier_rejects_two_created_contenders() {
3702        let mut receipt = test_exact_receipt();
3703        receipt
3704            .verify(&test_store_binding(), &test_device_binding())
3705            .expect("baseline exact receipt verifies");
3706        receipt.transcript.contenders[1].outcome = ProbeCreateOutcome::Created;
3707
3708        assert!(receipt
3709            .verify(&test_store_binding(), &test_device_binding())
3710            .is_err());
3711    }
3712
3713    #[test]
3714    fn custom_s3_origin_rejects_paths_and_canonicalizes_default_port() {
3715        assert_eq!(
3716            canonical_custom_s3_origin("HTTPS://Objects.Example:443").unwrap(),
3717            "https://objects.example"
3718        );
3719        assert!(canonical_custom_s3_origin("https://objects.example/").is_err());
3720        assert!(canonical_custom_s3_origin("https://objects.example/bucket").is_err());
3721    }
3722
3723    #[test]
3724    fn private_cloudkit_owner_exposes_its_exact_administrator_locator() {
3725        let binding = private_cloudkit_binding();
3726
3727        let locator = ProviderAccessLocator::for_current_administrator(&binding)
3728            .expect("private CloudKit owner exposes its exact administrator locator");
3729
3730        assert_eq!(
3731            locator,
3732            ProviderAccessLocator::CloudKitPrivateZoneOwner {
3733                owner_name: "private-owner".to_string(),
3734                zone_name: "private-zone".to_string(),
3735                owner_record_name: "current-user".to_string(),
3736            }
3737        );
3738        locator
3739            .validate_for(&binding.store, &binding.device)
3740            .expect("private CloudKit owner locator matches its binding");
3741    }
3742
3743    #[test]
3744    fn shared_cloudkit_participant_is_not_treated_as_the_zone_owner() {
3745        let mut binding = private_cloudkit_binding();
3746        binding.device.principal =
3747            crate::sync::storage::ProviderPrincipalId::CloudKitSharedZoneParticipant {
3748                record_name: "current-user".to_string(),
3749            };
3750
3751        assert!(ProviderAccessLocator::for_current_administrator(&binding).is_err());
3752    }
3753
3754    #[test]
3755    fn provider_admin_grants_coexist_and_replay_exactly() {
3756        let founder = admin_record(1, "founder");
3757        let mut state = ProviderAdminState::founder(founder.clone());
3758        let second = admin_record(2, "second");
3759        let change = set_change(&second, BTreeSet::new());
3760        state
3761            .apply(change.clone(), second.created_at.clone())
3762            .expect("a second administrator may coexist");
3763        state
3764            .apply(change, second.created_at.clone())
3765            .expect("an exact replay is idempotent");
3766        assert_eq!(state.active().len(), 2);
3767        assert!(state.tombstones().is_empty());
3768    }
3769
3770    #[test]
3771    fn provider_admin_rejects_conflicting_id_reuse() {
3772        let founder = admin_record(1, "founder");
3773        let mut state = ProviderAdminState::founder(founder);
3774        let second = admin_record(2, "second");
3775        state
3776            .apply(
3777                set_change(&second, BTreeSet::new()),
3778                second.created_at.clone(),
3779            )
3780            .unwrap();
3781        let mut conflicting = second.clone();
3782        conflicting.provider = ProviderDeviceBinding {
3783            principal: crate::sync::storage::ProviderPrincipalId::Aws {
3784                account_id: "999999999999".to_string(),
3785                principal: crate::sync::storage::AwsPrincipal::Root,
3786            },
3787        };
3788        assert_eq!(
3789            state.apply(
3790                set_change(&conflicting, BTreeSet::new()),
3791                conflicting.created_at.clone()
3792            ),
3793            Err(ProviderAdminReducerError::GrantIdReuse)
3794        );
3795    }
3796
3797    #[test]
3798    fn provider_admin_removal_and_replacement_retain_tombstones() {
3799        let founder = admin_record(1, "founder");
3800        let founder_id = founder.grant_id.clone();
3801        let mut state = ProviderAdminState::founder(founder);
3802        let replacement = admin_record(2, "replacement");
3803        state
3804            .apply(
3805                set_change(&replacement, BTreeSet::from([founder_id.clone()])),
3806                replacement.created_at.clone(),
3807            )
3808            .unwrap();
3809        assert!(state.records().contains_key(&founder_id));
3810        assert!(state.tombstones().contains(&founder_id));
3811        assert_eq!(
3812            state.apply(
3813                ProviderAdminChange::Remove {
3814                    removes: BTreeSet::from([replacement.grant_id.clone()]),
3815                },
3816                replacement.created_at.clone(),
3817            ),
3818            Err(ProviderAdminReducerError::NoEffectiveAdministrator)
3819        );
3820        assert!(state.records().contains_key(&replacement.grant_id));
3821        assert!(!state.tombstones().contains(&replacement.grant_id));
3822    }
3823
3824    #[test]
3825    fn provider_admin_replay_cannot_tombstone_a_newly_active_replacement() {
3826        let founder = admin_record(1, "founder");
3827        let mut state = ProviderAdminState::founder(founder);
3828        let second = admin_record(2, "second");
3829        state
3830            .apply(
3831                set_change(&second, BTreeSet::new()),
3832                second.created_at.clone(),
3833            )
3834            .unwrap();
3835        let third = admin_record(3, "third");
3836        state
3837            .apply(
3838                set_change(&third, BTreeSet::new()),
3839                third.created_at.clone(),
3840            )
3841            .unwrap();
3842        assert_eq!(
3843            state.apply(
3844                set_change(&second, BTreeSet::from([third.grant_id.clone()])),
3845                second.created_at.clone(),
3846            ),
3847            Err(ProviderAdminReducerError::UnknownReplacement)
3848        );
3849        assert!(state.active().contains(&third.grant_id));
3850    }
3851
3852    #[tokio::test]
3853    async fn database_probe_journal_rejects_a_skipped_progress_state() {
3854        let db = crate::sync::test_helpers::open_test_db();
3855        let probe_id = ProviderProbeId::from_bytes([44; 32]);
3856        let binding = crate::sync::storage::ResolvedProviderBinding {
3857            store: test_store_binding(),
3858            device: test_device_binding(),
3859        };
3860        let prepared = ProviderProbeJournalRecord::Exact(ExactProbeJournal {
3861            probe_id,
3862            binding,
3863            slot: ObjectSlot::logical("__coven_probe__/exact/journal".to_string()).unwrap(),
3864            lost_response_slot: ObjectSlot::logical(
3865                "__coven_probe__/lost-response/journal".to_string(),
3866            )
3867            .unwrap(),
3868            progress: ExactProbeProgress::Prepared,
3869        });
3870        assert_eq!(db.begin(prepared.clone()).await.unwrap(), prepared);
3871        let ProviderProbeJournalRecord::Exact(mut final_record) = prepared.clone() else {
3872            unreachable!()
3873        };
3874        final_record.progress = ExactProbeProgress::ReceiptReady {
3875            receipt: test_exact_receipt(),
3876        };
3877        let final_record = ProviderProbeJournalRecord::Exact(final_record);
3878        assert!(db.advance(&prepared, final_record).await.is_err());
3879        assert_eq!(db.load(probe_id).await.unwrap(), Some(prepared));
3880    }
3881
3882    fn admin_record(id: u8, label: &str) -> ProviderAdminGrantRecord {
3883        let root_object = crate::sync::storage::ExactObjectRef::new(
3884            crate::storage::cloud::ObjectSlot::logical(format!("roots/{label}")).unwrap(),
3885            1,
3886            ObjectHash::digest(&[id]),
3887        );
3888        let root = StoreRootRef {
3889            store_root_id: ObjectHash::digest(format!("{label} id").as_bytes()),
3890            store_root_hash: ObjectHash::digest(label.as_bytes()),
3891            object: root_object,
3892        };
3893        let registration: StoreDeviceRegistrationRef = serde_json::from_value(serde_json::json!({
3894            "device_id": ObjectHash::digest(&[id, 1]),
3895            "registration_hash": ObjectHash::digest(&[id, 2]),
3896            "object": {
3897                "slot": {"logical_key": format!("registrations/{label}"), "physical": {"kind": "logical_key"}},
3898                "stored_size": 1,
3899                "stored_hash": ObjectHash::digest(&[id, 3]),
3900            }
3901        }))
3902        .unwrap();
3903        ProviderAdminGrantRecord {
3904            grant_id: ProviderAdminGrantId(ObjectHash::digest(&[id, 4])),
3905            administrator: registration,
3906            provider: test_device_binding(),
3907            access: ProviderAccessLocator::S3SharedCredentialGeneration {
3908                generation: 1,
3909                access_key_id_hash: ObjectHash::digest(b"test access key"),
3910            },
3911            capability: ProviderCapabilityProof {
3912                exact_slots: test_exact_receipt(),
3913                serial_coordination: None,
3914            },
3915            created_at: ProviderAdminGrantOrigin::Founder { root },
3916        }
3917    }
3918
3919    fn set_change(
3920        record: &ProviderAdminGrantRecord,
3921        replaces: BTreeSet<ProviderAdminGrantId>,
3922    ) -> ProviderAdminChange {
3923        ProviderAdminChange::Set {
3924            administrator: record.administrator.clone(),
3925            provider: record.provider.clone(),
3926            access: record.access.clone(),
3927            capability: record.capability.clone(),
3928            grant_id: record.grant_id.clone(),
3929            replaces,
3930        }
3931    }
3932
3933    fn test_store_binding() -> crate::sync::storage::StoreProviderBinding {
3934        crate::sync::storage::StoreProviderBinding::S3 {
3935            endpoint: crate::sync::storage::S3EndpointBinding::Aws {
3936                partition: "aws".to_string(),
3937            },
3938            region: "us-east-1".to_string(),
3939            bucket: "bucket".to_string(),
3940            key_prefix: None,
3941        }
3942    }
3943
3944    fn private_cloudkit_binding() -> crate::sync::storage::ResolvedProviderBinding {
3945        crate::sync::storage::ResolvedProviderBinding {
3946            store: crate::sync::storage::StoreProviderBinding::CloudKit {
3947                container_id: "iCloud.example.coven".to_string(),
3948                environment: crate::sync::storage::CloudKitEnvironment::Development,
3949                owner_name: "private-owner".to_string(),
3950                zone_name: "private-zone".to_string(),
3951            },
3952            device: crate::sync::storage::ProviderDeviceBinding {
3953                principal: crate::sync::storage::ProviderPrincipalId::CloudKitPrivateZoneOwner {
3954                    record_name: "current-user".to_string(),
3955                },
3956            },
3957        }
3958    }
3959
3960    fn test_device_binding() -> crate::sync::storage::ProviderDeviceBinding {
3961        crate::sync::storage::ProviderDeviceBinding {
3962            principal: crate::sync::storage::ProviderPrincipalId::Aws {
3963                account_id: "123456789012".to_string(),
3964                principal: crate::sync::storage::AwsPrincipal::Root,
3965            },
3966        }
3967    }
3968
3969    fn test_exact_receipt() -> ExactSlotProbeReceipt {
3970        let probe_id = ProviderProbeId::from_bytes([7; 32]);
3971        let slot = crate::storage::cloud::ObjectSlot::logical("store-v1/probes/exact".to_string())
3972            .unwrap();
3973        let first = probe_payload(&probe_id, ProbePayloadLabel::ExactCreateFirst);
3974        let second = probe_payload(&probe_id, ProbePayloadLabel::ExactCreateSecond);
3975        let accepted = crate::sync::storage::ExactObjectRef::new(
3976            slot.clone(),
3977            first.len() as u64,
3978            ObjectHash::digest(&first),
3979        );
3980        let lost_slot =
3981            crate::storage::cloud::ObjectSlot::logical("store-v1/probes/lost".to_string()).unwrap();
3982        let lost_payload = probe_payload(&probe_id, ProbePayloadLabel::LostResponse);
3983        let lost_ref = crate::sync::storage::ExactObjectRef::new(
3984            lost_slot.clone(),
3985            lost_payload.len() as u64,
3986            ObjectHash::digest(&lost_payload),
3987        );
3988        let transcript = ExactSlotProbeTranscript {
3989            probe_id,
3990            logical_key: slot.logical_key().to_string(),
3991            slot,
3992            contenders: [
3993                ProbeCreateAttempt {
3994                    payload_hash: ObjectHash::digest(&first),
3995                    outcome: ProbeCreateOutcome::Created,
3996                },
3997                ProbeCreateAttempt {
3998                    payload_hash: ObjectHash::digest(&second),
3999                    outcome: ProbeCreateOutcome::RejectedOccupied,
4000                },
4001            ],
4002            accepted: accepted.clone(),
4003            full_read_hash: accepted.stored_hash(),
4004            range: ProbeRangeReceipt {
4005                start: PROBE_RANGE_START,
4006                end: PROBE_RANGE_END,
4007                bytes_hash: ObjectHash::digest(
4008                    &first[PROBE_RANGE_START as usize..PROBE_RANGE_END as usize],
4009                ),
4010            },
4011            delete_verified_absent: true,
4012            lost_response: LostResponseProbeReceipt {
4013                logical_key: "store-v1/probes/lost".to_string(),
4014                slot: lost_slot,
4015                payload_hash: ObjectHash::digest(&lost_payload),
4016                settled: lost_ref,
4017                readback_hash: ObjectHash::digest(&lost_payload),
4018                delete_verified_absent: true,
4019            },
4020        };
4021        ExactSlotProbeReceipt::from_transcript(
4022            transcript,
4023            &test_store_binding(),
4024            &test_device_binding(),
4025        )
4026    }
4027}