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