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