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