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 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#[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 async fn begin(
1852 &self,
1853 prepared: ProviderProbeJournalRecord,
1854 ) -> Result<ProviderProbeJournalRecord, StorageError>;
1855
1856 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}