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(¤t.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(¤t.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}