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