Skip to main content

coven_storage/provider_probe/
exact_slots.rs

1use super::*;
2
3impl ProviderProbeStorage {
4    pub async fn probe_exact_slots(
5        &self,
6        journal: &dyn ProviderProbeJournal,
7        probe_id: ProviderProbeId,
8        binding: &coven_protocol::objects::ResolvedProviderBinding,
9    ) -> Result<ExactSlotProbeReceipt, ProviderProbeError> {
10        let first = self.storage.as_ref();
11        let second = self.storage.as_ref();
12        binding.validate().map_err(ProviderProbeError::Storage)?;
13        let first_binding = first.provider_binding().await.map_err(StorageError::from)?;
14        let second_binding = second
15            .provider_binding()
16            .await
17            .map_err(StorageError::from)?;
18        if first_binding != *binding || second_binding != *binding {
19            return invalid("exact-slot probe clients do not match the receipt binding");
20        }
21        let id = hex::encode(probe_id.as_bytes());
22        let logical_key = format!("__coven_probe__/exact/{id}");
23        let conditional_logical_key = format!("__coven_probe__/conditional/{id}");
24        let lost_logical_key = format!("__coven_probe__/lost-response/{id}");
25        let mut durable = match journal.load(probe_id).await? {
26            Some(existing) => existing,
27            None => {
28                let allocated_slot = first
29                    .allocate_slot(&logical_key)
30                    .await
31                    .map_err(StorageError::from)?;
32                let allocated_lost_slot = first
33                    .allocate_slot(&lost_logical_key)
34                    .await
35                    .map_err(StorageError::from)?;
36                let allocated_conditional_slot = first
37                    .allocate_slot(&conditional_logical_key)
38                    .await
39                    .map_err(StorageError::from)?;
40                journal
41                    .begin(ProviderProbeJournalRecord::Exact(ExactProbeJournal {
42                        probe_id,
43                        binding: binding.clone(),
44                        slot: allocated_slot,
45                        conditional_slot: allocated_conditional_slot,
46                        lost_response_slot: allocated_lost_slot,
47                        progress: ExactProbeProgress::Prepared,
48                    }))
49                    .await?
50            }
51        };
52        let ProviderProbeJournalRecord::Exact(mut record) = durable.clone() else {
53            return invalid("exact probe id belongs to a different durable probe kind");
54        };
55        if record.probe_id != probe_id || record.binding != *binding {
56            return invalid("durable exact probe differs from its requested binding or id");
57        }
58        let slot = record.slot.clone();
59        let conditional_slot = record.conditional_slot.clone();
60        let lost_slot = record.lost_response_slot.clone();
61        if slot.logical_key() != logical_key
62            || conditional_slot.logical_key() != conditional_logical_key
63            || lost_slot.logical_key() != lost_logical_key
64        {
65            return invalid("exact-slot allocator changed the probe logical key");
66        }
67        let payloads = [
68            probe_payload(&probe_id, ProbePayloadLabel::ExactCreateFirst),
69            probe_payload(&probe_id, ProbePayloadLabel::ExactCreateSecond),
70        ];
71        if matches!(record.progress, ExactProbeProgress::Prepared) {
72            let (outcomes, _winner) = match first.read_at(&slot).await {
73                Err(CloudHomeError::NotFound(_)) => {
74                    let (left, right) = tokio::join!(
75                        create_exact_bytes(first, &slot, &payloads[0]),
76                        create_exact_bytes(second, &slot, &payloads[1]),
77                    );
78                    classify_exact_create_race(left, right)?
79                }
80                Ok(bytes) if bytes == payloads[0] => {
81                    require_occupied_rejection(
82                        create_exact_bytes(second, &slot, &payloads[1]).await,
83                    )?;
84                    (
85                        [
86                            ProbeCreateOutcome::Created,
87                            ProbeCreateOutcome::RejectedOccupied,
88                        ],
89                        0,
90                    )
91                }
92                Ok(bytes) if bytes == payloads[1] => {
93                    require_occupied_rejection(
94                        create_exact_bytes(first, &slot, &payloads[0]).await,
95                    )?;
96                    (
97                        [
98                            ProbeCreateOutcome::RejectedOccupied,
99                            ProbeCreateOutcome::Created,
100                        ],
101                        1,
102                    )
103                }
104                Ok(_) => return invalid("durable exact probe slot contains unknown bytes"),
105                Err(error) => return Err(ProviderProbeError::Storage(StorageError::from(error))),
106            };
107            advance_exact(
108                journal,
109                &mut durable,
110                &mut record,
111                ExactProbeProgress::Created { outcomes },
112            )
113            .await?;
114        }
115        let (outcomes, winner) = exact_race_state(&record.progress)?;
116        let (full, range) = if matches!(record.progress, ExactProbeProgress::Created { .. }) {
117            let full = first.read_at(&slot).await.map_err(StorageError::from)?;
118            if full != payloads[winner] {
119                return invalid("authoritative exact read does not match the create winner");
120            }
121            let range = first
122                .read_range_at(&slot, PROBE_RANGE_START, PROBE_RANGE_END)
123                .await
124                .map_err(StorageError::from)?;
125            if range != full[PROBE_RANGE_START as usize..PROBE_RANGE_END as usize] {
126                return invalid("exact range read does not match the authoritative full read");
127            }
128            (full, range)
129        } else {
130            (
131                payloads[winner].clone(),
132                payloads[winner][PROBE_RANGE_START as usize..PROBE_RANGE_END as usize].to_vec(),
133            )
134        };
135        if matches!(record.progress, ExactProbeProgress::Created { .. }) {
136            advance_exact(
137                journal,
138                &mut durable,
139                &mut record,
140                ExactProbeProgress::ReadsVerified { outcomes },
141            )
142            .await?;
143        }
144        let accepted =
145            ExactObjectRef::new(slot.clone(), full.len() as u64, ObjectHash::digest(&full));
146        if matches!(record.progress, ExactProbeProgress::ReadsVerified { .. }) {
147            let initial = probe_payload(&probe_id, ProbePayloadLabel::ConditionalInitial);
148            let conditional_payloads = [
149                probe_payload(&probe_id, ProbePayloadLabel::ConditionalFirst),
150                probe_payload(&probe_id, ProbePayloadLabel::ConditionalSecond),
151            ];
152            let mut current = match first.read_versioned_at(&conditional_slot).await {
153                Ok(current) => current,
154                Err(CloudHomeError::NotFound(_)) => {
155                    create_versioned_bytes(first, &conditional_slot, &initial).await?;
156                    first
157                        .read_versioned_at(&conditional_slot)
158                        .await
159                        .map_err(StorageError::from)?
160                }
161                Err(error) => return Err(ProviderProbeError::Storage(StorageError::from(error))),
162            };
163            if current.bytes != initial
164                && current.bytes != conditional_payloads[0]
165                && current.bytes != conditional_payloads[1]
166            {
167                return invalid("conditional-update probe slot contains unknown bytes");
168            }
169            let starting_payload_hash = ObjectHash::digest(&current.bytes);
170            let expected = current.version.clone();
171            let (left, right) = tokio::join!(
172                first.replace_at_if_version(
173                    &conditional_slot,
174                    &expected,
175                    conditional_payloads[0].clone(),
176                ),
177                second.replace_at_if_version(
178                    &conditional_slot,
179                    &expected,
180                    conditional_payloads[1].clone(),
181                ),
182            );
183            let (conditional_outcomes, conditional_winner) =
184                classify_conditional_update_race(left, right)?;
185            current = first
186                .read_versioned_at(&conditional_slot)
187                .await
188                .map_err(StorageError::from)?;
189            if current.bytes != conditional_payloads[conditional_winner] {
190                return invalid("conditional-update readback does not match its winning write");
191            }
192            let conditional = ConditionalUpdateProbeReceipt {
193                logical_key: conditional_logical_key.clone(),
194                slot: conditional_slot.clone(),
195                starting_payload_hash,
196                contenders: [
197                    ProbeConditionalAttempt {
198                        payload_hash: ObjectHash::digest(&conditional_payloads[0]),
199                        outcome: conditional_outcomes[0],
200                    },
201                    ProbeConditionalAttempt {
202                        payload_hash: ObjectHash::digest(&conditional_payloads[1]),
203                        outcome: conditional_outcomes[1],
204                    },
205                ],
206                accepted_payload_hash: ObjectHash::digest(&current.bytes),
207            };
208            advance_exact(
209                journal,
210                &mut durable,
211                &mut record,
212                ExactProbeProgress::ConditionalVerified {
213                    outcomes,
214                    conditional,
215                },
216            )
217            .await?;
218        }
219        let conditional = exact_conditional_evidence(&record.progress)?.clone();
220        if matches!(
221            record.progress,
222            ExactProbeProgress::ConditionalVerified { .. }
223        ) {
224            first
225                .delete_and_verify_absent(&slot)
226                .await
227                .map_err(StorageError::from)?;
228            delete_versioned_probe(first, &conditional_slot).await?;
229            advance_exact(
230                journal,
231                &mut durable,
232                &mut record,
233                ExactProbeProgress::PrimaryAbsent {
234                    outcomes,
235                    conditional: conditional.clone(),
236                },
237            )
238            .await?;
239        }
240        let lost_payload = probe_payload(&probe_id, ProbePayloadLabel::LostResponse);
241        if matches!(record.progress, ExactProbeProgress::PrimaryAbsent { .. }) {
242            match first.read_at(&lost_slot).await {
243                Ok(bytes) if bytes == lost_payload => {}
244                Ok(_) => return invalid("lost-response slot contains unknown bytes"),
245                Err(CloudHomeError::NotFound(_)) => {
246                    create_exact_bytes(first, &lost_slot, &lost_payload)
247                        .await
248                        .map_err(StorageError::from)?;
249                }
250                Err(error) => return Err(ProviderProbeError::Storage(StorageError::from(error))),
251            }
252            advance_exact(
253                journal,
254                &mut durable,
255                &mut record,
256                ExactProbeProgress::LostResponseCreated {
257                    outcomes,
258                    conditional: conditional.clone(),
259                },
260            )
261            .await?;
262        }
263        let lost_readback = if matches!(
264            record.progress,
265            ExactProbeProgress::LostResponseCreated { .. }
266        ) {
267            let readback = first
268                .read_at(&lost_slot)
269                .await
270                .map_err(StorageError::from)?;
271            if readback != lost_payload {
272                return invalid(
273                    "lost-response authoritative readback differs from committed bytes",
274                );
275            }
276            readback
277        } else {
278            lost_payload.clone()
279        };
280        let settled = ExactObjectRef::new(
281            lost_slot.clone(),
282            lost_readback.len() as u64,
283            ObjectHash::digest(&lost_readback),
284        );
285        if matches!(
286            record.progress,
287            ExactProbeProgress::LostResponseCreated { .. }
288        ) {
289            advance_exact(
290                journal,
291                &mut durable,
292                &mut record,
293                ExactProbeProgress::LostResponseReadVerified {
294                    outcomes,
295                    conditional: conditional.clone(),
296                },
297            )
298            .await?;
299        }
300        if matches!(
301            record.progress,
302            ExactProbeProgress::LostResponseReadVerified { .. }
303        ) {
304            first
305                .delete_and_verify_absent(&lost_slot)
306                .await
307                .map_err(StorageError::from)?;
308            advance_exact(
309                journal,
310                &mut durable,
311                &mut record,
312                ExactProbeProgress::Absent {
313                    outcomes,
314                    conditional: conditional.clone(),
315                },
316            )
317            .await?;
318        }
319
320        if let ExactProbeProgress::ReceiptReady { receipt } = &record.progress {
321            receipt.verify(&binding.store, &binding.device)?;
322            return Ok(receipt.clone());
323        }
324        let transcript = ExactSlotProbeTranscript {
325            probe_id,
326            logical_key,
327            slot,
328            contenders: [
329                ProbeCreateAttempt {
330                    payload_hash: ObjectHash::digest(&payloads[0]),
331                    outcome: outcomes[0],
332                },
333                ProbeCreateAttempt {
334                    payload_hash: ObjectHash::digest(&payloads[1]),
335                    outcome: outcomes[1],
336                },
337            ],
338            accepted,
339            full_read_hash: ObjectHash::digest(&full),
340            range: ProbeRangeReceipt {
341                start: PROBE_RANGE_START,
342                end: PROBE_RANGE_END,
343                bytes_hash: ObjectHash::digest(&range),
344            },
345            conditional,
346            lost_response: LostResponseProbeReceipt {
347                logical_key: lost_logical_key,
348                slot: lost_slot,
349                payload_hash: ObjectHash::digest(&lost_payload),
350                settled,
351                readback_hash: ObjectHash::digest(&lost_readback),
352            },
353        };
354        let receipt =
355            ExactSlotProbeReceipt::from_transcript(transcript, &binding.store, &binding.device);
356        receipt.verify(&binding.store, &binding.device)?;
357        advance_exact(
358            journal,
359            &mut durable,
360            &mut record,
361            ExactProbeProgress::ReceiptReady {
362                receipt: receipt.clone(),
363            },
364        )
365        .await?;
366        Ok(receipt)
367    }
368}
369
370fn exact_race_state(
371    progress: &ExactProbeProgress,
372) -> Result<([ProbeCreateOutcome; 2], usize), ProviderProbeError> {
373    let (outcomes, winner) = match progress {
374        ExactProbeProgress::Prepared => return invalid("exact probe has no durable create result"),
375        ExactProbeProgress::Created { outcomes }
376        | ExactProbeProgress::ReadsVerified { outcomes }
377        | ExactProbeProgress::ConditionalVerified { outcomes, .. }
378        | ExactProbeProgress::PrimaryAbsent { outcomes, .. }
379        | ExactProbeProgress::LostResponseCreated { outcomes, .. }
380        | ExactProbeProgress::LostResponseReadVerified { outcomes, .. }
381        | ExactProbeProgress::Absent { outcomes, .. } => {
382            let winner = outcomes
383                .iter()
384                .position(|outcome| *outcome == ProbeCreateOutcome::Created)
385                .ok_or_else(|| {
386                    ProviderProbeError::InvalidReceipt(
387                        "durable exact probe has no create winner".to_string(),
388                    )
389                })?;
390            (*outcomes, winner)
391        }
392        ExactProbeProgress::ReceiptReady { receipt } => {
393            let winner = receipt
394                .transcript
395                .contenders
396                .iter()
397                .position(|attempt| attempt.outcome == ProbeCreateOutcome::Created)
398                .ok_or_else(|| {
399                    ProviderProbeError::InvalidReceipt(
400                        "durable exact receipt has no create winner".to_string(),
401                    )
402                })?;
403            (
404                [
405                    receipt.transcript.contenders[0].outcome,
406                    receipt.transcript.contenders[1].outcome,
407                ],
408                winner,
409            )
410        }
411    };
412    if winner > 1 || outcomes[winner] != ProbeCreateOutcome::Created {
413        return invalid("durable exact probe has an invalid winner");
414    }
415    Ok((outcomes, winner))
416}
417
418fn exact_conditional_evidence(
419    progress: &ExactProbeProgress,
420) -> Result<&ConditionalUpdateProbeReceipt, ProviderProbeError> {
421    match progress {
422        ExactProbeProgress::ConditionalVerified { conditional, .. }
423        | ExactProbeProgress::PrimaryAbsent { conditional, .. }
424        | ExactProbeProgress::LostResponseCreated { conditional, .. }
425        | ExactProbeProgress::LostResponseReadVerified { conditional, .. }
426        | ExactProbeProgress::Absent { conditional, .. } => Ok(conditional),
427        ExactProbeProgress::ReceiptReady { receipt } => Ok(&receipt.transcript.conditional),
428        ExactProbeProgress::Prepared
429        | ExactProbeProgress::Created { .. }
430        | ExactProbeProgress::ReadsVerified { .. } => {
431            invalid("exact probe has no conditional-update evidence")
432        }
433    }
434}
435
436fn classify_conditional_update_race(
437    left: Result<ConditionalWriteOutcome, CloudHomeError>,
438    right: Result<ConditionalWriteOutcome, CloudHomeError>,
439) -> Result<([ProbeConditionalOutcome; 2], usize), ProviderProbeError> {
440    match (left, right) {
441        (
442            Ok(ConditionalWriteOutcome::Replaced(_)),
443            Ok(ConditionalWriteOutcome::VersionChanged),
444        ) => Ok((
445            [
446                ProbeConditionalOutcome::Replaced,
447                ProbeConditionalOutcome::RejectedRevision,
448            ],
449            0,
450        )),
451        (
452            Ok(ConditionalWriteOutcome::VersionChanged),
453            Ok(ConditionalWriteOutcome::Replaced(_)),
454        ) => Ok((
455            [
456                ProbeConditionalOutcome::RejectedRevision,
457                ProbeConditionalOutcome::Replaced,
458            ],
459            1,
460        )),
461        (left, right) => invalid(&format!(
462            "conditional-update race did not produce one replacement and one revision rejection: left={left:?}, right={right:?}"
463        )),
464    }
465}
466
467fn classify_exact_create_race(
468    left: Result<ExactCreateOutcome, CloudHomeError>,
469    right: Result<ExactCreateOutcome, CloudHomeError>,
470) -> Result<([ProbeCreateOutcome; 2], usize), ProviderProbeError> {
471    match (left, right) {
472        (
473            Ok(ExactCreateOutcome::Created),
474            Err(CloudHomeError::SlotCollision(_) | CloudHomeError::AlreadyExists(_)),
475        ) => Ok((
476            [
477                ProbeCreateOutcome::Created,
478                ProbeCreateOutcome::RejectedOccupied,
479            ],
480            0,
481        )),
482        (
483            Err(CloudHomeError::SlotCollision(_) | CloudHomeError::AlreadyExists(_)),
484            Ok(ExactCreateOutcome::Created),
485        ) => Ok((
486            [
487                ProbeCreateOutcome::RejectedOccupied,
488                ProbeCreateOutcome::Created,
489            ],
490            1,
491        )),
492        (left, right) => invalid(&format!(
493            "exact-slot race did not produce one create and one occupied rejection: left={left:?}, right={right:?}"
494        )),
495    }
496}
497
498fn require_occupied_rejection(
499    result: Result<ExactCreateOutcome, CloudHomeError>,
500) -> Result<(), ProviderProbeError> {
501    match result {
502        Err(CloudHomeError::SlotCollision(_) | CloudHomeError::AlreadyExists(_)) => Ok(()),
503        Ok(ExactCreateOutcome::Created) => {
504            invalid("settled exact probe contender unexpectedly created a second object")
505        }
506        result => invalid(&format!(
507            "settled exact probe contender was not rejected as occupied: result={result:?}"
508        )),
509    }
510}