Skip to main content

coven/storage/cloud/
cloudkit.rs

1//! CloudKit-backed `CloudHome` implementation.
2//!
3//! CloudKit's CKAsset has a 50MB limit, so large files are split into 10MB
4//! chunks stored as tokened part records plus a manifest record.
5//!
6//! The `CloudKitOps` trait defines synchronous record operations implemented by
7//! a host bridge to its CloudKit driver. `CloudKitCloudHome` wraps these ops,
8//! adds chunking logic, and implements `CloudHome`.
9
10use std::collections::HashSet;
11use std::sync::Arc;
12
13use async_trait::async_trait;
14use bytes::Bytes;
15
16use crate::id_provider::{IdRef, UuidProvider};
17
18use super::{
19    BlobBody, CloudAccessOutcome, CloudAccessState, CloudHeadCreateError, CloudHeadReplaceError,
20    CloudHeadStorage, CloudHeadVersion, CloudHome, CloudHomeError, CloudHomeJoinInfo,
21    CloudVersionedHead, ExactSlotStorage, ObjectSlot, PhysicalObjectLocator, RevokeOutcome,
22    UploadProgress,
23};
24
25const CHUNK_SIZE: usize = 10 * 1024 * 1024; // 10MB
26const CHUNK_MANIFEST_MAGIC: &[u8] = b"coven-cloudkit-chunk-manifest-v1\0";
27const CHUNK_MANIFEST_SUFFIX: &str = ".manifest";
28
29/// Synchronous interface for raw CloudKit record operations.
30/// Implemented by a host bridge to its platform CloudKit driver.
31/// Methods block the calling thread while CloudKit async operations complete.
32pub trait CloudKitOps: Send + Sync {
33    /// Stable CloudKit namespace and principal facts for the selected zone.
34    fn provider_identity(
35        &self,
36        scope: &CloudKitScope,
37    ) -> Result<CloudKitProviderIdentity, CloudHomeError>;
38
39    /// Fetch the accepted CKShare for a shared scope and return its exact
40    /// canonical record bytes plus the participant facts verified by the host.
41    fn accepted_read_write_share(
42        &self,
43        scope: &CloudKitScope,
44    ) -> Result<CloudKitAcceptedShareRecord, CloudHomeError>;
45
46    fn write_record(
47        &self,
48        scope: &CloudKitScope,
49        key: &str,
50        data: Vec<u8>,
51    ) -> Result<(), CloudHomeError>;
52    fn read_record(&self, scope: &CloudKitScope, key: &str) -> Result<Vec<u8>, CloudHomeError>;
53    fn list_records(
54        &self,
55        scope: &CloudKitScope,
56        prefix: &str,
57    ) -> Result<Vec<String>, CloudHomeError>;
58    fn delete_record(&self, scope: &CloudKitScope, key: &str) -> Result<(), CloudHomeError>;
59    fn record_exists(&self, scope: &CloudKitScope, key: &str) -> Result<bool, CloudHomeError>;
60    /// Read the exact CKRecord and return its opaque `recordChangeTag` with the
61    /// bytes. A later replacement must mutate this fetched record.
62    fn read_versioned_record(
63        &self,
64        scope: &CloudKitScope,
65        key: &str,
66    ) -> Result<CloudVersionedHead, CloudHomeError>;
67    fn create_record(
68        &self,
69        scope: &CloudKitScope,
70        key: &str,
71        data: Vec<u8>,
72    ) -> Result<CloudVersionedHead, CloudHeadCreateError>;
73    /// Replace by saving the previously fetched CKRecord with
74    /// `ifServerRecordUnchanged`; copying its tag into a new CKRecord is invalid.
75    fn replace_record(
76        &self,
77        scope: &CloudKitScope,
78        key: &str,
79        expected: &CloudHeadVersion,
80        data: Vec<u8>,
81    ) -> Result<CloudVersionedHead, CloudHeadReplaceError>;
82    /// Open a host-owned local staging batch. Staging never creates CloudKit
83    /// records; the host keeps payloads in temporary CKAsset files until commit.
84    fn begin_atomic_create(
85        &self,
86        scope: &CloudKitScope,
87    ) -> Result<CloudKitAtomicCreateBatch, CloudHomeError>;
88    /// Stage one bounded record payload in the host-owned batch.
89    fn stage_atomic_create_record(
90        &self,
91        scope: &CloudKitScope,
92        batch: &CloudKitAtomicCreateBatch,
93        record: CloudKitRecordCreate,
94    ) -> Result<(), CloudHomeError>;
95    /// Create every staged record as one atomic custom-zone modification. Every
96    /// record uses CloudKit's create-only save policy. A known precommit failure
97    /// leaves no record present. If the commit response is lost, the whole batch
98    /// may be present; preserve those records so the caller can read back every
99    /// requested key and settle the outcome.
100    /// Returned versions follow staging order when the response is received.
101    fn commit_atomic_create(
102        &self,
103        scope: &CloudKitScope,
104        batch: &CloudKitAtomicCreateBatch,
105    ) -> Result<Vec<CloudKitRecordVersion>, CloudHomeError>;
106    /// Discard host-local staging without deleting any CloudKit records the batch
107    /// may have committed. This is idempotent. On failure, return an error naming
108    /// the batch; the caller surfaces it and does not hide or retry it.
109    fn discard_atomic_create(
110        &self,
111        scope: &CloudKitScope,
112        batch: &CloudKitAtomicCreateBatch,
113    ) -> Result<(), CloudHomeError>;
114    /// Delete exactly these fetched record versions as one CloudKit atomic zone
115    /// modification. A changed or missing record fails the whole deletion.
116    fn delete_record_versions(
117        &self,
118        scope: &CloudKitScope,
119        records: &[CloudKitRecordVersion],
120    ) -> Result<(), CloudHomeError>;
121    fn share_for_member(
122        &self,
123        member_pubkey: &str,
124    ) -> Result<Option<CloudKitShare>, CloudHomeError>;
125    fn grant_share(&self, member_pubkey: &str) -> Result<CloudKitShare, CloudHomeError>;
126    fn revoke_share(&self, member_pubkey: &str) -> Result<(), CloudHomeError>;
127    fn accept_share(&self, share_url: &str) -> Result<CloudKitShare, CloudHomeError>;
128}
129
130#[derive(Clone, Debug, PartialEq, Eq)]
131pub struct CloudKitRecordVersion {
132    pub key: String,
133    pub version: CloudHeadVersion,
134}
135
136#[derive(Clone, Debug, PartialEq, Eq)]
137pub struct CloudKitRecordCreate {
138    pub key: String,
139    pub data: Vec<u8>,
140}
141
142#[derive(Clone, Debug, PartialEq, Eq, Hash)]
143pub struct CloudKitAtomicCreateBatch(String);
144
145impl CloudKitAtomicCreateBatch {
146    pub fn from_provider(value: String) -> Result<Self, CloudHomeError> {
147        if value.is_empty() {
148            return Err(CloudHomeError::Transport(
149                "CloudKit returned an empty atomic-create batch id".to_string(),
150            ));
151        }
152        Ok(Self(value))
153    }
154
155    pub fn as_provider(&self) -> &str {
156        &self.0
157    }
158}
159
160#[derive(Clone, Debug, PartialEq, Eq, Hash)]
161pub enum CloudKitScope {
162    Private,
163    Shared {
164        owner_name: String,
165        zone_name: String,
166    },
167}
168
169#[derive(Clone, Debug, PartialEq, Eq)]
170pub struct CloudKitProviderIdentity {
171    pub container_id: String,
172    pub environment: coven_core::sync::storage::CloudKitEnvironment,
173    pub owner_name: String,
174    pub zone_name: String,
175    pub current_user_record_name: String,
176}
177
178#[derive(Clone, Debug, PartialEq, Eq)]
179pub struct CloudKitShare {
180    pub share_url: String,
181    pub owner_name: String,
182    pub zone_name: String,
183}
184
185#[derive(Clone, Debug, PartialEq, Eq)]
186pub struct CloudKitAcceptedShareRecord {
187    pub share_record_name: String,
188    pub owner_name: String,
189    pub zone_name: String,
190    pub participant_record_name: String,
191    pub permission: CloudKitSharePermission,
192    pub acceptance: CloudKitShareAcceptance,
193    pub canonical_record: Vec<u8>,
194}
195
196#[derive(Clone, Debug, PartialEq, Eq)]
197pub enum CloudKitSharePermission {
198    ReadOnly,
199    ReadWrite,
200}
201
202#[derive(Clone, Debug, PartialEq, Eq)]
203pub enum CloudKitShareAcceptance {
204    Pending,
205    Accepted,
206}
207
208/// CloudKit-backed cloud home with automatic chunking for large files.
209#[derive(Clone)]
210pub(crate) struct CloudKitCloudHome {
211    ops: Arc<dyn CloudKitOps>,
212    ids: IdRef,
213    scope: CloudKitScope,
214}
215
216#[async_trait]
217impl CloudHeadStorage for CloudKitCloudHome {
218    async fn read_head(&self, key: &str) -> Result<CloudVersionedHead, CloudHomeError> {
219        let ops = self.ops.clone();
220        let scope = self.scope.clone();
221        let key = key.to_string();
222        blocking(move || ops.read_versioned_record(&scope, &key)).await
223    }
224
225    async fn create_head(
226        &self,
227        key: &str,
228        bytes: Vec<u8>,
229    ) -> Result<CloudVersionedHead, CloudHeadCreateError> {
230        let ops = self.ops.clone();
231        let scope = self.scope.clone();
232        let key = key.to_string();
233        tokio::task::spawn_blocking(move || ops.create_record(&scope, &key, bytes))
234            .await
235            .map_err(|error| {
236                CloudHeadCreateError::Storage(CloudHomeError::Transport(format!(
237                    "spawn_blocking failed: {error}"
238                )))
239            })?
240    }
241
242    async fn replace_head(
243        &self,
244        key: &str,
245        expected: &CloudHeadVersion,
246        bytes: Vec<u8>,
247    ) -> Result<CloudVersionedHead, CloudHeadReplaceError> {
248        let ops = self.ops.clone();
249        let scope = self.scope.clone();
250        let key = key.to_string();
251        let expected = expected.clone();
252        tokio::task::spawn_blocking(move || ops.replace_record(&scope, &key, &expected, bytes))
253            .await
254            .map_err(|error| {
255                CloudHeadReplaceError::Storage(CloudHomeError::Transport(format!(
256                    "spawn_blocking failed: {error}"
257                )))
258            })?
259    }
260
261    async fn delete_probe_head(&self, key: &str) -> Result<(), CloudHomeError> {
262        let ops = self.ops.clone();
263        let scope = self.scope.clone();
264        let key = key.to_string();
265        blocking(move || ops.delete_record(&scope, &key)).await
266    }
267}
268
269impl CloudKitCloudHome {
270    pub(crate) fn new_private(ops: Arc<dyn CloudKitOps>) -> Self {
271        Self::new_private_with_ids(ops, Arc::new(UuidProvider))
272    }
273
274    pub(crate) fn new_private_with_ids(ops: Arc<dyn CloudKitOps>, ids: IdRef) -> Self {
275        Self {
276            ops,
277            ids,
278            scope: CloudKitScope::Private,
279        }
280    }
281
282    pub(crate) fn new_shared(
283        ops: Arc<dyn CloudKitOps>,
284        owner_name: String,
285        zone_name: String,
286    ) -> Self {
287        Self::new_shared_with_ids(ops, Arc::new(UuidProvider), owner_name, zone_name)
288    }
289
290    pub(crate) fn new_shared_with_ids(
291        ops: Arc<dyn CloudKitOps>,
292        ids: IdRef,
293        owner_name: String,
294        zone_name: String,
295    ) -> Self {
296        Self {
297            ops,
298            ids,
299            scope: CloudKitScope::Shared {
300                owner_name,
301                zone_name,
302            },
303        }
304    }
305}
306
307pub(crate) async fn accept_share(
308    ops: Arc<dyn CloudKitOps>,
309    share_url: String,
310) -> Result<CloudKitShare, CloudHomeError> {
311    blocking(move || ops.accept_share(&share_url)).await
312}
313
314/// Run a synchronous CloudKit op on the blocking pool, mapping a join failure to a
315/// storage error. The Swift bridge methods block, so every `CloudHome` method
316/// wraps its call this way — one helper instead of the same `spawn_blocking(...)
317/// .await.map_err(...)` in each.
318async fn blocking<T, F>(f: F) -> Result<T, CloudHomeError>
319where
320    F: FnOnce() -> Result<T, CloudHomeError> + Send + 'static,
321    T: Send + 'static,
322{
323    tokio::task::spawn_blocking(f)
324        .await
325        .map_err(|e| CloudHomeError::Transport(format!("spawn_blocking failed: {e}")))?
326}
327
328struct BlockingState<T> {
329    result: std::sync::Mutex<Option<std::thread::Result<T>>>,
330    ready: std::sync::Condvar,
331    notify: tokio::sync::Notify,
332}
333
334struct BlockingCompletion<T> {
335    state: Arc<BlockingState<T>>,
336    consumed: bool,
337}
338
339impl<T> Drop for BlockingCompletion<T> {
340    fn drop(&mut self) {
341        if self.consumed {
342            return;
343        }
344        let mut result = self.state.result.lock().expect("lock blocking result");
345        while result.is_none() {
346            result = self
347                .state
348                .ready
349                .wait(result)
350                .expect("wait for blocking result");
351        }
352    }
353}
354
355async fn cancellation_safe_blocking<T, F>(f: F) -> Result<T, CloudHomeError>
356where
357    F: FnOnce() -> T + Send + 'static,
358    T: Send + 'static,
359{
360    let state = Arc::new(BlockingState {
361        result: std::sync::Mutex::new(None),
362        ready: std::sync::Condvar::new(),
363        notify: tokio::sync::Notify::new(),
364    });
365    let worker_state = state.clone();
366    tokio::task::spawn_blocking(move || {
367        let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(f));
368        worker_state
369            .result
370            .lock()
371            .expect("lock blocking result")
372            .replace(result);
373        worker_state.ready.notify_all();
374        worker_state.notify.notify_one();
375    });
376    let mut completion = BlockingCompletion {
377        state,
378        consumed: false,
379    };
380    let result = loop {
381        let notified = completion.state.notify.notified();
382        if let Some(result) = completion
383            .state
384            .result
385            .lock()
386            .expect("lock blocking result")
387            .take()
388        {
389            break result;
390        }
391        notified.await;
392    };
393    completion.consumed = true;
394    result.map_err(|_| CloudHomeError::Transport("CloudKit blocking task panicked".to_string()))
395}
396
397/// If a key ends with CloudKit layout metadata, strip that suffix to get the
398/// base object key.
399fn strip_part_suffix(key: &str) -> &str {
400    if let Some(base) = key.strip_suffix(CHUNK_MANIFEST_SUFFIX) {
401        return base;
402    }
403    if let Some(idx) = key.rfind(".part") {
404        let after = &key[idx + 5..];
405        let digits = after
406            .as_bytes()
407            .iter()
408            .take_while(|b| b.is_ascii_digit())
409            .count();
410        if digits > 0 && (digits == after.len() || after.as_bytes()[digits] == b'.') {
411            return &key[..idx];
412        }
413    }
414    key
415}
416
417#[derive(Clone, Debug, PartialEq, Eq)]
418struct ChunkManifest {
419    part_count: usize,
420    total_len: usize,
421    upload_id: String,
422}
423
424impl ChunkManifest {
425    fn new(total_len: usize, upload_id: String) -> Self {
426        Self {
427            part_count: total_len.div_ceil(CHUNK_SIZE),
428            total_len,
429            upload_id,
430        }
431    }
432}
433
434fn encode_chunk_manifest(manifest: ChunkManifest) -> Vec<u8> {
435    let mut encoded = CHUNK_MANIFEST_MAGIC.to_vec();
436    encoded.extend_from_slice(manifest.part_count.to_string().as_bytes());
437    encoded.push(b'\n');
438    encoded.extend_from_slice(manifest.total_len.to_string().as_bytes());
439    encoded.push(b'\n');
440    encoded.extend_from_slice(manifest.upload_id.as_bytes());
441    encoded.push(b'\n');
442    encoded
443}
444
445fn chunk_manifest_key(key: &str) -> String {
446    format!("{key}{CHUNK_MANIFEST_SUFFIX}")
447}
448
449fn chunk_part_key(key: &str, upload_id: &str, index: usize) -> String {
450    format!("{key}.part{index}.{upload_id}")
451}
452
453fn decode_chunk_manifest(data: &[u8]) -> Result<ChunkManifest, CloudHomeError> {
454    let body = data.strip_prefix(CHUNK_MANIFEST_MAGIC).ok_or_else(|| {
455        CloudHomeError::Transport("CloudKit chunk manifest missing magic".to_string())
456    })?;
457    let body = std::str::from_utf8(body).map_err(|e| {
458        CloudHomeError::Transport(format!("CloudKit chunk manifest is not UTF-8: {e}"))
459    })?;
460    let mut lines = body.lines();
461    let part_count = lines
462        .next()
463        .ok_or_else(|| {
464            CloudHomeError::Transport("CloudKit chunk manifest missing part count".to_string())
465        })?
466        .parse::<usize>()
467        .map_err(|e| {
468            CloudHomeError::Transport(format!(
469                "CloudKit chunk manifest part count is invalid: {e}"
470            ))
471        })?;
472    let total_len = lines
473        .next()
474        .ok_or_else(|| {
475            CloudHomeError::Transport("CloudKit chunk manifest missing total length".to_string())
476        })?
477        .parse::<usize>()
478        .map_err(|e| {
479            CloudHomeError::Transport(format!(
480                "CloudKit chunk manifest total length is invalid: {e}"
481            ))
482        })?;
483    let upload_id = lines
484        .next()
485        .ok_or_else(|| {
486            CloudHomeError::Transport("CloudKit chunk manifest missing upload id".to_string())
487        })?
488        .to_string();
489    if lines.next().is_some() {
490        return Err(CloudHomeError::Transport(
491            "CloudKit chunk manifest has extra fields".to_string(),
492        ));
493    }
494    if part_count == 0 || total_len == 0 {
495        return Err(CloudHomeError::Transport(
496            "CloudKit chunk manifest must describe a non-empty object".to_string(),
497        ));
498    }
499    if upload_id.is_empty()
500        || !upload_id
501            .as_bytes()
502            .iter()
503            .all(|b| b.is_ascii_alphanumeric() || *b == b'-' || *b == b'_')
504    {
505        return Err(CloudHomeError::Transport(
506            "CloudKit chunk manifest upload id is invalid".to_string(),
507        ));
508    }
509    if ChunkManifest::new(total_len, upload_id.clone()).part_count != part_count {
510        return Err(CloudHomeError::Transport(format!(
511            "CloudKit chunk manifest part count {part_count} does not match total length {total_len}"
512        )));
513    }
514    Ok(ChunkManifest {
515        part_count,
516        total_len,
517        upload_id,
518    })
519}
520
521struct CloudKitStagingCleanup {
522    ops: Arc<dyn CloudKitOps>,
523    scope: CloudKitScope,
524    batch: CloudKitAtomicCreateBatch,
525    armed: std::sync::atomic::AtomicBool,
526}
527
528impl CloudKitStagingCleanup {
529    fn disarm(&self) {
530        self.armed.store(false, std::sync::atomic::Ordering::SeqCst);
531    }
532
533    fn cleanup_failure(&self, operation: CloudHomeError) -> CloudHomeError {
534        self.disarm();
535        match self.ops.discard_atomic_create(&self.scope, &self.batch) {
536            Ok(()) => operation,
537            Err(cleanup) => CloudHomeError::CleanupFailed {
538                operation: Box::new(operation),
539                cleanup: Box::new(CloudHomeError::Transport(format!(
540                    "discard CloudKit atomic-create batch {:?}: {cleanup}",
541                    self.batch.as_provider()
542                ))),
543            },
544        }
545    }
546}
547
548impl Drop for CloudKitStagingCleanup {
549    fn drop(&mut self) {
550        if !*self.armed.get_mut() {
551            return;
552        }
553        if let Err(error) = self.ops.discard_atomic_create(&self.scope, &self.batch) {
554            tracing::error!(
555                batch = self.batch.as_provider(),
556                %error,
557                "CloudKit cancellation failed to discard atomic-create batch"
558            );
559            std::process::abort();
560        }
561    }
562}
563
564async fn begin_atomic_create(
565    ops: Arc<dyn CloudKitOps>,
566    scope: CloudKitScope,
567) -> Result<Arc<CloudKitStagingCleanup>, CloudHomeError> {
568    tokio::task::spawn_blocking(move || {
569        let batch = ops.begin_atomic_create(&scope)?;
570        Ok(Arc::new(CloudKitStagingCleanup {
571            ops,
572            scope,
573            batch,
574            armed: std::sync::atomic::AtomicBool::new(true),
575        }))
576    })
577    .await
578    .map_err(|error| {
579        CloudHomeError::Transport(format!(
580            "CloudKit atomic-create staging task failed: {error}"
581        ))
582    })?
583}
584
585async fn stage_atomic_create_record(
586    staging: Arc<CloudKitStagingCleanup>,
587    record: CloudKitRecordCreate,
588) -> Result<(), CloudHomeError> {
589    if record.data.len() > CHUNK_SIZE {
590        return Err(CloudHomeError::Configuration(format!(
591            "CloudKit staged record {:?} has {} bytes, above the {CHUNK_SIZE}-byte bound",
592            record.key,
593            record.data.len()
594        )));
595    }
596    tokio::task::spawn_blocking(move || {
597        staging
598            .ops
599            .stage_atomic_create_record(&staging.scope, &staging.batch, record)
600    })
601    .await
602    .map_err(|error| {
603        CloudHomeError::Transport(format!(
604            "CloudKit atomic-create staging task failed: {error}"
605        ))
606    })?
607}
608
609async fn commit_atomic_create(
610    staging: Arc<CloudKitStagingCleanup>,
611) -> Result<Vec<CloudKitRecordVersion>, CloudHomeError> {
612    tokio::task::spawn_blocking(move || {
613        let created = staging
614            .ops
615            .commit_atomic_create(&staging.scope, &staging.batch)?;
616        staging.disarm();
617        Ok(created)
618    })
619    .await
620    .map_err(|error| {
621        CloudHomeError::Transport(format!(
622            "CloudKit atomic-create commit task failed: {error}"
623        ))
624    })?
625}
626
627async fn authoritative_created_records(
628    ops: Arc<dyn CloudKitOps>,
629    scope: CloudKitScope,
630    keys: Vec<String>,
631) -> Result<Vec<CloudKitRecordVersion>, CloudHomeError> {
632    blocking(move || {
633        keys.into_iter()
634            .map(|key| {
635                let record = ops.read_versioned_record(&scope, &key).map_err(|error| {
636                    CloudHomeError::Transport(format!(
637                        "read committed CloudKit atomic-create record {key:?}: {error}"
638                    ))
639                })?;
640                Ok(CloudKitRecordVersion {
641                    key,
642                    version: record.version,
643                })
644            })
645            .collect()
646    })
647    .await
648}
649
650enum AtomicCreateReadback {
651    Created(Vec<CloudKitRecordVersion>),
652    Absent,
653}
654
655async fn settle_atomic_create_response_loss(
656    ops: Arc<dyn CloudKitOps>,
657    scope: CloudKitScope,
658    keys: Vec<String>,
659) -> Result<AtomicCreateReadback, CloudHomeError> {
660    blocking(move || {
661        let mut created = Vec::with_capacity(keys.len());
662        let mut missing = 0usize;
663        for key in keys {
664            match ops.read_versioned_record(&scope, &key) {
665                Ok(record) => created.push(CloudKitRecordVersion {
666                    key,
667                    version: record.version,
668                }),
669                Err(CloudHomeError::NotFound(_)) => missing += 1,
670                Err(error) => return Err(error),
671            }
672        }
673        match (created.is_empty(), missing) {
674            (true, _) => Ok(AtomicCreateReadback::Absent),
675            (false, 0) => Ok(AtomicCreateReadback::Created(created)),
676            (false, _) => Err(CloudHomeError::Transport(
677                "CloudKit atomic create exposed only part of its record batch".to_string(),
678            )),
679        }
680    })
681    .await
682}
683
684fn parse_chunk_key(key: &str, upload_id: &str) -> Result<Option<usize>, CloudHomeError> {
685    let Some((part_key, token)) = key.rsplit_once('.') else {
686        return Ok(None);
687    };
688    if token != upload_id {
689        return Ok(None);
690    }
691    let index = part_key
692        .rsplit_once(".part")
693        .and_then(|(_, suffix)| suffix.parse::<usize>().ok())
694        .ok_or_else(|| {
695            CloudHomeError::Transport(format!("chunk key {key:?} missing .part suffix"))
696        })?;
697    Ok(Some(index))
698}
699
700fn list_numbered_chunks(
701    ops: &dyn CloudKitOps,
702    scope: &CloudKitScope,
703    key: &str,
704    manifest: &ChunkManifest,
705) -> Result<Vec<(usize, String)>, CloudHomeError> {
706    let chunk_prefix = format!("{key}.part");
707    let mut numbered = Vec::new();
708    for chunk_key in ops.list_records(scope, &chunk_prefix)? {
709        let Some(index) = parse_chunk_key(&chunk_key, &manifest.upload_id)? else {
710            continue;
711        };
712        numbered.push((index, chunk_key));
713    }
714    numbered.sort_by_key(|(index, _)| *index);
715    Ok(numbered)
716}
717
718fn verify_chunk_manifest(
719    key: &str,
720    manifest: &ChunkManifest,
721    chunks: &[(usize, String)],
722) -> Result<(), CloudHomeError> {
723    if chunks.len() != manifest.part_count {
724        return Err(CloudHomeError::Transport(format!(
725            "CloudKit object {key} is incomplete: manifest expects {} parts, found {}",
726            manifest.part_count,
727            chunks.len()
728        )));
729    }
730    for (expected, (actual, _)) in chunks.iter().enumerate() {
731        if *actual != expected {
732            return Err(CloudHomeError::Transport(format!(
733                "CloudKit object {key} is incomplete: missing part {expected}"
734            )));
735        }
736    }
737    Ok(())
738}
739
740fn chunk_manifest_has_all_parts(manifest: &ChunkManifest, chunks: &[(usize, String)]) -> bool {
741    chunks.len() == manifest.part_count
742        && chunks
743            .iter()
744            .enumerate()
745            .all(|(expected, (actual, _))| *actual == expected)
746}
747
748fn read_chunk(
749    ops: &dyn CloudKitOps,
750    scope: &CloudKitScope,
751    key: &str,
752    manifest: &ChunkManifest,
753    index: usize,
754    chunk_key: &str,
755) -> Result<Vec<u8>, CloudHomeError> {
756    let chunk = ops.read_record(scope, chunk_key)?;
757    let expected_len = if index + 1 == manifest.part_count {
758        manifest.total_len - (CHUNK_SIZE * index)
759    } else {
760        CHUNK_SIZE
761    };
762    if chunk.len() != expected_len {
763        return Err(CloudHomeError::Transport(format!(
764            "CloudKit object {key} part {index} has {} bytes, expected {expected_len}",
765            chunk.len()
766        )));
767    }
768    Ok(chunk)
769}
770
771fn read_chunked_object(
772    ops: &dyn CloudKitOps,
773    scope: &CloudKitScope,
774    key: &str,
775    manifest: ChunkManifest,
776) -> Result<Vec<u8>, CloudHomeError> {
777    let chunks = list_numbered_chunks(ops, scope, key, &manifest)?;
778    verify_chunk_manifest(key, &manifest, &chunks)?;
779
780    let mut result = Vec::with_capacity(manifest.total_len);
781    for (index, chunk_key) in &chunks {
782        let chunk = read_chunk(ops, scope, key, &manifest, *index, chunk_key)?;
783        result.extend_from_slice(&chunk);
784    }
785    if result.len() != manifest.total_len {
786        return Err(CloudHomeError::Transport(format!(
787            "CloudKit object {key} assembled to {} bytes, expected {}",
788            result.len(),
789            manifest.total_len
790        )));
791    }
792    Ok(result)
793}
794
795fn delete_chunk_layout(
796    ops: &dyn CloudKitOps,
797    scope: &CloudKitScope,
798    key: &str,
799) -> Result<(), CloudHomeError> {
800    match ops.delete_record(scope, &chunk_manifest_key(key)) {
801        Ok(()) | Err(CloudHomeError::NotFound(_)) => {}
802        Err(e) => return Err(e),
803    }
804
805    let chunk_prefix = format!("{key}.part");
806    let chunks = ops.list_records(scope, &chunk_prefix)?;
807    for chunk_key in chunks {
808        match ops.delete_record(scope, &chunk_key) {
809            Ok(()) | Err(CloudHomeError::NotFound(_)) => {}
810            Err(e) => return Err(e),
811        }
812    }
813
814    Ok(())
815}
816
817fn delete_stale_chunk_records(
818    ops: &dyn CloudKitOps,
819    scope: &CloudKitScope,
820    key: &str,
821    upload_id: &str,
822) -> Result<(), CloudHomeError> {
823    let chunk_prefix = format!("{key}.part");
824    let chunks = ops.list_records(scope, &chunk_prefix)?;
825    for chunk_key in chunks {
826        if parse_chunk_key(&chunk_key, upload_id)?.is_some() {
827            continue;
828        }
829        match ops.delete_record(scope, &chunk_key) {
830            Ok(()) | Err(CloudHomeError::NotFound(_)) => {}
831            Err(e) => return Err(e),
832        }
833    }
834    Ok(())
835}
836
837fn delete_single_record(
838    ops: &dyn CloudKitOps,
839    scope: &CloudKitScope,
840    key: &str,
841) -> Result<(), CloudHomeError> {
842    match ops.delete_record(scope, key) {
843        Ok(()) | Err(CloudHomeError::NotFound(_)) => Ok(()),
844        Err(e) => Err(e),
845    }
846}
847
848/// Delete old single record and chunk records for a key.
849fn delete_all_variants(
850    ops: &dyn CloudKitOps,
851    scope: &CloudKitScope,
852    key: &str,
853) -> Result<(), CloudHomeError> {
854    delete_single_record(ops, scope, key)?;
855    delete_chunk_layout(ops, scope, key)
856}
857
858/// A [`PartSink`] over CloudKit's chunked record layout: each `send_part` writes
859/// one tokened part record (CKAsset caps at 50 MB, so a large blob is split),
860/// `finish` writes the `{key}.manifest` record that makes the object readable.
861/// Existing records stay readable until the manifest points at the new token.
862struct CloudKitPartSink {
863    ops: Arc<dyn CloudKitOps>,
864    scope: CloudKitScope,
865    key: String,
866    upload_id: String,
867    index: usize,
868    total_len: usize,
869    written_len: usize,
870    settled: Arc<std::sync::atomic::AtomicBool>,
871}
872
873impl CloudKitPartSink {
874    async fn abort(&mut self) -> Result<(), CloudHomeError> {
875        if self.settled.load(std::sync::atomic::Ordering::SeqCst) {
876            return Ok(());
877        }
878        let ops = self.ops.clone();
879        let scope = self.scope.clone();
880        let key = self.key.clone();
881        let upload_id = self.upload_id.clone();
882        let written_parts = self.index;
883        cancellation_safe_blocking(move || {
884            for i in 0..written_parts {
885                delete_single_record(&*ops, &scope, &chunk_part_key(&key, &upload_id, i))?;
886            }
887            Ok::<(), CloudHomeError>(())
888        })
889        .await??;
890        self.settled
891            .store(true, std::sync::atomic::Ordering::SeqCst);
892        Ok(())
893    }
894}
895
896fn combine_cloudkit_cleanup_failure(
897    operation: CloudHomeError,
898    cleanup: Result<(), CloudHomeError>,
899) -> CloudHomeError {
900    match cleanup {
901        Ok(()) => operation,
902        Err(cleanup) => CloudHomeError::CleanupFailed {
903            operation: Box::new(operation),
904            cleanup: Box::new(cleanup),
905        },
906    }
907}
908
909impl Drop for CloudKitPartSink {
910    fn drop(&mut self) {
911        if self.settled.load(std::sync::atomic::Ordering::SeqCst) {
912            return;
913        }
914        for index in 0..self.index {
915            let part_key = chunk_part_key(&self.key, &self.upload_id, index);
916            if let Err(error) = delete_single_record(&*self.ops, &self.scope, &part_key) {
917                tracing::error!(
918                    %error,
919                    key = %self.key,
920                    upload_id = %self.upload_id,
921                    "CloudKit cancellation failed to discard multipart part"
922                );
923                std::process::abort();
924            }
925        }
926    }
927}
928
929#[async_trait]
930impl super::PartSink for CloudKitPartSink {
931    fn part_size(&self) -> usize {
932        CHUNK_SIZE
933    }
934
935    async fn send_part(
936        &mut self,
937        part: bytes::Bytes,
938        _offset: u64,
939        _is_last: bool,
940    ) -> Result<(), CloudHomeError> {
941        let i = self.index;
942        self.index += 1;
943        self.written_len += part.len();
944        let chunk_key = chunk_part_key(&self.key, &self.upload_id, i);
945        let ops = self.ops.clone();
946        let scope = self.scope.clone();
947        cancellation_safe_blocking(move || ops.write_record(&scope, &chunk_key, part.to_vec()))
948            .await?
949    }
950
951    async fn abort(&mut self) -> Result<(), CloudHomeError> {
952        CloudKitPartSink::abort(self).await
953    }
954
955    async fn finish(mut self: Box<Self>) -> Result<(), CloudHomeError> {
956        let manifest = ChunkManifest::new(self.total_len, self.upload_id.clone());
957        if self.index != manifest.part_count || self.written_len != manifest.total_len {
958            let operation = CloudHomeError::Transport(format!(
959                "CloudKit multipart {} wrote {} parts/{} bytes, expected {} parts/{} bytes",
960                self.key, self.index, self.written_len, manifest.part_count, manifest.total_len
961            ));
962            let cleanup = self.abort().await;
963            return Err(combine_cloudkit_cleanup_failure(operation, cleanup));
964        }
965        let manifest_key = chunk_manifest_key(&self.key);
966        let manifest_data = encode_chunk_manifest(manifest);
967        let ops = self.ops.clone();
968        let scope = self.scope.clone();
969        let settled = self.settled.clone();
970        if let Err(operation) = cancellation_safe_blocking(move || {
971            let result = ops.write_record(&scope, &manifest_key, manifest_data);
972            if result.is_ok() {
973                settled.store(true, std::sync::atomic::Ordering::SeqCst);
974            }
975            result
976        })
977        .await?
978        {
979            let cleanup = self.abort().await;
980            return Err(combine_cloudkit_cleanup_failure(operation, cleanup));
981        }
982
983        // The manifest is published, so the object is now readable. The remaining
984        // steps only clean up records the previous object left; their failure fails
985        // loud but must not touch the parts just published.
986        let ops = self.ops.clone();
987        let scope = self.scope.clone();
988        let key = self.key.clone();
989        blocking(move || delete_single_record(&*ops, &scope, &key)).await?;
990
991        let ops = self.ops.clone();
992        let scope = self.scope.clone();
993        let key = self.key.clone();
994        let upload_id = self.upload_id.clone();
995        blocking(move || delete_stale_chunk_records(&*ops, &scope, &key, &upload_id)).await
996    }
997}
998
999#[async_trait]
1000impl CloudHome for CloudKitCloudHome {
1001    fn exact_slot_storage(self: Arc<Self>) -> Option<Arc<dyn ExactSlotStorage>> {
1002        Some(self)
1003    }
1004
1005    async fn put_object(&self, key: &str, data: Vec<u8>) -> Result<(), CloudHomeError> {
1006        let ops = self.ops.clone();
1007        let scope = self.scope.clone();
1008        let k = key.to_string();
1009        blocking(move || ops.write_record(&scope, &k, data)).await?;
1010
1011        let ops = self.ops.clone();
1012        let scope = self.scope.clone();
1013        let k = key.to_string();
1014        blocking(move || delete_chunk_layout(&*ops, &scope, &k)).await
1015    }
1016
1017    async fn open_multipart<'a>(
1018        &'a self,
1019        key: &str,
1020        total_len: u64,
1021    ) -> Result<super::BoxPartSink<'a>, CloudHomeError> {
1022        let total_len = usize::try_from(total_len).map_err(|_| {
1023            CloudHomeError::Transport(format!(
1024                "CloudKit object {key} is too large for this platform"
1025            ))
1026        })?;
1027        Ok(Box::new(CloudKitPartSink {
1028            ops: self.ops.clone(),
1029            scope: self.scope.clone(),
1030            key: key.to_string(),
1031            upload_id: self.ids.new_id(),
1032            index: 0,
1033            total_len,
1034            written_len: 0,
1035            settled: Arc::new(std::sync::atomic::AtomicBool::new(false)),
1036        }))
1037    }
1038
1039    fn multipart_threshold(&self) -> u64 {
1040        CHUNK_SIZE as u64
1041    }
1042
1043    async fn read(&self, key: &str) -> Result<Vec<u8>, CloudHomeError> {
1044        let ops = self.ops.clone();
1045        let scope = self.scope.clone();
1046        let key = key.to_string();
1047        blocking(move || {
1048            match ops.read_record(&scope, &key) {
1049                Ok(data) => return Ok(data),
1050                Err(CloudHomeError::NotFound(_)) => {}
1051                Err(e) => return Err(e),
1052            }
1053
1054            match ops.read_record(&scope, &chunk_manifest_key(&key)) {
1055                Ok(data) => {
1056                    let manifest = decode_chunk_manifest(&data)?;
1057                    return read_chunked_object(&*ops, &scope, &key, manifest);
1058                }
1059                Err(CloudHomeError::NotFound(_)) => {}
1060                Err(e) => return Err(e),
1061            }
1062
1063            let chunks = ops.list_records(&scope, &format!("{key}.part"))?;
1064            if chunks.is_empty() {
1065                return Err(CloudHomeError::NotFound(key));
1066            }
1067            Err(CloudHomeError::Transport(format!(
1068                "CloudKit object {key} has chunk records but no manifest"
1069            )))
1070        })
1071        .await
1072    }
1073
1074    async fn read_range(&self, key: &str, start: u64, end: u64) -> Result<Vec<u8>, CloudHomeError> {
1075        if end <= start {
1076            return Ok(Vec::new());
1077        }
1078
1079        let ops = self.ops.clone();
1080        let scope = self.scope.clone();
1081        let key = key.to_string();
1082        blocking(move || {
1083            let start = start as usize;
1084            let end = end as usize;
1085
1086            match ops.read_record(&scope, &key) {
1087                Ok(data) => {
1088                    if end > data.len() {
1089                        return Err(CloudHomeError::Transport(format!(
1090                            "range {start}..{end} exceeds file size {}",
1091                            data.len()
1092                        )));
1093                    }
1094                    return Ok(data[start..end].to_vec());
1095                }
1096                Err(CloudHomeError::NotFound(_)) => {}
1097                Err(e) => return Err(e),
1098            }
1099
1100            match ops.read_record(&scope, &chunk_manifest_key(&key)) {
1101                Ok(data) => {
1102                    let manifest = decode_chunk_manifest(&data)?;
1103                    let chunks = list_numbered_chunks(&*ops, &scope, &key, &manifest)?;
1104                    verify_chunk_manifest(&key, &manifest, &chunks)?;
1105                    if end > manifest.total_len {
1106                        return Err(CloudHomeError::Transport(format!(
1107                            "range {start}..{end} exceeds file size {}",
1108                            manifest.total_len
1109                        )));
1110                    }
1111
1112                    let first_chunk = start / CHUNK_SIZE;
1113                    let last_chunk = (end - 1) / CHUNK_SIZE;
1114                    let mut result = Vec::with_capacity(end - start);
1115                    for (i, chunk_key) in chunks
1116                        .iter()
1117                        .filter(|(i, _)| (first_chunk..=last_chunk).contains(i))
1118                    {
1119                        let chunk = read_chunk(&*ops, &scope, &key, &manifest, *i, chunk_key)?;
1120                        let chunk_start = i * CHUNK_SIZE;
1121                        let slice_start = if *i == first_chunk {
1122                            start - chunk_start
1123                        } else {
1124                            0
1125                        };
1126                        let slice_end = if *i == last_chunk {
1127                            end - chunk_start
1128                        } else {
1129                            chunk.len()
1130                        };
1131                        result.extend_from_slice(&chunk[slice_start..slice_end]);
1132                    }
1133                    return Ok(result);
1134                }
1135                Err(CloudHomeError::NotFound(_)) => {}
1136                Err(e) => return Err(e),
1137            }
1138
1139            let chunks = ops.list_records(&scope, &format!("{key}.part"))?;
1140            if chunks.is_empty() {
1141                return Err(CloudHomeError::NotFound(key));
1142            }
1143            Err(CloudHomeError::Transport(format!(
1144                "CloudKit object {key} has chunk records but no manifest"
1145            )))
1146        })
1147        .await
1148    }
1149
1150    async fn list(&self, prefix: &str) -> Result<Vec<String>, CloudHomeError> {
1151        let ops = self.ops.clone();
1152        let scope = self.scope.clone();
1153        let prefix = prefix.to_string();
1154        blocking(move || {
1155            let raw_keys = ops.list_records(&scope, &prefix)?;
1156
1157            // A base key exists only when its single record or its manifest is
1158            // present — the manifest is what makes a chunked object readable. Part
1159            // records with no manifest are an incomplete or aborted upload, which
1160            // `read` cannot assemble, so they are not reported.
1161            let present: HashSet<&str> = raw_keys.iter().map(String::as_str).collect();
1162            let mut base_keys: Vec<String> = raw_keys
1163                .iter()
1164                .map(|k| strip_part_suffix(k))
1165                .filter(|&base| {
1166                    present.contains(base) || present.contains(chunk_manifest_key(base).as_str())
1167                })
1168                .map(str::to_string)
1169                .collect();
1170            base_keys.sort();
1171            base_keys.dedup();
1172            Ok(base_keys)
1173        })
1174        .await
1175    }
1176
1177    async fn delete(&self, key: &str) -> Result<(), CloudHomeError> {
1178        let ops = self.ops.clone();
1179        let scope = self.scope.clone();
1180        let key = key.to_string();
1181        blocking(move || delete_all_variants(&*ops, &scope, &key)).await
1182    }
1183
1184    async fn exists(&self, key: &str) -> Result<bool, CloudHomeError> {
1185        let ops = self.ops.clone();
1186        let scope = self.scope.clone();
1187        let key = key.to_string();
1188        blocking(move || {
1189            if ops.record_exists(&scope, &key)? {
1190                return Ok(true);
1191            }
1192            let manifest = match ops.read_record(&scope, &chunk_manifest_key(&key)) {
1193                Ok(data) => decode_chunk_manifest(&data)?,
1194                Err(CloudHomeError::NotFound(_)) => return Ok(false),
1195                Err(e) => return Err(e),
1196            };
1197            let chunks = list_numbered_chunks(&*ops, &scope, &key, &manifest)?;
1198            Ok(chunk_manifest_has_all_parts(&manifest, &chunks))
1199        })
1200        .await
1201    }
1202
1203    async fn set_access(
1204        &self,
1205        desired: CloudAccessState,
1206    ) -> Result<CloudAccessOutcome, CloudHomeError> {
1207        let ops = self.ops.clone();
1208        match desired {
1209            CloudAccessState::Present { member_pubkey, .. } => {
1210                // CloudKit shares bind the joiner's identity at URL-accept time,
1211                // so no provider email is required.
1212                let lookup_ops = ops.clone();
1213                let lookup_member = member_pubkey.clone();
1214                let existing =
1215                    blocking(move || lookup_ops.share_for_member(&lookup_member)).await?;
1216                let expected = match existing {
1217                    Some(share) => share,
1218                    None => {
1219                        let grant_ops = ops.clone();
1220                        let grant_member = member_pubkey.clone();
1221                        blocking(move || grant_ops.grant_share(&grant_member)).await?
1222                    }
1223                };
1224                let verified = blocking(move || ops.share_for_member(&member_pubkey))
1225                    .await?
1226                    .ok_or_else(|| {
1227                        CloudHomeError::Transport(
1228                            "CloudKit member share is absent after setting it present".to_string(),
1229                        )
1230                    })?;
1231                if verified != expected {
1232                    return Err(CloudHomeError::Transport(
1233                        "CloudKit member share changed while verifying present access".to_string(),
1234                    ));
1235                }
1236                Ok(CloudAccessOutcome::Present(
1237                    CloudHomeJoinInfo::CloudKitShare {
1238                        share_url: verified.share_url,
1239                        owner_name: verified.owner_name,
1240                        zone_name: verified.zone_name,
1241                    },
1242                ))
1243            }
1244            CloudAccessState::Absent { member_pubkey, .. } => {
1245                let lookup_ops = ops.clone();
1246                let lookup_member = member_pubkey.clone();
1247                if blocking(move || lookup_ops.share_for_member(&lookup_member))
1248                    .await?
1249                    .is_some()
1250                {
1251                    let revoke_ops = ops.clone();
1252                    let revoke_member = member_pubkey.clone();
1253                    blocking(move || revoke_ops.revoke_share(&revoke_member)).await?;
1254                }
1255                if blocking(move || ops.share_for_member(&member_pubkey))
1256                    .await?
1257                    .is_some()
1258                {
1259                    return Err(CloudHomeError::Transport(
1260                        "CloudKit member share remains after setting access absent".to_string(),
1261                    ));
1262                }
1263                Ok(CloudAccessOutcome::Absent(RevokeOutcome::Revoked))
1264            }
1265        }
1266    }
1267}
1268
1269const EXACT_MANIFEST_MAGIC: &[u8] = b"coven-cloudkit-exact-manifest-v1\0";
1270
1271fn validate_cloudkit_slot(slot: &ObjectSlot) -> Result<(), CloudHomeError> {
1272    slot.validate()?;
1273    if slot.physical() != &PhysicalObjectLocator::LogicalKey {
1274        return Err(CloudHomeError::Configuration(format!(
1275            "CloudKit slot for {} must use its logical key",
1276            slot.logical_key()
1277        )));
1278    }
1279    Ok(())
1280}
1281
1282fn exact_part_key(logical_key: &str, index: usize) -> String {
1283    format!("{logical_key}.exact-part{index}")
1284}
1285
1286fn encode_exact_manifest(part_count: usize, total_len: usize) -> Vec<u8> {
1287    let mut bytes = EXACT_MANIFEST_MAGIC.to_vec();
1288    bytes.extend_from_slice(part_count.to_string().as_bytes());
1289    bytes.push(b'\n');
1290    bytes.extend_from_slice(total_len.to_string().as_bytes());
1291    bytes.push(b'\n');
1292    bytes
1293}
1294
1295fn decode_exact_manifest(bytes: &[u8]) -> Result<(usize, usize), CloudHomeError> {
1296    let text = std::str::from_utf8(bytes.strip_prefix(EXACT_MANIFEST_MAGIC).ok_or_else(|| {
1297        CloudHomeError::Transport("CloudKit exact object has an invalid manifest".to_string())
1298    })?)
1299    .map_err(|error| CloudHomeError::Transport(format!("CloudKit exact manifest: {error}")))?;
1300    let mut lines = text.lines();
1301    let part_count = lines
1302        .next()
1303        .ok_or_else(|| {
1304            CloudHomeError::Transport("CloudKit exact manifest omitted part count".to_string())
1305        })?
1306        .parse::<usize>()
1307        .map_err(|error| {
1308            CloudHomeError::Transport(format!("CloudKit exact manifest part count: {error}"))
1309        })?;
1310    let total_len = lines
1311        .next()
1312        .ok_or_else(|| {
1313            CloudHomeError::Transport("CloudKit exact manifest omitted length".to_string())
1314        })?
1315        .parse::<usize>()
1316        .map_err(|error| {
1317            CloudHomeError::Transport(format!("CloudKit exact manifest length: {error}"))
1318        })?;
1319    if lines.next().is_some() || part_count != total_len.div_ceil(CHUNK_SIZE) {
1320        return Err(CloudHomeError::Transport(
1321            "CloudKit exact manifest shape does not match its length".to_string(),
1322        ));
1323    }
1324    Ok((part_count, total_len))
1325}
1326
1327fn read_exact_cloudkit_object(
1328    ops: &dyn CloudKitOps,
1329    scope: &CloudKitScope,
1330    logical_key: &str,
1331) -> Result<(Vec<u8>, Vec<CloudKitRecordVersion>), CloudHomeError> {
1332    let manifest = ops.read_versioned_record(scope, logical_key)?;
1333    let (part_count, total_len) = decode_exact_manifest(&manifest.bytes)?;
1334    let mut bytes = Vec::with_capacity(total_len);
1335    let mut records = Vec::with_capacity(part_count + 1);
1336    records.push(CloudKitRecordVersion {
1337        key: logical_key.to_string(),
1338        version: manifest.version,
1339    });
1340    for index in 0..part_count {
1341        let key = exact_part_key(logical_key, index);
1342        let part = ops.read_versioned_record(scope, &key)?;
1343        let expected_len = if index + 1 == part_count {
1344            total_len - index * CHUNK_SIZE
1345        } else {
1346            CHUNK_SIZE
1347        };
1348        if part.bytes.len() != expected_len {
1349            return Err(CloudHomeError::Transport(format!(
1350                "CloudKit exact object {logical_key:?} part {index} has {} bytes, expected {expected_len}",
1351                part.bytes.len()
1352            )));
1353        }
1354        bytes.extend_from_slice(&part.bytes);
1355        records.push(CloudKitRecordVersion {
1356            key,
1357            version: part.version,
1358        });
1359    }
1360    Ok((bytes, records))
1361}
1362
1363#[async_trait]
1364impl ExactSlotStorage for CloudKitCloudHome {
1365    async fn provider_binding(
1366        &self,
1367    ) -> Result<coven_core::sync::storage::ResolvedProviderBinding, CloudHomeError> {
1368        use coven_core::sync::storage::{
1369            ProviderDeviceBinding, ProviderPrincipalId, ResolvedProviderBinding,
1370            StoreProviderBinding,
1371        };
1372
1373        let ops = self.ops.clone();
1374        let scope = self.scope.clone();
1375        let identity = blocking(move || ops.provider_identity(&scope)).await?;
1376        if identity.container_id.is_empty()
1377            || identity.owner_name.is_empty()
1378            || identity.zone_name.is_empty()
1379            || identity.current_user_record_name.is_empty()
1380        {
1381            return Err(CloudHomeError::Configuration(
1382                "CloudKit provider identity contains an empty stable identifier".to_string(),
1383            ));
1384        }
1385        if let CloudKitScope::Shared {
1386            owner_name,
1387            zone_name,
1388        } = &self.scope
1389        {
1390            if owner_name != &identity.owner_name || zone_name != &identity.zone_name {
1391                return Err(CloudHomeError::Configuration(format!(
1392                    "CloudKit provider identity resolved zone {}/{}, expected {owner_name}/{zone_name}",
1393                    identity.owner_name, identity.zone_name
1394                )));
1395            }
1396        }
1397        let principal = match &self.scope {
1398            CloudKitScope::Private => ProviderPrincipalId::CloudKitPrivateZoneOwner {
1399                record_name: identity.current_user_record_name,
1400            },
1401            CloudKitScope::Shared { .. } => ProviderPrincipalId::CloudKitSharedZoneParticipant {
1402                record_name: identity.current_user_record_name,
1403            },
1404        };
1405        Ok(ResolvedProviderBinding {
1406            store: StoreProviderBinding::CloudKit {
1407                container_id: identity.container_id,
1408                environment: identity.environment,
1409                owner_name: identity.owner_name,
1410                zone_name: identity.zone_name,
1411            },
1412            device: ProviderDeviceBinding { principal },
1413        })
1414    }
1415
1416    async fn cross_principal_evidence(
1417        &self,
1418    ) -> Result<coven_core::sync::provider::CrossPrincipalProviderEvidence, CloudHomeError> {
1419        use coven_core::sync::provider::{CloudKitAcceptedShare, CrossPrincipalProviderEvidence};
1420        use coven_core::sync::store_commit::ObjectHash;
1421
1422        let CloudKitScope::Shared {
1423            owner_name,
1424            zone_name,
1425        } = &self.scope
1426        else {
1427            return Err(CloudHomeError::Configuration(
1428                "CloudKit cross-principal evidence requires an accepted shared zone".to_string(),
1429            ));
1430        };
1431        let ops = self.ops.clone();
1432        let scope = self.scope.clone();
1433        let accepted = blocking(move || ops.accepted_read_write_share(&scope)).await?;
1434        let binding = self.provider_binding().await?;
1435        let coven_core::sync::storage::ProviderPrincipalId::CloudKitSharedZoneParticipant {
1436            record_name,
1437        } = binding.device.principal
1438        else {
1439            return Err(CloudHomeError::Configuration(
1440                "CloudKit adapter returned a non-CloudKit principal".to_string(),
1441            ));
1442        };
1443        if accepted.share_record_name.is_empty()
1444            || accepted.owner_name != *owner_name
1445            || accepted.zone_name != *zone_name
1446            || accepted.participant_record_name != record_name
1447            || accepted.permission != CloudKitSharePermission::ReadWrite
1448            || accepted.acceptance != CloudKitShareAcceptance::Accepted
1449            || accepted.canonical_record.is_empty()
1450        {
1451            return Err(CloudHomeError::Configuration(
1452                "CloudKit accepted share does not prove read-write participation in the selected zone"
1453                    .to_string(),
1454            ));
1455        }
1456        let share_slot = ObjectSlot::logical(format!(
1457            "__coven_cloudkit_share__/{}",
1458            hex::encode(ObjectHash::digest(accepted.share_record_name.as_bytes()).as_bytes())
1459        ))?;
1460        Ok(CrossPrincipalProviderEvidence::CloudKit(
1461            CloudKitAcceptedShare {
1462                share: coven_core::sync::storage::ExactObjectRef::new(
1463                    share_slot,
1464                    accepted.canonical_record.len() as u64,
1465                    ObjectHash::digest(&accepted.canonical_record),
1466                ),
1467                share_record_name: accepted.share_record_name,
1468                owner_name: accepted.owner_name,
1469                zone_name: accepted.zone_name,
1470                participant_record_name: accepted.participant_record_name,
1471            },
1472        ))
1473    }
1474
1475    async fn allocate_slot(&self, logical_key: &str) -> Result<ObjectSlot, CloudHomeError> {
1476        ObjectSlot::logical(logical_key.to_string())
1477    }
1478
1479    async fn create_at(
1480        &self,
1481        slot: &ObjectSlot,
1482        mut body: BlobBody,
1483        progress: &UploadProgress<'_>,
1484    ) -> Result<(), CloudHomeError> {
1485        validate_cloudkit_slot(slot)?;
1486        let total_len = usize::try_from(body.len()).map_err(|_| {
1487            CloudHomeError::Transport(format!(
1488                "CloudKit object {:?} is too large for this platform",
1489                slot.logical_key()
1490            ))
1491        })?;
1492        let part_count = total_len.div_ceil(CHUNK_SIZE);
1493        let staging = begin_atomic_create(self.ops.clone(), self.scope.clone()).await?;
1494        let mut requested_keys = Vec::with_capacity(part_count + 1);
1495        let mut written_len = 0usize;
1496        for index in 0..part_count {
1497            let part = match body.next_part(CHUNK_SIZE).await {
1498                Ok(Some(part)) => part,
1499                Ok(None) => {
1500                    return Err(staging.cleanup_failure(CloudHomeError::Transport(format!(
1501                        "CloudKit object {:?} ended after {written_len} of {total_len} bytes",
1502                        slot.logical_key()
1503                    ))))
1504                }
1505                Err(error) => return Err(staging.cleanup_failure(error)),
1506            };
1507            written_len += part.len();
1508            let key = exact_part_key(slot.logical_key(), index);
1509            if let Err(error) = stage_atomic_create_record(
1510                staging.clone(),
1511                CloudKitRecordCreate {
1512                    key: key.clone(),
1513                    data: part.to_vec(),
1514                },
1515            )
1516            .await
1517            {
1518                return Err(staging.cleanup_failure(error));
1519            }
1520            requested_keys.push(key);
1521        }
1522        match body.next_part(CHUNK_SIZE).await {
1523            Ok(None) if written_len == total_len => {}
1524            Ok(None) => {
1525                return Err(staging.cleanup_failure(CloudHomeError::Transport(format!(
1526                    "CloudKit object {:?} yielded {written_len} bytes, expected {total_len}",
1527                    slot.logical_key()
1528                ))))
1529            }
1530            Ok(Some(extra)) => {
1531                return Err(staging.cleanup_failure(CloudHomeError::Transport(format!(
1532                    "CloudKit object {:?} yielded at least {} bytes, expected {total_len}",
1533                    slot.logical_key(),
1534                    written_len + extra.len()
1535                ))))
1536            }
1537            Err(error) => return Err(staging.cleanup_failure(error)),
1538        }
1539        if let Err(error) = stage_atomic_create_record(
1540            staging.clone(),
1541            CloudKitRecordCreate {
1542                key: slot.logical_key().to_string(),
1543                data: encode_exact_manifest(part_count, total_len),
1544            },
1545        )
1546        .await
1547        {
1548            return Err(staging.cleanup_failure(error));
1549        }
1550        requested_keys.push(slot.logical_key().to_string());
1551        let created = match commit_atomic_create(staging.clone()).await {
1552            Ok(created) => created,
1553            Err(CloudHomeError::AlreadyExists(_)) => {
1554                return Err(staging.cleanup_failure(CloudHomeError::AlreadyExists(
1555                    slot.logical_key().to_string(),
1556                )));
1557            }
1558            Err(operation) => {
1559                match settle_atomic_create_response_loss(
1560                    self.ops.clone(),
1561                    self.scope.clone(),
1562                    requested_keys.clone(),
1563                )
1564                .await
1565                {
1566                    Ok(AtomicCreateReadback::Created(created)) => {
1567                        staging.disarm();
1568                        created
1569                    }
1570                    Ok(AtomicCreateReadback::Absent) => {
1571                        return Err(staging.cleanup_failure(operation))
1572                    }
1573                    Err(readback) => {
1574                        staging.disarm();
1575                        return Err(CloudHomeError::UnresolvedOutcome {
1576                            operation: Box::new(operation),
1577                            readback: Box::new(readback),
1578                        });
1579                    }
1580                }
1581            }
1582        };
1583        if created.len() != requested_keys.len()
1584            || created
1585                .iter()
1586                .zip(&requested_keys)
1587                .any(|(record, requested)| &record.key != requested)
1588        {
1589            authoritative_created_records(self.ops.clone(), self.scope.clone(), requested_keys)
1590                .await?;
1591        }
1592        progress(total_len as u64);
1593        Ok(())
1594    }
1595
1596    async fn read_at(&self, slot: &ObjectSlot) -> Result<Vec<u8>, CloudHomeError> {
1597        validate_cloudkit_slot(slot)?;
1598        let ops = self.ops.clone();
1599        let scope = self.scope.clone();
1600        let logical_key = slot.logical_key().to_string();
1601        blocking(move || {
1602            read_exact_cloudkit_object(&*ops, &scope, &logical_key).map(|value| value.0)
1603        })
1604        .await
1605    }
1606
1607    async fn read_range_at(
1608        &self,
1609        slot: &ObjectSlot,
1610        start: u64,
1611        end: u64,
1612    ) -> Result<Vec<u8>, CloudHomeError> {
1613        let bytes = self.read_at(slot).await?;
1614        let start = usize::try_from(start)
1615            .map_err(|_| CloudHomeError::Configuration("range start is too large".to_string()))?;
1616        let end = usize::try_from(end)
1617            .map_err(|_| CloudHomeError::Configuration("range end is too large".to_string()))?;
1618        bytes.get(start..end).map(<[u8]>::to_vec).ok_or_else(|| {
1619            CloudHomeError::Configuration(format!(
1620                "invalid range {start}..{end} for {} bytes",
1621                bytes.len()
1622            ))
1623        })
1624    }
1625
1626    async fn read_at_to_file(
1627        &self,
1628        slot: &ObjectSlot,
1629        destination: &std::path::Path,
1630    ) -> Result<(), super::CloudFileReadError> {
1631        let bytes = self.read_at(slot).await?;
1632        let stream: super::CloudObjectStream =
1633            Box::pin(futures_util::stream::once(
1634                async move { Ok(Bytes::from(bytes)) },
1635            ));
1636        super::write_cloud_object_stream(destination, stream)
1637            .await
1638            .map(drop)
1639    }
1640
1641    async fn delete_at(&self, slot: &ObjectSlot) -> Result<(), CloudHomeError> {
1642        validate_cloudkit_slot(slot)?;
1643        let ops = self.ops.clone();
1644        let scope = self.scope.clone();
1645        let logical_key = slot.logical_key().to_string();
1646        blocking(move || {
1647            let records = match read_exact_cloudkit_object(&*ops, &scope, &logical_key) {
1648                Ok((_, records)) => records,
1649                Err(CloudHomeError::NotFound(_)) => return Ok(()),
1650                Err(error) => return Err(error),
1651            };
1652            ops.delete_record_versions(&scope, &records)
1653        })
1654        .await
1655    }
1656}
1657
1658#[cfg(test)]
1659mod tests {
1660    use super::*;
1661    use crate::id_provider::SequentialIdProvider;
1662    use crate::storage::cloud::{no_progress, BlobBody};
1663    use std::collections::{HashMap, HashSet};
1664    use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
1665    use std::sync::Mutex;
1666
1667    fn exact_slot(key: &str) -> ObjectSlot {
1668        ObjectSlot::logical(key.to_string()).expect("valid exact slot")
1669    }
1670
1671    #[derive(Clone, Debug, PartialEq, Eq)]
1672    enum MockCall {
1673        Write(String),
1674        Read(String),
1675        List(String),
1676        Delete(String),
1677        Exists(String),
1678        Create(String),
1679        BeginBatch(String),
1680        Stage(String),
1681        CommitBatch(String),
1682        DiscardBatch(String),
1683        DeleteVersions(Vec<String>),
1684    }
1685
1686    struct PausedWrite {
1687        key: String,
1688        stored: Arc<std::sync::Barrier>,
1689        release: Arc<std::sync::Barrier>,
1690    }
1691
1692    struct MockCloudKitOps {
1693        store: Mutex<HashMap<(CloudKitScope, String), Vec<u8>>>,
1694        versions: Mutex<HashMap<(CloudKitScope, String), u64>>,
1695        calls: Mutex<Vec<MockCall>>,
1696        fail_deletes: Mutex<HashSet<String>>,
1697        fail_delete_once: Mutex<HashMap<String, usize>>,
1698        fail_writes: Mutex<HashSet<String>>,
1699        staged_batches: Mutex<HashMap<String, Vec<CloudKitRecordCreate>>>,
1700        next_batch: AtomicUsize,
1701        max_stage_payload: AtomicUsize,
1702        fail_discards: AtomicBool,
1703        lose_commit_response: AtomicBool,
1704        return_wrong_commit_keys: AtomicBool,
1705        pause_write_after_store: Mutex<Option<PausedWrite>>,
1706        record_exists_calls: AtomicUsize,
1707        grant_share_calls: AtomicUsize,
1708        revoke_share_calls: AtomicUsize,
1709        shares: Mutex<HashMap<String, CloudKitShare>>,
1710    }
1711
1712    impl MockCloudKitOps {
1713        fn new() -> Self {
1714            Self {
1715                store: Mutex::new(HashMap::new()),
1716                versions: Mutex::new(HashMap::new()),
1717                calls: Mutex::new(Vec::new()),
1718                fail_deletes: Mutex::new(HashSet::new()),
1719                fail_delete_once: Mutex::new(HashMap::new()),
1720                fail_writes: Mutex::new(HashSet::new()),
1721                staged_batches: Mutex::new(HashMap::new()),
1722                next_batch: AtomicUsize::new(0),
1723                max_stage_payload: AtomicUsize::new(0),
1724                fail_discards: AtomicBool::new(false),
1725                lose_commit_response: AtomicBool::new(false),
1726                return_wrong_commit_keys: AtomicBool::new(false),
1727                pause_write_after_store: Mutex::new(None),
1728                record_exists_calls: AtomicUsize::new(0),
1729                grant_share_calls: AtomicUsize::new(0),
1730                revoke_share_calls: AtomicUsize::new(0),
1731                shares: Mutex::new(HashMap::new()),
1732            }
1733        }
1734
1735        fn calls(&self) -> Vec<MockCall> {
1736            self.calls.lock().unwrap().clone()
1737        }
1738
1739        fn clear_calls(&self) {
1740            self.calls.lock().unwrap().clear();
1741        }
1742
1743        fn fail_delete(&self, key: &str) {
1744            self.fail_deletes.lock().unwrap().insert(key.to_string());
1745        }
1746
1747        fn fail_next_delete(&self, key: &str) {
1748            self.fail_delete_once
1749                .lock()
1750                .unwrap()
1751                .insert(key.to_string(), 1);
1752        }
1753
1754        fn fail_write(&self, key: &str) {
1755            self.fail_writes.lock().unwrap().insert(key.to_string());
1756        }
1757
1758        fn fail_discard(&self) {
1759            self.fail_discards.store(true, Ordering::SeqCst);
1760        }
1761
1762        fn lose_commit_response(&self) {
1763            self.lose_commit_response.store(true, Ordering::SeqCst);
1764        }
1765
1766        fn return_wrong_commit_keys(&self) {
1767            self.return_wrong_commit_keys.store(true, Ordering::SeqCst);
1768        }
1769
1770        fn pause_write_after_store(
1771            &self,
1772            key: &str,
1773        ) -> (Arc<std::sync::Barrier>, Arc<std::sync::Barrier>) {
1774            let stored = Arc::new(std::sync::Barrier::new(2));
1775            let release = Arc::new(std::sync::Barrier::new(2));
1776            let previous = self
1777                .pause_write_after_store
1778                .lock()
1779                .unwrap()
1780                .replace(PausedWrite {
1781                    key: key.to_string(),
1782                    stored: stored.clone(),
1783                    release: release.clone(),
1784                });
1785            assert!(previous.is_none(), "a CloudKit write is already paused");
1786            (stored, release)
1787        }
1788    }
1789
1790    impl CloudKitOps for MockCloudKitOps {
1791        fn provider_identity(
1792            &self,
1793            scope: &CloudKitScope,
1794        ) -> Result<CloudKitProviderIdentity, CloudHomeError> {
1795            let (owner_name, zone_name) = match scope {
1796                CloudKitScope::Private => ("private-owner", "private-zone"),
1797                CloudKitScope::Shared {
1798                    owner_name,
1799                    zone_name,
1800                } => (owner_name.as_str(), zone_name.as_str()),
1801            };
1802            Ok(CloudKitProviderIdentity {
1803                container_id: "iCloud.example.coven".to_string(),
1804                environment: coven_core::sync::storage::CloudKitEnvironment::Development,
1805                owner_name: owner_name.to_string(),
1806                zone_name: zone_name.to_string(),
1807                current_user_record_name: "current-user".to_string(),
1808            })
1809        }
1810
1811        fn accepted_read_write_share(
1812            &self,
1813            scope: &CloudKitScope,
1814        ) -> Result<CloudKitAcceptedShareRecord, CloudHomeError> {
1815            let CloudKitScope::Shared {
1816                owner_name,
1817                zone_name,
1818            } = scope
1819            else {
1820                return Err(CloudHomeError::NotFound(
1821                    "accepted CloudKit share".to_string(),
1822                ));
1823            };
1824            Ok(CloudKitAcceptedShareRecord {
1825                share_record_name: "accepted-share".to_string(),
1826                owner_name: owner_name.clone(),
1827                zone_name: zone_name.clone(),
1828                participant_record_name: "current-user".to_string(),
1829                permission: CloudKitSharePermission::ReadWrite,
1830                acceptance: CloudKitShareAcceptance::Accepted,
1831                canonical_record: b"canonical accepted CKShare".to_vec(),
1832            })
1833        }
1834
1835        fn write_record(
1836            &self,
1837            scope: &CloudKitScope,
1838            key: &str,
1839            data: Vec<u8>,
1840        ) -> Result<(), CloudHomeError> {
1841            self.calls
1842                .lock()
1843                .unwrap()
1844                .push(MockCall::Write(key.to_string()));
1845            if self.fail_writes.lock().unwrap().contains(key) {
1846                return Err(CloudHomeError::Transport(format!("write {key} failed")));
1847            }
1848            let record = (scope.clone(), key.to_string());
1849            self.store.lock().unwrap().insert(record.clone(), data);
1850            let mut versions = self.versions.lock().unwrap();
1851            let next = versions.get(&record).copied().unwrap_or(0) + 1;
1852            versions.insert(record, next);
1853            drop(versions);
1854            let pause = {
1855                let mut pause = self.pause_write_after_store.lock().unwrap();
1856                match pause.as_ref() {
1857                    Some(paused) if paused.key == key => pause.take(),
1858                    _ => None,
1859                }
1860            };
1861            if let Some(paused) = pause {
1862                assert_eq!(paused.key, key);
1863                paused.stored.wait();
1864                paused.release.wait();
1865            }
1866            Ok(())
1867        }
1868
1869        fn read_record(&self, scope: &CloudKitScope, key: &str) -> Result<Vec<u8>, CloudHomeError> {
1870            self.calls
1871                .lock()
1872                .unwrap()
1873                .push(MockCall::Read(key.to_string()));
1874            self.store
1875                .lock()
1876                .unwrap()
1877                .get(&(scope.clone(), key.to_string()))
1878                .cloned()
1879                .ok_or_else(|| CloudHomeError::NotFound(key.to_string()))
1880        }
1881
1882        fn list_records(
1883            &self,
1884            scope: &CloudKitScope,
1885            prefix: &str,
1886        ) -> Result<Vec<String>, CloudHomeError> {
1887            self.calls
1888                .lock()
1889                .unwrap()
1890                .push(MockCall::List(prefix.to_string()));
1891            let store = self.store.lock().unwrap();
1892            let mut keys: Vec<String> = store
1893                .keys()
1894                .filter(|(record_scope, key)| record_scope == scope && key.starts_with(prefix))
1895                .map(|(_, key)| key.clone())
1896                .collect();
1897            keys.sort();
1898            Ok(keys)
1899        }
1900
1901        fn delete_record(&self, scope: &CloudKitScope, key: &str) -> Result<(), CloudHomeError> {
1902            self.calls
1903                .lock()
1904                .unwrap()
1905                .push(MockCall::Delete(key.to_string()));
1906            if self.fail_deletes.lock().unwrap().contains(key) {
1907                return Err(CloudHomeError::Transport(format!("delete {key} failed")));
1908            }
1909            if let Some(remaining) = self.fail_delete_once.lock().unwrap().get_mut(key) {
1910                if *remaining > 0 {
1911                    *remaining -= 1;
1912                    return Err(CloudHomeError::Transport(format!("delete {key} failed")));
1913                }
1914            }
1915            self.store
1916                .lock()
1917                .unwrap()
1918                .remove(&(scope.clone(), key.to_string()));
1919            self.versions
1920                .lock()
1921                .unwrap()
1922                .remove(&(scope.clone(), key.to_string()));
1923            Ok(())
1924        }
1925
1926        fn record_exists(&self, scope: &CloudKitScope, key: &str) -> Result<bool, CloudHomeError> {
1927            self.record_exists_calls.fetch_add(1, Ordering::Relaxed);
1928            self.calls
1929                .lock()
1930                .unwrap()
1931                .push(MockCall::Exists(key.to_string()));
1932            Ok(self
1933                .store
1934                .lock()
1935                .unwrap()
1936                .contains_key(&(scope.clone(), key.to_string())))
1937        }
1938
1939        fn read_versioned_record(
1940            &self,
1941            scope: &CloudKitScope,
1942            key: &str,
1943        ) -> Result<CloudVersionedHead, CloudHomeError> {
1944            let record = (scope.clone(), key.to_string());
1945            let store = self.store.lock().unwrap();
1946            let versions = self.versions.lock().unwrap();
1947            Ok(CloudVersionedHead {
1948                bytes: store
1949                    .get(&record)
1950                    .cloned()
1951                    .ok_or_else(|| CloudHomeError::NotFound(key.to_string()))?,
1952                version: CloudHeadVersion::from_provider(
1953                    versions
1954                        .get(&record)
1955                        .copied()
1956                        .ok_or_else(|| CloudHomeError::NotFound(key.to_string()))?
1957                        .to_string(),
1958                )?,
1959            })
1960        }
1961
1962        fn create_record(
1963            &self,
1964            scope: &CloudKitScope,
1965            key: &str,
1966            data: Vec<u8>,
1967        ) -> Result<CloudVersionedHead, CloudHeadCreateError> {
1968            self.calls
1969                .lock()
1970                .unwrap()
1971                .push(MockCall::Create(key.to_string()));
1972            if self.fail_writes.lock().unwrap().contains(key) {
1973                return Err(CloudHeadCreateError::Storage(CloudHomeError::Transport(
1974                    format!("create {key} failed"),
1975                )));
1976            }
1977            let record = (scope.clone(), key.to_string());
1978            let mut store = self.store.lock().unwrap();
1979            let mut versions = self.versions.lock().unwrap();
1980            if store.contains_key(&record) {
1981                return Err(CloudHeadCreateError::AlreadyExists);
1982            }
1983            store.insert(record.clone(), data.clone());
1984            versions.insert(record, 1);
1985            Ok(CloudVersionedHead {
1986                bytes: data,
1987                version: CloudHeadVersion::from_provider("1".to_string())?,
1988            })
1989        }
1990
1991        fn begin_atomic_create(
1992            &self,
1993            _scope: &CloudKitScope,
1994        ) -> Result<CloudKitAtomicCreateBatch, CloudHomeError> {
1995            let batch = CloudKitAtomicCreateBatch::from_provider(format!(
1996                "batch-{}",
1997                self.next_batch.fetch_add(1, Ordering::SeqCst)
1998            ))?;
1999            self.calls
2000                .lock()
2001                .unwrap()
2002                .push(MockCall::BeginBatch(batch.as_provider().to_string()));
2003            self.staged_batches
2004                .lock()
2005                .unwrap()
2006                .insert(batch.as_provider().to_string(), Vec::new());
2007            Ok(batch)
2008        }
2009
2010        fn stage_atomic_create_record(
2011            &self,
2012            _scope: &CloudKitScope,
2013            batch: &CloudKitAtomicCreateBatch,
2014            record: CloudKitRecordCreate,
2015        ) -> Result<(), CloudHomeError> {
2016            self.calls
2017                .lock()
2018                .unwrap()
2019                .push(MockCall::Stage(record.key.clone()));
2020            self.max_stage_payload
2021                .fetch_max(record.data.len(), Ordering::SeqCst);
2022            self.staged_batches
2023                .lock()
2024                .unwrap()
2025                .get_mut(batch.as_provider())
2026                .ok_or_else(|| {
2027                    CloudHomeError::NotFound(format!(
2028                        "CloudKit staging batch {:?}",
2029                        batch.as_provider()
2030                    ))
2031                })?
2032                .push(record);
2033            Ok(())
2034        }
2035
2036        fn commit_atomic_create(
2037            &self,
2038            scope: &CloudKitScope,
2039            batch: &CloudKitAtomicCreateBatch,
2040        ) -> Result<Vec<CloudKitRecordVersion>, CloudHomeError> {
2041            self.calls
2042                .lock()
2043                .unwrap()
2044                .push(MockCall::CommitBatch(batch.as_provider().to_string()));
2045            let mut batches = self.staged_batches.lock().unwrap();
2046            let records = batches.get(batch.as_provider()).ok_or_else(|| {
2047                CloudHomeError::NotFound(format!(
2048                    "CloudKit staging batch {:?}",
2049                    batch.as_provider()
2050                ))
2051            })?;
2052            let fail_writes = self.fail_writes.lock().unwrap();
2053            let mut store = self.store.lock().unwrap();
2054            let mut versions = self.versions.lock().unwrap();
2055            for record in records {
2056                if fail_writes.contains(&record.key) {
2057                    return Err(CloudHomeError::Transport(format!(
2058                        "atomic create {:?} failed",
2059                        record.key
2060                    )));
2061                }
2062                if store.contains_key(&(scope.clone(), record.key.clone())) {
2063                    return Err(CloudHomeError::AlreadyExists(record.key.clone()));
2064                }
2065            }
2066            let records = batches
2067                .remove(batch.as_provider())
2068                .expect("validated CloudKit staging batch disappeared");
2069            let mut created = Vec::with_capacity(records.len());
2070            for record in records {
2071                let coordinate = (scope.clone(), record.key.clone());
2072                store.insert(coordinate.clone(), record.data);
2073                versions.insert(coordinate, 1);
2074                created.push(CloudKitRecordVersion {
2075                    key: record.key,
2076                    version: CloudHeadVersion::from_provider("1".to_string())?,
2077                });
2078            }
2079            if self.lose_commit_response.load(Ordering::SeqCst) {
2080                return Err(CloudHomeError::Transport(
2081                    "CloudKit commit response was lost".to_string(),
2082                ));
2083            }
2084            if self.return_wrong_commit_keys.load(Ordering::SeqCst) {
2085                for (index, record) in created.iter_mut().enumerate() {
2086                    record.key = format!("unexpected-returned-record-{index}");
2087                }
2088            }
2089            Ok(created)
2090        }
2091
2092        fn discard_atomic_create(
2093            &self,
2094            _scope: &CloudKitScope,
2095            batch: &CloudKitAtomicCreateBatch,
2096        ) -> Result<(), CloudHomeError> {
2097            self.calls
2098                .lock()
2099                .unwrap()
2100                .push(MockCall::DiscardBatch(batch.as_provider().to_string()));
2101            if self.fail_discards.load(Ordering::SeqCst) {
2102                return Err(CloudHomeError::Transport(format!(
2103                    "discard staging batch {:?} failed",
2104                    batch.as_provider()
2105                )));
2106            }
2107            self.staged_batches
2108                .lock()
2109                .unwrap()
2110                .remove(batch.as_provider());
2111            Ok(())
2112        }
2113
2114        fn delete_record_versions(
2115            &self,
2116            scope: &CloudKitScope,
2117            records: &[CloudKitRecordVersion],
2118        ) -> Result<(), CloudHomeError> {
2119            self.calls.lock().unwrap().push(MockCall::DeleteVersions(
2120                records.iter().map(|record| record.key.clone()).collect(),
2121            ));
2122            let fail_deletes = self.fail_deletes.lock().unwrap();
2123            let mut store = self.store.lock().unwrap();
2124            let mut versions = self.versions.lock().unwrap();
2125            for record in records {
2126                if fail_deletes.contains(&record.key) {
2127                    return Err(CloudHomeError::Transport(format!(
2128                        "delete {:?} failed",
2129                        record.key
2130                    )));
2131                }
2132                let storage_key = (scope.clone(), record.key.clone());
2133                let current = versions
2134                    .get(&storage_key)
2135                    .ok_or_else(|| CloudHomeError::NotFound(record.key.clone()))?;
2136                if current.to_string() != record.version.as_provider() {
2137                    return Err(CloudHomeError::Transport(format!(
2138                        "CloudKit record {:?} changed before exact deletion",
2139                        record.key
2140                    )));
2141                }
2142                if !store.contains_key(&storage_key) {
2143                    return Err(CloudHomeError::NotFound(record.key.clone()));
2144                }
2145            }
2146            for record in records {
2147                let storage_key = (scope.clone(), record.key.clone());
2148                store.remove(&storage_key);
2149                versions.remove(&storage_key);
2150            }
2151            Ok(())
2152        }
2153
2154        fn replace_record(
2155            &self,
2156            scope: &CloudKitScope,
2157            key: &str,
2158            expected: &CloudHeadVersion,
2159            data: Vec<u8>,
2160        ) -> Result<CloudVersionedHead, CloudHeadReplaceError> {
2161            let record = (scope.clone(), key.to_string());
2162            let mut store = self.store.lock().unwrap();
2163            let mut versions = self.versions.lock().unwrap();
2164            let current = versions
2165                .get(&record)
2166                .copied()
2167                .ok_or(CloudHeadReplaceError::VersionMismatch)?;
2168            if expected.as_provider() != current.to_string() {
2169                return Err(CloudHeadReplaceError::VersionMismatch);
2170            }
2171            let next = current + 1;
2172            store.insert(record.clone(), data.clone());
2173            versions.insert(record, next);
2174            Ok(CloudVersionedHead {
2175                bytes: data,
2176                version: CloudHeadVersion::from_provider(next.to_string())?,
2177            })
2178        }
2179
2180        fn grant_share(&self, member_pubkey: &str) -> Result<CloudKitShare, CloudHomeError> {
2181            self.grant_share_calls.fetch_add(1, Ordering::Relaxed);
2182            let share = CloudKitShare {
2183                share_url: format!("https://share.example/{member_pubkey}"),
2184                owner_name: "owner-name".to_string(),
2185                zone_name: "bae-store".to_string(),
2186            };
2187            self.shares
2188                .lock()
2189                .unwrap()
2190                .insert(member_pubkey.to_string(), share.clone());
2191            Ok(share)
2192        }
2193
2194        fn share_for_member(
2195            &self,
2196            member_pubkey: &str,
2197        ) -> Result<Option<CloudKitShare>, CloudHomeError> {
2198            Ok(self.shares.lock().unwrap().get(member_pubkey).cloned())
2199        }
2200
2201        fn revoke_share(&self, member_pubkey: &str) -> Result<(), CloudHomeError> {
2202            self.revoke_share_calls.fetch_add(1, Ordering::Relaxed);
2203            self.shares.lock().unwrap().remove(member_pubkey);
2204            Ok(())
2205        }
2206
2207        fn accept_share(&self, share_url: &str) -> Result<CloudKitShare, CloudHomeError> {
2208            Ok(CloudKitShare {
2209                share_url: share_url.to_string(),
2210                owner_name: "owner-name".to_string(),
2211                zone_name: "bae-store".to_string(),
2212            })
2213        }
2214    }
2215
2216    fn make_cloud_home() -> CloudKitCloudHome {
2217        CloudKitCloudHome::new_private_with_ids(
2218            Arc::new(MockCloudKitOps::new()),
2219            Arc::new(SequentialIdProvider::new("cloudkit-upload")),
2220        )
2221    }
2222
2223    fn make_cloud_home_with_ops() -> (CloudKitCloudHome, Arc<MockCloudKitOps>) {
2224        let ops = Arc::new(MockCloudKitOps::new());
2225        (
2226            CloudKitCloudHome::new_private_with_ids(
2227                ops.clone(),
2228                Arc::new(SequentialIdProvider::new("cloudkit-upload")),
2229            ),
2230            ops,
2231        )
2232    }
2233
2234    #[tokio::test]
2235    async fn provider_binding_uses_the_bridge_container_zone_and_current_user() {
2236        use coven_core::sync::storage::{ProviderPrincipalId, StoreProviderBinding};
2237        let (home, _) = make_cloud_home_with_ops();
2238
2239        let binding = ExactSlotStorage::provider_binding(&home)
2240            .await
2241            .expect("resolve CloudKit provider binding");
2242
2243        assert_eq!(
2244            binding.store,
2245            StoreProviderBinding::CloudKit {
2246                container_id: "iCloud.example.coven".to_string(),
2247                environment: coven_core::sync::storage::CloudKitEnvironment::Development,
2248                owner_name: "private-owner".to_string(),
2249                zone_name: "private-zone".to_string(),
2250            }
2251        );
2252        assert_eq!(
2253            binding.device.principal,
2254            ProviderPrincipalId::CloudKitPrivateZoneOwner {
2255                record_name: "current-user".to_string(),
2256            }
2257        );
2258    }
2259
2260    fn write_chunk_manifest(ops: &MockCloudKitOps, key: &str, total_len: usize) {
2261        write_chunk_manifest_with_upload_id(
2262            ops,
2263            key,
2264            total_len,
2265            "0123456789abcdef0123456789abcdef",
2266        );
2267    }
2268
2269    fn write_chunk_manifest_with_upload_id(
2270        ops: &MockCloudKitOps,
2271        key: &str,
2272        total_len: usize,
2273        upload_id: &str,
2274    ) {
2275        ops.write_record(
2276            &CloudKitScope::Private,
2277            &chunk_manifest_key(key),
2278            encode_chunk_manifest(ChunkManifest::new(total_len, upload_id.to_string())),
2279        )
2280        .unwrap();
2281    }
2282
2283    fn write_chunk_part(ops: &MockCloudKitOps, key: &str, index: usize, data: Vec<u8>) {
2284        ops.write_record(
2285            &CloudKitScope::Private,
2286            &chunk_part_key(key, "0123456789abcdef0123456789abcdef", index),
2287            data,
2288        )
2289        .unwrap();
2290    }
2291
2292    struct FailingBodyReader {
2293        emitted: bool,
2294    }
2295
2296    #[async_trait]
2297    impl crate::local_blob::PlaintextChunkReader for FailingBodyReader {
2298        async fn next_chunk(
2299            &mut self,
2300            _max: usize,
2301        ) -> Result<Vec<u8>, crate::local_blob::PlaintextChunkError> {
2302            if !self.emitted {
2303                self.emitted = true;
2304                return Ok(vec![7; CHUNK_SIZE]);
2305            }
2306            Err(crate::local_blob::PlaintextChunkError::Local(
2307                "injected body failure".to_string(),
2308            ))
2309        }
2310    }
2311
2312    struct PausedBodyReader {
2313        emitted: bool,
2314        waiting: Arc<tokio::sync::Notify>,
2315        release: Arc<tokio::sync::Notify>,
2316    }
2317
2318    #[async_trait]
2319    impl crate::local_blob::PlaintextChunkReader for PausedBodyReader {
2320        async fn next_chunk(
2321            &mut self,
2322            _max: usize,
2323        ) -> Result<Vec<u8>, crate::local_blob::PlaintextChunkError> {
2324            if !self.emitted {
2325                self.emitted = true;
2326                return Ok(vec![7; CHUNK_SIZE]);
2327            }
2328            self.waiting.notify_one();
2329            self.release.notified().await;
2330            Ok(vec![8])
2331        }
2332    }
2333
2334    #[tokio::test]
2335    async fn mutable_body_failure_reports_cleanup_failure_and_drop_retries_cleanup() {
2336        let (home, ops) = make_cloud_home_with_ops();
2337        let part_key = chunk_part_key("mutable/body-failure", "cloudkit-upload-0", 0);
2338        ops.fail_next_delete(&part_key);
2339        let reader = crate::local_blob::PlaintextReader::from_test_reader(FailingBodyReader {
2340            emitted: false,
2341        });
2342        let body = BlobBody::from_test_reader((CHUNK_SIZE + 1) as u64, reader);
2343
2344        let error = home
2345            .write("mutable/body-failure", body, &no_progress())
2346            .await
2347            .expect_err("body failure must report failed cleanup");
2348
2349        assert!(
2350            matches!(error, CloudHomeError::CleanupFailed { .. }),
2351            "{error}"
2352        );
2353        assert!(
2354            error.to_string().contains("injected body failure"),
2355            "{error}"
2356        );
2357        assert!(error.to_string().contains("delete"), "{error}");
2358        assert!(!ops
2359            .record_exists(&CloudKitScope::Private, &part_key)
2360            .expect("inspect canceled part"));
2361    }
2362
2363    #[tokio::test]
2364    async fn canceling_mutable_write_removes_every_staged_part() {
2365        let (home, ops) = make_cloud_home_with_ops();
2366        let waiting = Arc::new(tokio::sync::Notify::new());
2367        let release = Arc::new(tokio::sync::Notify::new());
2368        let reader = crate::local_blob::PlaintextReader::from_test_reader(PausedBodyReader {
2369            emitted: false,
2370            waiting: waiting.clone(),
2371            release: release.clone(),
2372        });
2373        let body = BlobBody::from_test_reader((CHUNK_SIZE + 1) as u64, reader);
2374        let write =
2375            tokio::spawn(async move { home.write("mutable/cancel", body, &no_progress()).await });
2376        waiting.notified().await;
2377
2378        write.abort();
2379        assert!(write.await.expect_err("write task canceled").is_cancelled());
2380        release.notify_waiters();
2381
2382        let part_key = chunk_part_key("mutable/cancel", "cloudkit-upload-0", 0);
2383        assert!(!ops
2384            .record_exists(&CloudKitScope::Private, &part_key)
2385            .expect("inspect canceled part"));
2386    }
2387
2388    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2389    async fn cancel_after_mutable_manifest_publish_preserves_the_committed_layout() {
2390        let (home, ops) = make_cloud_home_with_ops();
2391        let key = "mutable/published";
2392        let data = vec![9; CHUNK_SIZE + 1];
2393        let (stored, release) = ops.pause_write_after_store(&chunk_manifest_key(key));
2394        let write_home = home.clone();
2395        let write_data = data.clone();
2396        let write = tokio::spawn(async move {
2397            write_home
2398                .write(key, BlobBody::from_bytes(write_data), &no_progress())
2399                .await
2400        });
2401        tokio::task::spawn_blocking(move || stored.wait())
2402            .await
2403            .expect("wait for manifest publication");
2404
2405        write.abort();
2406        tokio::task::spawn_blocking(move || release.wait())
2407            .await
2408            .expect("release manifest publication");
2409        assert!(write.await.expect_err("write task canceled").is_cancelled());
2410
2411        assert_eq!(home.read(key).await.expect("read committed layout"), data);
2412        assert!(ops
2413            .record_exists(&CloudKitScope::Private, &chunk_manifest_key(key))
2414            .expect("inspect committed manifest"));
2415        assert_eq!(
2416            ops.list_records(&CloudKitScope::Private, &format!("{key}.part"))
2417                .expect("inspect committed parts")
2418                .len(),
2419            2
2420        );
2421    }
2422
2423    #[test]
2424    fn mutable_cancellation_cleanup_failure_terminates_the_process() {
2425        const CHILD: &str = "COVEN_CLOUDKIT_MUTABLE_CANCEL_CHILD";
2426        if std::env::var_os(CHILD).is_some() {
2427            let runtime = tokio::runtime::Runtime::new().expect("build child runtime");
2428            runtime.block_on(async {
2429                let (home, ops) = make_cloud_home_with_ops();
2430                let mut sink = home
2431                    .open_multipart("mutable/cancel", (CHUNK_SIZE + 1) as u64)
2432                    .await
2433                    .expect("open CloudKit multipart upload");
2434                sink.send_part(Bytes::from(vec![7; CHUNK_SIZE]), 0, false)
2435                    .await
2436                    .expect("write first multipart part");
2437                ops.fail_delete(&chunk_part_key("mutable/cancel", "cloudkit-upload-0", 0));
2438                drop(sink);
2439            });
2440            std::process::exit(0);
2441        }
2442
2443        let status = std::process::Command::new(
2444            std::env::current_exe().expect("locate CloudKit test executable"),
2445        )
2446        .arg("mutable_cancellation_cleanup_failure_terminates_the_process")
2447        .arg("--nocapture")
2448        .env(CHILD, "1")
2449        .status()
2450        .expect("run CloudKit mutable cancellation subprocess");
2451        assert!(!status.success(), "cancellation subprocess survived");
2452    }
2453
2454    #[tokio::test]
2455    async fn write_reports_progress_per_chunk_record() {
2456        use std::sync::atomic::{AtomicU64, Ordering};
2457        let ch = make_cloud_home();
2458        // 25 MB spans three records (10 + 10 + 5) so progress fires three
2459        // times, the last equalling the total.
2460        let total = 25 * 1024 * 1024u64;
2461        let data: Vec<u8> = vec![0u8; total as usize];
2462        let last = Arc::new(AtomicU64::new(0));
2463        let ticks = Arc::new(AtomicU64::new(0));
2464        let last2 = last.clone();
2465        let ticks2 = ticks.clone();
2466        let sink = move |n: u64| {
2467            last2.store(n, Ordering::Relaxed);
2468            ticks2.fetch_add(1, Ordering::Relaxed);
2469        };
2470        ch.write("chunked.bin", BlobBody::from_bytes(data), &sink)
2471            .await
2472            .unwrap();
2473        assert_eq!(last.load(Ordering::Relaxed), total);
2474        assert_eq!(ticks.load(Ordering::Relaxed), 3);
2475    }
2476
2477    #[tokio::test]
2478    async fn test_small_file_roundtrip() {
2479        let ch = make_cloud_home();
2480        let data = b"hello world".to_vec();
2481        ch.write(
2482            "small.bin",
2483            BlobBody::from_bytes(data.clone()),
2484            &no_progress(),
2485        )
2486        .await
2487        .unwrap();
2488        let read = ch.read("small.bin").await.unwrap();
2489        assert_eq!(read, data);
2490    }
2491
2492    #[tokio::test]
2493    async fn test_large_file_roundtrip() {
2494        let ch = make_cloud_home();
2495        // 25MB of data -- spans 3 chunks (10 + 10 + 5)
2496        let data: Vec<u8> = (0..25 * 1024 * 1024).map(|i| (i % 256) as u8).collect();
2497        ch.write(
2498            "large.bin",
2499            BlobBody::from_bytes(data.clone()),
2500            &no_progress(),
2501        )
2502        .await
2503        .unwrap();
2504        let read = ch.read("large.bin").await.unwrap();
2505        assert_eq!(read.len(), data.len());
2506        assert_eq!(read, data);
2507    }
2508
2509    #[tokio::test]
2510    async fn test_read_range_single() {
2511        let ch = make_cloud_home();
2512        ch.write(
2513            "range.bin",
2514            BlobBody::from_bytes(b"0123456789".to_vec()),
2515            &no_progress(),
2516        )
2517        .await
2518        .unwrap();
2519        let slice = ch.read_range("range.bin", 3, 7).await.unwrap();
2520        assert_eq!(slice, b"3456");
2521    }
2522
2523    #[tokio::test]
2524    async fn read_single_record_does_not_probe_existence() {
2525        let (ch, ops) = make_cloud_home_with_ops();
2526        let data = b"hello world".to_vec();
2527        ch.write(
2528            "single.bin",
2529            BlobBody::from_bytes(data.clone()),
2530            &no_progress(),
2531        )
2532        .await
2533        .unwrap();
2534
2535        let read = ch.read("single.bin").await.unwrap();
2536
2537        assert_eq!(read, data);
2538        assert_eq!(ops.record_exists_calls.load(Ordering::Relaxed), 0);
2539    }
2540
2541    #[tokio::test]
2542    async fn read_range_single_record_does_not_probe_existence() {
2543        let (ch, ops) = make_cloud_home_with_ops();
2544        ch.write(
2545            "single-range.bin",
2546            BlobBody::from_bytes(b"0123456789".to_vec()),
2547            &no_progress(),
2548        )
2549        .await
2550        .unwrap();
2551
2552        let read = ch.read_range("single-range.bin", 2, 6).await.unwrap();
2553
2554        assert_eq!(read, b"2345");
2555        assert_eq!(ops.record_exists_calls.load(Ordering::Relaxed), 0);
2556    }
2557
2558    #[tokio::test]
2559    async fn read_chunked_record_without_manifest_errors() {
2560        let (ch, ops) = make_cloud_home_with_ops();
2561        let first = vec![1u8; CHUNK_SIZE];
2562        let second = b"tail".to_vec();
2563        write_chunk_part(&ops, "chunked.bin", 0, first.clone());
2564        write_chunk_part(&ops, "chunked.bin", 1, second.clone());
2565
2566        let err = ch
2567            .read("chunked.bin")
2568            .await
2569            .expect_err("chunks without manifest must fail");
2570        let msg = err.to_string();
2571
2572        assert!(
2573            msg.contains("chunked.bin") && msg.contains("no manifest"),
2574            "unexpected error: {msg}"
2575        );
2576        assert!(!ch.exists("chunked.bin").await.unwrap());
2577    }
2578
2579    #[tokio::test]
2580    async fn list_omits_base_key_whose_manifest_is_absent() {
2581        let (ch, ops) = make_cloud_home_with_ops();
2582        // Part records with no manifest — an interrupted upload that never
2583        // published. `read` cannot assemble them, so `list` must not report them.
2584        write_chunk_part(&ops, "files/orphan.bin", 0, vec![1u8; CHUNK_SIZE]);
2585        write_chunk_part(&ops, "files/orphan.bin", 1, b"tail".to_vec());
2586        ch.write(
2587            "files/ok.bin",
2588            BlobBody::from_bytes(b"hi".to_vec()),
2589            &no_progress(),
2590        )
2591        .await
2592        .unwrap();
2593
2594        let keys = ch.list("files/").await.unwrap();
2595
2596        assert_eq!(keys, vec!["files/ok.bin".to_string()]);
2597    }
2598
2599    #[tokio::test]
2600    async fn multipart_part_failure_leaves_no_orphan_records_or_visibility() {
2601        let (ch, ops) = make_cloud_home_with_ops();
2602        // 25 MB spans three parts; fail the second part write mid-upload. The
2603        // upload id is the first id the sequential provider hands out.
2604        ops.fail_write(&chunk_part_key("orphan.bin", "cloudkit-upload-0", 1));
2605        let data: Vec<u8> = vec![0u8; 25 * 1024 * 1024];
2606
2607        let err = ch
2608            .write("orphan.bin", BlobBody::from_bytes(data), &no_progress())
2609            .await
2610            .expect_err("injected part write failure must fail the upload");
2611        assert!(err.to_string().contains("write"), "unexpected error: {err}");
2612
2613        assert!(!ch.exists("orphan.bin").await.unwrap());
2614        assert!(!ch
2615            .list("")
2616            .await
2617            .unwrap()
2618            .contains(&"orphan.bin".to_string()));
2619        assert!(
2620            ops.list_records(&CloudKitScope::Private, "orphan.bin.part")
2621                .unwrap()
2622                .is_empty(),
2623            "aborted upload must leave no part records"
2624        );
2625        assert!(!ops
2626            .record_exists(&CloudKitScope::Private, &chunk_manifest_key("orphan.bin"))
2627            .unwrap());
2628    }
2629
2630    #[tokio::test]
2631    async fn read_chunked_record_with_missing_manifest_part_errors() {
2632        let (ch, ops) = make_cloud_home_with_ops();
2633        let first = vec![1u8; CHUNK_SIZE];
2634        let second = vec![2u8; CHUNK_SIZE];
2635        let total_len = (CHUNK_SIZE * 2) + 4;
2636        write_chunk_manifest(&ops, "chunked.bin", total_len);
2637        write_chunk_part(&ops, "chunked.bin", 0, first);
2638        write_chunk_part(&ops, "chunked.bin", 1, second);
2639
2640        let err = ch
2641            .read("chunked.bin")
2642            .await
2643            .expect_err("missing manifest part must fail");
2644        let msg = err.to_string();
2645
2646        assert!(
2647            msg.contains("expects 3 parts") && msg.contains("found 2"),
2648            "unexpected error: {msg}"
2649        );
2650        assert!(!ch.exists("chunked.bin").await.unwrap());
2651    }
2652
2653    #[tokio::test]
2654    async fn read_range_chunked_rejects_range_past_manifest_length() {
2655        let ch = make_cloud_home();
2656        let data: Vec<u8> = vec![7u8; 15 * 1024 * 1024];
2657        ch.write(
2658            "range-limit.bin",
2659            BlobBody::from_bytes(data),
2660            &no_progress(),
2661        )
2662        .await
2663        .unwrap();
2664
2665        let err = ch
2666            .read_range("range-limit.bin", 0, (16 * 1024 * 1024) as u64)
2667            .await
2668            .expect_err("range past manifest length must fail");
2669        let msg = err.to_string();
2670
2671        assert!(msg.contains("exceeds file size"), "unexpected error: {msg}");
2672    }
2673
2674    #[tokio::test]
2675    async fn read_range_chunked_short_chunk_errors_instead_of_panicking() {
2676        let (ch, ops) = make_cloud_home_with_ops();
2677        let total_len = CHUNK_SIZE + 8;
2678        write_chunk_manifest(&ops, "short-tail.bin", total_len);
2679        write_chunk_part(&ops, "short-tail.bin", 0, vec![1u8; CHUNK_SIZE]);
2680        write_chunk_part(&ops, "short-tail.bin", 1, vec![2u8; 4]);
2681
2682        let err = ch
2683            .read_range("short-tail.bin", CHUNK_SIZE as u64, (CHUNK_SIZE + 8) as u64)
2684            .await
2685            .expect_err("short tail chunk must fail");
2686        let msg = err.to_string();
2687
2688        assert!(
2689            msg.contains("part 1") && msg.contains("expected 8"),
2690            "unexpected error: {msg}"
2691        );
2692    }
2693
2694    #[tokio::test]
2695    async fn test_read_range_chunked() {
2696        let ch = make_cloud_home();
2697        // Create data that spans 2 chunks: 15MB
2698        let data: Vec<u8> = (0..15 * 1024 * 1024).map(|i| (i % 256) as u8).collect();
2699        ch.write(
2700            "big.bin",
2701            BlobBody::from_bytes(data.clone()),
2702            &no_progress(),
2703        )
2704        .await
2705        .unwrap();
2706
2707        // Read a range that crosses the chunk boundary (last byte of chunk 0, first byte of chunk 1)
2708        let boundary = CHUNK_SIZE;
2709        let start = (boundary - 2) as u64;
2710        let end = (boundary + 3) as u64;
2711        let slice = ch.read_range("big.bin", start, end).await.unwrap();
2712        assert_eq!(slice.len(), 5);
2713        assert_eq!(slice, &data[start as usize..end as usize]);
2714    }
2715
2716    #[tokio::test]
2717    async fn test_list_deduplicates_chunks() {
2718        let ch = make_cloud_home();
2719        // Write a chunked file
2720        let data: Vec<u8> = vec![0u8; 25 * 1024 * 1024];
2721        ch.write(
2722            "files/album.flac",
2723            BlobBody::from_bytes(data),
2724            &no_progress(),
2725        )
2726        .await
2727        .unwrap();
2728
2729        // Also write a small file
2730        ch.write(
2731            "files/cover.jpg",
2732            BlobBody::from_bytes(b"img".to_vec()),
2733            &no_progress(),
2734        )
2735        .await
2736        .unwrap();
2737
2738        let keys = ch.list("files/").await.unwrap();
2739        assert_eq!(keys.len(), 2);
2740        assert!(keys.contains(&"files/album.flac".to_string()));
2741        assert!(keys.contains(&"files/cover.jpg".to_string()));
2742    }
2743
2744    #[tokio::test]
2745    async fn test_delete_removes_all_chunks() {
2746        let ch = make_cloud_home();
2747        let data: Vec<u8> = vec![0u8; 25 * 1024 * 1024];
2748        ch.write("to-delete.bin", BlobBody::from_bytes(data), &no_progress())
2749            .await
2750            .unwrap();
2751
2752        assert!(ch.exists("to-delete.bin").await.unwrap());
2753
2754        ch.delete("to-delete.bin").await.unwrap();
2755
2756        assert!(!ch.exists("to-delete.bin").await.unwrap());
2757
2758        // Verify the underlying ops store is empty of related keys
2759        let ops = &ch.ops;
2760        let keys = ops
2761            .list_records(&CloudKitScope::Private, "to-delete.bin")
2762            .unwrap();
2763        assert!(keys.is_empty());
2764    }
2765
2766    #[tokio::test]
2767    async fn test_overwrite_chunked_with_single() {
2768        let ch = make_cloud_home();
2769        // Write large file (chunked)
2770        let large_data: Vec<u8> = vec![0u8; 25 * 1024 * 1024];
2771        ch.write("file.bin", BlobBody::from_bytes(large_data), &no_progress())
2772            .await
2773            .unwrap();
2774
2775        // Overwrite with small file (single record)
2776        let small_data = b"small".to_vec();
2777        ch.write(
2778            "file.bin",
2779            BlobBody::from_bytes(small_data.clone()),
2780            &no_progress(),
2781        )
2782        .await
2783        .unwrap();
2784
2785        let read = ch.read("file.bin").await.unwrap();
2786        assert_eq!(read, small_data);
2787
2788        // Verify no chunk records remain
2789        let chunks = ch
2790            .ops
2791            .list_records(&CloudKitScope::Private, "file.bin.part")
2792            .unwrap();
2793        assert!(chunks.is_empty());
2794    }
2795
2796    #[tokio::test]
2797    async fn put_object_over_single_writes_without_deleting_base_first() {
2798        let (ch, ops) = make_cloud_home_with_ops();
2799        ch.write(
2800            "file.bin",
2801            BlobBody::from_bytes(b"old".to_vec()),
2802            &no_progress(),
2803        )
2804        .await
2805        .unwrap();
2806        ops.clear_calls();
2807
2808        ch.write(
2809            "file.bin",
2810            BlobBody::from_bytes(b"new".to_vec()),
2811            &no_progress(),
2812        )
2813        .await
2814        .unwrap();
2815
2816        let calls = ops.calls();
2817        assert_eq!(
2818            calls.first(),
2819            Some(&MockCall::Write("file.bin".to_string()))
2820        );
2821        assert!(
2822            !calls.contains(&MockCall::Delete("file.bin".to_string())),
2823            "single-record overwrite must not delete the base record: {calls:?}"
2824        );
2825        assert_eq!(ch.read("file.bin").await.unwrap(), b"new");
2826    }
2827
2828    #[tokio::test]
2829    async fn put_object_over_chunked_publishes_single_before_cleanup() {
2830        let (ch, ops) = make_cloud_home_with_ops();
2831        let large_data: Vec<u8> = vec![0u8; 25 * 1024 * 1024];
2832        ch.write("file.bin", BlobBody::from_bytes(large_data), &no_progress())
2833            .await
2834            .unwrap();
2835        ops.clear_calls();
2836
2837        ch.write(
2838            "file.bin",
2839            BlobBody::from_bytes(b"new".to_vec()),
2840            &no_progress(),
2841        )
2842        .await
2843        .unwrap();
2844
2845        let calls = ops.calls();
2846        assert_eq!(
2847            calls.first(),
2848            Some(&MockCall::Write("file.bin".to_string()))
2849        );
2850        assert_eq!(ch.read("file.bin").await.unwrap(), b"new");
2851        assert!(ch
2852            .ops
2853            .list_records(&CloudKitScope::Private, "file.bin.part")
2854            .unwrap()
2855            .is_empty());
2856        assert!(!ch
2857            .ops
2858            .record_exists(&CloudKitScope::Private, &chunk_manifest_key("file.bin"))
2859            .unwrap());
2860    }
2861
2862    #[tokio::test]
2863    async fn put_object_cleanup_failure_leaves_new_single_readable() {
2864        let (ch, ops) = make_cloud_home_with_ops();
2865        let large_data: Vec<u8> = vec![0u8; 15 * 1024 * 1024];
2866        ch.write("file.bin", BlobBody::from_bytes(large_data), &no_progress())
2867            .await
2868            .unwrap();
2869        let stale_chunk = ch
2870            .ops
2871            .list_records(&CloudKitScope::Private, "file.bin.part")
2872            .unwrap()
2873            .into_iter()
2874            .next()
2875            .expect("chunked setup writes a chunk");
2876        ops.fail_delete(&stale_chunk);
2877
2878        let err = ch
2879            .write(
2880                "file.bin",
2881                BlobBody::from_bytes(b"new".to_vec()),
2882                &no_progress(),
2883            )
2884            .await
2885            .expect_err("stale chunk cleanup failure must fail loud");
2886        let msg = err.to_string();
2887
2888        assert!(msg.contains("delete"), "unexpected error: {msg}");
2889        assert_eq!(ch.read("file.bin").await.unwrap(), b"new");
2890    }
2891
2892    #[tokio::test]
2893    async fn test_overwrite_single_with_chunked() {
2894        let (ch, ops) = make_cloud_home_with_ops();
2895        // Write small file
2896        ch.write(
2897            "file.bin",
2898            BlobBody::from_bytes(b"small".to_vec()),
2899            &no_progress(),
2900        )
2901        .await
2902        .unwrap();
2903        ops.clear_calls();
2904
2905        // Overwrite with large file (chunked)
2906        let large_data: Vec<u8> = vec![1u8; 25 * 1024 * 1024];
2907        ch.write(
2908            "file.bin",
2909            BlobBody::from_bytes(large_data.clone()),
2910            &no_progress(),
2911        )
2912        .await
2913        .unwrap();
2914
2915        let read = ch.read("file.bin").await.unwrap();
2916        assert_eq!(read, large_data);
2917
2918        let calls = ops.calls();
2919        let manifest_write = calls
2920            .iter()
2921            .position(|call| *call == MockCall::Write(chunk_manifest_key("file.bin")))
2922            .expect("chunked write publishes manifest");
2923        let base_delete = calls
2924            .iter()
2925            .position(|call| *call == MockCall::Delete("file.bin".to_string()))
2926            .expect("chunked write removes stale single base");
2927        assert!(
2928            manifest_write < base_delete,
2929            "chunk manifest must publish before stale base cleanup: {calls:?}"
2930        );
2931
2932        // The single-record base is replaced by the chunk layout.
2933        assert!(!ch
2934            .ops
2935            .record_exists(&CloudKitScope::Private, "file.bin")
2936            .unwrap());
2937        assert!(ch
2938            .ops
2939            .record_exists(&CloudKitScope::Private, &chunk_manifest_key("file.bin"))
2940            .unwrap());
2941    }
2942
2943    #[tokio::test]
2944    async fn chunked_over_longer_chunked_uses_new_token_before_stale_cleanup() {
2945        let (ch, ops) = make_cloud_home_with_ops();
2946        let old_data: Vec<u8> = vec![0u8; 25 * 1024 * 1024];
2947        ch.write("file.bin", BlobBody::from_bytes(old_data), &no_progress())
2948            .await
2949            .unwrap();
2950        let old_chunks = ops
2951            .list_records(&CloudKitScope::Private, "file.bin.part")
2952            .unwrap();
2953        ops.clear_calls();
2954
2955        let new_data: Vec<u8> = vec![1u8; 15 * 1024 * 1024];
2956        ch.write(
2957            "file.bin",
2958            BlobBody::from_bytes(new_data.clone()),
2959            &no_progress(),
2960        )
2961        .await
2962        .unwrap();
2963
2964        assert_eq!(ch.read("file.bin").await.unwrap(), new_data);
2965        let remaining_chunks = ch
2966            .ops
2967            .list_records(&CloudKitScope::Private, "file.bin.part")
2968            .unwrap();
2969        assert_eq!(remaining_chunks.len(), 2);
2970        assert!(
2971            old_chunks
2972                .iter()
2973                .all(|old| !remaining_chunks.iter().any(|new| new == old)),
2974            "old token chunks must be cleaned after new manifest publishes"
2975        );
2976    }
2977
2978    #[tokio::test]
2979    async fn test_exists() {
2980        let ch = make_cloud_home();
2981
2982        assert!(!ch.exists("nope.bin").await.unwrap());
2983
2984        ch.write(
2985            "yep.bin",
2986            BlobBody::from_bytes(b"data".to_vec()),
2987            &no_progress(),
2988        )
2989        .await
2990        .unwrap();
2991        assert!(ch.exists("yep.bin").await.unwrap());
2992
2993        // Chunked file
2994        let data: Vec<u8> = vec![0u8; 15 * 1024 * 1024];
2995        ch.write("chunked.bin", BlobBody::from_bytes(data), &no_progress())
2996            .await
2997            .unwrap();
2998        assert!(ch.exists("chunked.bin").await.unwrap());
2999    }
3000
3001    #[tokio::test]
3002    async fn test_read_range_empty_when_end_leq_start() {
3003        let ch = make_cloud_home();
3004        ch.write(
3005            "range.bin",
3006            BlobBody::from_bytes(b"0123456789".to_vec()),
3007            &no_progress(),
3008        )
3009        .await
3010        .unwrap();
3011
3012        // end == start returns empty
3013        let slice = ch.read_range("range.bin", 3, 3).await.unwrap();
3014        assert!(slice.is_empty());
3015
3016        // end < start returns empty
3017        let slice = ch.read_range("range.bin", 5, 2).await.unwrap();
3018        assert!(slice.is_empty());
3019
3020        // end == 0 returns empty (the underflow case)
3021        let slice = ch.read_range("range.bin", 0, 0).await.unwrap();
3022        assert!(slice.is_empty());
3023    }
3024
3025    #[tokio::test]
3026    async fn exact_bounded_records_are_create_only() {
3027        let (home, ops) = make_cloud_home_with_ops();
3028        let slot = exact_slot("copies/bounded");
3029        ExactSlotStorage::create_at(
3030            &home,
3031            &slot,
3032            BlobBody::from_bytes(b"first".to_vec()),
3033            &no_progress(),
3034        )
3035        .await
3036        .unwrap();
3037
3038        let collision = ExactSlotStorage::create_at(
3039            &home,
3040            &slot,
3041            BlobBody::from_bytes(b"second".to_vec()),
3042            &no_progress(),
3043        )
3044        .await
3045        .expect_err("an immutable record must never overwrite an existing key");
3046        assert!(matches!(collision, CloudHomeError::AlreadyExists(key) if key == "copies/bounded"));
3047        assert_eq!(
3048            ExactSlotStorage::read_at(&home, &slot).await.unwrap(),
3049            b"first"
3050        );
3051
3052        ops.write_record(
3053            &CloudKitScope::Private,
3054            "copies/bounded",
3055            b"replacement".to_vec(),
3056        )
3057        .unwrap();
3058        let changed = ExactSlotStorage::read_at(&home, &slot)
3059            .await
3060            .expect_err("an exact read must reject a replaced manifest");
3061        assert!(changed.to_string().contains("invalid manifest"));
3062    }
3063
3064    #[tokio::test]
3065    async fn exact_multipart_stages_one_bounded_part_at_a_time_and_manifest_last() {
3066        let (home, ops) = make_cloud_home_with_ops();
3067        let data = vec![7u8; CHUNK_SIZE + 13];
3068        let slot = exact_slot("copies/chunked");
3069
3070        ExactSlotStorage::create_at(
3071            &home,
3072            &slot,
3073            BlobBody::from_bytes(data.clone()),
3074            &no_progress(),
3075        )
3076        .await
3077        .unwrap();
3078
3079        assert_eq!(
3080            ops.calls(),
3081            vec![
3082                MockCall::BeginBatch("batch-0".to_string()),
3083                MockCall::Stage(exact_part_key("copies/chunked", 0)),
3084                MockCall::Stage(exact_part_key("copies/chunked", 1)),
3085                MockCall::Stage("copies/chunked".to_string()),
3086                MockCall::CommitBatch("batch-0".to_string()),
3087            ]
3088        );
3089        assert_eq!(ops.max_stage_payload.load(Ordering::SeqCst), CHUNK_SIZE);
3090        assert_eq!(ExactSlotStorage::read_at(&home, &slot).await.unwrap(), data);
3091
3092        let manifest_bytes = ops
3093            .read_record(&CloudKitScope::Private, "copies/chunked")
3094            .unwrap();
3095        assert_eq!(
3096            decode_exact_manifest(&manifest_bytes).unwrap(),
3097            (2, data.len())
3098        );
3099    }
3100
3101    #[tokio::test]
3102    async fn lost_atomic_commit_response_is_settled_by_readback() {
3103        let (home, ops) = make_cloud_home_with_ops();
3104        ops.lose_commit_response();
3105        let data = vec![2u8; CHUNK_SIZE + 1];
3106        let slot = exact_slot("copies/ambiguous");
3107
3108        ExactSlotStorage::create_at(
3109            &home,
3110            &slot,
3111            BlobBody::from_bytes(data.clone()),
3112            &no_progress(),
3113        )
3114        .await
3115        .expect("authoritative readback settles a committed create");
3116
3117        assert_eq!(ExactSlotStorage::read_at(&home, &slot).await.unwrap(), data);
3118    }
3119
3120    #[tokio::test]
3121    async fn concurrent_immutable_creates_have_one_winner() {
3122        let (home, _) = make_cloud_home_with_ops();
3123        let slot = exact_slot("copies/create-race");
3124        let left_progress = no_progress();
3125        let right_progress = no_progress();
3126
3127        let (left, right) = tokio::join!(
3128            ExactSlotStorage::create_at(
3129                &home,
3130                &slot,
3131                BlobBody::from_bytes(b"left".to_vec()),
3132                &left_progress,
3133            ),
3134            ExactSlotStorage::create_at(
3135                &home,
3136                &slot,
3137                BlobBody::from_bytes(b"right".to_vec()),
3138                &right_progress,
3139            ),
3140        );
3141
3142        assert!(matches!(
3143            (&left, &right),
3144            (Ok(()), Err(CloudHomeError::AlreadyExists(_)))
3145                | (Err(CloudHomeError::AlreadyExists(_)), Ok(()))
3146        ));
3147        let expected = if left.is_ok() {
3148            b"left".as_slice()
3149        } else {
3150            b"right".as_slice()
3151        };
3152        assert_eq!(
3153            ExactSlotStorage::read_at(&home, &slot).await.unwrap(),
3154            expected
3155        );
3156    }
3157
3158    #[tokio::test]
3159    async fn mismatched_commit_keys_are_checked_against_authoritative_records() {
3160        let (home, ops) = make_cloud_home_with_ops();
3161        ops.return_wrong_commit_keys();
3162        let data = vec![3u8; CHUNK_SIZE + 1];
3163        let slot = exact_slot("copies/locator-mismatch");
3164
3165        ExactSlotStorage::create_at(
3166            &home,
3167            &slot,
3168            BlobBody::from_bytes(data.clone()),
3169            &no_progress(),
3170        )
3171        .await
3172        .expect("authoritative reads verify committed records");
3173
3174        assert_eq!(ExactSlotStorage::read_at(&home, &slot).await.unwrap(), data);
3175    }
3176
3177    #[tokio::test]
3178    async fn immutable_atomic_multipart_failure_and_collision_create_no_partial_layout() {
3179        let (home, ops) = make_cloud_home_with_ops();
3180        let first_part = exact_part_key("copies/failed", 0);
3181        let second_part = exact_part_key("copies/failed", 1);
3182        ops.fail_write(&second_part);
3183        let data = vec![3u8; CHUNK_SIZE + 1];
3184        let slot = exact_slot("copies/failed");
3185
3186        let error =
3187            ExactSlotStorage::create_at(&home, &slot, BlobBody::from_bytes(data), &no_progress())
3188                .await
3189                .expect_err("part creation failure must abort the append");
3190        assert!(matches!(error, CloudHomeError::Transport(_)));
3191        assert!(!ops
3192            .record_exists(&CloudKitScope::Private, &first_part)
3193            .unwrap());
3194        assert!(!ops
3195            .record_exists(&CloudKitScope::Private, "copies/failed")
3196            .unwrap());
3197        assert!(ops.staged_batches.lock().unwrap().is_empty());
3198
3199        let (home, ops) = make_cloud_home_with_ops();
3200        let first_part = exact_part_key("copies/collision", 0);
3201        let second_part = exact_part_key("copies/collision", 1);
3202        ops.create_record(&CloudKitScope::Private, &second_part, b"existing".to_vec())
3203            .unwrap();
3204        let slot = exact_slot("copies/collision");
3205        let error = ExactSlotStorage::create_at(
3206            &home,
3207            &slot,
3208            BlobBody::from_bytes(vec![3u8; CHUNK_SIZE + 1]),
3209            &no_progress(),
3210        )
3211        .await
3212        .expect_err("a batch collision must reject the whole append");
3213        assert!(matches!(error, CloudHomeError::AlreadyExists(key) if key == "copies/collision"));
3214        assert!(!ops
3215            .record_exists(&CloudKitScope::Private, &first_part)
3216            .unwrap());
3217        assert!(!ops
3218            .record_exists(&CloudKitScope::Private, "copies/collision")
3219            .unwrap());
3220        assert!(ops
3221            .record_exists(&CloudKitScope::Private, &second_part)
3222            .unwrap());
3223        assert!(ops.staged_batches.lock().unwrap().is_empty());
3224    }
3225
3226    #[tokio::test]
3227    async fn immutable_staging_cleanup_failure_is_typed_and_remote_state_stays_empty() {
3228        let (home, ops) = make_cloud_home_with_ops();
3229        let second_part = exact_part_key("copies/discard", 1);
3230        ops.fail_write(&second_part);
3231        ops.fail_discard();
3232        let slot = exact_slot("copies/discard");
3233
3234        let error = ExactSlotStorage::create_at(
3235            &home,
3236            &slot,
3237            BlobBody::from_bytes(vec![4u8; CHUNK_SIZE + 1]),
3238            &no_progress(),
3239        )
3240        .await
3241        .expect_err("failed staging discard must be returned with the commit error");
3242
3243        assert!(matches!(error, CloudHomeError::CleanupFailed { .. }));
3244        assert!(error.to_string().contains("batch-0"), "{error}");
3245        assert!(ops.store.lock().unwrap().is_empty());
3246    }
3247
3248    #[tokio::test]
3249    async fn dropping_an_uncommitted_staging_batch_discards_host_local_payloads() {
3250        let (_home, ops) = make_cloud_home_with_ops();
3251        let staging = begin_atomic_create(ops.clone(), CloudKitScope::Private)
3252            .await
3253            .unwrap();
3254        stage_atomic_create_record(
3255            staging.clone(),
3256            CloudKitRecordCreate {
3257                key: "copies/cancelled.part0.upload".to_string(),
3258                data: vec![8u8; CHUNK_SIZE],
3259            },
3260        )
3261        .await
3262        .unwrap();
3263
3264        drop(staging);
3265
3266        assert!(ops.staged_batches.lock().unwrap().is_empty());
3267        assert!(ops.store.lock().unwrap().is_empty());
3268        assert!(ops
3269            .calls()
3270            .contains(&MockCall::DiscardBatch("batch-0".to_string())));
3271    }
3272
3273    #[test]
3274    fn cancellation_discard_failure_terminates_the_process() {
3275        const CHILD: &str = "COVEN_CLOUDKIT_CANCEL_DISCARD_ABORT_CHILD";
3276        if std::env::var_os(CHILD).is_some() {
3277            let runtime = tokio::runtime::Runtime::new().unwrap();
3278            runtime.block_on(async {
3279                let ops = Arc::new(MockCloudKitOps::new());
3280                let staging = begin_atomic_create(ops.clone(), CloudKitScope::Private)
3281                    .await
3282                    .unwrap();
3283                stage_atomic_create_record(
3284                    staging.clone(),
3285                    CloudKitRecordCreate {
3286                        key: "copies/cancelled.part0.upload".to_string(),
3287                        data: vec![8u8; CHUNK_SIZE],
3288                    },
3289                )
3290                .await
3291                .unwrap();
3292                ops.fail_discard();
3293                let started = Arc::new(std::sync::Barrier::new(2));
3294                let release = Arc::new(std::sync::Barrier::new(2));
3295                let worker_started = started.clone();
3296                let worker_release = release.clone();
3297                let owner = tokio::spawn(async move {
3298                    tokio::task::spawn_blocking(move || {
3299                        worker_started.wait();
3300                        worker_release.wait();
3301                        drop(staging);
3302                    })
3303                    .await
3304                    .unwrap();
3305                });
3306                started.wait();
3307                owner.abort();
3308                release.wait();
3309                tokio::time::sleep(std::time::Duration::from_secs(2)).await;
3310            });
3311            panic!("failed cancellation discard did not abort the process");
3312        }
3313
3314        let status = std::process::Command::new(std::env::current_exe().unwrap())
3315            .arg("cancellation_discard_failure_terminates_the_process")
3316            .arg("--nocapture")
3317            .env(CHILD, "1")
3318            .status()
3319            .expect("run CloudKit cancellation sabotage subprocess");
3320        assert!(
3321            !status.success(),
3322            "sabotage subprocess unexpectedly survived"
3323        );
3324    }
3325
3326    #[tokio::test]
3327    async fn exact_delete_removes_the_manifest_and_every_part() {
3328        let (home, ops) = make_cloud_home_with_ops();
3329        let slot = exact_slot("copies/delete");
3330        ExactSlotStorage::create_at(
3331            &home,
3332            &slot,
3333            BlobBody::from_bytes(vec![5u8; CHUNK_SIZE + 1]),
3334            &no_progress(),
3335        )
3336        .await
3337        .unwrap();
3338
3339        ExactSlotStorage::delete_at(&home, &slot)
3340            .await
3341            .expect("delete exact slot");
3342
3343        for key in [
3344            "copies/delete".to_string(),
3345            exact_part_key("copies/delete", 0),
3346            exact_part_key("copies/delete", 1),
3347        ] {
3348            assert!(!ops.record_exists(&CloudKitScope::Private, &key).unwrap());
3349        }
3350    }
3351
3352    #[tokio::test]
3353    async fn grant_access_returns_share_join_info_without_email() {
3354        let ch = make_cloud_home();
3355        // CloudKit shares bind identity at URL-accept time, so no invitee email
3356        // is supplied and the grant still succeeds.
3357        let join_info = ch
3358            .set_access(CloudAccessState::Present {
3359                member_pubkey: "member-pubkey".to_string(),
3360                provider_account_email: None,
3361            })
3362            .await
3363            .unwrap();
3364        assert_eq!(
3365            join_info,
3366            CloudAccessOutcome::Present(CloudHomeJoinInfo::CloudKitShare {
3367                share_url: "https://share.example/member-pubkey".to_string(),
3368                owner_name: "owner-name".to_string(),
3369                zone_name: "bae-store".to_string(),
3370            })
3371        );
3372    }
3373
3374    #[tokio::test]
3375    async fn revoke_access_unshares_and_reports_revoked() {
3376        let ch = make_cloud_home();
3377        let outcome = ch
3378            .set_access(CloudAccessState::Absent {
3379                member_pubkey: "member-pubkey".to_string(),
3380                provider_account_email: None,
3381            })
3382            .await
3383            .unwrap();
3384        // CloudKit removes the member's share participation, so it reports the
3385        // credential actually withdrawn rather than Unsupported.
3386        assert_eq!(outcome, CloudAccessOutcome::Absent(RevokeOutcome::Revoked));
3387    }
3388
3389    #[tokio::test]
3390    async fn repeated_present_access_reuses_the_verified_share() {
3391        let (home, ops) = make_cloud_home_with_ops();
3392        let desired = CloudAccessState::Present {
3393            member_pubkey: "member-pubkey".to_string(),
3394            provider_account_email: None,
3395        };
3396
3397        let first = home.set_access(desired.clone()).await.unwrap();
3398        let second = home.set_access(desired).await.unwrap();
3399
3400        assert_eq!(first, second);
3401        assert_eq!(ops.grant_share_calls.load(Ordering::Relaxed), 1);
3402    }
3403
3404    #[tokio::test]
3405    async fn repeated_absent_access_does_not_revoke_twice() {
3406        let (home, ops) = make_cloud_home_with_ops();
3407        home.set_access(CloudAccessState::Present {
3408            member_pubkey: "member-pubkey".to_string(),
3409            provider_account_email: None,
3410        })
3411        .await
3412        .unwrap();
3413        let desired = CloudAccessState::Absent {
3414            member_pubkey: "member-pubkey".to_string(),
3415            provider_account_email: None,
3416        };
3417
3418        home.set_access(desired.clone()).await.unwrap();
3419        home.set_access(desired).await.unwrap();
3420
3421        assert_eq!(ops.revoke_share_calls.load(Ordering::Relaxed), 1);
3422    }
3423
3424    #[test]
3425    fn test_strip_part_suffix() {
3426        assert_eq!(strip_part_suffix("file.bin.part0"), "file.bin");
3427        assert_eq!(strip_part_suffix("file.bin.part123"), "file.bin");
3428        assert_eq!(
3429            strip_part_suffix("file.bin.part123.0123456789abcdef0123456789abcdef"),
3430            "file.bin"
3431        );
3432        assert_eq!(strip_part_suffix("file.bin.manifest"), "file.bin");
3433        assert_eq!(strip_part_suffix("file.bin"), "file.bin");
3434        assert_eq!(strip_part_suffix("file.partition"), "file.partition");
3435        assert_eq!(strip_part_suffix("file.part"), "file.part"); // no digits after .part
3436    }
3437}