Skip to main content

coven_storage/remote/
storage_impl.rs

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            // A browsable home stores the plaintext in the clear, so its objects
1011            // carry no tags and a range read has nothing to check the provider's
1012            // answer against. Ranged reading is refused rather than served
1013            // unverified; the caller materializes the whole blob, where the row's
1014            // content hash can refuse it.
1015            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        // One ranged read of the prefix names the key and the chunk size; every
1028        // later range is arithmetic over the header it carries, so this is the
1029        // only request a range does not pay for.
1030        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
1079/// Reading a stored blob body into an unpublished sibling is a step of this
1080/// adapter's own verified download, not a capability the storage surface
1081/// offers: every caller reaches it through `stage_verified_blob_plaintext`.
1082impl 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}