Skip to main content

coven_storage/cloud/cloudkit/
exact.rs

1use super::chunking::*;
2use super::*;
3
4const EXACT_MANIFEST_MAGIC: &[u8] = b"coven-cloudkit-exact-manifest-v2\0";
5
6#[derive(Clone, Copy, Debug, PartialEq, Eq)]
7pub(crate) struct ExactManifest {
8    pub(crate) part_count: usize,
9    pub(crate) total_len: usize,
10    pub(crate) stored_hash: coven_protocol::store_commit::ObjectHash,
11}
12
13pub(crate) fn exact_part_key(logical_key: &str, index: usize) -> String {
14    format!("{logical_key}.exact-part{index}")
15}
16
17pub(crate) fn encode_exact_manifest(manifest: ExactManifest) -> Vec<u8> {
18    let mut bytes = EXACT_MANIFEST_MAGIC.to_vec();
19    bytes.extend_from_slice(manifest.part_count.to_string().as_bytes());
20    bytes.push(b'\n');
21    bytes.extend_from_slice(manifest.total_len.to_string().as_bytes());
22    bytes.push(b'\n');
23    bytes.extend_from_slice(manifest.stored_hash.to_string().as_bytes());
24    bytes.push(b'\n');
25    bytes
26}
27
28pub(crate) fn decode_exact_manifest(bytes: &[u8]) -> Result<ExactManifest, CloudHomeError> {
29    let text = std::str::from_utf8(bytes.strip_prefix(EXACT_MANIFEST_MAGIC).ok_or_else(|| {
30        CloudHomeError::Transport("CloudKit exact object has an invalid manifest".to_string())
31    })?)
32    .map_err(|error| CloudHomeError::transport("CloudKit exact manifest".to_string(), error))?;
33    let mut lines = text.lines();
34    let part_count = lines
35        .next()
36        .ok_or_else(|| {
37            CloudHomeError::Transport("CloudKit exact manifest omitted part count".to_string())
38        })?
39        .parse::<usize>()
40        .map_err(|error| {
41            CloudHomeError::transport("CloudKit exact manifest part count".to_string(), error)
42        })?;
43    let total_len = lines
44        .next()
45        .ok_or_else(|| {
46            CloudHomeError::Transport("CloudKit exact manifest omitted length".to_string())
47        })?
48        .parse::<usize>()
49        .map_err(|error| {
50            CloudHomeError::transport("CloudKit exact manifest length".to_string(), error)
51        })?;
52    let stored_hash = lines
53        .next()
54        .ok_or_else(|| {
55            CloudHomeError::Transport("CloudKit exact manifest omitted stored hash".to_string())
56        })?
57        .parse()
58        .map_err(|error| {
59            CloudHomeError::transport("CloudKit exact manifest stored hash".to_string(), error)
60        })?;
61    if lines.next().is_some() || part_count != total_len.div_ceil(CHUNK_SIZE) {
62        return Err(CloudHomeError::Transport(
63            "CloudKit exact manifest shape does not match its length".to_string(),
64        ));
65    }
66    Ok(ExactManifest {
67        part_count,
68        total_len,
69        stored_hash,
70    })
71}
72
73pub(crate) fn read_exact_cloudkit_object(
74    ops: &dyn CloudKitOps,
75    scope: &CloudKitScope,
76    logical_key: &str,
77) -> Result<(Vec<u8>, Vec<CloudKitRecordVersion>), CloudHomeError> {
78    let manifest = ops.read_versioned_record(scope, logical_key)?;
79    let manifest_data = decode_exact_manifest(&manifest.bytes)?;
80    let mut bytes = Vec::with_capacity(manifest_data.total_len);
81    let mut records = Vec::with_capacity(manifest_data.part_count + 1);
82    records.push(CloudKitRecordVersion {
83        key: logical_key.to_string(),
84        version: manifest.version,
85    });
86    for index in 0..manifest_data.part_count {
87        let key = exact_part_key(logical_key, index);
88        let part = read_exact_part(
89            ops,
90            scope,
91            logical_key,
92            manifest_data.part_count,
93            manifest_data.total_len,
94            index,
95            &key,
96        )?;
97        bytes.extend_from_slice(&part.bytes);
98        records.push(CloudKitRecordVersion {
99            key,
100            version: part.version,
101        });
102    }
103    Ok((bytes, records))
104}
105
106/// The plaintext length part `index` of an exact object carries: a full chunk
107/// for every part but the last, which holds the remainder.
108pub(crate) fn exact_part_len(part_count: usize, total_len: usize, index: usize) -> usize {
109    if index + 1 == part_count {
110        total_len - index * CHUNK_SIZE
111    } else {
112        CHUNK_SIZE
113    }
114}
115
116/// Read one part record and refuse a length its manifest does not assign it. A
117/// part that is short is not the part the manifest describes, so splicing it
118/// would silently serve the wrong bytes at every later offset.
119pub(crate) fn read_exact_part(
120    ops: &dyn CloudKitOps,
121    scope: &CloudKitScope,
122    logical_key: &str,
123    part_count: usize,
124    total_len: usize,
125    index: usize,
126    key: &str,
127) -> Result<CloudVersionedObject, CloudHomeError> {
128    let part = ops.read_versioned_record(scope, key)?;
129    let expected_len = exact_part_len(part_count, total_len, index);
130    if part.bytes.len() != expected_len {
131        return Err(CloudHomeError::Transport(format!(
132            "CloudKit exact object {logical_key:?} part {index} has {} bytes, expected {expected_len}",
133            part.bytes.len()
134        )));
135    }
136    Ok(part)
137}
138
139/// Read one byte range of an exact CloudKit object, fetching only the part
140/// records that cover it.
141///
142/// The whole-object sibling is [`read_exact_cloudkit_object`]. Both read the
143/// same manifest, but this one never touches a part the range does not reach —
144/// which is what makes a ranged read of a blob cost the range. Reading the whole
145/// object and slicing would answer correctly and cost the object, so the
146/// caller's O(range) guarantee lives or dies here.
147pub(crate) fn read_exact_cloudkit_range(
148    ops: &dyn CloudKitOps,
149    scope: &CloudKitScope,
150    logical_key: &str,
151    start: usize,
152    end: usize,
153) -> Result<Vec<u8>, CloudHomeError> {
154    let manifest = ops.read_versioned_record(scope, logical_key)?;
155    let manifest_data = decode_exact_manifest(&manifest.bytes)?;
156    if end > manifest_data.total_len {
157        return Err(CloudHomeError::Transport(format!(
158            "range {start}..{end} exceeds CloudKit exact object {logical_key:?} size {}",
159            manifest_data.total_len
160        )));
161    }
162    let first = start / CHUNK_SIZE;
163    let last = (end - 1) / CHUNK_SIZE;
164    if last >= manifest_data.part_count {
165        return Err(CloudHomeError::Transport(format!(
166            "range {start}..{end} needs part {last} of CloudKit exact object {logical_key:?}, which has {}",
167            manifest_data.part_count
168        )));
169    }
170    let mut bytes = Vec::with_capacity(end - start);
171    for index in first..=last {
172        let key = exact_part_key(logical_key, index);
173        let part = read_exact_part(
174            ops,
175            scope,
176            logical_key,
177            manifest_data.part_count,
178            manifest_data.total_len,
179            index,
180            &key,
181        )?;
182        let part_start = index * CHUNK_SIZE;
183        let from = start.saturating_sub(part_start);
184        let to = (end - part_start).min(part.bytes.len());
185        bytes.extend_from_slice(&part.bytes[from..to]);
186    }
187    Ok(bytes)
188}
189
190#[async_trait]
191impl ExactSlotStorage for CloudKitCloudHome {
192    async fn provider_binding(
193        &self,
194    ) -> Result<coven_protocol::objects::ResolvedProviderBinding, CloudHomeError> {
195        use coven_protocol::objects::{
196            ProviderDeviceBinding, ProviderPrincipalId, ResolvedProviderBinding,
197            StoreProviderBinding,
198        };
199
200        let ops = self.ops.clone();
201        let scope = self.scope.clone();
202        let identity = blocking(move || ops.provider_identity(&scope)).await?;
203        if identity.container_id.is_empty()
204            || identity.owner_name.is_empty()
205            || identity.zone_name.is_empty()
206            || identity.current_user_record_name.is_empty()
207        {
208            return Err(CloudHomeError::Configuration(
209                "CloudKit provider identity contains an empty stable identifier".to_string(),
210            ));
211        }
212        if let CloudKitScope::Shared {
213            owner_name,
214            zone_name,
215        } = &self.scope
216        {
217            if owner_name != &identity.owner_name || zone_name != &identity.zone_name {
218                return Err(CloudHomeError::Configuration(format!(
219                    "CloudKit provider identity resolved zone {}/{}, expected {owner_name}/{zone_name}",
220                    identity.owner_name, identity.zone_name
221                )));
222            }
223        }
224        let principal = match &self.scope {
225            CloudKitScope::Private => ProviderPrincipalId::CloudKitPrivateZoneOwner {
226                record_name: identity.current_user_record_name,
227            },
228            CloudKitScope::Shared { .. } => ProviderPrincipalId::CloudKitSharedZoneParticipant {
229                record_name: identity.current_user_record_name,
230            },
231        };
232        Ok(ResolvedProviderBinding {
233            store: StoreProviderBinding::CloudKit {
234                container_id: identity.container_id,
235                environment: identity.environment,
236                owner_name: identity.owner_name,
237                zone_name: identity.zone_name,
238            },
239            device: ProviderDeviceBinding { principal },
240        })
241    }
242
243    async fn cross_principal_evidence(
244        &self,
245    ) -> Result<coven_protocol::provider::CrossPrincipalProviderEvidence, CloudHomeError> {
246        use coven_protocol::provider::{CloudKitAcceptedShare, CrossPrincipalProviderEvidence};
247        use coven_protocol::store_commit::ObjectHash;
248
249        let CloudKitScope::Shared {
250            owner_name,
251            zone_name,
252        } = &self.scope
253        else {
254            return Err(CloudHomeError::Configuration(
255                "CloudKit cross-principal evidence requires an accepted shared zone".to_string(),
256            ));
257        };
258        let ops = self.ops.clone();
259        let scope = self.scope.clone();
260        let accepted = blocking(move || ops.accepted_read_write_share(&scope)).await?;
261        let binding = self.provider_binding().await?;
262        let coven_protocol::objects::ProviderPrincipalId::CloudKitSharedZoneParticipant {
263            record_name,
264        } = binding.device.principal
265        else {
266            return Err(CloudHomeError::Configuration(
267                "CloudKit adapter returned a non-CloudKit principal".to_string(),
268            ));
269        };
270        if accepted.share_record_name.is_empty()
271            || accepted.owner_name != *owner_name
272            || accepted.zone_name != *zone_name
273            || accepted.participant_record_name != record_name
274            || accepted.permission != CloudKitSharePermission::ReadWrite
275            || accepted.acceptance != CloudKitShareAcceptance::Accepted
276            || accepted.canonical_record.is_empty()
277        {
278            return Err(CloudHomeError::Configuration(
279                "CloudKit accepted share does not prove read-write participation in the selected zone"
280                    .to_string(),
281            ));
282        }
283        let share_slot = ObjectSlot::logical(format!(
284            "__coven_cloudkit_share__/{}",
285            hex::encode(ObjectHash::digest(accepted.share_record_name.as_bytes()).as_bytes())
286        ))?;
287        Ok(CrossPrincipalProviderEvidence::CloudKit(
288            CloudKitAcceptedShare {
289                share: coven_protocol::objects::ExactObjectRef::new(
290                    share_slot,
291                    accepted.canonical_record.len() as u64,
292                    ObjectHash::digest(&accepted.canonical_record),
293                ),
294                share_record_name: accepted.share_record_name,
295                owner_name: accepted.owner_name,
296                zone_name: accepted.zone_name,
297                participant_record_name: accepted.participant_record_name,
298            },
299        ))
300    }
301
302    async fn create_at(
303        &self,
304        upload: &crate::cloud::ExactUpload<'_>,
305        control: &crate::cloud::UploadControl,
306    ) -> Result<crate::cloud::ExactCreateOutcome, CloudHomeError> {
307        if matches!(
308            self.exact_upload_verification,
309            coven_foundation::config::ExactUploadVerification::UploadChecksum
310        ) {
311            return Err(CloudHomeError::Configuration(
312                "CloudKit does not accept a caller-supplied upload checksum".to_string(),
313            ));
314        }
315        let slot = upload.object().slot();
316        let mut body = upload.body().await?;
317        slot.require_logical_key_for("CloudKit")?;
318        let total_len = usize::try_from(body.len()).map_err(|_| {
319            CloudHomeError::Transport(format!(
320                "CloudKit object {:?} is too large for this platform",
321                slot.logical_key()
322            ))
323        })?;
324        let part_count = total_len.div_ceil(CHUNK_SIZE);
325        let staging = self.begin_atomic_create().await?;
326        let mut requested_keys = Vec::with_capacity(part_count + 1);
327        let mut written_len = 0usize;
328        for index in 0..part_count {
329            control.wait_until_resumed().await;
330            let part = match body.next_part(CHUNK_SIZE).await {
331                Ok(Some(part)) => part,
332                Ok(None) => {
333                    return Err(staging.cleanup_failure(CloudHomeError::Transport(format!(
334                        "CloudKit object {:?} ended after {written_len} of {total_len} bytes",
335                        slot.logical_key()
336                    ))))
337                }
338                Err(error) => return Err(staging.cleanup_failure(error)),
339            };
340            written_len += part.len();
341            let key = exact_part_key(slot.logical_key(), index);
342            if let Err(error) = staging
343                .clone()
344                .stage_record(CloudKitRecordCreate {
345                    key: key.clone(),
346                    data: part.to_vec(),
347                })
348                .await
349            {
350                return Err(staging.cleanup_failure(error));
351            }
352            requested_keys.push(key);
353        }
354        match body.next_part(CHUNK_SIZE).await {
355            Ok(None) if written_len == total_len => {}
356            Ok(None) => {
357                return Err(staging.cleanup_failure(CloudHomeError::Transport(format!(
358                    "CloudKit object {:?} yielded {written_len} bytes, expected {total_len}",
359                    slot.logical_key()
360                ))))
361            }
362            Ok(Some(extra)) => {
363                return Err(staging.cleanup_failure(CloudHomeError::Transport(format!(
364                    "CloudKit object {:?} yielded at least {} bytes, expected {total_len}",
365                    slot.logical_key(),
366                    written_len + extra.len()
367                ))))
368            }
369            Err(error) => return Err(staging.cleanup_failure(error)),
370        }
371        if let Err(error) = staging
372            .clone()
373            .stage_record(CloudKitRecordCreate {
374                key: slot.logical_key().to_string(),
375                data: encode_exact_manifest(ExactManifest {
376                    part_count,
377                    total_len,
378                    stored_hash: upload.object().stored_hash(),
379                }),
380            })
381            .await
382        {
383            return Err(staging.cleanup_failure(error));
384        }
385        requested_keys.push(slot.logical_key().to_string());
386        let outcome = match staging.clone().commit().await {
387            Ok(created) => {
388                if created.len() != requested_keys.len()
389                    || created
390                        .iter()
391                        .zip(&requested_keys)
392                        .any(|(record, requested)| &record.key != requested)
393                {
394                    self.exact_manifest(slot).await?;
395                }
396                crate::cloud::ExactCreateOutcome::Created
397            }
398            Err(CloudHomeError::AlreadyExists(_)) => {
399                let collision = staging.cleanup_failure(CloudHomeError::AlreadyExists(
400                    slot.logical_key().to_string(),
401                ));
402                if !matches!(collision, CloudHomeError::AlreadyExists(_)) {
403                    return Err(collision);
404                }
405                return match self.verify_exact_upload(upload, false).await {
406                    Ok(()) => Ok(crate::cloud::ExactCreateOutcome::AlreadyPresent),
407                    Err(CloudHomeError::NotFound(_)) => Err(collision),
408                    Err(CloudHomeError::AlreadyExists(_)) => Err(collision),
409                    Err(slot_collision @ CloudHomeError::SlotCollision(_)) => Err(slot_collision),
410                    Err(settlement) => Err(CloudHomeError::UnresolvedOutcome {
411                        operation: Box::new(collision),
412                        settlement: Box::new(settlement),
413                    }),
414                };
415            }
416            Err(operation) => {
417                match self
418                    .settle_atomic_create_response_loss(slot.logical_key().to_string())
419                    .await
420                {
421                    Ok(AtomicCreateReadback::Created) => {
422                        staging.disarm();
423                        crate::cloud::ExactCreateOutcome::AlreadyPresent
424                    }
425                    Ok(AtomicCreateReadback::Absent) => {
426                        return Err(staging.cleanup_failure(operation))
427                    }
428                    Err(readback) => {
429                        staging.disarm();
430                        return Err(CloudHomeError::UnresolvedOutcome {
431                            operation: Box::new(operation),
432                            settlement: Box::new(readback),
433                        });
434                    }
435                }
436            }
437        };
438        self.verify_exact_upload(
439            upload,
440            matches!(outcome, crate::cloud::ExactCreateOutcome::Created),
441        )
442        .await?;
443        control.report(total_len as u64);
444        Ok(outcome)
445    }
446
447    async fn list_slots(&self, prefix: &str) -> Result<Vec<ObjectSlot>, CloudHomeError> {
448        crate::cloud::logical_slots(CloudHome::list(self, prefix).await?)
449    }
450
451    async fn create_versioned_at(
452        &self,
453        upload: &crate::cloud::ExactUpload<'_>,
454        control: &crate::cloud::UploadControl,
455    ) -> Result<crate::cloud::ExactCreateOutcome, CloudHomeError> {
456        let slot = upload.object().slot();
457        slot.require_logical_key_for("CloudKit")?;
458        let bytes = upload.body().await?.collect().await?;
459        if bytes.len() > CHUNK_SIZE {
460            return Err(CloudHomeError::Configuration(format!(
461                "CloudKit versioned record {:?} has {} bytes, above the {CHUNK_SIZE}-byte record bound",
462                slot.logical_key(),
463                bytes.len()
464            )));
465        }
466        let staging = self.begin_atomic_create().await?;
467        if let Err(error) = staging
468            .clone()
469            .stage_record(CloudKitRecordCreate {
470                key: slot.logical_key().to_string(),
471                data: bytes.clone(),
472            })
473            .await
474        {
475            return Err(staging.cleanup_failure(error));
476        }
477        let outcome = match staging.clone().commit().await {
478            Ok(created) if created.len() == 1 && created[0].key == slot.logical_key() => {
479                crate::cloud::ExactCreateOutcome::Created
480            }
481            Ok(_) => {
482                let observed = self.read_versioned_at(slot).await?;
483                if observed.bytes != bytes {
484                    return Err(CloudHomeError::SlotCollision(
485                        slot.logical_key().to_string(),
486                    ));
487                }
488                crate::cloud::ExactCreateOutcome::Created
489            }
490            Err(CloudHomeError::AlreadyExists(_)) => {
491                let collision = staging.cleanup_failure(CloudHomeError::AlreadyExists(
492                    slot.logical_key().to_string(),
493                ));
494                if !matches!(collision, CloudHomeError::AlreadyExists(_)) {
495                    return Err(collision);
496                }
497                let observed = self.read_versioned_at(slot).await?;
498                if observed.bytes != bytes {
499                    return Err(CloudHomeError::SlotCollision(
500                        slot.logical_key().to_string(),
501                    ));
502                }
503                crate::cloud::ExactCreateOutcome::AlreadyPresent
504            }
505            Err(operation) => match self.read_versioned_at(slot).await {
506                Ok(observed) if observed.bytes == bytes => {
507                    staging.disarm();
508                    crate::cloud::ExactCreateOutcome::AlreadyPresent
509                }
510                Ok(_) => {
511                    staging.disarm();
512                    return Err(CloudHomeError::SlotCollision(
513                        slot.logical_key().to_string(),
514                    ));
515                }
516                Err(CloudHomeError::NotFound(_)) => {
517                    return Err(staging.cleanup_failure(operation));
518                }
519                Err(settlement) => {
520                    staging.disarm();
521                    return Err(CloudHomeError::UnresolvedOutcome {
522                        operation: Box::new(operation),
523                        settlement: Box::new(settlement),
524                    });
525                }
526            },
527        };
528        control.report(bytes.len() as u64);
529        Ok(outcome)
530    }
531
532    async fn read_at(&self, slot: &ObjectSlot) -> Result<Vec<u8>, CloudHomeError> {
533        slot.require_logical_key_for("CloudKit")?;
534        let ops = self.ops.clone();
535        let scope = self.scope.clone();
536        let logical_key = slot.logical_key().to_string();
537        blocking(move || {
538            read_exact_cloudkit_object(&*ops, &scope, &logical_key).map(|value| value.0)
539        })
540        .await
541    }
542
543    async fn read_versioned_at(
544        &self,
545        slot: &ObjectSlot,
546    ) -> Result<crate::cloud::CloudVersionedObject, CloudHomeError> {
547        slot.require_logical_key_for("CloudKit")?;
548        let ops = self.ops.clone();
549        let scope = self.scope.clone();
550        let logical_key = slot.logical_key().to_string();
551        blocking(move || ops.read_versioned_record(&scope, &logical_key)).await
552    }
553
554    async fn replace_at_if_version(
555        &self,
556        slot: &ObjectSlot,
557        expected: &crate::cloud::CloudObjectVersion,
558        bytes: Vec<u8>,
559    ) -> Result<crate::cloud::ConditionalWriteOutcome, CloudHomeError> {
560        slot.require_logical_key_for("CloudKit")?;
561        let ops = self.ops.clone();
562        let scope = self.scope.clone();
563        let logical_key = slot.logical_key().to_string();
564        let expected = expected.clone();
565        cancellation_safe_blocking(move || {
566            ops.replace_record_if_version(&scope, &logical_key, &expected, bytes)
567        })
568        .await?
569    }
570
571    async fn delete_versioned_at(&self, slot: &ObjectSlot) -> Result<(), CloudHomeError> {
572        slot.require_logical_key_for("CloudKit")?;
573        let ops = self.ops.clone();
574        let scope = self.scope.clone();
575        let logical_key = slot.logical_key().to_string();
576        cancellation_safe_blocking(move || {
577            match ops.delete_record(&scope, &logical_key) {
578                Ok(()) | Err(CloudHomeError::NotFound(_)) => {}
579                Err(error) => return Err(error),
580            }
581            if ops.record_exists(&scope, &logical_key)? {
582                return Err(CloudHomeError::Transport(format!(
583                    "CloudKit versioned record {logical_key:?} remains after deletion"
584                )));
585            }
586            Ok(())
587        })
588        .await?
589    }
590
591    async fn read_range_at(
592        &self,
593        slot: &ObjectSlot,
594        start: u64,
595        end: u64,
596    ) -> Result<Vec<u8>, CloudHomeError> {
597        slot.require_logical_key_for("CloudKit")?;
598        let start = usize::try_from(start)
599            .map_err(|_| CloudHomeError::Configuration("range start is too large".to_string()))?;
600        let end = usize::try_from(end)
601            .map_err(|_| CloudHomeError::Configuration("range end is too large".to_string()))?;
602        if end < start {
603            return Err(CloudHomeError::Configuration(format!(
604                "invalid range {start}..{end}"
605            )));
606        }
607        if end == start {
608            return Ok(Vec::new());
609        }
610        let ops = self.ops.clone();
611        let scope = self.scope.clone();
612        let logical_key = slot.logical_key().to_string();
613        blocking(move || read_exact_cloudkit_range(&*ops, &scope, &logical_key, start, end)).await
614    }
615
616    async fn read_at_to_file(
617        &self,
618        slot: &ObjectSlot,
619        destination: &std::path::Path,
620        progress: crate::cloud::DownloadProgress,
621    ) -> Result<(), crate::cloud::CloudFileReadError> {
622        let bytes = self.read_at(slot).await?;
623        let stream: crate::cloud::CloudObjectStream = Box::pin(futures_util::stream::once(
624            async move { Ok(Bytes::from(bytes)) },
625        ));
626        crate::cloud::write_cloud_object_stream(destination, stream, progress)
627            .await
628            .map(drop)
629    }
630
631    async fn delete_at(&self, slot: &ObjectSlot) -> Result<(), CloudHomeError> {
632        slot.require_logical_key_for("CloudKit")?;
633        let ops = self.ops.clone();
634        let scope = self.scope.clone();
635        let logical_key = slot.logical_key().to_string();
636        blocking(move || {
637            let records = match read_exact_cloudkit_object(&*ops, &scope, &logical_key) {
638                Ok((_, records)) => records,
639                Err(CloudHomeError::NotFound(_)) => return Ok(()),
640                Err(error) => return Err(error),
641            };
642            ops.delete_record_versions(&scope, &records)
643        })
644        .await
645    }
646}