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