1use super::blob_io::*;
2use super::cipher::*;
3use super::*;
4
5const STORE_KEY_CONFIRMATION_PLAINTEXT: &[u8] = b"coven store key confirmation v1";
6
7fn store_key_confirmation_aad(
8 creation_id: coven_protocol::store_commit::StoreCreationId,
9) -> Vec<u8> {
10 format!("coven.store-key-confirmation.v1\0{creation_id}").into_bytes()
11}
12
13macro_rules! store_blob_protection {
14 ($storage:expr) => {{
15 let cipher = $storage.cipher.read().unwrap();
16 $storage
17 .pending_rotation
18 .check(cipher.current_generation())?;
19 match &*cipher {
20 CloudCipher::Encrypted(encryption) => {
21 coven_protocol::objects::BlobSpoolProtection::Opaque(encryption.clone())
22 }
23 CloudCipher::Plaintext => coven_protocol::objects::BlobSpoolProtection::Browsable,
24 }
25 }};
26}
27
28#[async_trait]
29impl CloudSyncObjectStorage for CloudSyncConnection {
30 fn blob_path_scheme(&self) -> BlobPathScheme {
31 self.blob_path_scheme()
32 }
33
34 async fn probe_provider(&self) -> Result<(), StorageError> {
35 self.probe().await.map_err(Into::into)
36 }
37
38 fn provider_requests(
39 &self,
40 ) -> Option<std::sync::Arc<dyn coven_foundation::stage_timing::ProviderRequests>> {
41 CloudSyncConnection::provider_requests(self)
42 }
43
44 async fn set_member_access(
45 &self,
46 state: crate::cloud::CloudAccessState,
47 ) -> Result<crate::cloud::CloudAccessOutcome, StorageError> {
48 self.home.set_access(state).await.map_err(Into::into)
49 }
50
51 async fn read_blob_tombstone(
52 &self,
53 stored: &coven_protocol::blob::locator::StoredBlobRef,
54 ) -> Result<Option<Vec<u8>>, StorageError> {
55 let key = crate::blob_tombstone_key(stored, self.cipher_suffix());
56 let stored_bytes = match self.home.read(&key).await.map_err(StorageError::from) {
57 Ok(bytes) => bytes,
58 Err(StorageError::NotFound(_)) => return Ok(None),
59 Err(error) => return Err(error),
60 };
61 let aad_context = cloud_aad_context(&self.store_id, &key);
62 self.open_stored_data(stored_bytes, &aad_context)
63 .map(Some)
64 .map_err(|source| StorageError::Decryption {
65 context: format!("blob tombstone {key}"),
66 source,
67 })
68 }
69
70 async fn write_blob_tombstone(
71 &self,
72 stored: &coven_protocol::blob::locator::StoredBlobRef,
73 plaintext: Vec<u8>,
74 ) -> Result<(), StorageError> {
75 let key = crate::blob_tombstone_key(stored, self.cipher_suffix());
76 let aad_context = cloud_aad_context(&self.store_id, &key);
77 let stored_bytes = self.seal_stored_data(plaintext, &aad_context)?;
78 self.home
79 .write(
80 &key,
81 BlobBody::from_bytes(stored_bytes),
82 &crate::cloud::no_progress(),
83 )
84 .await
85 .map_err(Into::into)
86 }
87
88 async fn list_blob_tombstones(&self) -> Result<Vec<crate::ListedBlobTombstone>, StorageError> {
89 let suffix = self.cipher_suffix();
90 let keys = self
91 .home
92 .list(crate::BLOB_TOMBSTONE_PREFIX)
93 .await
94 .map_err(StorageError::from)?;
95 let mut listed = Vec::with_capacity(keys.len());
96 for key in keys {
97 let object_id = key
98 .strip_suffix(suffix)
99 .and_then(|key| key.strip_prefix(crate::BLOB_TOMBSTONE_PREFIX))
100 .and_then(|encoded| encoded.parse().ok());
101 let Some(object_id) = object_id else {
102 listed.push(crate::ListedBlobTombstone::InvalidKey { provider_key: key });
103 continue;
104 };
105 let stored_bytes = self.home.read(&key).await.map_err(StorageError::from)?;
106 let aad_context = cloud_aad_context(&self.store_id, &key);
107 match self.open_stored_data(stored_bytes, &aad_context) {
108 Ok(plaintext) => listed.push(crate::ListedBlobTombstone::Opened {
109 object_id,
110 plaintext,
111 }),
112 Err(error) => listed.push(crate::ListedBlobTombstone::InvalidBody {
113 provider_key: key,
114 source: std::sync::Arc::new(error),
115 }),
116 }
117 }
118 Ok(listed)
119 }
120
121 async fn blob_tombstone_exists(
122 &self,
123 stored: &coven_protocol::blob::locator::StoredBlobRef,
124 ) -> Result<bool, StorageError> {
125 let key = crate::blob_tombstone_key(stored, self.cipher_suffix());
126 self.home.exists(&key).await.map_err(Into::into)
127 }
128
129 async fn delete_blob_tombstone(
130 &self,
131 stored: &coven_protocol::blob::locator::StoredBlobRef,
132 ) -> Result<(), StorageError> {
133 let key = crate::blob_tombstone_key(stored, self.cipher_suffix());
134 self.home.delete(&key).await.map_err(Into::into)
135 }
136
137 #[cfg(any(test, feature = "test-utils"))]
138 async fn read_provider_bytes_for_test(&self, key: &str) -> Result<Vec<u8>, StorageError> {
139 self.home.read(key).await.map_err(Into::into)
140 }
141
142 #[cfg(any(test, feature = "test-utils"))]
143 async fn write_provider_bytes_for_test(
144 &self,
145 key: &str,
146 bytes: Vec<u8>,
147 ) -> Result<(), StorageError> {
148 self.home
149 .write(
150 key,
151 BlobBody::from_bytes(bytes),
152 &crate::cloud::no_progress(),
153 )
154 .await
155 .map_err(Into::into)
156 }
157
158 #[cfg(any(test, feature = "test-utils"))]
159 async fn list_provider_keys_for_test(&self, prefix: &str) -> Result<Vec<String>, StorageError> {
160 self.home.list(prefix).await.map_err(Into::into)
161 }
162
163 #[cfg(any(test, feature = "test-utils"))]
164 async fn provider_key_exists_for_test(&self, key: &str) -> Result<bool, StorageError> {
165 self.home.exists(key).await.map_err(Into::into)
166 }
167
168 async fn reserve_cross_principal_response_slot(
169 &self,
170 probe_id: coven_protocol::provider::ProviderProbeId,
171 ) -> Result<ObjectSlot, coven_protocol::provider::ProviderProbeError> {
172 self.provider_probes
173 .reserve_cross_principal_response_slot(probe_id)
174 .await
175 }
176
177 async fn prepare_cross_principal_challenge(
178 &self,
179 publication_journal: &dyn coven_protocol::provider::DeviceJoinChallengePublicationJournal,
180 probe_id: coven_protocol::provider::ProviderProbeId,
181 store: &coven_protocol::StoreProviderBinding,
182 context: &coven_protocol::provider::CrossPrincipalChallengeContext,
183 administrator_signer: &dyn coven_keys::keys::DeviceSigningAuthority,
184 ) -> Result<
185 coven_protocol::provider::CrossPrincipalProbeChallenge,
186 coven_protocol::provider::ProviderProbeError,
187 > {
188 self.provider_probes
189 .prepare_cross_principal_challenge(
190 publication_journal,
191 probe_id,
192 store,
193 context,
194 administrator_signer,
195 )
196 .await
197 }
198
199 async fn settle_cross_principal_challenge(
200 &self,
201 publication_journal: &dyn coven_protocol::provider::DeviceJoinChallengePublicationJournal,
202 authorization: &coven_protocol::provider::DeviceJoinChallengePublicationAuthorization,
203 challenge: &coven_protocol::provider::CrossPrincipalProbeChallenge,
204 context: &coven_protocol::provider::CrossPrincipalChallengeContext,
205 store: &coven_protocol::StoreProviderBinding,
206 ) -> Result<
207 coven_protocol::provider::CrossPrincipalProbeChallenge,
208 coven_protocol::provider::ProviderProbeError,
209 > {
210 self.provider_probes
211 .settle_cross_principal_challenge(
212 publication_journal,
213 authorization,
214 challenge,
215 context,
216 store,
217 )
218 .await
219 }
220
221 async fn create_cross_principal_response(
222 &self,
223 challenge: &coven_protocol::provider::CrossPrincipalProbeChallenge,
224 context: &coven_protocol::provider::CrossPrincipalResponseContext,
225 store: &coven_protocol::StoreProviderBinding,
226 administrator_signing_pubkey: &str,
227 peer_signer: &coven_keys::keys::UserKeypair,
228 ) -> Result<
229 coven_protocol::provider::CrossPrincipalProbeResponse,
230 coven_protocol::provider::ProviderProbeError,
231 > {
232 self.provider_probes
233 .create_cross_principal_response(
234 challenge,
235 context,
236 store,
237 administrator_signing_pubkey,
238 peer_signer,
239 )
240 .await
241 }
242
243 async fn complete_cross_principal_probe(
244 &self,
245 journal: &dyn coven_protocol::provider::ProviderProbeJournal,
246 challenge: &coven_protocol::provider::CrossPrincipalProbeChallenge,
247 response: &coven_protocol::provider::CrossPrincipalProbeResponse,
248 context: &coven_protocol::provider::CrossPrincipalResponseContext,
249 store: &coven_protocol::StoreProviderBinding,
250 administrator_signer: &dyn coven_keys::keys::DeviceSigningAuthority,
251 peer_signing_pubkey: &str,
252 ) -> Result<
253 coven_protocol::provider::CrossPrincipalProbeReceipt,
254 coven_protocol::provider::ProviderProbeError,
255 > {
256 self.provider_probes
257 .complete_cross_principal_probe(
258 journal,
259 challenge,
260 response,
261 context,
262 store,
263 administrator_signer,
264 peer_signing_pubkey,
265 )
266 .await
267 }
268
269 async fn probe_exact_slots(
270 &self,
271 journal: &dyn coven_protocol::provider::ProviderProbeJournal,
272 probe_id: coven_protocol::provider::ProviderProbeId,
273 binding: &ResolvedProviderBinding,
274 ) -> Result<
275 coven_protocol::provider::ExactSlotProbeReceipt,
276 coven_protocol::provider::ProviderProbeError,
277 > {
278 self.provider_probes
279 .probe_exact_slots(journal, probe_id, binding)
280 .await
281 }
282
283 async fn observe_exact_slot(
284 &self,
285 slot: &ObjectSlot,
286 ) -> Result<Option<ExactObjectRef>, StorageError> {
287 self.home.observe_at(slot).await.map_err(Into::into)
288 }
289
290 async fn delete_exact_slot_and_verify_absent(
291 &self,
292 slot: &ObjectSlot,
293 ) -> Result<(), StorageError> {
294 self.home
295 .delete_and_verify_absent(slot)
296 .await
297 .map_err(Into::into)
298 }
299
300 fn store_blob_key_fingerprint(
301 &self,
302 ) -> Result<Option<coven_keys::encryption::KeyFingerprint>, StorageError> {
303 let cipher = self.cipher.read().unwrap();
304 self.pending_rotation.check(cipher.current_generation())?;
305 Ok(match &*cipher {
306 CloudCipher::Encrypted(encryption) => Some(encryption.seal_key_fingerprint()),
307 CloudCipher::Plaintext => None,
308 })
309 }
310
311 fn create_store_key_confirmation(
312 &self,
313 creation_id: coven_protocol::store_commit::StoreCreationId,
314 ) -> Result<coven_protocol::store_commit::StoreKeyConfirmation, StorageError> {
315 let cipher = self.cipher.read().unwrap();
316 self.pending_rotation.check(cipher.current_generation())?;
317 Ok(match &*cipher {
318 CloudCipher::Encrypted(_) => {
319 coven_protocol::store_commit::StoreKeyConfirmation::Opaque(cipher.seal(
320 STORE_KEY_CONFIRMATION_PLAINTEXT.to_vec(),
321 &store_key_confirmation_aad(creation_id),
322 ))
323 }
324 CloudCipher::Plaintext => {
325 coven_protocol::store_commit::StoreKeyConfirmation::NotRequired
326 }
327 })
328 }
329
330 fn verify_store_key_confirmation(
331 &self,
332 creation_id: coven_protocol::store_commit::StoreCreationId,
333 confirmation: &coven_protocol::store_commit::StoreKeyConfirmation,
334 ) -> Result<(), StorageError> {
335 let cipher = self.cipher.read().unwrap();
336 self.pending_rotation.check(cipher.current_generation())?;
337 match (&*cipher, confirmation) {
338 (
339 CloudCipher::Plaintext,
340 coven_protocol::store_commit::StoreKeyConfirmation::NotRequired,
341 ) => Ok(()),
342 (
343 CloudCipher::Encrypted(_),
344 coven_protocol::store_commit::StoreKeyConfirmation::Opaque(sealed),
345 ) => {
346 let opened = cipher
347 .open(sealed.clone(), &store_key_confirmation_aad(creation_id))
348 .map_err(|source| StorageError::Decryption {
349 context: "Store key confirmation".to_string(),
350 source,
351 })?;
352 if opened == STORE_KEY_CONFIRMATION_PLAINTEXT {
353 Ok(())
354 } else {
355 Err(StorageError::InvalidContent(
356 "Store key confirmation plaintext differs".to_string(),
357 ))
358 }
359 }
360 (
361 CloudCipher::Plaintext,
362 coven_protocol::store_commit::StoreKeyConfirmation::Opaque(_),
363 ) => Err(StorageError::InvalidContent(
364 "opaque Store root opened through browsable storage".to_string(),
365 )),
366 (
367 CloudCipher::Encrypted(_),
368 coven_protocol::store_commit::StoreKeyConfirmation::NotRequired,
369 ) => Err(StorageError::InvalidContent(
370 "browsable Store root opened through opaque storage".to_string(),
371 )),
372 }
373 }
374
375 async fn provider_binding(&self) -> Result<ResolvedProviderBinding, StorageError> {
376 self.home.provider_binding().await.map_err(Into::into)
377 }
378
379 async fn allocate_protocol_slot(
380 &self,
381 context: &ProtocolObjectContext,
382 semantic_prefix: &str,
383 extension: &str,
384 ) -> Result<ObjectSlot, StorageError> {
385 context.validate_path(semantic_prefix)?;
386 context.validate_extension(extension)?;
387 Ok(self
388 .home
389 .allocate_slot(&format!("{semantic_prefix}{extension}"))
390 .await?)
391 }
392
393 fn prepare_protocol_object(
394 &self,
395 context: &ProtocolObjectContext,
396 slot: ObjectSlot,
397 semantic_prefix: &str,
398 data: Vec<u8>,
399 ) -> Result<PreparedExactObject, StorageError> {
400 context.validate_slot(&slot, semantic_prefix)?;
401 let aad = protocol_object_aad_context(context, semantic_prefix);
402 let stored = self.seal_protocol_data(context, data, &aad)?;
403 let reference = ExactObjectRef::new(
404 slot,
405 stored.len() as u64,
406 coven_protocol::store_commit::ObjectHash::digest(&stored),
407 );
408 PreparedExactObject::new(reference, stored)
409 }
410
411 async fn open_prepared_protocol_object(
412 &self,
413 context: &ProtocolObjectContext,
414 prepared: &PreparedExactObject,
415 semantic_prefix: &str,
416 ) -> Result<Vec<u8>, StorageError> {
417 context.validate_reference(prepared.reference(), semantic_prefix)?;
418 let aad = protocol_object_aad_context(context, semantic_prefix);
419 let object = prepared.reference().clone();
420 let stored = prepared.stored_bytes().to_vec();
421 self.verify_and_open_protocol_data(
422 "verify and open prepared protocol object",
423 context,
424 object,
425 stored,
426 aad,
427 )
428 .await
429 }
430
431 async fn create_protocol_object(
432 &self,
433 prepared: &PreparedExactObject,
434 ) -> Result<(), StorageError> {
435 let upload =
436 crate::cloud::ExactUpload::from_bytes(prepared.reference(), prepared.stored_bytes())?;
437 self.home
438 .create_at(
439 &upload,
440 &crate::cloud::UploadControl::running(crate::cloud::no_progress()),
441 )
442 .await
443 .map(drop)
444 .map_err(Into::into)
445 }
446
447 async fn create_versioned_protocol_record(
448 &self,
449 context: &ProtocolObjectContext,
450 prepared: &PreparedExactObject,
451 semantic_prefix: &str,
452 expected: &[u8],
453 ) -> Result<crate::cloud::CloudObjectVersion, StorageError> {
454 self.verify_prepared_protocol_object(context, prepared, semantic_prefix, expected)
455 .await?;
456 let upload =
457 crate::cloud::ExactUpload::from_bytes(prepared.reference(), prepared.stored_bytes())?;
458 self.home
459 .create_versioned_at(
460 &upload,
461 &crate::cloud::UploadControl::running(crate::cloud::no_progress()),
462 )
463 .await?;
464 let versioned = self
465 .home
466 .read_versioned_at(prepared.reference().slot())
467 .await?;
468 prepared.reference().verify(&versioned.bytes)?;
469 let aad = protocol_object_aad_context(context, semantic_prefix);
470 let opened = self
471 .verify_and_open_protocol_data(
472 "verify and open created versioned protocol record",
473 context,
474 prepared.reference().clone(),
475 versioned.bytes,
476 aad,
477 )
478 .await?;
479 if opened != expected {
480 return Err(StorageError::PreparedObjectMismatch(
481 prepared.reference().slot().logical_key().to_string(),
482 ));
483 }
484 Ok(versioned.version)
485 }
486
487 async fn read_protocol_object(
488 &self,
489 context: &ProtocolObjectContext,
490 object: &ExactObjectRef,
491 semantic_prefix: &str,
492 ) -> Result<Vec<u8>, StorageError> {
493 context.validate_reference(object, semantic_prefix)?;
494 let stored = self.home.read_at(object.slot()).await?;
495 let aad = protocol_object_aad_context(context, semantic_prefix);
496 let object = object.clone();
497 self.verify_and_open_protocol_data(
498 "verify and open protocol object",
499 context,
500 object,
501 stored,
502 aad,
503 )
504 .await
505 }
506
507 async fn read_protocol_object_with_progress(
508 &self,
509 context: &ProtocolObjectContext,
510 object: &ExactObjectRef,
511 semantic_prefix: &str,
512 progress: crate::cloud::DownloadProgress,
513 ) -> Result<Vec<u8>, StorageError> {
514 context.validate_reference(object, semantic_prefix)?;
515 let temporary = tempfile::tempdir().map_err(StorageError::Io)?;
516 let stored_path = temporary.path().join("protocol-object");
517 self.home
518 .read_at_to_file(object.slot(), &stored_path, progress)
519 .await
520 .map_err(map_cloud_file_read_error)?;
521 let stored = tokio::fs::read(&stored_path)
522 .await
523 .map_err(StorageError::Io)?;
524 let aad = protocol_object_aad_context(context, semantic_prefix);
525 self.verify_and_open_protocol_data(
526 "verify and open streamed protocol object",
527 context,
528 object.clone(),
529 stored,
530 aad,
531 )
532 .await
533 }
534
535 async fn read_versioned_protocol_record(
536 &self,
537 context: &ProtocolObjectContext,
538 slot: &ObjectSlot,
539 semantic_prefix: &str,
540 ) -> Result<(Vec<u8>, crate::cloud::CloudObjectVersion), StorageError> {
541 context.validate_slot(slot, semantic_prefix)?;
542 let versioned = self.home.read_versioned_at(slot).await?;
543 let aad = protocol_object_aad_context(context, semantic_prefix);
544 let object = ExactObjectRef::new(
545 slot.clone(),
546 versioned.bytes.len() as u64,
547 coven_protocol::store_commit::ObjectHash::digest(&versioned.bytes),
548 );
549 let opened = self
550 .verify_and_open_protocol_data(
551 "verify and open versioned protocol record",
552 context,
553 object,
554 versioned.bytes,
555 aad,
556 )
557 .await?;
558 Ok((opened, versioned.version))
559 }
560
561 async fn replace_protocol_record_if_version(
562 &self,
563 context: &ProtocolObjectContext,
564 slot: &ObjectSlot,
565 semantic_prefix: &str,
566 expected: &crate::cloud::CloudObjectVersion,
567 data: Vec<u8>,
568 ) -> Result<crate::cloud::ConditionalWriteOutcome, StorageError> {
569 context.validate_slot(slot, semantic_prefix)?;
570 let aad = protocol_object_aad_context(context, semantic_prefix);
571 let stored = self.seal_protocol_data(context, data, &aad)?;
572 self.home
573 .replace_at_if_version(slot, expected, stored)
574 .await
575 .map_err(Into::into)
576 }
577
578 async fn list_protocol_slots(
579 &self,
580 context: &ProtocolObjectContext,
581 listing_prefix: &str,
582 ) -> Result<Vec<ObjectSlot>, StorageError> {
583 let listed = self.home.list_slots(listing_prefix).await?;
584 Ok(listed
585 .into_iter()
586 .filter(|slot| context.semantic_prefix_of(slot).is_some())
587 .collect())
588 }
589
590 async fn read_protocol_slot(
591 &self,
592 context: &ProtocolObjectContext,
593 slot: &ObjectSlot,
594 semantic_prefix: &str,
595 ) -> Result<(Vec<u8>, ExactObjectRef), StorageError> {
596 let (opened, prepared) = self
597 .read_prepared_protocol_slot(context, slot, semantic_prefix)
598 .await?;
599 Ok((opened, prepared.reference().clone()))
600 }
601
602 async fn read_prepared_protocol_slot(
603 &self,
604 context: &ProtocolObjectContext,
605 slot: &ObjectSlot,
606 semantic_prefix: &str,
607 ) -> Result<(Vec<u8>, PreparedExactObject), StorageError> {
608 context.validate_slot(slot, semantic_prefix)?;
609 let stored = self.home.read_at(slot).await?;
610 let aad = protocol_object_aad_context(context, semantic_prefix);
611 let slot = slot.clone();
612 self.identify_and_open_protocol_data(context, slot, stored, aad)
613 .await
614 }
615
616 async fn delete_protocol_object(&self, object: &ExactObjectRef) -> Result<(), StorageError> {
617 match self.home.read_at(object.slot()).await {
618 Err(crate::cloud::CloudHomeError::NotFound(_)) => return Ok(()),
619 Err(error) => return Err(error.into()),
620 Ok(stored)
621 if stored.len() as u64 != object.stored_size()
622 || coven_protocol::store_commit::ObjectHash::digest(&stored)
623 != object.stored_hash() =>
624 {
625 return Err(StorageError::SlotCollision(format!(
626 "exact delete target {} contains different bytes",
627 object.slot().logical_key()
628 )));
629 }
630 Ok(_) => {}
631 }
632 let delete_error = self.home.delete_at(object.slot()).await.err();
633 if delete_error
634 .as_ref()
635 .is_some_and(|error| !error.is_retryable())
636 {
637 return Err(delete_error.expect("delete error exists").into());
638 }
639 match self.home.read_at(object.slot()).await {
640 Err(crate::cloud::CloudHomeError::NotFound(_)) => Ok(()),
641 Err(readback) => match delete_error {
642 Some(operation) => Err(StorageError::UnresolvedOutcome {
643 operation: Box::new(operation.into()),
644 settlement: Box::new(readback.into()),
645 }),
646 None => Err(readback.into()),
647 },
648 Ok(_) => match delete_error {
649 Some(error) => Err(error.into()),
650 None => Err(StorageError::Storage(format!(
651 "exact object remains after delete: {}",
652 object.slot().logical_key()
653 ))),
654 },
655 }
656 }
657
658 async fn allocate_blob_slot(
659 &self,
660 locator: &coven_protocol::blob::locator::BlobLocator,
661 authority: &coven_protocol::objects::BlobWriteAuthority<'_>,
662 ) -> Result<ObjectSlot, StorageError> {
663 self.validate_blob_locator_home(locator)?;
664 self.validate_blob_append_authority(locator, authority)
665 .await?;
666 Ok(self.home.allocate_slot(&locator.semantic_key()).await?)
667 }
668
669 async fn seal_blob_to_spool(
670 &self,
671 locator: &coven_protocol::blob::locator::BlobLocator,
672 authority: &coven_protocol::objects::BlobWriteAuthority<'_>,
673 protection: coven_protocol::objects::BlobSpoolProtection,
674 plaintext_file: &Path,
675 spool: coven_foundation::local_file::AtomicStagedFile,
676 progress: crate::cloud::PreparationProgress,
677 ) -> Result<coven_protocol::objects::BlobSpoolWrite, StorageError> {
678 let spool_file = spool.destination().to_path_buf();
679 self.validate_blob_locator_home(locator)?;
680 self.validate_blob_append_authority(locator, authority)
681 .await?;
682 match tokio::fs::metadata(&spool_file).await {
683 Ok(metadata) => {
684 if !metadata.is_file() {
685 return Err(StorageError::LocalFilesystem(
686 coven_foundation::atomic_file::FileError::NotFile {
687 subject: "blob spool path",
688 path: spool_file,
689 },
690 ));
691 }
692 let (stored_size, stored_hash) = crate::local_file::exact_file_facts(&spool_file)
693 .await
694 .map_err(StorageError::LocalFilesystem)?;
695 let object = ExactObjectRef::new(
696 ObjectSlot::logical(locator.semantic_key())?,
697 stored_size,
698 stored_hash,
699 );
700 let blob =
701 coven_protocol::blob::locator::StoredBlobRef::new(locator.clone(), object)?;
702 let mut reader =
703 ExactBlobPlaintextReader::new(&spool_file, &self.store_id, &blob, protection)
704 .await?;
705 loop {
706 let chunk = coven_foundation::local_file::PlaintextChunkReader::next_chunk(
707 &mut reader,
708 1 << 20,
709 )
710 .await
711 .map_err(StorageError::from)?;
712 if chunk.is_empty() {
713 break;
714 }
715 }
716 progress(locator.plaintext_size());
717 return Ok(coven_protocol::objects::BlobSpoolWrite::Reused);
718 }
719 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
720 Err(error) => {
721 return Err(StorageError::LocalFilesystem(
722 coven_foundation::atomic_file::FileError::Path {
723 operation: "inspect blob spool",
724 path: spool_file,
725 source: error,
726 },
727 ));
728 }
729 }
730
731 let retry_protection = protection.clone();
732 let body = match (locator, protection) {
733 (
734 coven_protocol::blob::locator::BlobLocator::Opaque {
735 scope,
736 key_fingerprint,
737 ..
738 },
739 coven_protocol::objects::BlobSpoolProtection::Opaque(encryption),
740 ) => {
741 if encryption.seal_key_fingerprint() != *key_fingerprint {
742 return Err(StorageError::InvalidContent(format!(
743 "blob locator key fingerprint {key_fingerprint} differs from the supplied audience key {}",
744 encryption.seal_key_fingerprint()
745 )));
746 }
747 let aad = cloud_aad_context(&self.store_id, &locator.semantic_key());
748 CloudCipher::Encrypted(encryption)
749 .open_exact_body(
750 scope.clone(),
751 plaintext_file,
752 &aad,
753 self.blob_chunking.chunk(),
754 locator.plaintext_size(),
755 locator.plaintext_hash(),
756 progress,
757 )
758 .await
759 .map_err(StorageError::LocalFilesystem)?
760 }
761 (
762 coven_protocol::blob::locator::BlobLocator::Browsable { .. },
763 coven_protocol::objects::BlobSpoolProtection::Browsable,
764 ) => CloudCipher::Plaintext
765 .open_exact_body(
766 coven_protocol::blob::BlobScope::Master,
767 plaintext_file,
768 &[],
769 self.blob_chunking.chunk(),
770 locator.plaintext_size(),
771 locator.plaintext_hash(),
772 progress,
773 )
774 .await
775 .map_err(StorageError::LocalFilesystem)?,
776 (coven_protocol::blob::locator::BlobLocator::Opaque { .. }, _) => {
777 return Err(StorageError::Configuration(
778 "opaque blob locator requires audience encryption".to_string(),
779 ));
780 }
781 (coven_protocol::blob::locator::BlobLocator::Browsable { .. }, _) => {
782 return Err(StorageError::Configuration(
783 "browsable blob locator cannot use audience encryption".to_string(),
784 ));
785 }
786 };
787 let expected_size = body.len();
788 let stream = futures_util::stream::try_unfold(body, |mut body| async move {
789 match body.next_part(1 << 20).await? {
790 Some(chunk) => Ok::<_, crate::cloud::CloudHomeError>(Some((chunk, body))),
791 None => Ok::<_, crate::cloud::CloudHomeError>(None),
792 }
793 });
794 let (staged, written) =
795 spool
796 .write_byte_stream(Box::pin(stream))
797 .await
798 .map_err(|error| match error {
799 coven_foundation::local_file::ByteStreamWriteError::Source(error) => {
800 error.into()
801 }
802 coven_foundation::local_file::ByteStreamWriteError::SourceCleanup {
803 source,
804 cleanup,
805 } => StorageError::CleanupFailed {
806 operation: Box::new(source.into()),
807 cleanup: Box::new(StorageError::LocalFilesystem(cleanup)),
808 },
809 coven_foundation::local_file::ByteStreamWriteError::Local(error) => {
810 StorageError::LocalFilesystem(error)
811 }
812 })?;
813 if written != expected_size {
814 return Err(StorageError::InvalidContent(format!(
815 "blob spool {} contains {written} stored bytes, expected {expected_size}",
816 spool_file.display()
817 )));
818 }
819 match staged.commit_new().await {
820 Ok(()) => Ok(coven_protocol::objects::BlobSpoolWrite::Created),
821 Err(coven_foundation::local_file::CommitNewFileError::DestinationExists(_)) => {
822 let (stored_size, stored_hash) = crate::local_file::exact_file_facts(&spool_file)
823 .await
824 .map_err(StorageError::LocalFilesystem)?;
825 let object = ExactObjectRef::new(
826 ObjectSlot::logical(locator.semantic_key())?,
827 stored_size,
828 stored_hash,
829 );
830 let blob =
831 coven_protocol::blob::locator::StoredBlobRef::new(locator.clone(), object)?;
832 let mut reader = ExactBlobPlaintextReader::new(
833 &spool_file,
834 &self.store_id,
835 &blob,
836 retry_protection,
837 )
838 .await?;
839 loop {
840 let chunk = coven_foundation::local_file::PlaintextChunkReader::next_chunk(
841 &mut reader,
842 1 << 20,
843 )
844 .await
845 .map_err(StorageError::from)?;
846 if chunk.is_empty() {
847 break;
848 }
849 }
850 Ok(coven_protocol::objects::BlobSpoolWrite::Reused)
851 }
852 Err(error) => Err(StorageError::CommitNewFile(error)),
853 }
854 }
855
856 async fn seal_store_blob_to_spool(
857 &self,
858 locator: &coven_protocol::blob::locator::BlobLocator,
859 authority: &coven_protocol::objects::BlobWriteAuthority<'_>,
860 plaintext_file: &Path,
861 spool: coven_foundation::local_file::AtomicStagedFile,
862 progress: crate::cloud::PreparationProgress,
863 ) -> Result<coven_protocol::objects::BlobSpoolWrite, StorageError> {
864 let protection = store_blob_protection!(self);
865 self.seal_blob_to_spool(
866 locator,
867 authority,
868 protection,
869 plaintext_file,
870 spool,
871 progress,
872 )
873 .await
874 }
875
876 async fn prepare_blob_object(
877 &self,
878 locator: &coven_protocol::blob::locator::BlobLocator,
879 authority: &coven_protocol::objects::BlobWriteAuthority<'_>,
880 slot: ObjectSlot,
881 stored_file: &Path,
882 ) -> Result<coven_protocol::blob::locator::StoredBlobRef, StorageError> {
883 self.validate_blob_locator_home(locator)?;
884 self.validate_blob_append_authority(locator, authority)
885 .await?;
886 let expected = locator.semantic_key();
887 if slot.logical_key() != expected {
888 return Err(StorageError::Parse(format!(
889 "blob slot {:?} does not match locator key {expected:?}",
890 slot.logical_key()
891 )));
892 }
893 let (stored_size, stored_hash) = crate::local_file::exact_file_facts(stored_file)
894 .await
895 .map_err(StorageError::LocalFilesystem)?;
896 coven_protocol::blob::locator::StoredBlobRef::new(
897 locator.clone(),
898 ExactObjectRef::new(slot, stored_size, stored_hash),
899 )
900 .map_err(StorageError::from)
901 }
902
903 async fn create_blob_object_from_file(
904 &self,
905 blob: &coven_protocol::blob::locator::StoredBlobRef,
906 authority: &coven_protocol::objects::BlobWriteAuthority<'_>,
907 stored_file: &Path,
908 control: &crate::cloud::UploadControl,
909 ) -> Result<(), StorageError> {
910 let locator = blob.locator();
911 let object = blob.object();
912 self.validate_blob_locator_home(locator)?;
913 self.validate_blob_append_authority(locator, authority)
914 .await?;
915 let expected = locator.semantic_key();
916 if object.slot().logical_key() != expected {
917 return Err(StorageError::Parse(format!(
918 "blob object {:?} does not match locator key {expected:?}",
919 object.slot().logical_key()
920 )));
921 }
922 let upload = crate::cloud::ExactUpload::from_file(object, stored_file).await?;
923 self.home
924 .create_at(&upload, control)
925 .await
926 .map(drop)
927 .map_err(Into::into)
928 }
929
930 async fn verify_blob_object(
931 &self,
932 blob: &coven_protocol::blob::locator::StoredBlobRef,
933 ) -> Result<(), StorageError> {
934 self.validate_blob_locator_home(blob.locator())?;
935 let expected = blob.locator().semantic_key();
936 if blob.object().slot().logical_key() != expected {
937 return Err(StorageError::Parse(format!(
938 "blob object {:?} does not match locator key {expected:?}",
939 blob.object().slot().logical_key()
940 )));
941 }
942 let stored = self.home.read_at(blob.object().slot()).await?;
943 blob.object().verify(&stored)
944 }
945
946 async fn stage_verified_blob_plaintext(
947 &self,
948 blob: &coven_protocol::blob::locator::StoredBlobRef,
949 protection: coven_protocol::objects::BlobSpoolProtection,
950 mut plaintext: coven_foundation::local_file::AtomicStagedFile,
951 progress: crate::cloud::DownloadProgress,
952 ) -> Result<coven_foundation::local_file::AtomicStagedFile, StorageError> {
953 let stored_destination = plaintext
954 .destination()
955 .with_extension("coven-stored-download");
956 let stored_stage = plaintext
957 .stage_peer(&stored_destination)
958 .await
959 .map_err(StorageError::LocalFilesystem)?;
960 let stored = self
961 .stage_exact_blob_download(blob, stored_stage, progress)
962 .await?;
963 let mut reader =
964 ExactBlobPlaintextReader::new(stored.path(), &self.store_id, blob, protection).await?;
965 let written =
966 plaintext
967 .write_plaintext(&mut reader)
968 .await
969 .map_err(|error| match error {
970 coven_foundation::local_file::StreamWriteError::Source(error) => error.into(),
971 coven_foundation::local_file::StreamWriteError::Local(error) => {
972 StorageError::LocalFilesystem(error)
973 }
974 })?;
975 if written != blob.locator().plaintext_size() {
976 return Err(StorageError::InvalidContent(format!(
977 "blob {} plaintext stage contains {written} bytes, expected {}",
978 blob.locator().locator_hash(),
979 blob.locator().plaintext_size()
980 )));
981 }
982 Ok(plaintext)
983 }
984
985 async fn stage_verified_store_blob_plaintext(
986 &self,
987 blob: &coven_protocol::blob::locator::StoredBlobRef,
988 plaintext: coven_foundation::local_file::AtomicStagedFile,
989 progress: crate::cloud::DownloadProgress,
990 ) -> Result<coven_foundation::local_file::AtomicStagedFile, StorageError> {
991 let protection = store_blob_protection!(self);
992 self.stage_verified_blob_plaintext(blob, protection, plaintext, progress)
993 .await
994 }
995
996 async fn open_blob_range_reader(
997 &self,
998 blob: &coven_protocol::blob::locator::StoredBlobRef,
999 protection: coven_protocol::objects::BlobSpoolProtection,
1000 ) -> Result<BlobRangeReader, StorageError> {
1001 let locator = blob.locator();
1002 self.validate_blob_locator_home(locator)?;
1003 let slot = blob.object().slot().clone();
1004 let (scope, key_fingerprint) = match locator {
1005 coven_protocol::blob::locator::BlobLocator::Opaque {
1006 scope,
1007 key_fingerprint,
1008 ..
1009 } => (scope, key_fingerprint),
1010 coven_protocol::blob::locator::BlobLocator::Browsable { .. } => {
1016 return Err(StorageError::Configuration(format!(
1017 "blob {} is stored in the clear, which has no per-range verification",
1018 locator.locator_hash()
1019 )));
1020 }
1021 };
1022 let coven_protocol::objects::BlobSpoolProtection::Opaque(master) = protection else {
1023 return Err(StorageError::Configuration(
1024 "opaque blob locator requires audience encryption".to_string(),
1025 ));
1026 };
1027 let prefix = self
1031 .home
1032 .read_range_at(&slot, 0, (KeyTag::LEN + SEALED_BLOB_HEADER_LEN) as u64)
1033 .await
1034 .map_err(StorageError::from)?;
1035 let opener = verified_sealed_blob_opener(
1036 &prefix,
1037 blob,
1038 key_fingerprint,
1039 scope,
1040 &master,
1041 &cloud_aad_context(&self.store_id, &locator.semantic_key()),
1042 )?;
1043 Ok(BlobRangeReader::new(
1044 self.home.clone(),
1045 slot,
1046 opener,
1047 locator.plaintext_size(),
1048 self.blob_chunking.window(),
1049 ))
1050 }
1051
1052 async fn open_store_blob_range_reader(
1053 &self,
1054 blob: &coven_protocol::blob::locator::StoredBlobRef,
1055 ) -> Result<BlobRangeReader, StorageError> {
1056 let protection = store_blob_protection!(self);
1057 self.open_blob_range_reader(blob, protection).await
1058 }
1059
1060 async fn delete_blob_object(
1061 &self,
1062 blob: &coven_protocol::blob::locator::StoredBlobRef,
1063 ) -> Result<(), StorageError> {
1064 let locator = blob.locator();
1065 let object = blob.object();
1066 self.validate_blob_locator_home(locator)?;
1067 let expected = locator.semantic_key();
1068 if object.slot().logical_key() != expected {
1069 return Err(StorageError::Parse(format!(
1070 "blob object {:?} does not match locator key {expected:?}",
1071 object.slot().logical_key()
1072 )));
1073 }
1074 self.delete_protocol_object(object).await?;
1075 Ok(())
1076 }
1077}
1078
1079impl CloudSyncConnection {
1083 async fn stage_exact_blob_download(
1084 &self,
1085 blob: &coven_protocol::blob::locator::StoredBlobRef,
1086 mut staged: coven_foundation::local_file::AtomicStagedFile,
1087 progress: crate::cloud::DownloadProgress,
1088 ) -> Result<coven_foundation::local_file::AtomicStagedFile, StorageError> {
1089 let locator = blob.locator();
1090 let object = blob.object();
1091 self.validate_blob_locator_home(locator)?;
1092 let expected = locator.semantic_key();
1093 if object.slot().logical_key() != expected {
1094 return Err(StorageError::Parse(format!(
1095 "blob object {:?} does not match locator key {expected:?}",
1096 object.slot().logical_key()
1097 )));
1098 }
1099 self.home
1100 .read_at_to_file(
1101 object.slot(),
1102 staged.path_for_atomic_replacement(),
1103 progress,
1104 )
1105 .await
1106 .map_err(map_cloud_file_read_error)?;
1107 {
1108 let (size, digest) = coven_foundation::local_file::file_facts(staged.path())
1109 .await
1110 .map_err(StorageError::LocalFilesystem)?;
1111 object.verify_stored_facts(
1112 staged.path(),
1113 size,
1114 coven_protocol::store_commit::ObjectHash::from_digest(digest),
1115 )?;
1116 }
1117 Ok(staged)
1118 }
1119}