Skip to main content

coven_storage/
provider_probe.rs

1//! Cross-principal provider probe execution: reserving, creating, and
2//! settling exact probe slots on the primary and peer provider storage, over
3//! the probe transcript model in [`coven_protocol::provider`].
4
5use std::sync::Arc;
6
7use crate::cloud::{
8    CloudHomeError, ConditionalWriteOutcome, ExactCloudHome, ExactCreateOutcome, ExactSlotStorage,
9    ExactUpload,
10};
11use coven_keys::keys::UserKeypair;
12use coven_protocol::objects::{ExactObjectRef, ObjectSlot, StorageError};
13use coven_protocol::provider::*;
14use coven_protocol::provider::{
15    advance_cross_completion, advance_exact, cross_challenge_hash, cross_response_hash, invalid,
16    validate_cross_provider_evidence, validate_cross_provider_evidence_context,
17};
18use coven_protocol::store_commit::ObjectHash;
19use coven_protocol::StoreProviderBinding;
20
21mod exact_slots;
22
23pub struct ProviderProbeStorage {
24    storage: Arc<dyn ExactCloudHome>,
25}
26
27impl ProviderProbeStorage {
28    pub fn new(storage: Arc<dyn ExactCloudHome>) -> Self {
29        Self { storage }
30    }
31
32    pub async fn reserve_cross_principal_response_slot(
33        &self,
34        probe_id: ProviderProbeId,
35    ) -> Result<ObjectSlot, ProviderProbeError> {
36        let logical = cross_peer_logical_key(probe_id);
37        let slot = self
38            .storage
39            .allocate_slot(&logical)
40            .await
41            .map_err(StorageError::from)?;
42        if slot.logical_key() != logical {
43            return invalid("cross-principal response slot changed its logical key");
44        }
45        Ok(slot)
46    }
47
48    pub async fn prepare_cross_principal_challenge(
49        &self,
50        publication_journal: &dyn DeviceJoinChallengePublicationJournal,
51        probe_id: ProviderProbeId,
52        store: &StoreProviderBinding,
53        context: &CrossPrincipalChallengeContext,
54        administrator_signer: &dyn coven_keys::keys::DeviceSigningAuthority,
55    ) -> Result<CrossPrincipalProbeChallenge, ProviderProbeError> {
56        let administrator_live = self
57            .storage
58            .provider_binding()
59            .await
60            .map_err(StorageError::from)?;
61        if administrator_live.store != *store
62            || administrator_live.device != context.administrator_binding
63        {
64            return invalid("cross-principal administrator does not match the challenge context");
65        }
66        validate_cross_provider_evidence_context(store, context)?;
67        let suffix = hex::encode(probe_id.as_bytes());
68        let administrator_key = format!("__coven_probe__/cross/{suffix}/administrator");
69        let administrator_slot = self
70            .storage
71            .allocate_slot(&administrator_key)
72            .await
73            .map_err(StorageError::from)?;
74        if administrator_slot.logical_key() != administrator_key {
75            return invalid("cross-principal administrator slot changed its logical key");
76        }
77        let administrator_payload = probe_payload(&probe_id, ProbePayloadLabel::CrossAdministrator);
78        let administrator_object = ProbeExactObjectReceipt {
79            slot: administrator_slot.clone(),
80            payload_hash: ObjectHash::digest(&administrator_payload),
81            object: ExactObjectRef::new(
82                administrator_slot,
83                administrator_payload.len() as u64,
84                ObjectHash::digest(&administrator_payload),
85            ),
86        };
87        let conditional_slot = self
88            .storage
89            .allocate_slot(&cross_conditional_logical_key(probe_id))
90            .await
91            .map_err(StorageError::from)?;
92        let unsigned = CrossPrincipalProbeChallenge {
93            probe_id,
94            administrator_object,
95            conditional_slot,
96            challenge_hash: ObjectHash::digest(&[]),
97            administrator_signature: String::new(),
98        };
99        let challenge_hash = cross_challenge_hash(store, context, &unsigned);
100        let challenge = CrossPrincipalProbeChallenge {
101            challenge_hash,
102            administrator_signature: hex::encode(
103                administrator_signer.sign(challenge_hash.as_bytes()),
104            ),
105            ..unsigned
106        };
107        challenge.verify(context, store, &administrator_signer.public_key_hex())?;
108        let durable = publication_journal.prepare(&challenge).await?;
109        if durable.challenge != challenge {
110            return invalid("durable cross-principal challenge differs from its prepared bytes");
111        }
112        Ok(challenge)
113    }
114
115    pub async fn settle_cross_principal_challenge(
116        &self,
117        publication_journal: &dyn DeviceJoinChallengePublicationJournal,
118        authorization: &DeviceJoinChallengePublicationAuthorization,
119        challenge: &CrossPrincipalProbeChallenge,
120        context: &CrossPrincipalChallengeContext,
121        store: &StoreProviderBinding,
122    ) -> Result<CrossPrincipalProbeChallenge, ProviderProbeError> {
123        let live = self
124            .storage
125            .provider_binding()
126            .await
127            .map_err(StorageError::from)?;
128        if live.store != *store || live.device != context.administrator_binding {
129            return invalid("cross-principal administrator does not match the published challenge");
130        }
131        publication_journal
132            .claim_published(authorization, challenge)
133            .await?;
134        let payload = probe_payload(&challenge.probe_id, ProbePayloadLabel::CrossAdministrator);
135        self.settle_exact_create(&challenge.administrator_object.slot, &payload)
136            .await?;
137        let observed = self
138            .storage
139            .read_at(&challenge.administrator_object.slot)
140            .await
141            .map_err(StorageError::from)?;
142        if observed != payload {
143            return invalid("published cross-principal challenge differs from its signed bytes");
144        }
145        let initial = probe_payload(&challenge.probe_id, ProbePayloadLabel::ConditionalInitial);
146        let current = match self
147            .storage
148            .read_versioned_at(&challenge.conditional_slot)
149            .await
150        {
151            Ok(current) => current,
152            Err(CloudHomeError::NotFound(_)) => {
153                create_versioned_bytes(
154                    self.storage.as_ref(),
155                    &challenge.conditional_slot,
156                    &initial,
157                )
158                .await?;
159                self.storage
160                    .read_versioned_at(&challenge.conditional_slot)
161                    .await
162                    .map_err(StorageError::from)?
163            }
164            Err(error) => return Err(StorageError::from(error).into()),
165        };
166        let peer = probe_payload(&challenge.probe_id, ProbePayloadLabel::ConditionalFirst);
167        if current.bytes != initial && current.bytes != peer {
168            return invalid("published cross-principal conditional slot contains unknown bytes");
169        }
170        Ok(challenge.clone())
171    }
172
173    pub async fn create_cross_principal_response(
174        &self,
175        challenge: &CrossPrincipalProbeChallenge,
176        context: &CrossPrincipalResponseContext,
177        store: &StoreProviderBinding,
178        administrator_signing_pubkey: &str,
179        peer_signer: &UserKeypair,
180    ) -> Result<CrossPrincipalProbeResponse, ProviderProbeError> {
181        challenge.verify(&context.challenge, store, administrator_signing_pubkey)?;
182        let peer_pubkey = coven_keys::keys::public_key_hex(peer_signer);
183        if context.challenge.member_pubkey != peer_pubkey {
184            return invalid("cross-principal peer signer is not the joining member");
185        }
186        let live = self
187            .storage
188            .provider_binding()
189            .await
190            .map_err(StorageError::from)?;
191        if live.store != *store || live.device != context.challenge.peer_binding {
192            return invalid("cross-principal peer does not match the response context");
193        }
194        let evidence = self
195            .storage
196            .cross_principal_evidence()
197            .await
198            .map_err(StorageError::from)?;
199        validate_cross_provider_evidence(
200            store,
201            &context.challenge.administrator_binding,
202            &context.challenge.peer_binding,
203            &evidence,
204        )?;
205        let expected_peer_key = cross_peer_logical_key(challenge.probe_id);
206        if context.response_slot.logical_key() != expected_peer_key {
207            return invalid("cross-principal response slot uses the wrong logical key");
208        }
209        let administrator_payload =
210            probe_payload(&challenge.probe_id, ProbePayloadLabel::CrossAdministrator);
211        let administrator_read = self
212            .storage
213            .read_at(&challenge.administrator_object.slot)
214            .await
215            .map_err(StorageError::from)?;
216        if administrator_read != administrator_payload {
217            return invalid("peer read bytes differ from the signed cross-principal challenge");
218        }
219        let conditional_start = match self.storage.read_at(&context.response_slot).await {
220            Ok(bytes) => {
221                let start: CrossPrincipalConditionalStart =
222                    serde_json::from_slice(&bytes).map_err(StorageError::from)?;
223                if start.challenge_hash != challenge.challenge_hash
224                    || start.canonical_bytes() != bytes
225                {
226                    return invalid(
227                        "durable peer observation differs from its challenge or canonical bytes",
228                    );
229                }
230                start
231            }
232            Err(CloudHomeError::NotFound(_)) => {
233                let current = self
234                    .storage
235                    .read_versioned_at(&challenge.conditional_slot)
236                    .await
237                    .map_err(StorageError::from)?;
238                let initial =
239                    probe_payload(&challenge.probe_id, ProbePayloadLabel::ConditionalInitial);
240                if current.bytes != initial {
241                    return invalid("peer conditional observation does not start at the administrator's initial bytes");
242                }
243                let start = CrossPrincipalConditionalStart {
244                    challenge_hash: challenge.challenge_hash,
245                    version: current.version,
246                };
247                self.settle_exact_create(&context.response_slot, &start.canonical_bytes())
248                    .await?;
249                start
250            }
251            Err(error) => return Err(StorageError::from(error).into()),
252        };
253        let peer_payload = conditional_start.canonical_bytes();
254        let peer_read = self
255            .storage
256            .read_at(&context.response_slot)
257            .await
258            .map_err(StorageError::from)?;
259        if peer_read != peer_payload {
260            return invalid("peer response readback differs from its prepared observation");
261        }
262        let replacement = probe_payload(&challenge.probe_id, ProbePayloadLabel::ConditionalFirst);
263        // Retry the same expected revision. A lost successful response settles
264        // through exact readback; it never authorizes a write against a new revision.
265        self.storage
266            .replace_at_if_version(
267                &challenge.conditional_slot,
268                &conditional_start.version,
269                replacement.clone(),
270            )
271            .await
272            .map_err(StorageError::from)?;
273        let accepted = self
274            .storage
275            .read_versioned_at(&challenge.conditional_slot)
276            .await
277            .map_err(StorageError::from)?;
278        if accepted.bytes != replacement || accepted.version == conditional_start.version {
279            return invalid("peer conditional readback does not establish its replacement");
280        }
281        let peer_object = ProbeExactObjectReceipt {
282            slot: context.response_slot.clone(),
283            payload_hash: ObjectHash::digest(&peer_payload),
284            object: ExactObjectRef::new(
285                context.response_slot.clone(),
286                peer_payload.len() as u64,
287                ObjectHash::digest(&peer_payload),
288            ),
289        };
290        let unsigned = CrossPrincipalProbeResponse {
291            conditional_start,
292            provider_evidence: evidence,
293            peer_object,
294            peer_read_administrator_hash: ObjectHash::digest(&administrator_read),
295            response_hash: ObjectHash::digest(&[]),
296            peer_signature: String::new(),
297        };
298        let response_hash = cross_response_hash(store, context, challenge, &unsigned);
299        let response = CrossPrincipalProbeResponse {
300            response_hash,
301            peer_signature: hex::encode(peer_signer.sign(response_hash.as_bytes())),
302            ..unsigned
303        };
304        response.verify(
305            challenge,
306            context,
307            store,
308            administrator_signing_pubkey,
309            &peer_pubkey,
310        )?;
311        Ok(response)
312    }
313
314    pub async fn complete_cross_principal_probe(
315        &self,
316        journal: &dyn ProviderProbeJournal,
317        challenge: &CrossPrincipalProbeChallenge,
318        response: &CrossPrincipalProbeResponse,
319        context: &CrossPrincipalResponseContext,
320        store: &StoreProviderBinding,
321        administrator_signer: &dyn coven_keys::keys::DeviceSigningAuthority,
322        peer_signing_pubkey: &str,
323    ) -> Result<CrossPrincipalProbeReceipt, ProviderProbeError> {
324        let administrator_pubkey = administrator_signer.public_key_hex();
325        challenge.verify(&context.challenge, store, &administrator_pubkey)?;
326        response.verify(
327            challenge,
328            context,
329            store,
330            &administrator_pubkey,
331            peer_signing_pubkey,
332        )?;
333        let live = self
334            .storage
335            .provider_binding()
336            .await
337            .map_err(StorageError::from)?;
338        if live.store != *store || live.device != context.challenge.administrator_binding {
339            return invalid("cross-principal administrator does not match the completion context");
340        }
341        let prepared =
342            ProviderProbeJournalRecord::CrossPrincipal(CrossPrincipalCompletionJournal {
343                probe_id: challenge.probe_id,
344                store: store.clone(),
345                context: context.clone(),
346                challenge: challenge.clone(),
347                response: response.clone(),
348                progress: CrossPrincipalCompletionProgress::Prepared,
349            });
350        let mut durable = match journal.load(challenge.probe_id).await? {
351            Some(existing) => existing,
352            None => journal.begin(prepared).await?,
353        };
354        let ProviderProbeJournalRecord::CrossPrincipal(mut record) = durable.clone() else {
355            return invalid("cross-principal probe id belongs to another durable probe kind");
356        };
357        if record.probe_id != challenge.probe_id
358            || record.store != *store
359            || record.context != *context
360            || record.challenge != *challenge
361            || record.response != *response
362        {
363            return invalid("durable cross-principal completion differs from the requested proof");
364        }
365        if let CrossPrincipalCompletionProgress::ReceiptReady { receipt } = &record.progress {
366            receipt.verify(context, store, &administrator_pubkey, peer_signing_pubkey)?;
367            return Ok(receipt.clone());
368        }
369        let peer_payload = response.conditional_start.canonical_bytes();
370        if matches!(record.progress, CrossPrincipalCompletionProgress::Prepared) {
371            let observed = self
372                .storage
373                .read_at(&response.peer_object.slot)
374                .await
375                .map_err(StorageError::from)?;
376            if observed != peer_payload {
377                return invalid("administrator read differs from the signed peer response");
378            }
379            let peer_replacement =
380                probe_payload(&challenge.probe_id, ProbePayloadLabel::ConditionalFirst);
381            let administrator_replacement =
382                probe_payload(&challenge.probe_id, ProbePayloadLabel::ConditionalSecond);
383            let before = self
384                .storage
385                .read_versioned_at(&challenge.conditional_slot)
386                .await
387                .map_err(StorageError::from)?;
388            if before.bytes != peer_replacement
389                || before.version == response.conditional_start.version
390            {
391                return invalid(
392                    "administrator does not observe the peer's conditional replacement",
393                );
394            }
395            let rejected = self
396                .storage
397                .replace_at_if_version(
398                    &challenge.conditional_slot,
399                    &response.conditional_start.version,
400                    administrator_replacement.clone(),
401                )
402                .await
403                .map_err(StorageError::from)?;
404            if rejected != ConditionalWriteOutcome::VersionChanged {
405                return invalid("administrator replaced a revision already consumed by the peer");
406            }
407            let after = self
408                .storage
409                .read_versioned_at(&challenge.conditional_slot)
410                .await
411                .map_err(StorageError::from)?;
412            if after != before {
413                return invalid("rejected administrator update changed the exact peer readback");
414            }
415            let conditional = ConditionalUpdateProbeReceipt {
416                logical_key: cross_conditional_logical_key(challenge.probe_id),
417                slot: challenge.conditional_slot.clone(),
418                starting_payload_hash: ObjectHash::digest(&probe_payload(
419                    &challenge.probe_id,
420                    ProbePayloadLabel::ConditionalInitial,
421                )),
422                contenders: [
423                    ProbeConditionalAttempt {
424                        payload_hash: ObjectHash::digest(&peer_replacement),
425                        outcome: ProbeConditionalOutcome::Replaced,
426                    },
427                    ProbeConditionalAttempt {
428                        payload_hash: ObjectHash::digest(&administrator_replacement),
429                        outcome: ProbeConditionalOutcome::RejectedRevision,
430                    },
431                ],
432                accepted_payload_hash: ObjectHash::digest(&after.bytes),
433            };
434            advance_cross_completion(
435                journal,
436                &mut durable,
437                &mut record,
438                CrossPrincipalCompletionProgress::ReadsVerified {
439                    administrator_read_peer_hash: ObjectHash::digest(&observed),
440                    conditional,
441                },
442            )
443            .await?;
444        }
445        let (administrator_read_peer_hash, conditional) =
446            cross_completion_evidence(&record.progress)?;
447        let conditional = conditional.clone();
448        if matches!(
449            record.progress,
450            CrossPrincipalCompletionProgress::ReadsVerified { .. }
451        ) {
452            delete_versioned_probe(self.storage.as_ref(), &challenge.conditional_slot).await?;
453            self.storage
454                .delete_and_verify_absent(&response.peer_object.slot)
455                .await
456                .map_err(StorageError::from)?;
457            advance_cross_completion(
458                journal,
459                &mut durable,
460                &mut record,
461                CrossPrincipalCompletionProgress::ResponseObjectsAbsent {
462                    administrator_read_peer_hash,
463                    conditional: conditional.clone(),
464                },
465            )
466            .await?;
467        }
468        if matches!(
469            record.progress,
470            CrossPrincipalCompletionProgress::ResponseObjectsAbsent { .. }
471        ) {
472            self.storage
473                .delete_and_verify_absent(&challenge.administrator_object.slot)
474                .await
475                .map_err(StorageError::from)?;
476            advance_cross_completion(
477                journal,
478                &mut durable,
479                &mut record,
480                CrossPrincipalCompletionProgress::Absent {
481                    administrator_read_peer_hash,
482                    conditional: conditional.clone(),
483                },
484            )
485            .await?;
486        }
487        let transcript = CrossPrincipalProbeTranscript {
488            challenge: challenge.clone(),
489            response: response.clone(),
490            administrator_read_peer_hash,
491            conditional,
492        };
493        let receipt =
494            CrossPrincipalProbeReceipt::signed(transcript, context, store, administrator_signer)?;
495        advance_cross_completion(
496            journal,
497            &mut durable,
498            &mut record,
499            CrossPrincipalCompletionProgress::ReceiptReady {
500                receipt: receipt.clone(),
501            },
502        )
503        .await?;
504        Ok(receipt)
505    }
506
507    async fn settle_exact_create(
508        &self,
509        slot: &ObjectSlot,
510        payload: &[u8],
511    ) -> Result<(), ProviderProbeError> {
512        match self.storage.read_at(slot).await {
513            Ok(bytes) if bytes == payload => Ok(()),
514            Ok(_) => invalid("durable provider probe slot contains different bytes"),
515            Err(CloudHomeError::NotFound(_)) => {
516                create_exact_bytes(self.storage.as_ref(), slot, payload)
517                    .await
518                    .map(drop)
519                    .map_err(StorageError::from)
520                    .map_err(ProviderProbeError::Storage)
521            }
522            Err(error) => Err(ProviderProbeError::Storage(StorageError::from(error))),
523        }
524    }
525}
526
527fn cross_completion_evidence(
528    progress: &CrossPrincipalCompletionProgress,
529) -> Result<(ObjectHash, &ConditionalUpdateProbeReceipt), ProviderProbeError> {
530    match progress {
531        CrossPrincipalCompletionProgress::ReadsVerified {
532            administrator_read_peer_hash,
533            conditional,
534        }
535        | CrossPrincipalCompletionProgress::ResponseObjectsAbsent {
536            administrator_read_peer_hash,
537            conditional,
538        }
539        | CrossPrincipalCompletionProgress::Absent {
540            administrator_read_peer_hash,
541            conditional,
542        } => Ok((*administrator_read_peer_hash, conditional)),
543        CrossPrincipalCompletionProgress::Prepared
544        | CrossPrincipalCompletionProgress::ReceiptReady { .. } => {
545            invalid("cross-principal completion has no durable administrator read")
546        }
547    }
548}
549
550async fn create_exact_bytes(
551    storage: &dyn ExactSlotStorage,
552    slot: &ObjectSlot,
553    bytes: &[u8],
554) -> Result<ExactCreateOutcome, CloudHomeError> {
555    let object = ExactObjectRef::new(slot.clone(), bytes.len() as u64, ObjectHash::digest(bytes));
556    let upload = ExactUpload::from_bytes(&object, bytes).map_err(CloudHomeError::from)?;
557    storage
558        .create_at(
559            &upload,
560            &crate::cloud::UploadControl::running(crate::cloud::no_progress()),
561        )
562        .await
563}
564
565async fn create_versioned_bytes(
566    storage: &dyn ExactSlotStorage,
567    slot: &ObjectSlot,
568    bytes: &[u8],
569) -> Result<(), ProviderProbeError> {
570    let object = ExactObjectRef::new(slot.clone(), bytes.len() as u64, ObjectHash::digest(bytes));
571    let upload = ExactUpload::from_bytes(&object, bytes)?;
572    storage
573        .create_versioned_at(
574            &upload,
575            &crate::cloud::UploadControl::running(crate::cloud::no_progress()),
576        )
577        .await
578        .map(drop)
579        .map_err(StorageError::from)
580        .map_err(ProviderProbeError::Storage)
581}
582
583async fn delete_versioned_probe(
584    storage: &dyn ExactSlotStorage,
585    slot: &ObjectSlot,
586) -> Result<(), ProviderProbeError> {
587    storage
588        .delete_versioned_at(slot)
589        .await
590        .map_err(StorageError::from)?;
591    match storage.read_versioned_at(slot).await {
592        Err(CloudHomeError::NotFound(_)) => Ok(()),
593        Ok(_) => invalid("conditional-update probe record remains after deletion"),
594        Err(error) => Err(StorageError::from(error).into()),
595    }
596}
597
598#[cfg(test)]
599#[path = "provider_probe_tests.rs"]
600mod tests;
601
602#[cfg(test)]
603#[path = "provider_probe_s3_tests.rs"]
604mod s3_tests;