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