1use std::collections::HashMap;
7use std::sync::{Arc, Mutex};
8
9use coven_database::Database;
10use coven_foundation::store_dir::StoreDir;
11use coven_keys::encryption::MasterKeyring;
12use coven_keys::keys::{KeyError, MasterKeyCustody, UserKeypair};
13#[cfg(test)]
14use coven_protocol::store_commit::ObjectHash;
15use coven_storage::CloudSyncObjectStorage;
16
17#[cfg(test)]
18mod uploaded_changeset;
19
20pub use coven_database::synthetic_store::*;
23pub use coven_foundation::store_dir::temp_store_dir;
24pub use coven_storage::cloud::test_utils::{test_cloud_home, test_cloud_home_with_binding};
25
26pub fn staged_snapshot_image(bytes: &[u8]) -> coven_database::SnapshotDatabaseImage {
27 let file = tempfile::NamedTempFile::new().expect("create staged snapshot fixture path");
28 let path = file.path().to_path_buf();
29 file.close().expect("release staged snapshot fixture path");
30 coven_database::SnapshotDatabaseImage::create(path, bytes)
31 .expect("write staged snapshot fixture")
32}
33
34#[cfg(test)]
35pub fn test_cache_locator_hash(label: &str) -> ObjectHash {
36 ObjectHash::digest(label.as_bytes())
37}
38
39#[derive(Clone, Default)]
46pub struct TestCustody {
47 value: Arc<Mutex<Option<String>>>,
48 fail: Arc<std::sync::atomic::AtomicBool>,
49}
50
51impl TestCustody {
52 pub fn set_initial_key(&self, key: [u8; 32]) {
53 *self.value.lock().unwrap() = Some(
54 MasterKeyring::from(coven_keys::encryption::EncryptionService::from_key(key))
55 .to_serialized(),
56 );
57 }
58
59 pub fn stored_key(&self) -> Option<String> {
60 self.value.lock().unwrap().clone()
61 }
62
63 pub fn fail_writes(&self) {
65 self.fail.store(true, std::sync::atomic::Ordering::SeqCst);
66 }
67
68 pub fn allow_writes(&self) {
70 self.fail.store(false, std::sync::atomic::Ordering::SeqCst);
71 }
72}
73
74impl MasterKeyCustody for TestCustody {
75 fn unlock(&self) -> Result<Option<MasterKeyring>, KeyError> {
76 self.value
77 .lock()
78 .unwrap()
79 .as_deref()
80 .map(MasterKeyring::from_serialized)
81 .transpose()
82 .map_err(KeyError::Encryption)
83 }
84
85 fn persist(&self, keyring: &MasterKeyring) -> Result<(), KeyError> {
86 if self.fail.load(std::sync::atomic::Ordering::SeqCst) {
87 return Err(KeyError::Custody {
88 operation: "persist",
89 source: Box::new(std::io::Error::other("forced keyring write failure")),
90 });
91 }
92 *self.value.lock().unwrap() = Some(keyring.to_serialized());
93 Ok(())
94 }
95
96 fn forget(&self) -> Result<(), KeyError> {
97 *self.value.lock().unwrap() = None;
98 Ok(())
99 }
100}
101
102pub fn copy_payload_files(
109 from: &coven_foundation::store_dir::StoreDir,
110 to: &coven_foundation::store_dir::StoreDir,
111) {
112 let source = from.payload_spool_dir();
113 let destination = to.payload_spool_dir();
114 match std::fs::metadata(&source) {
115 Ok(metadata) if metadata.is_dir() => {}
116 Ok(_) => panic!("payload file path is not a directory: {}", source.display()),
117 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return,
118 Err(error) => panic!(
119 "inspect payload file directory {}: {error}",
120 source.display()
121 ),
122 }
123 std::fs::create_dir_all(&destination).expect("create the copied payload spool directory");
124 for entry in std::fs::read_dir(&source).expect("read the payload spool being copied") {
125 let entry = entry.expect("payload spool entry");
126 std::fs::copy(entry.path(), destination.join(entry.file_name()))
127 .expect("copy one payload into the copied store directory");
128 }
129}
130
131pub fn pubkey_hex(kp: &UserKeypair) -> String {
134 coven_keys::keys::public_key_hex(kp)
135}
136
137pub fn user_keypair_from_seed(seed: [u8; 32]) -> UserKeypair {
139 let signing_key = ed25519_dalek::SigningKey::from_bytes(&seed);
140 UserKeypair::from_signing_key_bytes(&signing_key.to_keypair_bytes())
141 .expect("seed-derived signing key is valid")
142}
143
144pub struct TestDropboxAccessAdministrator {
148 pub namespace_id: String,
149}
150
151#[async_trait::async_trait]
152impl crate::sync::store::DeviceProviderAccessAdministrator for TestDropboxAccessAdministrator {
153 async fn grant_member_access(
154 &self,
155 _member_pubkey: &str,
156 _provider_account_email: Option<&str>,
157 peer: &coven_protocol::objects::ProviderDeviceBinding,
158 ) -> Result<coven_protocol::provider::ProviderAccessLocator, crate::sync::store::DeviceJoinError>
159 {
160 let coven_protocol::objects::ProviderPrincipalId::Dropbox { account_id } = &peer.principal
161 else {
162 return Err(crate::sync::store::DeviceJoinError::Provider(
163 "test Dropbox access administrator received a non-Dropbox peer".to_string(),
164 ));
165 };
166 Ok(
167 coven_protocol::provider::ProviderAccessLocator::DropboxSharedFolderMember {
168 namespace_id: self.namespace_id.clone(),
169 account_id: account_id.clone(),
170 },
171 )
172 }
173}
174
175pub struct CrossPrincipalTestDevice {
176 storage: std::sync::Arc<dyn coven_storage::CloudSyncObjectStorage>,
177 access_administrator: TestDropboxAccessAdministrator,
178}
179
180impl CrossPrincipalTestDevice {
181 pub async fn pending_device_join_observation(
182 &self,
183 pending: &crate::sync::store::DeviceJoinJournalDatabase,
184 root: &coven_protocol::store_commit::StoreRootRef,
185 attempt_id: coven_protocol::store_commit::DeviceJoinAttemptId,
186 ) -> Result<crate::sync::store::PendingDeviceJoinObservation<'_>, TestError> {
187 crate::sync::store::PendingDeviceJoinObservation::open(
188 pending,
189 &self.storage,
190 root,
191 attempt_id,
192 )
193 .await
194 .map_err(TestError::from)
195 }
196
197 #[cfg(test)]
205 pub async fn install_store_snapshot<'a>(
206 &'a self,
207 store_dir: &'a coven_foundation::store_dir::StoreDir,
208 root: &coven_protocol::store_commit::StoreRootRef,
209 membership: &coven_protocol::membership::MembershipChain,
210 identity: &UserKeypair,
211 device_id: String,
212 binary_schema_version: u32,
213 synced_tables: Vec<coven_protocol::synced_schema::SyncedTable>,
214 migrations: &[coven_database::Migration],
215 ) -> Result<crate::sync::store::RestoringStore<'a>, TestError> {
216 let history_verifier = crate::sync::store::HistoryConstructionAuthority::for_snapshot()
217 .open_pinned(self.storage.as_ref(), root)
218 .await
219 .map_err(crate::sync::store::SnapshotError::from)?;
220 let cancel = tokio::sync::watch::channel(false).1;
221 store_dir.ensure_created()?;
222 Ok(crate::sync::store::PreparedSnapshotBootstrap::prepare(
223 &self.storage,
224 history_verifier,
225 &coven_protocol::membership::MembershipFloor(membership.head_refs().to_vec()),
226 binary_schema_version,
227 &store_dir.db_path(),
228 identity,
229 std::sync::Arc::new(|_| {}),
230 &cancel,
231 )
232 .await?
233 .install(
234 store_dir,
235 synced_tables,
236 coven_protocol::blob::BLOB_TOMBSTONE_GRACE,
237 coven_protocol::blob::TransferLimits::one_at_a_time(),
238 device_id,
239 std::sync::Arc::new(coven_foundation::clock::SystemClock),
240 migrations,
241 coven_database::CovenMigrationPolicy::ApplyPending,
242 None,
243 )
244 .await?)
245 }
246
247 pub async fn open_pending_device_join(
248 &self,
249 pending: &crate::sync::store::DeviceJoinJournalDatabase,
250 identity: &UserKeypair,
251 offer: coven_protocol::store_commit::device_join_exchange::DeviceJoinOffer,
252 ) -> Result<crate::sync::store::PendingDeviceJoinAuthority<'_>, TestError> {
253 let observation = self
254 .pending_device_join_observation(pending, &offer.store_root, offer.attempt_id)
255 .await?;
256 crate::sync::store::PendingDeviceJoinAuthority::open(observation, identity, offer)
257 .await
258 .map_err(TestError::from)
259 }
260
261 pub async fn authorize_device_provider_access(
262 &self,
263 owner: &TestDevice,
264 request: coven_protocol::store_commit::device_join_exchange::DeviceProviderAccessRequest,
265 ) -> Result<
266 coven_protocol::store_commit::device_join_exchange::DeviceProviderAdmissionApproval,
267 crate::sync::DeviceJoinError,
268 > {
269 owner
270 .authorize_device_provider_access(request, Some(&self.access_administrator))
271 .await
272 }
273}
274
275pub struct TestStore {
276 home: std::sync::Arc<coven_storage::cloud::test_utils::InMemoryCloudHome>,
277 provider_requests: Option<std::sync::Arc<dyn coven_foundation::stage_timing::ProviderRequests>>,
280 storage: std::sync::Arc<coven_storage::CloudSyncConnection>,
281 root: coven_protocol::store_commit::StoreRootRef,
282 signer: UserKeypair,
283 founder: TestDevice,
284 producers: Arc<tokio::sync::Mutex<TestStoreProducers>>,
285}
286
287pub type TestStoreParts = (Arc<TestStore>, Arc<coven_storage::CloudSyncConnection>);
288
289#[derive(Debug)]
290pub struct TestError(Box<TestErrorCause>);
291
292impl TestError {
293 pub(crate) fn invariant(message: impl Into<String>) -> Self {
294 Self(Box::new(TestErrorCause::Invariant(message.into())))
295 }
296
297 #[cfg(test)]
298 pub(crate) fn initialization_source(
299 &self,
300 ) -> Option<&crate::sync::store::StoreInitializationError> {
301 match self.0.as_ref() {
302 TestErrorCause::Initialization(error) => Some(error),
303 _ => None,
304 }
305 }
306}
307
308impl std::fmt::Display for TestError {
309 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
310 self.0.fmt(formatter)
311 }
312}
313
314impl std::error::Error for TestError {
315 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
316 Some(self.0.as_ref())
317 }
318}
319
320#[derive(Debug, thiserror::Error)]
321enum TestErrorCause {
322 #[error("test fixture invariant failed: {0}")]
323 Invariant(String),
324 #[error(transparent)]
325 Database(#[from] coven_database::DbError),
326 #[error(transparent)]
327 HostWrite(#[from] coven_database::HostWriteError<coven_database::DbError>),
328 #[error(transparent)]
329 Storage(#[from] coven_protocol::objects::StorageError),
330 #[error(transparent)]
331 Protocol(#[from] coven_protocol::store_commit::StoreProtocolError),
332 #[error(transparent)]
333 Initialization(#[from] crate::sync::store::StoreInitializationError),
334 #[error(transparent)]
335 Store(#[from] crate::sync::store::StoreError),
336 #[error(transparent)]
337 Registration(#[from] crate::sync::store::StoreRegistrationError),
338 #[error(transparent)]
339 WriterAuthorization(#[from] crate::sync::store::StoreWriterAuthorizationError),
340 #[error(transparent)]
341 Pull(#[from] crate::sync::store::StorePullError),
342 #[error(transparent)]
343 Cycle(#[from] crate::sync::cycle::SyncCycleFailure),
344 #[error(transparent)]
345 SyncInitialization(#[from] crate::sync::cycle::InitSyncError),
346 #[error(transparent)]
347 Membership(#[from] crate::sync::store::MembershipOpsError),
348 #[error(transparent)]
349 MembershipMutation(#[from] crate::sync::store::MembershipMutationError),
350 #[error(transparent)]
351 DeviceJoin(#[from] crate::sync::store::DeviceJoinError),
352 #[error(transparent)]
353 OwnerPromotion(#[from] crate::sync::store::OwnerPromotionError),
354 #[error(transparent)]
355 Snapshot(#[from] crate::sync::store::SnapshotError),
356 #[error(transparent)]
357 Acknowledgement(#[from] crate::sync::store::StoreAckError),
358 #[error(transparent)]
359 PublishedBlobDrop(#[from] crate::sync::store::blob::PublishedBlobDropError),
360 #[error(transparent)]
361 Encryption(#[from] coven_keys::encryption::EncryptionError),
362 #[error(transparent)]
363 Key(#[from] coven_keys::keys::KeyError),
364 #[error(transparent)]
365 File(#[from] coven_foundation::atomic_file::FileError),
366 #[error(transparent)]
367 Json(#[from] serde_json::Error),
368 #[error(transparent)]
369 Io(#[from] std::io::Error),
370 #[error(transparent)]
371 TestPull(#[from] TestPullError),
372 #[error(transparent)]
373 DatabaseOpen(#[from] coven_database::OpenError),
374}
375
376macro_rules! test_error_from {
377 ($source:ty, $variant:ident) => {
378 impl From<$source> for TestError {
379 fn from(source: $source) -> Self {
380 Self(Box::new(TestErrorCause::$variant(source)))
381 }
382 }
383 };
384}
385
386test_error_from!(coven_database::DbError, Database);
387test_error_from!(
388 coven_database::HostWriteError<coven_database::DbError>,
389 HostWrite
390);
391test_error_from!(coven_protocol::objects::StorageError, Storage);
392test_error_from!(coven_protocol::store_commit::StoreProtocolError, Protocol);
393test_error_from!(crate::sync::store::StoreInitializationError, Initialization);
394test_error_from!(crate::sync::store::StoreError, Store);
395test_error_from!(crate::sync::store::StoreRegistrationError, Registration);
396test_error_from!(
397 crate::sync::store::StoreWriterAuthorizationError,
398 WriterAuthorization
399);
400test_error_from!(crate::sync::store::StorePullError, Pull);
401test_error_from!(crate::sync::cycle::SyncCycleFailure, Cycle);
402test_error_from!(crate::sync::cycle::InitSyncError, SyncInitialization);
403test_error_from!(crate::sync::store::MembershipOpsError, Membership);
404test_error_from!(
405 crate::sync::store::MembershipMutationError,
406 MembershipMutation
407);
408test_error_from!(crate::sync::store::DeviceJoinError, DeviceJoin);
409test_error_from!(crate::sync::store::OwnerPromotionError, OwnerPromotion);
410test_error_from!(crate::sync::store::SnapshotError, Snapshot);
411test_error_from!(crate::sync::store::StoreAckError, Acknowledgement);
412test_error_from!(
413 crate::sync::store::blob::PublishedBlobDropError,
414 PublishedBlobDrop
415);
416test_error_from!(coven_keys::encryption::EncryptionError, Encryption);
417test_error_from!(coven_keys::keys::KeyError, Key);
418test_error_from!(coven_foundation::atomic_file::FileError, File);
419test_error_from!(serde_json::Error, Json);
420test_error_from!(std::io::Error, Io);
421test_error_from!(TestPullError, TestPull);
422test_error_from!(coven_database::OpenError, DatabaseOpen);
423
424#[derive(Debug, thiserror::Error)]
428pub enum TestPullError {
429 #[error("open Store: {0}")]
430 Open(#[source] crate::sync::store::StoreInitializationError),
431 #[error("authorize Store writer: {0}")]
432 Authorize(#[from] crate::sync::store::StoreWriterAuthorizationError),
433 #[error("pull: {0}")]
434 Pull(#[from] crate::sync::cycle::SyncCycleFailure),
435}
436
437mod test_device {
438 use super::*;
439
440 pub struct TestDeviceSigningAuthority {
441 registration: coven_protocol::store_commit::ReferencedStoreDeviceRegistration,
442 device_signer: UserKeypair,
443 }
444
445 impl TestDeviceSigningAuthority {
446 pub fn registration_ref(
447 &self,
448 ) -> &coven_protocol::store_commit::StoreDeviceRegistrationRef {
449 self.registration.reference()
450 }
451
452 pub fn registration(&self) -> &coven_protocol::store_commit::StoreDeviceRegistration {
453 self.registration.value()
454 }
455
456 pub fn referenced_registration(
457 &self,
458 ) -> &coven_protocol::store_commit::ReferencedStoreDeviceRegistration {
459 &self.registration
460 }
461
462 pub fn sign_provider_admission_approval_without_shape_validation_for_test(
463 &self,
464 request: coven_protocol::store_commit::device_join_exchange::DeviceProviderAccessRequest,
465 admission: coven_protocol::store_commit::device_join_exchange::DeviceProviderAdmission,
466 ) -> coven_protocol::store_commit::device_join_exchange::DeviceProviderAdmissionApproval
467 {
468 coven_protocol::store_commit::device_join_exchange::DeviceProviderAdmissionApproval::signed_without_shape_validation_for_test(
469 request,
470 admission,
471 &self.device_signer,
472 )
473 }
474
475 pub fn sign_reclaim_receipt_for_test(
476 &self,
477 store_root_hash: coven_protocol::store_commit::ObjectHash,
478 authorization: coven_protocol::reclaim::ReclaimAuthorizationRef,
479 provider_admin_state: coven_protocol::circle_control::StoreMembershipStateRef,
480 provider_admin_grant: coven_protocol::provider::ProviderAdminGrantId,
481 ) -> Result<
482 coven_protocol::reclaim::ReclaimReceipt,
483 coven_protocol::store_commit::StoreProtocolError,
484 > {
485 coven_protocol::reclaim::ReclaimReceipt::signed(
486 store_root_hash,
487 authorization,
488 provider_admin_state,
489 provider_admin_grant,
490 self.registration.reference().clone(),
491 self.registration.value(),
492 &self.device_signer,
493 )
494 }
495 }
496
497 #[derive(Clone)]
498 pub struct TestDevice {
499 db: coven_database::StoreDatabase,
500 store: std::sync::Arc<crate::sync::store::Store>,
501 store_dir: StoreDir,
502 device_id: String,
503 storage: std::sync::Arc<coven_storage::CloudSyncConnection>,
504 identity: UserKeypair,
505 settled: std::sync::Arc<crate::sync::store::SettledCycle>,
508 }
509
510 impl TestDevice {
511 pub async fn query_test_text(&self, sql: &str) -> String {
520 self.db
521 .test_query_optional_text(sql.to_string())
522 .await
523 .expect("test text query failed")
524 .expect("test text query matched no row")
525 }
526
527 pub async fn test_row_exists(&self, sql: &str) -> bool {
528 self.db
529 .test_query_optional_text(format!("SELECT 'found' FROM ({sql}) LIMIT 1"))
530 .await
531 .expect("test row-existence query failed")
532 .is_some()
533 }
534
535 pub async fn latest_local_store_device_registration(
536 &self,
537 ) -> Result<Option<coven_database::DurableDeviceRegistration>, coven_database::DbError>
538 {
539 self.db.latest_local_store_device_registration().await
540 }
541
542 pub async fn replay_row_count_for_test(
543 &self,
544 table: &str,
545 ) -> Result<i64, coven_database::DbError> {
546 self.db
547 .replay_row_count_for_test(
548 self.store.root_ref_for_test().clone(),
549 table.to_string(),
550 )
551 .await
552 }
553
554 pub fn device_id(&self) -> String {
555 self.device_id.clone()
556 }
557
558 pub fn host_write_blob_staging(&self) -> crate::sync::store::HostWriteBlobStaging {
559 self.store
560 .host_write_blob_staging(tokio::runtime::Handle::current())
561 }
562
563 pub async fn create(
564 db: &Database,
565 store_dir: StoreDir,
566 storage: std::sync::Arc<coven_storage::CloudSyncConnection>,
567 founder_timestamp: &str,
568 identity: UserKeypair,
569 ) -> Result<Self, crate::sync::store::StoreInitializationError> {
570 Self::create_with_database(
571 coven_database::StoreDatabase::new(db),
572 store_dir,
573 storage,
574 founder_timestamp,
575 identity,
576 )
577 .await
578 }
579
580 pub async fn create_with_database(
581 database: coven_database::StoreDatabase,
582 store_dir: StoreDir,
583 storage: std::sync::Arc<coven_storage::CloudSyncConnection>,
584 founder_timestamp: &str,
585 identity: UserKeypair,
586 ) -> Result<Self, crate::sync::store::StoreInitializationError> {
587 database.assert_owns_payload_directory_for_test(&store_dir);
588 let initialized = crate::sync::store::Store::create(
589 database.clone(),
590 storage.clone(),
591 store_dir.clone(),
592 founder_timestamp,
593 &identity,
594 Some(coven_keys::encryption::EncryptionService::from_key(
595 [42; 32],
596 )),
597 )
598 .await?;
599 let (store, device_id) = initialized.into_parts();
600 Ok(Self {
601 db: database,
602 store: std::sync::Arc::new(store),
603 store_dir,
604 device_id,
605 storage,
606 identity,
607 settled: std::sync::Arc::default(),
608 })
609 }
610
611 pub async fn open_with_database(
612 database: coven_database::StoreDatabase,
613 store_dir: StoreDir,
614 storage: std::sync::Arc<coven_storage::CloudSyncConnection>,
615 root: &coven_protocol::store_commit::StoreRootRef,
616 identity: &UserKeypair,
617 ) -> Result<Self, crate::sync::store::StoreInitializationError> {
618 database.assert_owns_payload_directory_for_test(&store_dir);
619 let initialized = crate::sync::store::Store::open(
620 database.clone(),
621 storage.clone(),
622 store_dir.clone(),
623 root,
624 identity,
625 Some(coven_keys::encryption::EncryptionService::from_key(
626 [42; 32],
627 )),
628 )
629 .await?;
630 let (store, device_id) = initialized.into_parts();
631 Ok(Self {
632 db: database,
633 store: std::sync::Arc::new(store),
634 store_dir,
635 device_id,
636 storage,
637 identity: identity.clone(),
638 settled: std::sync::Arc::default(),
639 })
640 }
641
642 pub async fn activate_joined(
643 observer: Self,
644 joining_database: coven_database::StoreDatabase,
645 joining_store_dir: StoreDir,
646 joining_identity: &UserKeypair,
647 published_at: &str,
648 storage: std::sync::Arc<coven_storage::CloudSyncConnection>,
649 ) -> Result<Self, TestError> {
650 Self::activate_joined_with_clock(
651 observer,
652 joining_database,
653 joining_store_dir,
654 joining_identity,
655 published_at,
656 storage,
657 Arc::new(coven_foundation::clock::SystemClock),
658 )
659 .await
660 }
661
662 pub fn activate_joined_with_clock<'a>(
663 observer: Self,
664 joining_database: coven_database::StoreDatabase,
665 joining_store_dir: StoreDir,
666 joining_identity: &'a UserKeypair,
667 published_at: &'a str,
668 storage: std::sync::Arc<coven_storage::CloudSyncConnection>,
669 clock: coven_foundation::clock::ClockRef,
670 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<Self, TestError>> + Send + 'a>>
671 {
672 Box::pin(async move {
673 let activated_database = joining_database.clone();
674 let snapshot = observer.ensure_device_join_snapshot_for_test();
675 snapshot.await?;
676 let pending_dir = tempfile::tempdir()?;
677 let pending = crate::sync::store::DeviceJoinJournalDatabase::open_for_test(
678 pending_dir.path().join("pending-device-join.sqlite"),
679 )?;
680 let joining_pubkey = pubkey_hex(joining_identity);
681 let begin = observer.begin_device_join(&joining_pubkey);
682 let offer = begin.await?;
683 let mut pending_join = observer
684 .open_pending_device_join_for_test(&pending, joining_identity, offer)
685 .await?;
686 let access_request = pending_join.prepare_provider_access_request().await?;
687 let authorize = observer.authorize_device_provider_access(access_request, None);
688 let approval = Box::pin(authorize).await?;
689 let registration_request =
690 pending_join.prepare_registration_request(approval).await?;
691 let activate = observer.activate_same_principal_join_for_test(registration_request);
692 let join = activate.await?;
693 let routing_encryption =
694 coven_keys::encryption::EncryptionService::from_key([42; 32]);
695 let object_storage: std::sync::Arc<dyn coven_storage::CloudSyncObjectStorage> =
696 storage.clone();
697 let progress: crate::sync::JoiningDeviceJoinProgressObserver =
698 std::sync::Arc::new(|_| {});
699 let cancel = tokio::sync::watch::channel(false).1;
700 let device_id = join
701 .bootstrap
702 .bootstrap
703 .request
704 .expected_registration()
705 .device_id
706 .to_string();
707 let history = crate::sync::store::HistoryConstructionAuthority::admission()
708 .open_pinned(
709 object_storage.as_ref(),
710 &join.installation.authority.store_root,
711 )
712 .await
713 .map_err(crate::sync::store::SnapshotError::from)?;
714 let accepted_membership = history
715 .load_accepted_membership_authority(
716 &join.installation.bootstrap.membership.0,
717 None,
718 )
719 .await
720 .map_err(crate::sync::store::SnapshotError::from)?;
721 let prepared = crate::sync::store::PreparedDeviceJoinSnapshot::prepare(
722 &object_storage,
723 (*join.installation).clone(),
724 &accepted_membership,
725 joining_database.schema_version(),
726 &joining_store_dir.db_path(),
727 &progress,
728 &cancel,
729 )
730 .await?;
731 let installed = prepared.install(
732 joining_database.synced_tables_for_test(),
733 joining_database.blob_tombstone_grace(),
734 joining_database.transfer_limits(),
735 device_id,
736 clock,
737 &coven_database::synthetic_store::test_migrations(),
738 coven_database::CovenMigrationPolicy::ApplyPending,
739 &routing_encryption,
740 )?;
741 drop(pending_join);
742 let completion =
743 crate::sync::store::PendingDeviceJoinAuthority::prepare_same_principal_completion(
744 &pending,
745 &object_storage,
746 &joining_store_dir,
747 joining_identity,
748 join,
749 installed,
750 published_at,
751 Some(&routing_encryption),
752 Some(accepted_membership.chain().clone()),
753 );
754 let completion = Box::pin(completion).await?;
755 completion.complete().await?;
756 let database_image =
757 coven_database::DatabaseImageTest::open(&joining_store_dir.db_path())?
758 .into_bytes()?;
759 joining_database
760 .replace_with_database_image_for_test(database_image)
761 .await?;
762 Self::load_with_database(
763 activated_database,
764 storage,
765 joining_identity.clone(),
766 joining_store_dir,
767 )
768 .await
769 .map_err(TestError::from)
770 })
771 }
772
773 #[cfg(test)]
777 #[allow(clippy::too_many_arguments)]
778 pub async fn activate_joined_from_snapshot(
779 observer: Self,
780 joining_store_dir: StoreDir,
781 joining_identity: &UserKeypair,
782 published_at: &str,
783 storage: std::sync::Arc<coven_storage::CloudSyncConnection>,
784 synced_tables: Vec<coven_protocol::synced_schema::SyncedTable>,
785 migrations: Vec<coven_database::Migration>,
786 binary_schema_version: u32,
787 ) -> Result<Self, TestError> {
788 observer.ensure_device_join_snapshot_for_test().await?;
789 let pending_dir = tempfile::tempdir()?;
790 let pending = crate::sync::store::DeviceJoinJournalDatabase::open_for_test(
791 pending_dir.path().join("pending-device-join.sqlite"),
792 )?;
793 let offer = observer
794 .begin_device_join(&pubkey_hex(joining_identity))
795 .await?;
796 let mut pending_join = observer
797 .open_pending_device_join_for_test(&pending, joining_identity, offer.clone())
798 .await?;
799 let access_request = pending_join.prepare_provider_access_request().await?;
800 let approval = observer
801 .authorize_device_provider_access(access_request, None)
802 .await?;
803 let registration_request = pending_join.prepare_registration_request(approval).await?;
804 let join = observer
805 .activate_same_principal_join_for_test(registration_request)
806 .await?;
807 drop(pending_join);
808 let routing_encryption = coven_keys::encryption::EncryptionService::from_key([42; 32]);
809 let history_verifier = crate::sync::store::HistoryConstructionAuthority::for_snapshot()
810 .open_pinned(storage.as_ref(), &offer.store_root)
811 .await
812 .map_err(crate::sync::store::SnapshotError::from)?;
813 let cancel = tokio::sync::watch::channel(false).1;
814 joining_store_dir.ensure_created()?;
815 let open_synced_tables = synced_tables.clone();
816 let joined_device_id = offer.attempt_id.to_string();
817 let membership = observer.membership().await?;
818 let storage_object: std::sync::Arc<dyn coven_storage::CloudSyncObjectStorage> =
819 storage.clone();
820 let restoring = crate::sync::store::PreparedSnapshotBootstrap::prepare(
821 &storage_object,
822 history_verifier,
823 &coven_protocol::membership::MembershipFloor(membership.head_refs().to_vec()),
824 binary_schema_version,
825 &joining_store_dir.db_path(),
826 joining_identity,
827 std::sync::Arc::new(|_| {}),
828 &cancel,
829 )
830 .await?
831 .install(
832 &joining_store_dir,
833 synced_tables,
834 coven_protocol::blob::BLOB_TOMBSTONE_GRACE,
835 coven_protocol::blob::TransferLimits::one_at_a_time(),
836 offer.attempt_id.to_string(),
837 std::sync::Arc::new(coven_foundation::clock::SystemClock),
838 &migrations,
839 coven_database::CovenMigrationPolicy::ApplyPending,
840 Some(&routing_encryption),
841 )
842 .await?;
843 let mut joining = restoring.begin_device_join(&pending, offer).await?;
844 joining
845 .bootstrap(
846 join.bootstrap.clone(),
847 published_at,
848 Some(&routing_encryption),
849 )
850 .await?;
851 joining.complete(join.activation).await?;
852 drop(joining);
857 open_joined_test_device(
858 joining_store_dir,
859 joining_identity,
860 storage,
861 joined_device_id,
862 open_synced_tables,
863 &migrations,
864 )
865 .await
866 }
867
868 pub async fn latest_local_store_snapshot_for_test(
869 &self,
870 ) -> Result<Option<coven_database::PublishedStoreSnapshot>, TestError> {
871 Ok(self.db.latest_local_store_snapshot().await?)
872 }
873
874 pub fn publish_snapshot_generation_for_test(
877 &self,
878 ) -> std::pin::Pin<
879 Box<
880 dyn std::future::Future<
881 Output = Result<coven_database::PublishedStoreSnapshot, TestError>,
882 > + Send
883 + '_,
884 >,
885 > {
886 Box::pin(async move {
887 let routing_encryption =
888 coven_keys::encryption::EncryptionService::from_key([42; 32]);
889 let mut writer = self.authorize_writer().await?;
890 writer.seed_retained_history().await?;
891 let mut snapshots = writer.snapshots();
892 let capture = snapshots.capture_snapshot_cut(Some(&routing_encryption));
893 let cut = capture.await?;
894 let publish = snapshots.push_snapshot_cut(cut, "2026-07-16T00:00:00Z".into());
895 let metadata = publish.await?;
896 self.publish_acknowledgement(metadata.coverage.clone())
897 .await?;
898 self.db.latest_local_store_snapshot().await?.ok_or_else(|| {
899 TestError::invariant("the published generation is absent".to_string())
900 })
901 })
902 }
903
904 pub async fn ensure_device_join_snapshot_for_test(&self) -> Result<(), TestError> {
905 if self.db.latest_local_store_snapshot().await?.is_some() {
906 return Ok(());
907 }
908 self.publish_snapshot_generation_for_test().await?;
909 Ok(())
910 }
911
912 pub async fn load(
913 db: &Database,
914 store_dir: StoreDir,
915 storage: std::sync::Arc<coven_storage::CloudSyncConnection>,
916 identity: UserKeypair,
917 ) -> Result<Self, crate::sync::store::StoreError> {
918 Self::load_with_database(
919 coven_database::StoreDatabase::new(db),
920 storage,
921 identity,
922 store_dir,
923 )
924 .await
925 }
926
927 pub async fn load_with_database(
928 database: coven_database::StoreDatabase,
929 storage: std::sync::Arc<coven_storage::CloudSyncConnection>,
930 identity: UserKeypair,
931 store_dir: StoreDir,
932 ) -> Result<Self, crate::sync::store::StoreError> {
933 database.assert_owns_payload_directory_for_test(&store_dir);
934 let store = crate::sync::store::Store::load(
935 database.clone(),
936 storage.clone(),
937 store_dir.clone(),
938 identity.clone(),
939 Some(coven_keys::encryption::EncryptionService::from_key(
940 [42; 32],
941 )),
942 )
943 .await?;
944 let device_id = database
945 .get_protocol_state(coven_database::LOCAL_DEVICE_ID_STATE_KEY)
946 .await?
947 .ok_or(crate::sync::store::StoreError::MissingState {
948 key: coven_database::LOCAL_DEVICE_ID_STATE_KEY,
949 })?;
950 Ok(Self {
951 db: database,
952 store: std::sync::Arc::new(store),
953 store_dir,
954 device_id,
955 storage,
956 identity,
957 settled: std::sync::Arc::default(),
958 })
959 }
960
961 pub fn adopt_key_rotation(
962 &self,
963 encryption: &coven_keys::encryption::EncryptionService,
964 custody: &dyn coven_keys::keys::MasterKeyCustody,
965 ) -> Result<String, coven_keys::keys::KeyError> {
966 self.storage
967 .adopt_key_rotation_for_test(encryption, custody)
968 }
969
970 pub fn store_root(&self) -> &coven_protocol::store_commit::StoreRootRef {
971 self.store.store_root()
972 }
973
974 pub async fn authorize_writer(
975 &self,
976 ) -> Result<crate::sync::store::AuthorizedWriterOperation<'_>, crate::sync::store::StoreError>
977 {
978 self.store
979 .authorize_writer()
980 .await
981 .map_err(crate::sync::store::StoreError::from)
982 }
983
984 pub async fn execute_unscoped_host_sql_for_test(
985 &self,
986 sql: String,
987 ) -> Result<(), coven_database::HostWriteError<coven_database::DbError>> {
988 self.store.execute_unscoped_host_sql_for_test(sql).await
989 }
990
991 pub async fn membership_for_test(
992 &self,
993 ) -> Result<coven_protocol::membership::MembershipChain, crate::sync::store::StoreError>
994 {
995 self.store.membership_for_test().await
996 }
997
998 pub async fn latest_local_store_position(
999 &self,
1000 ) -> Result<
1001 Option<coven_protocol::store_commit::StoreBatchCommitRef>,
1002 crate::sync::store::StoreError,
1003 > {
1004 self.store.latest_local_store_position().await
1005 }
1006
1007 pub async fn load_commit_for_test(
1008 &self,
1009 reference: &coven_protocol::store_commit::StoreBatchCommitRef,
1010 ) -> Result<
1011 coven_protocol::store_commit::VerifiedStoreBatchCommit,
1012 crate::sync::store::StoreError,
1013 > {
1014 self.store.load_commit_for_test(reference).await
1015 }
1016
1017 pub async fn load_membership_head_for_test(
1018 &self,
1019 reference: &coven_protocol::membership::MembershipHeadRef,
1020 ) -> Result<coven_protocol::membership::AuthorHead, crate::sync::store::StoreError>
1021 {
1022 self.store.load_membership_head_for_test(reference).await
1023 }
1024
1025 pub async fn load_exact_materialized_commit(
1026 &self,
1027 stream_id: &str,
1028 sequence: u64,
1029 ) -> Result<
1030 Option<(
1031 coven_protocol::store_commit::StoreBatchCommitRef,
1032 coven_protocol::store_commit::VerifiedStoreBatchCommit,
1033 )>,
1034 crate::sync::store::StoreError,
1035 > {
1036 self.store
1037 .load_exact_materialized_commit(stream_id, sequence)
1038 .await
1039 }
1040
1041 pub fn device_join_transport(
1042 &self,
1043 ) -> crate::sync::store::device_join::transport::StoreDeviceJoinTransport<'_> {
1044 self.store.device_join_transport()
1045 }
1046
1047 pub fn circles(&self) -> crate::sync::store::StoreCircleCommands<'_> {
1048 self.store.circles()
1049 }
1050
1051 pub async fn circle_epoch_access(
1052 &self,
1053 circle_id: coven_protocol::circle::CircleId,
1054 expected_control: coven_protocol::circle::CircleControlCoord,
1055 ) -> Result<
1056 Option<coven_protocol::circle_activation::CircleEpochAccess>,
1057 coven_database::DbError,
1058 > {
1059 self.store
1060 .circle_epoch_access(circle_id, expected_control)
1061 .await
1062 }
1063
1064 pub async fn discard_blocked_write(
1065 &self,
1066 write_id: coven_protocol::write::WriteId,
1067 ) -> Result<Vec<coven_protocol::write::WriteId>, crate::sync::store::StoreError> {
1068 let routing_encryption = coven_keys::encryption::EncryptionService::from_key([42; 32]);
1069 self.store
1070 .discard_blocked_write(write_id, Some(&routing_encryption))
1071 .await
1072 }
1073
1074 pub async fn restore_membership(
1075 &self,
1076 ) -> Result<
1077 crate::sync::store::StoreRestoreMembership,
1078 crate::sync::store::MembershipOpsError,
1079 > {
1080 self.store.restore_membership().await
1081 }
1082
1083 pub async fn owner_recovery_for_test(
1084 &self,
1085 ) -> Result<crate::sync::store::RestoringStore<'_>, TestError> {
1086 self.store.owner_recovery_for_test().await
1087 }
1088
1089 pub async fn begin_device_join(
1090 &self,
1091 member_pubkey: &str,
1092 ) -> Result<crate::sync::DeviceJoinOffer, crate::sync::DeviceJoinError> {
1093 self.store.begin_device_join(member_pubkey).await
1094 }
1095
1096 pub async fn begin_owner_promotion_for_device(
1097 &self,
1098 device_id: coven_protocol::StoreDeviceId,
1099 ) -> Result<
1100 coven_protocol::store_commit::OwnerPromotionRequest,
1101 crate::sync::store::OwnerPromotionError,
1102 > {
1103 self.store.begin_owner_promotion_for_device(device_id).await
1104 }
1105
1106 pub async fn begin_owner_promotion(
1107 &self,
1108 member_registration: coven_protocol::store_commit::StoreDeviceRegistrationRef,
1109 ) -> Result<
1110 coven_protocol::store_commit::OwnerPromotionRequest,
1111 crate::sync::store::OwnerPromotionError,
1112 > {
1113 self.store.begin_owner_promotion(member_registration).await
1114 }
1115
1116 pub async fn accept_owner_promotion(
1117 &self,
1118 request: coven_protocol::store_commit::OwnerPromotionRequest,
1119 ) -> Result<
1120 coven_protocol::store_commit::OwnerPromotionAcceptance,
1121 crate::sync::store::OwnerPromotionError,
1122 > {
1123 self.store.accept_owner_promotion(request).await
1124 }
1125
1126 pub async fn finalize_owner_promotion(
1127 &self,
1128 encryption: &coven_keys::encryption::EncryptionService,
1129 acceptance: coven_protocol::store_commit::OwnerPromotionAcceptance,
1130 ) -> Result<
1131 coven_protocol::circle_control::StoreMembershipStateRef,
1132 crate::sync::store::OwnerPromotionError,
1133 > {
1134 self.store
1135 .finalize_owner_promotion(encryption, acceptance)
1136 .await
1137 }
1138
1139 pub async fn blob_key_fingerprint_for_test(
1140 &self,
1141 authority: &coven_protocol::blob::RowBlobAuthority,
1142 stored: &coven_protocol::blob::locator::StoredBlobRef,
1143 ) -> Result<Option<coven_keys::encryption::KeyFingerprint>, TestError> {
1144 self.store
1145 .blob_key_fingerprint_for_test(authority, stored)
1146 .await
1147 .map_err(TestError::from)
1148 }
1149
1150 pub async fn announcement_stream_id_for_test(
1151 &self,
1152 ) -> Result<coven_protocol::membership::AuthorStreamId, crate::sync::store::StoreError>
1153 {
1154 self.store.announcement_stream_id_for_test().await
1155 }
1156
1157 pub async fn owner_promotion_target_for_test(
1158 &self,
1159 ) -> Result<
1160 coven_protocol::store_commit::StoreDeviceRegistrationRef,
1161 crate::sync::store::StoreError,
1162 > {
1163 self.store.owner_promotion_target_for_test().await
1164 }
1165
1166 pub async fn resign_snapshot_meta_for_test(
1167 &self,
1168 meta: coven_protocol::store_commit::SnapshotMeta,
1169 ) -> Result<coven_protocol::store_commit::SnapshotMeta, crate::sync::store::StoreError>
1170 {
1171 self.store.resign_snapshot_meta_for_test(meta).await
1172 }
1173
1174 pub async fn parse_local_snapshot_meta_for_test(
1175 &self,
1176 bytes: &[u8],
1177 reference: &coven_protocol::store_commit::StoreSnapshotRef,
1178 ) -> Result<coven_protocol::store_commit::SnapshotMeta, crate::sync::store::StoreError>
1179 {
1180 self.store
1181 .parse_local_snapshot_meta_for_test(bytes, reference)
1182 .await
1183 }
1184
1185 pub async fn prepare_operation_plan_for_test(
1186 &self,
1187 ) -> Result<crate::sync::store::StoreOperationCommitPlan, crate::sync::store::StoreError>
1188 {
1189 self.store.prepare_operation_plan_for_test().await
1190 }
1191
1192 pub async fn authorize_retained_outbound_for_test(
1193 &self,
1194 order: &coven_protocol::store_commit::StoreCommitOrder,
1195 candidate_membership_heads: &[coven_protocol::membership::MembershipHeadRef],
1196 ) -> Result<crate::sync::store::MergeOutboundAuthorization, crate::sync::store::StoreError>
1197 {
1198 self.store
1199 .authorize_retained_outbound_for_test(order, candidate_membership_heads)
1200 .await
1201 }
1202
1203 pub async fn complete_revoke_rotation_adoption_for_test(
1204 &self,
1205 pending_rotation: &dyn coven_storage::CloudSyncRotationStateAccess,
1206 adopted_generation: u64,
1207 ) -> Result<(), crate::sync::store::MembershipMutationError> {
1208 self.store
1209 .complete_revoke_rotation_adoption_for_test(pending_rotation, adopted_generation)
1210 .await
1211 }
1212
1213 pub async fn retained_merge_replay_inputs_for_test(
1214 &self,
1215 ) -> Result<Vec<coven_database::OwnedVerifiedMergeMaterialization>, coven_database::DbError>
1216 {
1217 self.store.retained_merge_replay_inputs_for_test().await
1218 }
1219
1220 pub async fn resolved_store_device_state_for_test(
1221 &self,
1222 reference: &coven_protocol::store_commit::StoreDeviceStateRef,
1223 ) -> Result<coven_protocol::store_commit::ResolvedStoreDeviceState, coven_database::DbError>
1224 {
1225 self.store
1226 .resolved_store_device_state_for_test(reference)
1227 .await
1228 }
1229
1230 pub async fn retained_merge_materialization_for_test(
1231 &self,
1232 reference: coven_protocol::store_commit::StoreBatchCommitRef,
1233 ) -> Result<coven_database::OwnedVerifiedMergeMaterialization, coven_database::DbError>
1234 {
1235 self.store
1236 .retained_merge_materialization_for_test(reference)
1237 .await
1238 }
1239
1240 pub async fn load_membership_at_exact_heads_for_test(
1241 &self,
1242 heads: &[coven_protocol::membership::MembershipHeadRef],
1243 ) -> Result<coven_protocol::membership::MembershipChain, crate::sync::store::StoreError>
1244 {
1245 self.store
1246 .load_membership_at_exact_heads_for_test(heads)
1247 .await
1248 }
1249
1250 pub async fn project_membership_for_test(
1251 &self,
1252 candidate_heads: &[coven_protocol::membership::MembershipHeadRef],
1253 ) -> Result<coven_protocol::membership::MembershipChain, crate::sync::store::StoreError>
1254 {
1255 self.store
1256 .project_membership_for_test(candidate_heads)
1257 .await
1258 }
1259
1260 pub async fn assert_deep_membership_projection_for_test(
1261 &self,
1262 heads: &[coven_protocol::membership::MembershipHeadRef],
1263 ) -> Result<(), crate::sync::store::StoreError> {
1264 self.store
1265 .assert_deep_membership_projection_for_test(heads)
1266 .await
1267 }
1268
1269 pub async fn load_registration_for_test(
1270 &self,
1271 reference: &coven_protocol::store_commit::StoreDeviceRegistrationRef,
1272 ) -> Result<
1273 coven_protocol::store_commit::StoreDeviceRegistration,
1274 crate::sync::store::StoreError,
1275 > {
1276 self.store.load_registration_for_test(reference).await
1277 }
1278
1279 pub async fn verify_installable_snapshots_for_test(
1280 &self,
1281 snapshots: &[coven_database::PublishedStoreSnapshot],
1282 ) -> Result<(), crate::sync::store::StoreError> {
1283 self.store
1284 .verify_installable_snapshots_for_test(snapshots)
1285 .await
1286 }
1287
1288 pub async fn open_circle_package_for_test(
1289 &self,
1290 access: &coven_protocol::circle_activation::CircleEpochAccess,
1291 commit: &coven_protocol::store_commit::VerifiedStoreBatchCommit,
1292 reference: &coven_protocol::store_commit::CirclePackageRef,
1293 ) -> Result<Vec<u8>, crate::sync::store::StoreError> {
1294 self.store
1295 .open_circle_package_for_test(access, commit, reference)
1296 .await
1297 }
1298
1299 #[allow(clippy::too_many_arguments)]
1300 pub async fn pull_readiness_for_test(
1301 &self,
1302 coverage: &coven_protocol::store_commit::CommitFrontier,
1303 frontier: &std::collections::BTreeMap<
1304 String,
1305 coven_protocol::store_commit::StoreBatchCommitRef,
1306 >,
1307 commit_ref: &coven_protocol::store_commit::StoreBatchCommitRef,
1308 commit: &coven_protocol::store_commit::StoreBatchCommit,
1309 ) -> Result<crate::sync::store::Readiness, crate::sync::store::StorePullError> {
1310 self.store
1311 .pull_readiness_for_test(coverage, frontier, commit_ref, commit)
1312 .await
1313 }
1314
1315 pub async fn verified_merge_membership_prefix_for_test(
1316 &self,
1317 references: impl IntoIterator<Item = coven_protocol::store_commit::StoreBatchCommitRef>,
1318 predecessors: impl IntoIterator<Item = coven_protocol::store_commit::StoreBatchCommitRef>,
1319 ) -> Result<
1320 crate::sync::store::VerifiedMergeMembershipPrefix,
1321 crate::sync::store::StorePullError,
1322 > {
1323 self.store
1324 .verified_merge_membership_prefix_for_test(references, predecessors)
1325 .await
1326 }
1327
1328 pub async fn retained_merge_history_frontier_for_test(
1329 &self,
1330 references: Vec<coven_protocol::store_commit::StoreBatchCommitRef>,
1331 ) -> Result<Vec<coven_database::RetainedMergeHistoryCheckpoint>, coven_database::DbError>
1332 {
1333 self.store
1334 .retained_merge_history_frontier_for_test(references)
1335 .await
1336 }
1337
1338 pub async fn verified_circle_activation_for_test(
1339 &self,
1340 circle_id: coven_protocol::circle::CircleId,
1341 control: coven_protocol::circle::CircleControlCoord,
1342 ) -> Result<
1343 Option<coven_protocol::circle_activation::VerifiedCircleReference>,
1344 coven_database::DbError,
1345 > {
1346 self.store
1347 .verified_circle_activation_for_test(circle_id, control)
1348 .await
1349 }
1350
1351 pub async fn finalized_circle_close_outcome_for_test(
1352 &self,
1353 circle_id: coven_protocol::circle::CircleId,
1354 ) -> Result<
1355 coven_protocol::circle::CircleEpochCloseOutcome,
1356 crate::sync::store::CircleOperationError,
1357 > {
1358 self.store
1359 .finalized_circle_close_outcome_for_test(circle_id)
1360 .await
1361 }
1362
1363 pub async fn circle_package_is_retained_for_replay_for_test(
1364 &self,
1365 target: coven_protocol::store_commit::CirclePackageRef,
1366 activation: coven_protocol::store_commit::StoreBatchCommitRef,
1367 ) -> Result<bool, coven_database::DbError> {
1368 self.store
1369 .circle_package_is_retained_for_replay_for_test(target, activation)
1370 .await
1371 }
1372
1373 pub async fn load_circle_acknowledgement_for_test(
1374 &self,
1375 reference: &coven_protocol::store_commit::CircleAckRef,
1376 ) -> Result<coven_protocol::store_commit::CircleAck, crate::sync::store::StoreAckError>
1377 {
1378 self.store
1379 .load_circle_acknowledgement_for_test(reference)
1380 .await
1381 }
1382
1383 pub async fn load_applicable_circle_packages_for_test(
1384 &self,
1385 verified: &coven_protocol::store_commit::VerifiedStoreBatchCommit,
1386 activations: &[&coven_protocol::circle_activation::VerifiedCircleActivations],
1387 author: &coven_protocol::store_commit::StoreDeviceRegistration,
1388 local_store_membership: coven_protocol::membership::LocalStoreMembership,
1389 ) -> Result<
1390 Vec<crate::sync::store::LoadedCirclePackage>,
1391 crate::sync::store::CirclePackageReadError,
1392 > {
1393 self.store
1394 .load_applicable_circle_packages_for_test(
1395 verified,
1396 activations,
1397 author,
1398 local_store_membership,
1399 )
1400 .await
1401 }
1402
1403 pub fn protocol_root_for_test(&self) -> &coven_protocol::store_commit::StoreProtocolRoot {
1404 self.store.protocol_root_for_test()
1405 }
1406
1407 pub async fn prepare_acknowledgement_activation_for_test(
1408 &self,
1409 acknowledgement: coven_protocol::store_commit::StoreAckRef,
1410 object: coven_protocol::objects::ExactProtocolObject<
1411 coven_protocol::store_commit::StoreAck,
1412 >,
1413 candidate: coven_protocol::prepared_commit::PreparedStoreOperationCommit,
1414 ) -> Result<(), coven_database::DbError> {
1415 self.store
1416 .prepare_acknowledgement_activation_for_test(acknowledgement, object, candidate)
1417 .await
1418 }
1419
1420 pub async fn prepare_merge_history_successor_for_test(
1421 &self,
1422 verified_commit: &coven_protocol::store_commit::VerifiedStoreBatchCommit,
1423 recovery_author: Option<&coven_protocol::store_commit::StoreDeviceRegistrationRef>,
1424 evidence: crate::sync::store::MergeHistorySuccessorEvidence,
1425 ) -> Result<crate::sync::store::PreparedMergeHistorySuccessor, crate::sync::store::StoreError>
1426 {
1427 self.store
1428 .prepare_merge_history_successor_for_test(
1429 verified_commit,
1430 recovery_author,
1431 evidence,
1432 )
1433 .await
1434 }
1435
1436 pub async fn prepare_device_join_bootstrap_for_test(
1437 &self,
1438 bootstrap_cut: &coven_protocol::store_commit::StoreHistoryCut,
1439 attempt_activation: &coven_protocol::store_commit::StoreBatchCommitRef,
1440 membership_state: &coven_protocol::circle_control::StoreMembershipStateRef,
1441 installed: &coven_protocol::store_commit::CommitFrontier,
1442 ) -> Result<coven_database::DeviceJoinBootstrapPlan, crate::sync::store::StoreError>
1443 {
1444 self.store
1445 .prepare_device_join_bootstrap_for_test(
1446 bootstrap_cut,
1447 attempt_activation,
1448 membership_state,
1449 installed,
1450 )
1451 .await
1452 }
1453
1454 pub async fn load_store_package_for_test(
1455 &self,
1456 reference: &coven_protocol::store_commit::StoreBatchCommitRef,
1457 ) -> Result<
1458 Option<coven_protocol::objects::VerifiedObject<Vec<u8>>>,
1459 crate::sync::store::StoreError,
1460 > {
1461 self.store.load_store_package_for_test(reference).await
1462 }
1463
1464 pub async fn load_store_ack_for_test(
1465 &self,
1466 reference: &coven_protocol::store_commit::StoreAckRef,
1467 registration: &coven_protocol::store_commit::StoreDeviceRegistration,
1468 ) -> Result<coven_protocol::store_commit::StoreAck, crate::sync::store::StoreError>
1469 {
1470 self.store
1471 .load_store_ack_for_test(reference, registration)
1472 .await
1473 }
1474
1475 #[allow(clippy::too_many_arguments)]
1476 pub async fn remove_member(
1477 &self,
1478 public_key_hex: &str,
1479 encryption: &coven_keys::encryption::EncryptionService,
1480 master_keys: &dyn coven_keys::keys::MasterKeyCustody,
1481 cipher: &dyn coven_storage::CloudSyncCipherStateAccess,
1482 pending_rotation: &dyn coven_storage::CloudSyncRotationStateAccess,
1483 ) -> Result<String, crate::sync::store::MembershipOpsError> {
1484 self.store
1485 .remove_member(
1486 public_key_hex,
1487 encryption,
1488 master_keys,
1489 cipher,
1490 pending_rotation,
1491 )
1492 .await
1493 }
1494
1495 pub async fn authorize_device_provider_access(
1496 &self,
1497 request: coven_protocol::store_commit::device_join_exchange::DeviceProviderAccessRequest,
1498 access_administrator: Option<
1499 &dyn crate::sync::store::DeviceProviderAccessAdministrator,
1500 >,
1501 ) -> Result<
1502 coven_protocol::store_commit::device_join_exchange::DeviceProviderAdmissionApproval,
1503 crate::sync::DeviceJoinError,
1504 > {
1505 self.store
1506 .authorize_device_provider_access(request, access_administrator)
1507 .await
1508 }
1509
1510 pub async fn publish_device_provider_challenge(
1511 &self,
1512 bootstrap: coven_protocol::store_commit::device_join_exchange::ProvisionalDeviceBootstrap,
1513 ) -> Result<
1514 coven_protocol::store_commit::device_join_exchange::ProviderReadyDeviceBootstrap,
1515 crate::sync::DeviceJoinError,
1516 > {
1517 self.store
1518 .publish_device_provider_challenge(bootstrap)
1519 .await
1520 }
1521
1522 pub async fn complete_device_provider_admission(
1523 &self,
1524 readiness: coven_protocol::store_commit::device_join_exchange::DeviceJoinReadiness,
1525 ) -> Result<
1526 coven_protocol::store_commit::device_join_exchange::DeviceProviderAdmissionCompletion,
1527 crate::sync::DeviceJoinError,
1528 > {
1529 self.store
1530 .complete_device_provider_admission(readiness)
1531 .await
1532 }
1533
1534 pub async fn abandon_device_join(
1535 &self,
1536 offer: coven_protocol::store_commit::device_join_exchange::DeviceJoinOffer,
1537 ) -> Result<
1538 coven_protocol::store_commit::device_join_exchange::DeviceJoinAbandonment,
1539 crate::sync::DeviceJoinError,
1540 > {
1541 self.store.abandon_device_join(offer).await
1542 }
1543
1544 pub async fn accept_device_registration_request(
1545 &self,
1546 request: coven_protocol::store_commit::device_join_exchange::DeviceRegistrationRequest,
1547 ) -> Result<
1548 coven_protocol::store_commit::device_join_exchange::ProvisionalDeviceBootstrap,
1549 crate::sync::DeviceJoinError,
1550 > {
1551 self.store.accept_device_registration_request(request).await
1552 }
1553
1554 pub async fn activate_same_principal_join_for_test(
1555 &self,
1556 request: coven_protocol::store_commit::device_join_exchange::DeviceRegistrationRequest,
1557 ) -> Result<
1558 coven_protocol::store_commit::device_join_exchange::SamePrincipalDeviceJoin,
1559 crate::sync::DeviceJoinError,
1560 > {
1561 let mut writer = self
1562 .store
1563 .authorize_writer()
1564 .await
1565 .map_err(crate::sync::DeviceJoinError::from)?;
1566 writer
1567 .join_operation()
1568 .activate_same_principal_join(request)
1569 .await
1570 }
1571
1572 pub async fn finalize_device_join(
1573 &self,
1574 completion: coven_protocol::store_commit::device_join_exchange::DeviceProviderAdmissionCompletion,
1575 ) -> Result<
1576 coven_protocol::store_commit::device_join_exchange::DeviceJoinActivation,
1577 crate::sync::DeviceJoinError,
1578 > {
1579 self.store.finalize_device_join(completion).await
1580 }
1581
1582 pub async fn device_exclusion_operations_for_test(
1583 &self,
1584 ) -> Result<
1585 Vec<crate::sync::store::StoreDeviceExclusionOperationInfo>,
1586 crate::sync::store::StoreDeviceExclusionError,
1587 > {
1588 self.store.device_exclusion_operations_for_test().await
1589 }
1590
1591 pub async fn stage_uploaded_device_exclusion_proposal_for_test(
1592 &self,
1593 ) -> Result<
1594 coven_protocol::store_commit::StoreDeviceExclusionProposalRef,
1595 crate::sync::store::StoreDeviceExclusionError,
1596 > {
1597 self.store
1598 .stage_uploaded_device_exclusion_proposal_for_test()
1599 .await
1600 }
1601
1602 pub async fn propose_device_exclusion(
1603 &self,
1604 target: &coven_protocol::store_commit::StoreDeviceRegistrationRef,
1605 ) -> Result<
1606 crate::sync::store::StoreDeviceExclusionResult,
1607 crate::sync::store::StoreDeviceExclusionError,
1608 > {
1609 self.store.propose_device_exclusion(target).await
1610 }
1611
1612 pub async fn cancel_device_exclusion(
1613 &self,
1614 proposal: &coven_protocol::store_commit::StoreDeviceExclusionProposalRef,
1615 ) -> Result<
1616 crate::sync::store::StoreDeviceExclusionResult,
1617 crate::sync::store::StoreDeviceExclusionError,
1618 > {
1619 self.store.cancel_device_exclusion(proposal).await
1620 }
1621
1622 pub async fn finalize_device_exclusion(
1623 &self,
1624 proposal: &coven_protocol::store_commit::StoreDeviceExclusionProposalRef,
1625 ) -> Result<
1626 crate::sync::store::StoreDeviceExclusionResult,
1627 crate::sync::store::StoreDeviceExclusionError,
1628 > {
1629 self.store.finalize_device_exclusion(proposal).await
1630 }
1631
1632 pub async fn pending_device_join_observation_for_test(
1633 &self,
1634 pending: &crate::sync::store::DeviceJoinJournalDatabase,
1635 offer: &coven_protocol::store_commit::device_join_exchange::DeviceJoinOffer,
1636 ) -> Result<crate::sync::store::PendingDeviceJoinObservation<'_>, TestError> {
1637 self.store
1638 .pending_device_join_observation_for_test(pending, offer)
1639 .await
1640 .map_err(TestError::from)
1641 }
1642
1643 pub async fn open_pending_device_join_for_test(
1644 &self,
1645 pending: &crate::sync::store::DeviceJoinJournalDatabase,
1646 identity: &UserKeypair,
1647 offer: coven_protocol::store_commit::device_join_exchange::DeviceJoinOffer,
1648 ) -> Result<crate::sync::store::PendingDeviceJoinAuthority<'_>, TestError> {
1649 self.store
1650 .open_pending_device_join_for_test(pending, identity, offer)
1651 .await
1652 .map_err(TestError::from)
1653 }
1654
1655 pub async fn prepare_snapshot_bootstrap_for_test(
1656 &self,
1657 membership_floor: &coven_protocol::membership::MembershipFloor,
1658 binary_schema_version: u32,
1659 target_path: &std::path::Path,
1660 restorer_identity: &UserKeypair,
1661 ) -> Result<
1662 crate::sync::store::PreparedSnapshotBootstrap<'_>,
1663 crate::sync::store::SnapshotError,
1664 > {
1665 self.store
1666 .prepare_snapshot_bootstrap_for_test(
1667 membership_floor,
1668 binary_schema_version,
1669 target_path,
1670 restorer_identity,
1671 )
1672 .await
1673 }
1674
1675 #[allow(clippy::too_many_arguments)]
1676 pub async fn admit_member(
1677 &self,
1678 member_pubkey: &str,
1679 member_email: Option<&str>,
1680 role: coven_protocol::membership::MemberRole,
1681 encryption: &coven_keys::encryption::EncryptionService,
1682 store_id: &str,
1683 store_name: &str,
1684 ) -> Result<crate::sync::store::MemberAdmission, crate::sync::store::MembershipOpsError>
1685 {
1686 self.store
1687 .admit_member(
1688 member_pubkey,
1689 member_email,
1690 role,
1691 encryption,
1692 store_id,
1693 store_name,
1694 )
1695 .await
1696 }
1697
1698 pub async fn drain_uploads(
1699 &self,
1700 clock: &dyn coven_foundation::clock::Clock,
1701 routing_encryption: Option<&coven_keys::encryption::EncryptionService>,
1702 observer: Option<&dyn coven_protocol::blob::BlobTransitionObserver>,
1703 ) -> Result<crate::blob::DrainOutcome, TestError> {
1704 self.store
1705 .authorize_writer()
1706 .await
1707 .map_err(TestError::from)?
1708 .drain_uploads(clock, routing_encryption, observer)
1709 .await
1710 .map_err(TestError::from)
1711 }
1712
1713 pub async fn publish_pending_store_database(&self) -> Result<bool, TestError> {
1714 let mut writer = self.store.authorize_writer().await?;
1715 let prepared = writer.prepare_pending_store_write().await?;
1716 let published = writer.drain_store_writes().await?;
1717 if published > 0 {
1718 let through_sequence = self
1719 .latest_local_store_position()
1720 .await?
1721 .expect("published Store write has no local position")
1722 .coord
1723 .sequence();
1724 crate::sync::test_owner_graph::TestOwnerGraph::new(
1725 self.db.clone(),
1726 self.store_dir.clone(),
1727 )
1728 .drain_published_blob_drop_intents(through_sequence)
1729 .await?;
1730 coven_database::LocalBlobCleanup::new(&self.db)
1731 .drain()
1732 .await?;
1733 }
1734 Ok(prepared || published > 0)
1735 }
1736
1737 pub async fn publish_fixture_position(&self, note_id: &str) -> u64 {
1738 self.db
1739 .insert_fixture_position_for_test(note_id)
1740 .await
1741 .expect("insert fixture Store position");
1742 assert!(self
1743 .publish_pending_store_database()
1744 .await
1745 .expect("publish fixture Store position"));
1746 self.latest_local_store_position()
1747 .await
1748 .expect("read fixture Store position")
1749 .expect("fixture Store write has an exact position")
1750 .coord
1751 .sequence()
1752 }
1753
1754 pub async fn publish_exact_remote_blob_binding(
1755 &self,
1756 root_id: &str,
1757 row_id: &str,
1758 bytes: &[u8],
1759 ) -> coven_protocol::blob::locator::StoredBlobRef {
1760 let local = self
1761 .db
1762 .row_blob_ref("note_photos", row_id)
1763 .await
1764 .expect("load exact Local row blob reference");
1765 let source = self
1766 .store_dir
1767 .local_blob_path(&local.blob().namespace, &local.blob().id)
1768 .expect("resolve host blob source");
1769 coven_foundation::local_file::AtomicStagedFile::write_for_test(&source, bytes)
1770 .await
1771 .expect("write host blob source");
1772 crate::sync::test_owner_graph::TestOwnerGraph::new(
1773 self.db.clone(),
1774 self.store_dir.clone(),
1775 )
1776 .make_remote("notes", root_id, "Notes Root", false)
1777 .await
1778 .expect("start exact make_remote");
1779 let clock = coven_foundation::clock::FixedClock(
1780 chrono::DateTime::parse_from_rfc3339("2024-06-01T01:00:00Z")
1781 .expect("valid exact blob publication time")
1782 .with_timezone(&chrono::Utc),
1783 );
1784 let outcome = self
1785 .drain_uploads(&clock, None, None)
1786 .await
1787 .expect("drain exact blob upload");
1788 assert_eq!(outcome.uploaded(), 1);
1789 assert!(self
1790 .publish_pending_store_database()
1791 .await
1792 .expect("publish exact remote blob binding"));
1793 self.db
1794 .row_blob_ref("note_photos", row_id)
1795 .await
1796 .expect("load exact Remote row blob reference")
1797 .stored()
1798 .cloned()
1799 .expect("Remote row owns an exact stored blob reference")
1800 }
1801
1802 pub async fn activated_store_device_registration_for_test(
1803 &self,
1804 reference: coven_protocol::store_commit::StoreDeviceRegistrationRef,
1805 ) -> Result<
1806 coven_protocol::store_commit::ReferencedStoreDeviceRegistration,
1807 coven_database::DbError,
1808 > {
1809 self.db.activated_store_device_registration(reference).await
1810 }
1811
1812 pub fn schema_version(&self) -> u32 {
1813 self.db.schema_version()
1814 }
1815
1816 pub async fn device_authority_for_test(
1817 &self,
1818 ) -> Result<TestDeviceSigningAuthority, TestError> {
1819 let registration = self
1820 .db
1821 .activated_store_device_registration_records()
1822 .await?
1823 .into_iter()
1824 .find(|registration| registration.value().device_id.to_string() == self.device_id)
1825 .ok_or_else(|| {
1826 TestError::invariant("test device registration is not active".to_string())
1827 })?;
1828 let device_signer = registration.value().device_signer(&self.identity)?;
1829 Ok(TestDeviceSigningAuthority {
1830 registration,
1831 device_signer,
1832 })
1833 }
1834
1835 #[cfg(test)]
1836 pub async fn prepare_uploaded_changeset_for_test(
1837 &self,
1838 sequence: u64,
1839 changeset: Vec<u8>,
1840 ) -> Result<coven_database::PreparedStoreWriteCommit, TestError> {
1841 let before = self.latest_local_store_position().await?;
1842 let expected = before
1843 .as_ref()
1844 .map_or(1, |reference| reference.coord.sequence() + 1);
1845 if sequence != expected {
1846 return Err(TestError::invariant(format!(
1847 "test producer expected sequence {expected}, got {sequence}"
1848 )));
1849 }
1850 self.db.enqueue_store_changeset_for_test(changeset).await?;
1851 self.prepare_uploaded_pending_write_for_test().await
1852 }
1853
1854 #[cfg(test)]
1855 pub async fn prepare_uploaded_pending_write_for_test(
1856 &self,
1857 ) -> Result<coven_database::PreparedStoreWriteCommit, TestError> {
1858 let mut writer = self.authorize_writer().await?;
1859 if !writer.prepare_pending_store_write().await? {
1860 return Err(TestError::invariant(
1861 "test changeset did not prepare a Store commit",
1862 ));
1863 }
1864 let pending = self
1865 .db
1866 .oldest_prepared_store_write()
1867 .await?
1868 .ok_or_else(|| TestError::invariant("prepared test changeset is absent"))?;
1869 let (uploaded, _resume) = self.db.arm_test_pause(
1870 coven_database::DatabaseTestPoint::StoreWriteCommitUploaded {
1871 write_id: pending.commit.value.write_id.clone(),
1872 },
1873 );
1874 {
1875 let publication = writer.drain_store_writes();
1876 tokio::pin!(publication);
1877 tokio::time::timeout(std::time::Duration::from_secs(10), async {
1878 tokio::select! {
1879 _ = uploaded.notified() => {},
1880 result = &mut publication => panic!("publication returned before commit upload: {result:?}"),
1881 }
1882 }).await.expect("upload the package and commit before interrupting publication");
1883 }
1884 Ok(pending)
1885 }
1886
1887 pub async fn publish_changeset_for_test(
1888 &self,
1889 sequence: u64,
1890 changeset: Vec<u8>,
1891 schema_version: u32,
1892 ) -> Result<coven_protocol::store_commit::StoreBatchCommitRef, TestError> {
1893 if schema_version != self.db.schema_version() {
1894 return Err(TestError::invariant(format!(
1895 "test changeset schema version {schema_version} differs from producer schema {}",
1896 self.db.schema_version()
1897 )));
1898 }
1899 let before = self.latest_local_store_position().await?;
1900 let expected = before
1901 .as_ref()
1902 .map_or(1, |reference| reference.coord.sequence().saturating_add(1));
1903 if sequence != expected {
1904 return Err(TestError::invariant(format!(
1905 "test producer expected sequence {expected}, got {sequence}"
1906 )));
1907 }
1908 self.db.enqueue_store_changeset_for_test(changeset).await?;
1909 let mut writer = self.authorize_writer().await?;
1910 let routing_encryption = coven_keys::encryption::EncryptionService::from_key([42; 32]);
1911 let published = writer
1912 .publish_pending_store_writes(Some(&routing_encryption))
1913 .await?;
1914 if published == 0 {
1915 return Err(TestError::invariant(
1916 "test changeset did not prepare a Store commit".to_string(),
1917 ));
1918 }
1919 writer.latest_local_store_position().await?.ok_or_else(|| {
1920 TestError::invariant("published test changeset has no Store position".to_string())
1921 })
1922 }
1923
1924 pub async fn publish_changeset_after_for_test(
1925 &self,
1926 changeset: Vec<u8>,
1927 previous_sequence: u64,
1928 ) -> Result<coven_protocol::store_commit::StoreBatchCommitRef, TestError> {
1929 let before = self.latest_local_store_position().await?;
1930 let actual_previous_sequence = before
1931 .as_ref()
1932 .map_or(0, |position| position.coord.sequence());
1933 if actual_previous_sequence != previous_sequence {
1934 return Err(TestError::invariant(format!(
1935 "Store position is {actual_previous_sequence}, expected {previous_sequence}"
1936 )));
1937 }
1938 self.db.enqueue_store_changeset_for_test(changeset).await?;
1939 let mut writer = self.store.authorize_writer().await?;
1940 if !writer.prepare_pending_store_write().await? {
1941 return Err(TestError::invariant(
1942 "test changeset did not prepare a Store commit".to_string(),
1943 ));
1944 }
1945 writer.drain_store_writes().await?;
1946 writer.latest_local_store_position().await?.ok_or_else(|| {
1947 TestError::invariant("published test changeset has no Store position".to_string())
1948 })
1949 }
1950
1951 pub async fn create_exact_opaque_blob(
1952 &self,
1953 namespace: &str,
1954 id: &str,
1955 bytes: &[u8],
1956 ) -> coven_protocol::blob::locator::StoredBlobRef {
1957 let registration = self
1958 .db
1959 .local_blob_write_authority()
1960 .await
1961 .expect("load exact blob write authority");
1962 let authority = coven_protocol::objects::BlobWriteAuthority::new(®istration);
1963 let protection = coven_keys::encryption::EncryptionService::from_key([42; 32]);
1964 let locator = coven_protocol::blob::locator::BlobLocator::opaque(
1965 namespace,
1966 id,
1967 authority.reference.clone(),
1968 coven_protocol::blob::locator::RemoteAudience::Store,
1969 coven_protocol::blob::BlobScope::Master,
1970 protection.seal_key_fingerprint(),
1971 bytes.len() as u64,
1972 coven_protocol::store_commit::ObjectHash::digest(bytes),
1973 )
1974 .expect("build exact blob locator");
1975 let temp = tempfile::tempdir().expect("create exact blob spool directory");
1976 let plaintext = temp.path().join("plaintext");
1977 let spool = temp.path().join("stored");
1978 coven_foundation::local_file::AtomicStagedFile::write_for_test(&plaintext, bytes)
1979 .await
1980 .expect("write exact blob plaintext");
1981 let slot = self
1982 .storage
1983 .allocate_blob_slot(&locator, &authority)
1984 .await
1985 .expect("allocate exact blob slot");
1986 let spool_stage = self
1987 .store_dir
1988 .stage_atomic_file(&spool)
1989 .await
1990 .expect("create exact blob spool stage");
1991 self.storage
1992 .seal_blob_to_spool(
1993 &locator,
1994 &authority,
1995 coven_protocol::objects::BlobSpoolProtection::Opaque(protection),
1996 &plaintext,
1997 spool_stage,
1998 coven_storage::cloud::no_preparation_progress(),
1999 )
2000 .await
2001 .expect("seal exact blob");
2002 let stored = self
2003 .storage
2004 .prepare_blob_object(&locator, &authority, slot, &spool)
2005 .await
2006 .expect("prepare exact blob object");
2007 let control =
2008 coven_storage::cloud::UploadControl::running(coven_storage::cloud::no_progress());
2009 self.storage
2010 .create_blob_object_from_file(&stored, &authority, &spool, &control)
2011 .await
2012 .expect("create exact blob object");
2013 stored
2014 }
2015
2016 pub async fn create_exact_browsable_blob(
2017 &self,
2018 namespace: &str,
2019 id: &str,
2020 cloud_path: &str,
2021 bytes: &[u8],
2022 ) -> coven_protocol::blob::locator::StoredBlobRef {
2023 let registration = self
2024 .db
2025 .local_blob_write_authority()
2026 .await
2027 .expect("load browsable blob write authority");
2028 let authority = coven_protocol::objects::BlobWriteAuthority::new(®istration);
2029 let locator = coven_protocol::blob::locator::BlobLocator::browsable(
2030 namespace,
2031 id,
2032 authority.reference.clone(),
2033 cloud_path,
2034 bytes.len() as u64,
2035 coven_protocol::store_commit::ObjectHash::digest(bytes),
2036 )
2037 .expect("build browsable blob locator");
2038 let temp = tempfile::tempdir().expect("create browsable blob spool directory");
2039 let plaintext = temp.path().join("plaintext");
2040 let spool = temp.path().join("stored");
2041 coven_foundation::local_file::AtomicStagedFile::write_for_test(&plaintext, bytes)
2042 .await
2043 .expect("write browsable blob plaintext");
2044 let slot = self
2045 .storage
2046 .allocate_blob_slot(&locator, &authority)
2047 .await
2048 .expect("allocate browsable blob slot");
2049 let spool_stage = self
2050 .store_dir
2051 .stage_atomic_file(&spool)
2052 .await
2053 .expect("create browsable blob spool stage");
2054 self.storage
2055 .seal_blob_to_spool(
2056 &locator,
2057 &authority,
2058 coven_protocol::objects::BlobSpoolProtection::Browsable,
2059 &plaintext,
2060 spool_stage,
2061 coven_storage::cloud::no_preparation_progress(),
2062 )
2063 .await
2064 .expect("stage browsable blob");
2065 let stored = self
2066 .storage
2067 .prepare_blob_object(&locator, &authority, slot, &spool)
2068 .await
2069 .expect("prepare browsable blob object");
2070 let control =
2071 coven_storage::cloud::UploadControl::running(coven_storage::cloud::no_progress());
2072 self.storage
2073 .create_blob_object_from_file(&stored, &authority, &spool, &control)
2074 .await
2075 .expect("create browsable blob object");
2076 stored
2077 }
2078
2079 pub async fn run_cycle(
2080 &self,
2081 observer: Option<&dyn coven_protocol::blob::BlobTransitionObserver>,
2082 ) -> Result<crate::sync::cycle::SyncCycleResult, crate::sync::cycle::SyncCycleFailure>
2083 {
2084 self.run_cycle_with(&coven_foundation::clock::SystemClock, None, observer)
2085 .await
2086 }
2087
2088 pub async fn run_cycle_with(
2089 &self,
2090 clock: &dyn coven_foundation::clock::Clock,
2091 master_keys: Option<std::sync::Arc<dyn coven_keys::keys::MasterKeyCustody>>,
2092 observer: Option<&dyn coven_protocol::blob::BlobTransitionObserver>,
2093 ) -> Result<crate::sync::cycle::SyncCycleResult, crate::sync::cycle::SyncCycleFailure>
2094 {
2095 self.run_cycle_with_storage(
2096 self.store.clone(),
2097 self.storage.clone(),
2098 clock,
2099 master_keys,
2100 observer,
2101 )
2102 .await
2103 }
2104
2105 pub async fn run_cycle_with_interceptor<I>(
2106 &self,
2107 clock: &dyn coven_foundation::clock::Clock,
2108 master_keys: Option<std::sync::Arc<dyn coven_keys::keys::MasterKeyCustody>>,
2109 observer: Option<&dyn coven_protocol::blob::BlobTransitionObserver>,
2110 interceptor: I,
2111 ) -> Result<crate::sync::cycle::SyncCycleResult, crate::sync::cycle::SyncCycleFailure>
2112 where
2113 I: super::StorageInterceptor + 'static,
2114 {
2115 let storage = std::sync::Arc::new(super::InterceptedStorage::new(
2116 self.storage.clone(),
2117 interceptor,
2118 ));
2119 let store_storage: std::sync::Arc<dyn coven_storage::CloudSyncObjectStorage> =
2120 storage.clone();
2121 let store = std::sync::Arc::new(self.store.with_test_storage(store_storage));
2122 self.run_cycle_with_storage(store, storage, clock, master_keys, observer)
2123 .await
2124 }
2125
2126 async fn run_cycle_with_storage<S>(
2127 &self,
2128 store: std::sync::Arc<crate::sync::store::Store>,
2129 storage: std::sync::Arc<S>,
2130 clock: &dyn coven_foundation::clock::Clock,
2131 master_keys: Option<std::sync::Arc<dyn coven_keys::keys::MasterKeyCustody>>,
2132 observer: Option<&dyn coven_protocol::blob::BlobTransitionObserver>,
2133 ) -> Result<crate::sync::cycle::SyncCycleResult, crate::sync::cycle::SyncCycleFailure>
2134 where
2135 S: crate::sync::cycle::CloudSyncCycleConnection + 'static,
2136 {
2137 let components = crate::sync::cycle::SyncComponents::from_retained_test_device(
2138 store,
2139 self.db.clone(),
2140 self.store_dir.clone(),
2141 storage,
2142 self.storage.store_id().to_string(),
2143 self.device_id.clone(),
2144 master_keys.unwrap_or_else(|| std::sync::Arc::new(super::TestCustody::default())),
2145 self.settled.clone(),
2146 );
2147 components
2148 .run_cycle(
2149 clock,
2150 observer,
2151 coven_foundation::config::Config::DEFAULT_SNAPSHOT_COMMIT_THRESHOLD,
2152 )
2153 .await
2154 }
2155
2156 pub fn current_keyring_for_test(&self) -> Option<coven_storage::CloudKeyringFacts> {
2157 self.storage.keyring_facts_for_test()
2158 }
2159
2160 pub fn mark_rotation_committed_for_test(
2161 &self,
2162 generation: u64,
2163 ) -> Result<(), coven_storage::RotationStateError> {
2164 self.storage.mark_rotation_committed_for_test(generation)
2165 }
2166
2167 pub fn pending_rotation_generation_for_test(&self) -> Option<u64> {
2168 self.storage.pending_rotation_generation_for_test()
2169 }
2170
2171 pub fn clear_rotation_gate_for_test(&self) {
2172 self.storage.clear_rotation_gate_for_test();
2173 }
2174
2175 pub async fn create_circle(
2176 &self,
2177 metadata_stamp: &str,
2178 name: &str,
2179 ) -> Result<coven_protocol::CircleId, crate::sync::store::CircleOperationError> {
2180 self.store
2181 .circles()
2182 .create_circle(metadata_stamp, name)
2183 .await
2184 }
2185
2186 async fn circle_writer(
2189 &self,
2190 ) -> Result<
2191 crate::sync::store::AuthorizedWriterOperation<'_>,
2192 crate::sync::store::CircleOperationError,
2193 > {
2194 self.store
2195 .authorize_writer()
2196 .await
2197 .map_err(crate::sync::store::CircleOperationError::from)
2198 }
2199
2200 #[cfg(test)]
2201 pub(crate) async fn prepare_circle_operation(
2202 &self,
2203 metadata_stamp: &str,
2204 name: &str,
2205 ) -> Result<
2206 crate::sync::store::circles::PreparedCircleJournal,
2207 crate::sync::store::CircleOperationError,
2208 > {
2209 self.circle_writer()
2210 .await?
2211 .circles()
2212 .prepare_create_for_test(metadata_stamp, name)
2213 .await
2214 }
2215
2216 pub async fn publish_circle_epoch_close_response(
2217 &self,
2218 ) -> Result<(), crate::sync::store::CircleOperationError> {
2219 self.circle_writer()
2220 .await?
2221 .circles()
2222 .publish_circle_epoch_close_responses()
2223 .await
2224 }
2225
2226 pub async fn publish_circle_operation(
2227 &self,
2228 operation_id: &coven_protocol::circle::CircleOperationId,
2229 ) -> Result<(), crate::sync::store::CircleOperationError> {
2230 let routing_key = coven_protocol::circle::derive_row_routing_key(
2231 &coven_keys::encryption::EncryptionService::from_key([42; 32]),
2232 self.store.store_root().store_root_hash,
2233 )
2234 .expect("derive Circle test routing key");
2235 self.circle_writer()
2236 .await?
2237 .circles()
2238 .publish_prepared_operation_for_test(operation_id, Some(&routing_key))
2239 .await
2240 }
2241
2242 pub async fn resume_circle_operations(
2243 &self,
2244 ) -> Result<(), crate::sync::store::CircleOperationError> {
2245 let routing_key = coven_protocol::circle::derive_row_routing_key(
2246 &coven_keys::encryption::EncryptionService::from_key([42; 32]),
2247 self.store.store_root().store_root_hash,
2248 )
2249 .expect("derive Circle test routing key");
2250 self.circle_writer()
2251 .await?
2252 .circles()
2253 .resume_circle_operations(Some(&routing_key))
2254 .await
2255 }
2256
2257 pub async fn retry_circle_operation(
2258 &self,
2259 operation_id: &coven_protocol::circle::CircleOperationId,
2260 ) -> Result<(), crate::sync::store::CircleOperationError> {
2261 self.store
2262 .circles()
2263 .retry_circle_operation(
2264 operation_id,
2265 Some(&coven_keys::encryption::EncryptionService::from_key(
2266 [42; 32],
2267 )),
2268 )
2269 .await
2270 }
2271
2272 pub async fn prepare_pending_store_write(
2273 &self,
2274 ) -> Result<bool, crate::sync::store::StoreError> {
2275 self.store
2276 .authorize_writer()
2277 .await
2278 .map_err(crate::sync::store::StoreError::from)?
2279 .prepare_pending_store_write()
2280 .await
2281 }
2282
2283 #[cfg(test)]
2284 pub async fn prepare_blocked_transfer_candidate(
2285 &self,
2286 label: &str,
2287 ) -> coven_protocol::write::WriteId {
2288 let statement = format!(
2289 "INSERT INTO notes (id, title, body, shared, _updated_at, created_at) \
2290 VALUES ('{label}', 'pending', NULL, 1, \
2291 '0000000002000-0000-{label}', '2026-07-18')"
2292 );
2293 self.db
2294 .run_host_store_write_for_test(None, None, move |transaction| {
2295 transaction
2296 .execute_batch(&statement)
2297 .map_err(coven_database::DbError::from)
2298 })
2299 .await
2300 .expect("capture transfer candidate host write");
2301 assert!(self
2302 .prepare_pending_store_write()
2303 .await
2304 .expect("prepare transfer candidate"));
2305 let candidate = self
2306 .db
2307 .oldest_prepared_store_write()
2308 .await
2309 .expect("load transfer candidate")
2310 .expect("transfer candidate exists");
2311 let write_id = candidate.commit.value.write_id.clone();
2312 self.db
2313 .set_write_status(
2314 &write_id,
2315 coven_protocol::write::WriteStatus::Blocked(
2316 coven_protocol::write::WriteBlock::InvalidProtocolState {
2317 reason: "exercise restored author-exclusion evidence".to_string(),
2318 },
2319 ),
2320 )
2321 .await
2322 .expect("block transfer candidate");
2323 write_id
2324 }
2325
2326 #[cfg(test)]
2327 pub async fn prepare_store_operation_plan_for_test(
2328 &self,
2329 ) -> Result<crate::sync::store::StoreOperationCommitPlan, crate::sync::store::StoreError>
2330 {
2331 self.store
2332 .authorize_writer()
2333 .await
2334 .map_err(crate::sync::store::StoreError::from)?
2335 .prepare_plan()
2336 .await
2337 }
2338
2339 pub async fn drain_store_writes(&self) -> Result<u64, crate::sync::store::StoreError> {
2340 self.store
2341 .authorize_writer()
2342 .await
2343 .map_err(crate::sync::store::StoreError::from)?
2344 .drain_store_writes()
2345 .await
2346 }
2347
2348 pub async fn reclaim_packages(
2349 &self,
2350 ) -> Result<crate::sync::store::StoreReclaimResult, crate::sync::store::StoreReclaimError>
2351 {
2352 self.store
2353 .authorize_writer()
2354 .await
2355 .map_err(crate::sync::store::StoreReclaimError::from)?
2356 .reclaim_packages(&crate::sync::store::SettledCycle::default())
2357 .await
2358 }
2359
2360 pub async fn prepare_peer_exclusion(
2361 &self,
2362 target: &coven_protocol::store_commit::StoreDeviceRegistrationRef,
2363 ) -> coven_protocol::store_commit::StoreDeviceExclusionProposalRef {
2364 let proposal = match self
2365 .propose_device_exclusion(target)
2366 .await
2367 .expect("propose peer exclusion")
2368 {
2369 crate::sync::store::StoreDeviceExclusionResult::ProposalActivated {
2370 proposal,
2371 ..
2372 } => proposal,
2373 result => panic!("unexpected exclusion proposal result: {result:?}"),
2374 };
2375 let frontier = coven_protocol::store_commit::CommitFrontier::from_refs(
2376 self.db
2377 .materialized_frontier()
2378 .await
2379 .expect("read proposal frontier"),
2380 )
2381 .expect("exact proposal frontier");
2382 let (_, state) = self
2383 .db
2384 .store_device_state_for_history_cut(&coven_protocol::store_commit::StoreHistoryCut(
2385 frontier.0,
2386 ))
2387 .await
2388 .expect("read accepted proposal state");
2389 let record = &state.devices[&target.device_id];
2390 assert!(matches!(
2391 record.status,
2392 coven_protocol::store_commit::StoreDeviceStatus::Active
2393 ));
2394 assert!(matches!(record.proposals.get(&proposal.proposal_id),
2395 Some(coven_protocol::store_commit::StoreDeviceProposalState::Pending { proposal: actual }) if actual == &proposal));
2396 proposal
2397 }
2398
2399 pub async fn activate_peer_exclusion(
2400 &self,
2401 proposal: &coven_protocol::store_commit::StoreDeviceExclusionProposalRef,
2402 ) -> coven_protocol::store_commit::StoreDeviceExclusionRef {
2403 let result = self
2404 .finalize_device_exclusion(proposal)
2405 .await
2406 .expect("finalize peer exclusion");
2407 let crate::sync::store::StoreDeviceExclusionResult::OutcomeActivated {
2408 outcome:
2409 coven_protocol::store_commit::StoreDeviceExclusionOutcomeRef::Excluded(exclusion),
2410 ..
2411 } = result
2412 else {
2413 panic!("unexpected exclusion result: {result:?}")
2414 };
2415 let frontier = coven_protocol::store_commit::CommitFrontier::from_refs(
2416 self.db
2417 .materialized_frontier()
2418 .await
2419 .expect("read exclusion frontier"),
2420 )
2421 .expect("exact exclusion frontier");
2422 let (_, state) = self
2423 .db
2424 .store_device_state_for_history_cut(&coven_protocol::store_commit::StoreHistoryCut(
2425 frontier.0,
2426 ))
2427 .await
2428 .expect("read accepted exclusion state");
2429 assert!(matches!(&state.devices[&proposal.target.device_id].status,
2430 coven_protocol::store_commit::StoreDeviceStatus::Inactive { terminals } if terminals.contains(&exclusion)));
2431 exclusion
2432 }
2433
2434 pub async fn finalize_peer_exclusion(
2435 &self,
2436 target: &coven_protocol::store_commit::StoreDeviceRegistrationRef,
2437 ) -> coven_protocol::store_commit::StoreDeviceExclusionRef {
2438 let proposal = self.prepare_peer_exclusion(target).await;
2439 self.activate_peer_exclusion(&proposal).await
2440 }
2441
2442 pub async fn prepare_circle_object(
2443 &self,
2444 context: &coven_protocol::objects::ProtocolObjectContext,
2445 semantic_prefix: &str,
2446 extension: &str,
2447 bytes: Vec<u8>,
2448 ) -> Result<
2449 coven_protocol::objects::PreparedExactObject,
2450 crate::sync::store::CircleOperationError,
2451 > {
2452 self.circle_writer()
2453 .await?
2454 .circles()
2455 .prepare_circle_object_for_test(context, semantic_prefix, extension, bytes)
2456 .await
2457 }
2458
2459 pub async fn prepare_circle_object_at(
2460 &self,
2461 context: &coven_protocol::objects::ProtocolObjectContext,
2462 slot: coven_protocol::objects::ObjectSlot,
2463 semantic_prefix: &str,
2464 bytes: Vec<u8>,
2465 ) -> Result<
2466 coven_protocol::objects::PreparedExactObject,
2467 crate::sync::store::CircleOperationError,
2468 > {
2469 self.circle_writer()
2470 .await?
2471 .circles()
2472 .prepare_circle_object_at_for_test(context, slot, semantic_prefix, bytes)
2473 }
2474
2475 pub async fn prepare_circle_activation_objects(
2476 &self,
2477 draft: coven_protocol::circle::CircleTransitionDraft,
2478 history: &crate::sync::store::CircleTransitionHistory,
2479 candidate_family: coven_protocol::store_commit::CandidateFamilyId,
2480 ) -> Result<
2481 (
2482 coven_protocol::circle::PreparedCircleTransition,
2483 coven_protocol::store_commit::CircleActivationObjects,
2484 std::collections::BTreeMap<String, coven_protocol::objects::PreparedExactObject>,
2485 Option<coven_protocol::objects::ExactObjectRef>,
2486 Vec<coven_protocol::store_commit::StreamActivation>,
2487 ),
2488 crate::sync::store::CircleOperationError,
2489 > {
2490 self.circle_writer()
2491 .await?
2492 .circles()
2493 .prepare_circle_activation_objects_for_test(draft, history, candidate_family)
2494 .await
2495 }
2496
2497 pub async fn sign_circle_commit(
2498 &self,
2499 old_commit: &coven_protocol::store_commit::StoreBatchCommit,
2500 coord: coven_protocol::store_commit::StoreCommitCoord,
2501 reference: coven_protocol::store_commit::CircleControlRef,
2502 stream_activations: Vec<coven_protocol::store_commit::StreamActivation>,
2503 ) -> Result<
2504 coven_protocol::store_commit::StoreBatchCommit,
2505 crate::sync::store::CircleOperationError,
2506 > {
2507 self.circle_writer()
2508 .await?
2509 .circles()
2510 .sign_circle_commit_for_test(old_commit, coord, reference, stream_activations)
2511 }
2512
2513 pub async fn rename_circle(
2514 &self,
2515 metadata_stamp: &str,
2516 circle_id: coven_protocol::CircleId,
2517 name: &str,
2518 ) -> Result<(), crate::sync::store::CircleOperationError> {
2519 self.store
2520 .circles()
2521 .rename_circle(metadata_stamp, circle_id, name)
2522 .await
2523 }
2524
2525 pub async fn delete_circle(
2526 &self,
2527 circle_id: coven_protocol::CircleId,
2528 ) -> Result<(), crate::sync::store::CircleOperationError> {
2529 self.store.circles().delete_circle(circle_id).await
2530 }
2531
2532 pub async fn load_circle_activations(
2533 &self,
2534 commit_ref: &coven_protocol::store_commit::StoreBatchCommitRef,
2535 commit: &coven_protocol::store_commit::StoreBatchCommit,
2536 author: &coven_protocol::store_commit::StoreDeviceRegistration,
2537 ) -> Result<
2538 coven_protocol::circle_activation::VerifiedCircleActivations,
2539 crate::sync::store::CircleOperationError,
2540 > {
2541 let routing_key = coven_protocol::circle::derive_row_routing_key(
2542 &coven_keys::encryption::EncryptionService::from_key([42; 32]),
2543 commit.store_root_hash,
2544 )
2545 .expect("derive Circle test routing key");
2546 self.store
2547 .load_circle_activations_for_test(commit_ref, commit, author, Some(&routing_key))
2548 .await
2549 }
2550
2551 pub async fn circle_blob_opening_error(
2552 &self,
2553 authority: &coven_protocol::blob::RowBlobAuthority,
2554 stored: &coven_protocol::blob::locator::StoredBlobRef,
2555 ) -> crate::sync::store::StoreError {
2556 match self
2557 .store
2558 .blob_key_fingerprint_for_test(authority, stored)
2559 .await
2560 {
2561 Ok(_) => panic!("invalid Circle blob authority must fail"),
2562 Err(error) => error,
2563 }
2564 }
2565
2566 pub async fn load_circle_snapshot_refs(
2567 &self,
2568 circle_id: coven_protocol::CircleId,
2569 access: &coven_protocol::circle_activation::CircleEpochAccess,
2570 ) -> Result<
2571 Vec<(
2572 coven_protocol::store_commit::CircleSnapshotRef,
2573 coven_protocol::store_commit::CircleSnapshotMeta,
2574 )>,
2575 TestError,
2576 > {
2577 self.store
2578 .authorize_writer()
2579 .await?
2580 .circles()
2581 .snapshots()
2582 .load_circle_snapshot_refs_for_test(circle_id, access)
2583 .await
2584 .map_err(TestError::from)
2585 }
2586
2587 pub async fn membership(
2588 &self,
2589 ) -> Result<coven_protocol::membership::MembershipChain, TestError> {
2590 self.store
2591 .membership_for_test()
2592 .await
2593 .map_err(TestError::from)
2594 }
2595
2596 pub fn protocol_root(&self) -> &coven_protocol::store_commit::StoreProtocolRoot {
2597 self.store.protocol_root_for_test()
2598 }
2599
2600 #[cfg(test)]
2601 pub async fn prepare_wrapped_key(
2602 &self,
2603 recipient: &str,
2604 value: coven_protocol::wrapped_store_key::WrappedStoreKey,
2605 ) -> Result<coven_protocol::wrapped_store_key::PreparedWrappedStoreKey, TestError> {
2606 self.store
2607 .prepare_wrapped_key_for_test(recipient, value)
2608 .await
2609 }
2610
2611 #[cfg(test)]
2612 pub async fn membership_keyring_facts(&self) -> Result<([u8; 32], usize), TestError> {
2613 self.store.membership_keyring_facts_for_test().await
2614 }
2615
2616 pub async fn publish_snapshot(
2617 &self,
2618 db_image: Vec<u8>,
2619 coverage: coven_protocol::store_commit::CommitFrontier,
2620 ) -> Result<coven_protocol::store_commit::SnapshotMeta, TestError> {
2621 self.publish_snapshot_at(db_image, coverage, "2026-07-16T00:00:00Z")
2622 .await
2623 .map_err(TestError::from)
2624 }
2625
2626 pub async fn publish_snapshot_at(
2627 &self,
2628 db_image: Vec<u8>,
2629 coverage: coven_protocol::store_commit::CommitFrontier,
2630 created_at: &str,
2631 ) -> Result<coven_protocol::store_commit::SnapshotMeta, crate::sync::store::SnapshotError>
2632 {
2633 self.store
2634 .publish_snapshot_for_test(
2635 coven_database::CreatedSnapshot::new(
2636 staged_snapshot_image(&db_image),
2637 Vec::new(),
2638 ),
2639 coverage,
2640 created_at.to_string(),
2641 )
2642 .await
2643 }
2644
2645 pub async fn resume_snapshot_publication(
2646 &self,
2647 ) -> Result<
2648 Option<coven_protocol::store_commit::SnapshotMeta>,
2649 crate::sync::store::SnapshotError,
2650 > {
2651 self.store
2652 .authorize_writer()
2653 .await
2654 .map_err(crate::sync::store::SnapshotError::from)?
2655 .resume_snapshot_publication()
2656 .await
2657 }
2658
2659 pub fn publish_acknowledgement(
2662 &self,
2663 frontier: coven_protocol::store_commit::CommitFrontier,
2664 ) -> std::pin::Pin<
2665 Box<
2666 dyn std::future::Future<
2667 Output = Result<Option<coven_database::AdvancedReplayBaseline>, TestError>,
2668 > + Send
2669 + '_,
2670 >,
2671 > {
2672 Box::pin(async move {
2673 let staged = self
2674 .store
2675 .stage_acknowledgement_for_test(frontier, "2026-07-16T00:00:01Z".to_string())
2676 .await?;
2677 let expected = u64::from(staged.acknowledgement.is_some());
2678 let published = self.store.drain_acknowledgements_for_test().await?;
2679 if published != expected {
2680 return Err(TestError::invariant(format!(
2681 "acknowledgement fixture published {published} statements for {expected} changed assertions"
2682 )));
2683 }
2684 match self.stand_on_accepted_snapshot().await? {
2685 crate::sync::store::ReplayBaselineAdvance::Advanced(advanced) => {
2686 Ok(Some(advanced))
2687 }
2688 crate::sync::store::ReplayBaselineAdvance::Declined(_) => Ok(None),
2689 }
2690 })
2691 }
2692
2693 pub async fn stage_acknowledgement(
2694 &self,
2695 frontier: coven_protocol::store_commit::CommitFrontier,
2696 sync_time: String,
2697 ) -> Result<Option<coven_protocol::store_commit::StoreAck>, TestError> {
2698 self.store
2699 .stage_acknowledgement_for_test(frontier, sync_time)
2700 .await
2701 .map(|staged| staged.acknowledgement)
2702 .map_err(TestError::from)
2703 }
2704
2705 pub async fn publish_acknowledgement_without_advancing(
2708 &self,
2709 frontier: coven_protocol::store_commit::CommitFrontier,
2710 ) -> Result<(), TestError> {
2711 self.store
2712 .stage_acknowledgement_for_test(frontier, "2026-07-16T00:00:03Z".to_string())
2713 .await?
2714 .acknowledgement
2715 .ok_or_else(|| {
2716 TestError::invariant(
2717 "the fixture acknowledgement asserted nothing new".to_string(),
2718 )
2719 })?;
2720 let published = self.store.drain_acknowledgements_for_test().await?;
2721 if published != 1 {
2722 return Err(TestError::invariant(format!(
2723 "fixture published {published} acknowledgements instead of one"
2724 )));
2725 }
2726 Ok(())
2727 }
2728
2729 pub async fn stand_on_accepted_snapshot(
2731 &self,
2732 ) -> Result<crate::sync::store::ReplayBaselineAdvance, TestError> {
2733 self.store
2734 .stand_on_accepted_snapshot_for_test()
2735 .await
2736 .map_err(TestError::from)
2737 }
2738
2739 pub async fn advance_baseline_by_acknowledging(
2742 &self,
2743 frontier: coven_protocol::store_commit::CommitFrontier,
2744 ) -> Result<Option<coven_database::AdvancedReplayBaseline>, TestError> {
2745 self.store
2746 .stage_acknowledgement_for_test(frontier, "2026-07-16T00:00:02Z".to_string())
2747 .await?;
2748 self.store.drain_acknowledgements_for_test().await?;
2749 match self.stand_on_accepted_snapshot().await? {
2750 crate::sync::store::ReplayBaselineAdvance::Advanced(advanced) => Ok(Some(advanced)),
2751 crate::sync::store::ReplayBaselineAdvance::Declined(_) => Ok(None),
2752 }
2753 }
2754
2755 pub async fn materialized_frontier(
2756 &self,
2757 ) -> Result<
2758 std::collections::BTreeMap<String, coven_protocol::store_commit::StoreBatchCommitRef>,
2759 TestError,
2760 > {
2761 self.db
2762 .materialized_frontier()
2763 .await
2764 .map_err(TestError::from)
2765 }
2766
2767 pub async fn drain_acknowledgements(&self) -> Result<u64, TestError> {
2768 self.store
2769 .drain_acknowledgements_for_test()
2770 .await
2771 .map_err(TestError::from)
2772 }
2773
2774 #[cfg(test)]
2775 pub async fn stage_acknowledgement_exact(
2776 &self,
2777 frontier: coven_protocol::store_commit::CommitFrontier,
2778 sync_time: String,
2779 ) -> Result<Option<coven_protocol::store_commit::StoreAck>, crate::sync::store::StoreAckError>
2780 {
2781 self.store
2782 .stage_acknowledgement_for_test(frontier, sync_time)
2783 .await
2784 .map(|staged| staged.acknowledgement)
2785 }
2786
2787 #[cfg(test)]
2788 pub async fn acknowledgement_frontier(
2789 &self,
2790 ) -> Result<coven_protocol::store_commit::CommitFrontier, crate::sync::store::StoreAckError>
2791 {
2792 coven_protocol::store_commit::CommitFrontier::from_refs(
2793 self.db.materialized_frontier().await?,
2794 )
2795 .map_err(crate::sync::store::StoreAckError::Protocol)
2796 }
2797
2798 #[cfg(test)]
2804 pub async fn stage_current_acknowledgement(
2805 &self,
2806 sync_time: &str,
2807 ) -> Result<coven_protocol::store_commit::StoreAck, crate::sync::store::StoreAckError>
2808 {
2809 Ok(self
2810 .stage_current_acknowledgement_if_new(sync_time)
2811 .await?
2812 .expect("the standing acknowledgement no longer holds"))
2813 }
2814
2815 #[cfg(test)]
2816 pub async fn stage_current_acknowledgement_if_new(
2817 &self,
2818 sync_time: &str,
2819 ) -> Result<Option<coven_protocol::store_commit::StoreAck>, crate::sync::store::StoreAckError>
2820 {
2821 let frontier = self.acknowledgement_frontier().await?;
2822 self.stage_acknowledgement_exact(frontier, sync_time.to_string())
2823 .await
2824 }
2825
2826 #[cfg(any(test, feature = "test-utils"))]
2827 pub fn typed_device_id(&self) -> coven_protocol::store_commit::StoreDeviceId {
2828 self.device_id
2829 .parse()
2830 .expect("TestDevice retains a valid Store device id")
2831 }
2832
2833 #[cfg(test)]
2834 pub async fn prepare_acknowledgement_candidate_for_test(
2835 &self,
2836 outbound: &coven_database::OutboundStoreAck,
2837 ) -> coven_protocol::prepared_commit::PreparedStoreOperationCommit {
2838 let mut writer = self
2839 .authorize_writer()
2840 .await
2841 .expect("authorize acknowledgement writer");
2842 let plan = writer
2843 .prepare_plan()
2844 .await
2845 .expect("prepare acknowledgement activation");
2846 plan.validate_acknowledgement(&outbound.ack.value)
2847 .expect("acknowledgement matches activation predecessor");
2848 let candidate = writer
2849 .prepare_candidate(
2850 &plan,
2851 crate::sync::store::StoreOperationBatch::Acknowledgement {
2852 reference: outbound.reference.clone(),
2853 value: outbound.ack.value.clone(),
2854 circle_acknowledgements: outbound.circle_acknowledgements.clone(),
2855 },
2856 )
2857 .await
2858 .expect("prepare acknowledgement candidate");
2859 self.prepare_acknowledgement_activation_for_test(
2860 outbound.reference.clone(),
2861 outbound.ack.clone(),
2862 candidate.clone(),
2863 )
2864 .await
2865 .expect("persist acknowledgement candidate");
2866 candidate
2867 }
2868
2869 #[cfg(test)]
2870 pub async fn drain_acknowledgements_exact(
2871 &self,
2872 ) -> Result<u64, crate::sync::store::StoreAckError> {
2873 self.store.drain_acknowledgements_for_test().await
2874 }
2875
2876 #[cfg(test)]
2877 pub async fn stage_circle_acknowledgements(
2878 &self,
2879 frontier: &coven_protocol::store_commit::CommitFrontier,
2880 sync_time: &str,
2881 ) -> Result<(), crate::sync::store::StoreAckError> {
2882 self.store
2883 .stage_circle_acknowledgements_for_test(frontier, sync_time)
2884 .await
2885 }
2886
2887 pub async fn load_commit_ancestry_until(
2888 &self,
2889 start: coven_protocol::store_commit::StoreBatchCommitRef,
2890 coverage: &coven_protocol::store_commit::CommitFrontier,
2891 ) -> Result<
2892 Vec<(
2893 coven_protocol::store_commit::StoreBatchCommitRef,
2894 coven_protocol::store_commit::VerifiedStoreBatchCommit,
2895 )>,
2896 TestError,
2897 > {
2898 self.store
2899 .load_commit_ancestry_until_for_test(start, coverage)
2900 .await
2901 .map_err(TestError::from)
2902 }
2903
2904 pub async fn export_activated_device_continuation(
2905 &self,
2906 ) -> Result<coven_protocol::recovery::ActivatedContinuation, TestError> {
2907 self.store
2908 .export_activated_device_continuation_for_test()
2909 .await
2910 .map_err(TestError::from)
2911 }
2912
2913 pub async fn latest_store_position(
2914 &self,
2915 ) -> Result<Option<coven_protocol::store_commit::StoreBatchCommitRef>, TestError> {
2916 self.store
2917 .latest_local_store_position()
2918 .await
2919 .map_err(TestError::from)
2920 }
2921
2922 pub async fn pull_store(
2923 &self,
2924 ) -> Result<
2925 (
2926 std::collections::BTreeMap<String, u64>,
2927 crate::sync::store::StorePullResult,
2928 ),
2929 TestPullError,
2930 > {
2931 let routing_encryption = coven_keys::encryption::EncryptionService::from_key([42; 32]);
2932 self.pull_store_with_encryption(&routing_encryption).await
2933 }
2934
2935 pub async fn pull_store_and_fill_eager(
2939 &self,
2940 ) -> Result<
2941 (
2942 std::collections::BTreeMap<String, u64>,
2943 crate::sync::store::StorePullResult,
2944 ),
2945 TestPullError,
2946 > {
2947 let pulled = self.pull_store().await?;
2948 crate::sync::test_owner_graph::TestOwnerGraph::new(
2949 self.db.clone(),
2950 self.store_dir.clone(),
2951 )
2952 .fill_eager_cache(self.storage.clone())
2953 .await
2954 .expect("fill the eager cache behind the pull");
2955 Ok(pulled)
2956 }
2957
2958 pub async fn pull_store_with_encryption(
2959 &self,
2960 routing_encryption: &coven_keys::encryption::EncryptionService,
2961 ) -> Result<
2962 (
2963 std::collections::BTreeMap<String, u64>,
2964 crate::sync::store::StorePullResult,
2965 ),
2966 TestPullError,
2967 > {
2968 let mut authorization = self.store.authorize_writer().await?;
2969 let result = authorization.pull(Some(routing_encryption)).await?;
2970 let sequences = result
2971 .frontier
2972 .iter()
2973 .map(|(stream, reference)| (stream.clone(), reference.coord.sequence()))
2974 .collect();
2975 Ok((sequences, result))
2976 }
2977 }
2978}
2979
2980pub use test_device::{TestDevice, TestDeviceSigningAuthority};
2981
2982#[cfg(test)]
2990pub struct JoinedTestStore {
2991 database: coven_database::StoreDatabase,
2992}
2993
2994#[cfg(test)]
2995impl JoinedTestStore {
2996 pub async fn latest_local_store_device_registration(
2997 &self,
2998 ) -> Result<Option<coven_database::DurableDeviceRegistration>, coven_database::DbError> {
2999 self.database.latest_local_store_device_registration().await
3000 }
3001
3002 pub async fn query_test_text(&self, sql: &str) -> String {
3003 self.database
3004 .test_query_optional_text(sql.to_string())
3005 .await
3006 .expect("test text query failed")
3007 .expect("test text query matched no row")
3008 }
3009}
3010
3011#[cfg(test)]
3013fn open_joined_test_store(
3014 store_dir: &StoreDir,
3015 device_id: String,
3016 synced_tables: Vec<coven_protocol::synced_schema::SyncedTable>,
3017 migrations: &[coven_database::Migration],
3018) -> Result<Database, TestError> {
3019 Ok(Database::open(
3020 &store_dir.db_path(),
3021 synced_tables,
3022 coven_protocol::blob::BLOB_TOMBSTONE_GRACE,
3023 coven_protocol::blob::TransferLimits::one_at_a_time(),
3024 device_id,
3025 std::sync::Arc::new(coven_foundation::clock::SystemClock),
3026 coven_database::CovenMigrationPolicy::ApplyPending,
3027 migrations,
3028 )?)
3029}
3030
3031#[cfg(test)]
3038async fn open_joined_test_device(
3039 store_dir: StoreDir,
3040 identity: &UserKeypair,
3041 storage: std::sync::Arc<coven_storage::CloudSyncConnection>,
3042 device_id: String,
3043 synced_tables: Vec<coven_protocol::synced_schema::SyncedTable>,
3044 migrations: &[coven_database::Migration],
3045) -> Result<TestDevice, TestError> {
3046 let database = open_joined_test_store(&store_dir, device_id, synced_tables, migrations)?;
3047 TestDevice::load(&database, store_dir, storage, identity.clone())
3048 .await
3049 .map_err(TestError::from)
3050}
3051
3052struct TestStoreProducers {
3053 unassigned: Option<TestDevice>,
3054 by_name: HashMap<String, TestDevice>,
3055}
3056
3057impl TestStore {
3058 pub fn root(&self) -> coven_protocol::store_commit::StoreRootRef {
3059 self.root.clone()
3060 }
3061
3062 pub async fn execute_unscoped_host_sql_for_test(
3063 &self,
3064 sql: impl Into<String>,
3065 ) -> Result<(), coven_database::HostWriteError<coven_database::DbError>> {
3066 self.founder
3067 .execute_unscoped_host_sql_for_test(sql.into())
3068 .await
3069 }
3070
3071 pub async fn bind_founder_device(
3072 &self,
3073 database: &Database,
3074 store_dir: StoreDir,
3075 ) -> Result<TestDevice, crate::sync::store::StoreError> {
3076 self.bind_device_in(database, store_dir, &self.signer).await
3077 }
3078
3079 pub async fn open_store_with_identity(
3080 &self,
3081 database: &Database,
3082 store_dir: StoreDir,
3083 identity: &UserKeypair,
3084 ) -> Result<crate::sync::store::Store, crate::sync::store::StoreInitializationError> {
3085 self.open_store_with_storage(
3086 coven_database::StoreDatabase::new(database),
3087 self.storage.clone(),
3088 store_dir,
3089 identity,
3090 )
3091 .await
3092 }
3093
3094 pub async fn open_store_with_storage(
3095 &self,
3096 database: coven_database::StoreDatabase,
3097 storage: Arc<dyn coven_storage::CloudSyncObjectStorage>,
3098 store_dir: StoreDir,
3099 identity: &UserKeypair,
3100 ) -> Result<crate::sync::store::Store, crate::sync::store::StoreInitializationError> {
3101 crate::sync::store::Store::open(
3102 database,
3103 storage,
3104 store_dir,
3105 &self.root,
3106 identity,
3107 Some(coven_keys::encryption::EncryptionService::from_key(
3108 [42; 32],
3109 )),
3110 )
3111 .await
3112 .map(|initialized| initialized.into_parts().0)
3113 }
3114
3115 pub async fn open_founder_store_with_storage(
3116 &self,
3117 database: coven_database::StoreDatabase,
3118 storage: Arc<dyn coven_storage::CloudSyncObjectStorage>,
3119 store_dir: StoreDir,
3120 ) -> Result<crate::sync::store::Store, crate::sync::store::StoreInitializationError> {
3121 self.open_store_with_storage(database, storage, store_dir, &self.signer)
3122 .await
3123 }
3124
3125 pub fn tombstone_deletions(&self) -> Vec<String> {
3126 self.home.deletes_seen()
3127 }
3128
3129 pub fn tombstone_provider_key(
3130 &self,
3131 stored: &coven_protocol::blob::locator::StoredBlobRef,
3132 ) -> String {
3133 coven_storage::blob_tombstone_key(
3134 stored,
3135 coven_storage::CloudSyncCipherStateAccess::suffix(self.storage.as_ref()),
3136 )
3137 }
3138
3139 pub fn stored_tombstone_bytes(&self, key: &str) -> Option<Vec<u8>> {
3140 let stored = self.home.get(key)?;
3141 let aad_context = coven_storage::cloud_aad_context(self.storage.store_id(), key);
3142 coven_storage::CloudSyncCipherStateAccess::open(self.storage.as_ref(), stored, &aad_context)
3143 .ok()
3144 }
3145
3146 pub async fn plant_tombstone_bytes(
3147 &self,
3148 key: &str,
3149 bytes: Vec<u8>,
3150 ) -> Result<(), coven_protocol::objects::StorageError> {
3151 let aad_context = coven_storage::cloud_aad_context(self.storage.store_id(), key);
3152 let stored = coven_storage::CloudSyncCipherStateAccess::seal(
3153 self.storage.as_ref(),
3154 bytes,
3155 &aad_context,
3156 );
3157 self.storage
3158 .write_provider_bytes_for_test(key, stored)
3159 .await
3160 }
3161
3162 pub async fn plant_tombstone(&self, tombstone: &crate::blob::delete::BlobTombstoneJson) {
3166 let key = self.tombstone_provider_key(&tombstone.stored);
3167 let bytes = serde_json::to_vec(tombstone).expect("serialize tombstone");
3168 self.plant_tombstone_bytes(&key, bytes)
3169 .await
3170 .expect("plant tombstone");
3171 }
3172
3173 pub fn fail_exact_delete_on_call(&self, call: usize) {
3174 self.home.fail_exact_delete_on_call(call);
3175 }
3176
3177 pub fn fail_nth_exact_delete_of(
3178 &self,
3179 slots: &[&coven_protocol::objects::ObjectSlot],
3180 call: usize,
3181 ) {
3182 self.home.fail_nth_exact_delete_of(slots, call);
3183 }
3184
3185 pub fn sort_provider_listings(&self) {
3186 self.home.sort_listings();
3187 }
3188
3189 pub fn provider_object_is_absent(&self, logical_key: &str) -> bool {
3190 self.home.get(logical_key).is_none()
3191 }
3192
3193 pub fn arm_provider_write_failures(&self) {
3194 self.home.arm_write_failures();
3195 }
3196
3197 pub fn fail_exact_create_before_call(&self, call: usize) {
3198 self.home.fail_exact_create_before_call(call);
3199 }
3200
3201 pub fn exact_creates(&self) -> Vec<coven_protocol::objects::ObjectSlot> {
3202 self.home.exact_creates()
3203 }
3204
3205 pub fn clear_exact_creates(&self) {
3206 self.home.clear_exact_creates();
3207 }
3208
3209 pub fn fail_exact_create_after_call(&self, call: usize) {
3210 self.home.fail_exact_create_after_call(call);
3211 }
3212
3213 pub fn pause_after_exact_create_call(
3214 &self,
3215 call: usize,
3216 ) -> (
3217 std::sync::Arc<tokio::sync::Notify>,
3218 std::sync::Arc<tokio::sync::Notify>,
3219 ) {
3220 self.home.pause_after_exact_create_call(call)
3221 }
3222
3223 pub async fn pull_with_storage_for_test(
3224 &self,
3225 database: &Database,
3226 storage: Arc<dyn coven_storage::CloudSyncObjectStorage>,
3227 store_dir: &StoreDir,
3228 routing_encryption: Option<&coven_keys::encryption::EncryptionService>,
3229 ) -> Result<crate::sync::store::StorePullResult, crate::sync::cycle::SyncCycleFailure> {
3230 let store = crate::sync::store::Store::load(
3231 coven_database::StoreDatabase::new(database),
3232 storage,
3233 store_dir.clone(),
3234 self.signer.clone(),
3235 routing_encryption.cloned(),
3236 )
3237 .await
3238 .map_err(|error| crate::sync::cycle::SyncCycleFailure::operation("load Store", error))?;
3239 store
3240 .authorize_writer()
3241 .await
3242 .map_err(|error| {
3243 crate::sync::cycle::SyncCycleFailure::operation("authorize Store writer", error)
3244 })?
3245 .pull(routing_encryption)
3246 .await
3247 }
3248
3249 pub async fn founder_recovery_authority(
3250 &self,
3251 ) -> coven_protocol::recovery::OwnerRecoveryAuthority {
3252 let protocol_root = self.founder.protocol_root_for_test();
3253 let owner_grant = protocol_root.descriptor.founder_grant.clone();
3254 let activation = coven_protocol::store_commit::OwnerRecoveryActivationId::derive(
3255 &self.root,
3256 &coven_keys::keys::public_key_hex(&self.signer),
3257 &owner_grant,
3258 &protocol_root.descriptor.founder_recovery,
3259 )
3260 .expect("derive founder recovery activation");
3261 coven_protocol::recovery::OwnerRecoveryAuthority {
3262 owner_identity_secret: hex::encode(self.signer.to_keypair_bytes()),
3263 owner_grant: owner_grant.clone(),
3264 recovery: coven_protocol::store_commit::OwnerRecoveryCursor {
3265 owner_grant,
3266 position: coven_protocol::store_commit::OwnerRecoveryPosition::BeforeFirst {
3267 activation,
3268 },
3269 },
3270 published_at: "2026-07-17T00:00:00Z".to_string(),
3271 }
3272 }
3273
3274 pub async fn create_circle(
3275 &self,
3276 metadata_stamp: &str,
3277 name: &str,
3278 ) -> Result<coven_protocol::CircleId, crate::sync::store::CircleOperationError> {
3279 self.founder.create_circle(metadata_stamp, name).await
3280 }
3281
3282 pub async fn run_founder_cycle(
3283 &self,
3284 observer: Option<&dyn coven_protocol::blob::BlobTransitionObserver>,
3285 ) -> Result<crate::sync::cycle::SyncCycleResult, crate::sync::cycle::SyncCycleFailure> {
3286 self.founder.run_cycle(observer).await
3287 }
3288
3289 pub async fn publish_fixture_position(&self, note_id: &str) -> u64 {
3290 self.founder.publish_fixture_position(note_id).await
3291 }
3292
3293 pub async fn create_exact_opaque_blob(
3294 &self,
3295 namespace: &str,
3296 id: &str,
3297 bytes: &[u8],
3298 ) -> coven_protocol::blob::locator::StoredBlobRef {
3299 self.founder
3300 .create_exact_opaque_blob(namespace, id, bytes)
3301 .await
3302 }
3303
3304 pub async fn create_exact_browsable_blob(
3305 &self,
3306 namespace: &str,
3307 id: &str,
3308 cloud_path: &str,
3309 bytes: &[u8],
3310 ) -> coven_protocol::blob::locator::StoredBlobRef {
3311 self.founder
3312 .create_exact_browsable_blob(namespace, id, cloud_path, bytes)
3313 .await
3314 }
3315
3316 pub async fn publish_exact_remote_blob_binding(
3317 &self,
3318 root_id: &str,
3319 row_id: &str,
3320 bytes: &[u8],
3321 ) -> coven_protocol::blob::locator::StoredBlobRef {
3322 self.founder
3323 .publish_exact_remote_blob_binding(root_id, row_id, bytes)
3324 .await
3325 }
3326
3327 pub async fn pull_into_result(
3328 &self,
3329 db: &Database,
3330 store_dir: &StoreDir,
3331 ) -> Result<
3332 (
3333 std::collections::BTreeMap<String, u64>,
3334 crate::sync::store::StorePullResult,
3335 ),
3336 TestPullError,
3337 > {
3338 let device = Box::pin(self.open_into(db, store_dir.clone()))
3339 .await
3340 .map_err(TestPullError::Open)?;
3341 device.pull_store().await
3342 }
3343
3344 pub async fn pull_into(
3345 &self,
3346 db: &Database,
3347 store_dir: &StoreDir,
3348 ) -> (
3349 std::collections::BTreeMap<String, u64>,
3350 crate::sync::store::StorePullResult,
3351 ) {
3352 self.pull_into_result(db, store_dir)
3353 .await
3354 .expect("pull exact test Store")
3355 }
3356
3357 pub async fn pull_and_fill_into(
3362 &self,
3363 db: &Database,
3364 store_dir: &StoreDir,
3365 ) -> (
3366 std::collections::BTreeMap<String, u64>,
3367 crate::sync::store::StorePullResult,
3368 ) {
3369 let pulled = self.pull_into(db, store_dir).await;
3370 crate::sync::test_owner_graph::TestOwnerGraph::new(
3371 coven_database::StoreDatabase::new(db),
3372 store_dir.clone(),
3373 )
3374 .fill_eager_cache(self.storage.clone())
3375 .await
3376 .expect("fill the eager cache behind the pull");
3377 pulled
3378 }
3379
3380 pub async fn promote_active_member_fixture(
3381 &self,
3382 owner_db: &Database,
3383 owner_db_store_dir: StoreDir,
3384 member_db: &Database,
3385 member_db_store_dir: StoreDir,
3386 owner: &UserKeypair,
3387 member: &UserKeypair,
3388 encryption: &coven_keys::encryption::EncryptionService,
3389 ) -> Result<coven_protocol::circle_control::StoreMembershipStateRef, TestError> {
3390 let owner_device = self
3391 .bind_device_in(owner_db, owner_db_store_dir.clone(), owner)
3392 .await?;
3393 let member_device = self
3394 .bind_device_in(member_db, member_db_store_dir.clone(), member)
3395 .await?;
3396 let request = owner_device
3397 .begin_owner_promotion_for_device(member_device.typed_device_id())
3398 .await?;
3399 let acceptance = member_device.accept_owner_promotion(request).await?;
3400 let finalized = owner_device
3401 .finalize_owner_promotion(encryption, acceptance)
3402 .await?;
3403 let (_, pull) = member_device.pull_store_with_encryption(encryption).await?;
3404 if !pull.held_positions.is_empty() {
3405 return Err(TestError::invariant(format!(
3406 "Owner promotion pull held signed positions: {:?}",
3407 pull.held_positions
3408 )));
3409 }
3410 Ok(finalized)
3411 }
3412
3413 pub async fn create(
3414 db: &Database,
3415 store_dir: StoreDir,
3416 store_id: &str,
3417 signer: UserKeypair,
3418 home: Arc<coven_storage::cloud::test_utils::InMemoryCloudHome>,
3419 ) -> Result<Arc<Self>, TestError> {
3420 Box::pin(Self::create_with_protection(
3421 db,
3422 store_dir,
3423 store_id,
3424 signer,
3425 home,
3426 coven_storage::CloudCipher::Encrypted(
3427 coven_keys::encryption::EncryptionService::from_key([42; 32]),
3428 ),
3429 coven_storage::BlobPathScheme::Hashed,
3430 ))
3431 .await
3432 }
3433
3434 pub async fn create_with_connection(
3435 db: &Database,
3436 store_dir: StoreDir,
3437 store_id: &str,
3438 signer: UserKeypair,
3439 home: Arc<coven_storage::cloud::test_utils::InMemoryCloudHome>,
3440 ) -> Result<TestStoreParts, TestError> {
3441 Box::pin(Self::create_with_protection_database(
3442 coven_database::StoreDatabase::new(db),
3443 store_dir,
3444 store_id,
3445 signer,
3446 home,
3447 coven_storage::CloudCipher::Encrypted(
3448 coven_keys::encryption::EncryptionService::from_key([42; 32]),
3449 ),
3450 coven_storage::BlobPathScheme::Hashed,
3451 ))
3452 .await
3453 }
3454
3455 pub async fn create_encrypted(
3456 db: &Database,
3457 store_dir: StoreDir,
3458 store_id: &str,
3459 signer: UserKeypair,
3460 home: Arc<coven_storage::cloud::test_utils::InMemoryCloudHome>,
3461 encryption: coven_keys::encryption::EncryptionService,
3462 ) -> Result<Arc<Self>, TestError> {
3463 Self::create_with_protection(
3464 db,
3465 store_dir,
3466 store_id,
3467 signer,
3468 home,
3469 coven_storage::CloudCipher::Encrypted(encryption),
3470 coven_storage::BlobPathScheme::Hashed,
3471 )
3472 .await
3473 }
3474
3475 pub async fn create_encrypted_with_connection(
3476 db: &Database,
3477 store_dir: StoreDir,
3478 store_id: &str,
3479 signer: UserKeypair,
3480 home: Arc<coven_storage::cloud::test_utils::InMemoryCloudHome>,
3481 encryption: coven_keys::encryption::EncryptionService,
3482 ) -> Result<TestStoreParts, TestError> {
3483 Self::create_with_protection_database(
3484 coven_database::StoreDatabase::new(db),
3485 store_dir,
3486 store_id,
3487 signer,
3488 home,
3489 coven_storage::CloudCipher::Encrypted(encryption),
3490 coven_storage::BlobPathScheme::Hashed,
3491 )
3492 .await
3493 }
3494
3495 pub async fn create_with_database(
3496 database: coven_database::StoreDatabase,
3497 store_dir: StoreDir,
3498 store_id: &str,
3499 signer: UserKeypair,
3500 home: Arc<coven_storage::cloud::test_utils::InMemoryCloudHome>,
3501 ) -> Result<Arc<Self>, TestError> {
3502 Box::pin(Self::create_with_protection_database(
3503 database,
3504 store_dir,
3505 store_id,
3506 signer,
3507 home,
3508 coven_storage::CloudCipher::Encrypted(
3509 coven_keys::encryption::EncryptionService::from_key([42; 32]),
3510 ),
3511 coven_storage::BlobPathScheme::Hashed,
3512 ))
3513 .await
3514 .map(|(store, _)| store)
3515 }
3516
3517 pub async fn create_browsable(
3522 db: &Database,
3523 store_dir: StoreDir,
3524 store_id: &str,
3525 signer: UserKeypair,
3526 home: Arc<coven_storage::cloud::test_utils::InMemoryCloudHome>,
3527 ) -> Result<Arc<Self>, TestError> {
3528 Box::pin(Self::create_with_protection(
3529 db,
3530 store_dir,
3531 store_id,
3532 signer,
3533 home,
3534 coven_storage::CloudCipher::Plaintext,
3535 coven_storage::BlobPathScheme::Plain,
3536 ))
3537 .await
3538 }
3539
3540 pub async fn create_browsable_with_connection(
3541 db: &Database,
3542 store_dir: StoreDir,
3543 store_id: &str,
3544 signer: UserKeypair,
3545 home: Arc<coven_storage::cloud::test_utils::InMemoryCloudHome>,
3546 ) -> Result<TestStoreParts, TestError> {
3547 Box::pin(Self::create_with_protection_database(
3548 coven_database::StoreDatabase::new(db),
3549 store_dir,
3550 store_id,
3551 signer,
3552 home,
3553 coven_storage::CloudCipher::Plaintext,
3554 coven_storage::BlobPathScheme::Plain,
3555 ))
3556 .await
3557 }
3558
3559 async fn create_with_protection(
3560 db: &Database,
3561 store_dir: StoreDir,
3562 store_id: &str,
3563 signer: UserKeypair,
3564 home: std::sync::Arc<coven_storage::cloud::test_utils::InMemoryCloudHome>,
3565 cipher: coven_storage::CloudCipher,
3566 blob_paths: coven_storage::BlobPathScheme,
3567 ) -> Result<Arc<Self>, TestError> {
3568 Self::create_with_protection_database(
3569 coven_database::StoreDatabase::new(db),
3570 store_dir,
3571 store_id,
3572 signer,
3573 home,
3574 cipher,
3575 blob_paths,
3576 )
3577 .await
3578 .map(|(store, _)| store)
3579 }
3580
3581 async fn create_with_protection_database(
3582 database: coven_database::StoreDatabase,
3583 store_dir: StoreDir,
3584 store_id: &str,
3585 signer: UserKeypair,
3586 home: std::sync::Arc<coven_storage::cloud::test_utils::InMemoryCloudHome>,
3587 cipher: coven_storage::CloudCipher,
3588 blob_paths: coven_storage::BlobPathScheme,
3589 ) -> Result<TestStoreParts, TestError> {
3590 let counted: std::sync::Arc<dyn coven_storage::ExactCloudHome> =
3594 std::sync::Arc::new(coven_storage::cloud::CountingCloudHome::new(home.clone()));
3595 let provider_requests = coven_storage::cloud::CloudHome::provider_requests(&*counted);
3596 let storage = std::sync::Arc::new(coven_storage::CloudSyncConnection::new(
3597 counted,
3598 cipher,
3599 blob_paths,
3600 store_id,
3601 signer.clone(),
3602 ));
3603 let founder = TestDevice::create_with_database(
3604 database,
3605 store_dir,
3606 storage.clone(),
3607 store_id,
3608 signer.clone(),
3609 )
3610 .await?;
3611 let root = founder.store_root().clone();
3612 let store = Arc::new(Self {
3613 home,
3614 provider_requests,
3615 storage: storage.clone(),
3616 root,
3617 signer,
3618 founder: founder.clone(),
3619 producers: Arc::new(tokio::sync::Mutex::new(TestStoreProducers {
3620 unassigned: Some(founder),
3621 by_name: HashMap::new(),
3622 })),
3623 });
3624 Ok((store, storage))
3625 }
3626
3627 pub fn delay_exact_full_reads(&self, delay: std::time::Duration) {
3632 self.home.delay_exact_full_reads(delay);
3633 }
3634
3635 pub fn exact_full_read_max_inflight(&self) -> usize {
3636 self.home.exact_full_read_max_inflight()
3637 }
3638
3639 pub fn exact_reads(&self) -> Vec<coven_protocol::objects::ObjectSlot> {
3640 self.home.exact_reads()
3641 }
3642
3643 pub fn clear_exact_reads(&self) {
3644 self.home.clear_exact_reads();
3645 }
3646
3647 pub fn provider_requests_issued(&self) -> u64 {
3648 self.provider_requests
3649 .as_ref()
3650 .expect("test Store home is counted")
3651 .issued()
3652 }
3653
3654 pub fn protocol_founder_pubkey(&self) -> String {
3655 coven_keys::keys::public_key_hex(&self.signer)
3656 }
3657
3658 pub async fn create_exact_protocol_object(
3659 &self,
3660 context: &coven_protocol::objects::ProtocolObjectContext,
3661 semantic_prefix: &str,
3662 extension: &str,
3663 bytes: &[u8],
3664 ) -> Result<coven_protocol::objects::ExactObjectRef, TestError> {
3665 let slot = self
3666 .storage
3667 .allocate_protocol_slot(context, semantic_prefix, extension)
3668 .await?;
3669 let prepared =
3670 self.storage
3671 .prepare_protocol_object(context, slot, semantic_prefix, bytes.to_vec())?;
3672 self.storage.create_protocol_object(&prepared).await?;
3673 Ok(prepared.reference().clone())
3674 }
3675
3676 pub async fn publish_exact_protocol_object(
3679 &self,
3680 object: &coven_protocol::objects::ExactObjectRef,
3681 bytes: Vec<u8>,
3682 ) -> Result<(), coven_protocol::objects::StorageError> {
3683 let prepared = coven_protocol::objects::PreparedExactObject::new(object.clone(), bytes)?;
3684 self.storage.create_protocol_object(&prepared).await
3685 }
3686
3687 pub async fn publish_prepared_protocol_object(
3688 &self,
3689 prepared: &coven_protocol::objects::PreparedExactObject,
3690 ) -> Result<(), coven_protocol::objects::StorageError> {
3691 self.storage.create_protocol_object(prepared).await
3692 }
3693
3694 pub async fn read_exact_protocol_object(
3695 &self,
3696 context: &coven_protocol::objects::ProtocolObjectContext,
3697 object: &coven_protocol::objects::ExactObjectRef,
3698 semantic_prefix: &str,
3699 ) -> Result<Vec<u8>, coven_protocol::objects::StorageError> {
3700 self.storage
3701 .read_protocol_object(context, object, semantic_prefix)
3702 .await
3703 }
3704
3705 pub async fn contains_blob_object(&self, reference: &coven_protocol::blob::RowBlobRef) -> bool {
3706 match reference.stored() {
3707 Some(stored) => self
3708 .contains_stored_blob_object(stored)
3709 .await
3710 .unwrap_or_else(|error| panic!("verify exact blob object: {error}")),
3711 None => false,
3712 }
3713 }
3714
3715 pub async fn contains_stored_blob_object(
3716 &self,
3717 stored: &coven_protocol::blob::locator::StoredBlobRef,
3718 ) -> Result<bool, coven_protocol::objects::StorageError> {
3719 match self.storage.verify_blob_object(stored).await {
3720 Ok(()) => Ok(true),
3721 Err(coven_protocol::objects::StorageError::NotFound(_)) => Ok(false),
3722 Err(error) => Err(error),
3723 }
3724 }
3725
3726 pub async fn contains_blob_tombstone(
3727 &self,
3728 stored: &coven_protocol::blob::locator::StoredBlobRef,
3729 ) -> Result<bool, coven_storage::cloud::CloudHomeError> {
3730 let key = coven_storage::blob_tombstone_key(
3731 stored,
3732 coven_storage::CloudSyncCipherStateAccess::suffix(self.storage.as_ref()),
3733 );
3734 coven_storage::cloud::CloudHome::exists(self.home.as_ref(), &key).await
3735 }
3736
3737 pub async fn contains_membership_rollup(
3738 &self,
3739 rollup: &coven_protocol::store_commit::MembershipRollupRef,
3740 ) -> Result<bool, TestError> {
3741 let context = coven_protocol::objects::ProtocolObjectContext::signed_plaintext(
3742 self.root.store_root_hash,
3743 coven_protocol::objects::ProtocolObjectDomain::StoreMembershipRollup,
3744 );
3745 let prefix = coven_protocol::store_commit::semantic_prefix_from_exact_object(
3746 &rollup.object,
3747 coven_protocol::objects::ProtectedObjectDomain::StoreMembershipRollup.extension(),
3748 )?;
3749 match self
3750 .storage
3751 .read_protocol_object(&context, &rollup.object, &prefix)
3752 .await
3753 {
3754 Ok(_) => Ok(true),
3755 Err(coven_protocol::objects::StorageError::NotFound(_)) => Ok(false),
3756 Err(error) => Err(TestError::from(error)),
3757 }
3758 }
3759
3760 pub async fn contains_circle_snapshot_image(
3761 &self,
3762 circle_id: coven_protocol::circle::CircleId,
3763 meta: &coven_protocol::store_commit::CircleSnapshotMeta,
3764 ) -> Result<bool, TestError> {
3765 let access = self
3766 .founder
3767 .circle_epoch_access(circle_id, meta.control.clone())
3768 .await?
3769 .ok_or_else(|| {
3770 TestError::invariant(
3771 "the Circle snapshot control has no retained access".to_string(),
3772 )
3773 })?;
3774 let context = access.protocol_context(
3775 self.root.store_root_hash,
3776 coven_protocol::objects::ProtocolObjectDomain::CircleSnapshotImage,
3777 );
3778 let prefix = coven_protocol::store_commit::semantic_prefix_from_exact_object(
3779 &meta.bootstrap.image.object,
3780 coven_protocol::objects::ProtectedObjectDomain::CircleSnapshotImage.extension(),
3781 )?;
3782 match self
3783 .storage
3784 .read_protocol_object(&context, &meta.bootstrap.image.object, &prefix)
3785 .await
3786 {
3787 Ok(_) => Ok(true),
3788 Err(coven_protocol::objects::StorageError::NotFound(_)) => Ok(false),
3789 Err(error) => Err(TestError::from(error)),
3790 }
3791 }
3792
3793 pub async fn circle_package_in(
3794 &self,
3795 commit_ref: &coven_protocol::store_commit::StoreBatchCommitRef,
3796 ) -> coven_protocol::store_commit::CirclePackageRef {
3797 let commit = self
3798 .founder
3799 .load_commit_for_test(commit_ref)
3800 .await
3801 .expect("load the exact Circle package commit");
3802 let [package] = commit.value().circle_packages() else {
3803 panic!("the commit must carry exactly one Circle package");
3804 };
3805 package.clone()
3806 }
3807
3808 pub async fn circle_package_object_present(
3809 &self,
3810 package: &coven_protocol::store_commit::CirclePackageRef,
3811 activation: &coven_protocol::store_commit::StoreBatchCommitRef,
3812 ) -> bool {
3813 let access = self
3814 .founder
3815 .circle_epoch_access(package.circle_id, package.control.clone())
3816 .await
3817 .expect("resolve Circle package access")
3818 .expect("the package's control stays retained after its epoch closed");
3819 let context = access.protocol_context(
3820 self.root.store_root_hash,
3821 coven_protocol::objects::ProtocolObjectDomain::CirclePackage,
3822 );
3823 let prefix = coven_protocol::store_commit::circle_package_semantic_prefix(
3824 package.circle_id,
3825 package.package.candidate_family,
3826 &activation.coord.stream_id.to_string(),
3827 activation.coord.sequence(),
3828 package.package.content_hash,
3829 );
3830 match self
3831 .storage
3832 .read_protocol_object(&context, &package.package.object, &prefix)
3833 .await
3834 {
3835 Ok(_) => true,
3836 Err(coven_protocol::objects::StorageError::NotFound(_)) => false,
3837 Err(error) => panic!("read the exact Circle package object: {error}"),
3838 }
3839 }
3840
3841 pub async fn overwrite_membership_head(
3842 &self,
3843 reference: &coven_protocol::membership::MembershipHeadRef,
3844 head: &coven_protocol::membership::AuthorHead,
3845 ) {
3846 self.storage
3847 .delete_protocol_object(&reference.object)
3848 .await
3849 .expect("delete exact head before replacement");
3850 let context = coven_protocol::objects::ProtocolObjectContext::signed_plaintext(
3851 self.root.store_root_hash,
3852 coven_protocol::objects::ProtocolObjectDomain::StoreMembershipHead,
3853 );
3854 let prefix = coven_protocol::store_commit::membership_head_slot_prefix(
3855 &reference.coord.author_pubkey,
3856 &reference.coord.author_owner_grant,
3857 reference.coord.stream_id,
3858 reference.coord.seq,
3859 );
3860 let prepared = self
3861 .storage
3862 .prepare_protocol_object(
3863 &context,
3864 reference.object.slot().clone(),
3865 &prefix,
3866 serde_json::to_vec(head).expect("serialize replacement head"),
3867 )
3868 .expect("prepare replacement head");
3869 self.storage
3870 .create_protocol_object(&prepared)
3871 .await
3872 .expect("write replacement head");
3873 }
3874
3875 pub async fn delete_membership_head_for_test(
3876 &self,
3877 reference: &coven_protocol::membership::MembershipHeadRef,
3878 ) -> Result<(), coven_protocol::objects::StorageError> {
3879 self.storage.delete_protocol_object(&reference.object).await
3880 }
3881
3882 pub async fn pending_device_join_observation(
3883 &self,
3884 pending: &crate::sync::store::DeviceJoinJournalDatabase,
3885 offer: &coven_protocol::store_commit::device_join_exchange::DeviceJoinOffer,
3886 ) -> Result<crate::sync::store::PendingDeviceJoinObservation<'_>, TestError> {
3887 self.founder
3888 .pending_device_join_observation_for_test(pending, offer)
3889 .await
3890 }
3891
3892 pub async fn open_pending_device_join(
3893 &self,
3894 pending: &crate::sync::store::DeviceJoinJournalDatabase,
3895 identity: &UserKeypair,
3896 offer: coven_protocol::store_commit::device_join_exchange::DeviceJoinOffer,
3897 ) -> Result<crate::sync::store::PendingDeviceJoinAuthority<'_>, TestError> {
3898 self.founder
3899 .open_pending_device_join_for_test(pending, identity, offer)
3900 .await
3901 }
3902
3903 pub async fn prepare_snapshot_bootstrap<'a>(
3904 &'a self,
3905 membership_floor: &coven_protocol::membership::MembershipFloor,
3906 binary_schema_version: u32,
3907 target_path: &std::path::Path,
3908 restorer_identity: &UserKeypair,
3909 ) -> Result<crate::sync::store::PreparedSnapshotBootstrap<'a>, crate::sync::store::SnapshotError>
3910 {
3911 self.founder
3912 .prepare_snapshot_bootstrap_for_test(
3913 membership_floor,
3914 binary_schema_version,
3915 target_path,
3916 restorer_identity,
3917 )
3918 .await
3919 }
3920
3921 pub async fn bind_device_in(
3922 &self,
3923 db: &Database,
3924 store_dir: StoreDir,
3925 identity: &UserKeypair,
3926 ) -> Result<TestDevice, crate::sync::store::StoreError> {
3927 TestDevice::load_with_database(
3928 coven_database::StoreDatabase::new(db),
3929 std::sync::Arc::new(self.storage.connection_for_test_identity(identity.clone())),
3930 identity.clone(),
3931 store_dir,
3932 )
3933 .await
3934 }
3935
3936 pub async fn bind_device(
3937 &self,
3938 db: &Database,
3939 store_dir: StoreDir,
3940 identity: &UserKeypair,
3941 ) -> Result<TestDevice, crate::sync::store::StoreError> {
3942 self.bind_device_in(db, store_dir, identity).await
3943 }
3944
3945 pub async fn drain_uploads(
3946 &self,
3947 database: &coven_database::StoreDatabase,
3948 store_dir: &coven_foundation::store_dir::StoreDir,
3949 clock: &dyn coven_foundation::clock::Clock,
3950 routing_encryption: Option<&coven_keys::encryption::EncryptionService>,
3951 observer: Option<&dyn coven_protocol::blob::BlobTransitionObserver>,
3952 ) -> Result<crate::blob::DrainOutcome, TestError> {
3953 let store = self
3954 .bind_store_device(database, store_dir.clone(), &self.signer)
3955 .await?;
3956 store
3957 .drain_uploads(clock, routing_encryption, observer)
3958 .await
3959 }
3960
3961 pub async fn activate_joined_device(
3962 &self,
3963 observer_db: &Database,
3964 observer_store_dir: StoreDir,
3965 joining_db: &Database,
3966 joining_store_dir: StoreDir,
3967 joining_identity: &UserKeypair,
3968 published_at: &str,
3969 ) -> Result<TestDevice, TestError> {
3970 self.activate_joined_device_with_clock(
3971 observer_db,
3972 observer_store_dir,
3973 joining_db,
3974 joining_store_dir,
3975 joining_identity,
3976 published_at,
3977 Arc::new(coven_foundation::clock::SystemClock),
3978 )
3979 .await
3980 }
3981
3982 #[allow(clippy::too_many_arguments)]
3983 pub fn activate_joined_device_with_clock<'a>(
3984 &'a self,
3985 observer_db: &'a Database,
3986 observer_store_dir: StoreDir,
3987 joining_db: &'a Database,
3988 joining_store_dir: StoreDir,
3989 joining_identity: &'a UserKeypair,
3990 published_at: &'a str,
3991 clock: coven_foundation::clock::ClockRef,
3992 ) -> std::pin::Pin<
3993 Box<dyn std::future::Future<Output = Result<TestDevice, TestError>> + Send + 'a>,
3994 > {
3995 Box::pin(async move {
3996 let binding = self.bind_device_in(observer_db, observer_store_dir, &self.signer);
3997 let observer = binding.await?;
3998 TestDevice::activate_joined_with_clock(
3999 observer,
4000 coven_database::StoreDatabase::new(joining_db),
4001 joining_store_dir,
4002 joining_identity,
4003 published_at,
4004 std::sync::Arc::new(
4005 self.storage
4006 .connection_for_test_identity(joining_identity.clone()),
4007 ),
4008 clock,
4009 )
4010 .await
4011 })
4012 }
4013
4014 #[cfg(test)]
4018 #[allow(clippy::too_many_arguments)]
4019 pub async fn activate_joined_device_from_snapshot(
4020 &self,
4021 observer_db: &Database,
4022 observer_store_dir: StoreDir,
4023 joining_store_dir: StoreDir,
4024 joining_identity: &UserKeypair,
4025 published_at: &str,
4026 synced_tables: Vec<coven_protocol::synced_schema::SyncedTable>,
4027 migrations: Vec<coven_database::Migration>,
4028 binary_schema_version: u32,
4029 ) -> Result<TestDevice, TestError> {
4030 let observer = self
4031 .bind_device_in(observer_db, observer_store_dir, &self.signer)
4032 .await?;
4033 TestDevice::activate_joined_from_snapshot(
4034 observer,
4035 joining_store_dir,
4036 joining_identity,
4037 published_at,
4038 std::sync::Arc::new(
4039 self.storage
4040 .connection_for_test_identity(joining_identity.clone()),
4041 ),
4042 synced_tables,
4043 migrations,
4044 binary_schema_version,
4045 )
4046 .await
4047 }
4048
4049 pub async fn bind_store_device(
4050 &self,
4051 database: &coven_database::StoreDatabase,
4052 store_dir: StoreDir,
4053 identity: &UserKeypair,
4054 ) -> Result<TestDevice, TestError> {
4055 if identity.public_key() != self.signer.public_key() {
4056 return Err(TestError::invariant(
4057 "custom Store database binding requires the founder identity".to_string(),
4058 ));
4059 }
4060 TestDevice::load_with_database(
4061 database.clone(),
4062 std::sync::Arc::new(self.storage.connection_for_test_identity(identity.clone())),
4063 identity.clone(),
4064 store_dir,
4065 )
4066 .await
4067 .map_err(TestError::from)
4068 }
4069
4070 pub async fn admit_member(
4071 &self,
4072 db: &Database,
4073 store_dir: StoreDir,
4074 identity: &UserKeypair,
4075 member_pubkey: &str,
4076 member_email: Option<&str>,
4077 role: coven_protocol::membership::MemberRole,
4078 encryption: &coven_keys::encryption::EncryptionService,
4079 store_name: &str,
4080 ) -> Result<crate::sync::store::MemberAdmission, crate::sync::store::MembershipOpsError> {
4081 let device = self
4082 .bind_device_in(db, store_dir, identity)
4083 .await
4084 .map_err(crate::sync::store::MembershipOpsError::Store)?;
4085 device
4086 .admit_member(
4087 member_pubkey,
4088 member_email,
4089 role,
4090 encryption,
4091 self.storage.store_id(),
4092 store_name,
4093 )
4094 .await
4095 }
4096
4097 pub async fn admit_and_activate_peer(
4098 &self,
4099 observer_db: &Database,
4100 observer_db_store_dir: StoreDir,
4101 peer_db: &Database,
4102 peer_db_store_dir: StoreDir,
4103 peer: &UserKeypair,
4104 ) -> Result<TestDevice, TestError> {
4105 self.admit_member(
4106 observer_db,
4107 observer_db_store_dir.clone(),
4108 &self.signer,
4109 &pubkey_hex(peer),
4110 None,
4111 coven_protocol::membership::MemberRole::Member,
4112 &coven_keys::encryption::EncryptionService::from_key([42; 32]),
4113 "Test Store",
4114 )
4115 .await?;
4116 self.activate_joined_device(
4117 observer_db,
4118 observer_db_store_dir.clone(),
4119 peer_db,
4120 peer_db_store_dir.clone(),
4121 peer,
4122 "2026-07-16T00:00:00Z",
4123 )
4124 .await
4125 }
4126
4127 pub async fn remove_member(
4128 &self,
4129 db: &Database,
4130 store_dir: StoreDir,
4131 identity: &UserKeypair,
4132 member_pubkey: &str,
4133 encryption: &coven_keys::encryption::EncryptionService,
4134 master_keys: &dyn coven_keys::keys::MasterKeyCustody,
4135 ) -> Result<String, crate::sync::store::MembershipOpsError> {
4136 let device = self
4137 .bind_device_in(db, store_dir, identity)
4138 .await
4139 .map_err(crate::sync::store::MembershipOpsError::Store)?;
4140 device
4141 .remove_member(
4142 member_pubkey,
4143 encryption,
4144 master_keys,
4145 self.storage.as_ref(),
4146 self.storage.as_ref(),
4147 )
4148 .await
4149 }
4150
4151 pub async fn device_id(&self, name: &str) -> Result<String, TestError> {
4152 self.ensure_producer_registered(name).await?;
4153 let producers = self.producers.lock().await;
4154 Ok(producers
4155 .by_name
4156 .get(name)
4157 .expect("registered test producer exists")
4158 .device_id())
4159 }
4160
4161 pub async fn latest_store_position(
4162 &self,
4163 ) -> Result<Option<coven_protocol::store_commit::StoreBatchCommitRef>, TestError> {
4164 self.founder.latest_store_position().await
4165 }
4166
4167 pub async fn load_commit_for_test(
4168 &self,
4169 reference: &coven_protocol::store_commit::StoreBatchCommitRef,
4170 ) -> Result<coven_protocol::store_commit::VerifiedStoreBatchCommit, TestError> {
4171 self.founder
4172 .load_commit_for_test(reference)
4173 .await
4174 .map_err(TestError::from)
4175 }
4176
4177 pub async fn load_membership_head_for_test(
4178 &self,
4179 reference: &coven_protocol::membership::MembershipHeadRef,
4180 ) -> Result<coven_protocol::membership::AuthorHead, TestError> {
4181 self.founder
4182 .load_membership_head_for_test(reference)
4183 .await
4184 .map_err(TestError::from)
4185 }
4186
4187 pub async fn load_registration_for_test(
4188 &self,
4189 reference: &coven_protocol::store_commit::StoreDeviceRegistrationRef,
4190 ) -> Result<coven_protocol::store_commit::StoreDeviceRegistration, TestError> {
4191 self.founder
4192 .load_registration_for_test(reference)
4193 .await
4194 .map_err(TestError::from)
4195 }
4196
4197 pub async fn load_store_package_for_test(
4198 &self,
4199 reference: &coven_protocol::store_commit::StoreBatchCommitRef,
4200 ) -> Result<Option<coven_protocol::objects::VerifiedObject<Vec<u8>>>, TestError> {
4201 self.founder
4202 .load_store_package_for_test(reference)
4203 .await
4204 .map_err(TestError::from)
4205 }
4206
4207 pub async fn prepare_founder_store_partition_blob_for_test(
4208 &self,
4209 fact: &coven_database::StoreWriteBlobFact,
4210 authority: &coven_protocol::objects::BlobWriteAuthority<'_>,
4211 ) -> Result<(), crate::sync::store::StoreError> {
4212 self.founder
4213 .authorize_writer()
4214 .await?
4215 .prepare_store_partition_blob(fact, authority)
4216 .await
4217 .map(|_| ())
4218 }
4219
4220 pub async fn next_commit_sequence(&self, name: &str) -> Result<u64, TestError> {
4221 self.ensure_producer_registered(name).await?;
4222 let producer = {
4223 let producers = self.producers.lock().await;
4224 producers
4225 .by_name
4226 .get(name)
4227 .expect("registered test producer exists")
4228 .clone()
4229 };
4230 producer
4231 .latest_local_store_position()
4232 .await?
4233 .map_or(Ok(1), |reference| {
4234 reference.coord.sequence().checked_add(1).ok_or_else(|| {
4235 TestError::invariant("test producer sequence exhausted u64".to_string())
4236 })
4237 })
4238 }
4239
4240 pub async fn founder_device_authority(&self) -> Result<TestDeviceSigningAuthority, TestError> {
4241 self.founder.device_authority_for_test().await
4242 }
4243
4244 async fn ensure_producer_registered(&self, name: &str) -> Result<(), TestError> {
4245 {
4246 let producers = self.producers.lock().await;
4247 if producers.by_name.contains_key(name) {
4248 return Ok(());
4249 }
4250 }
4251
4252 let unassigned = {
4253 let mut producers = self.producers.lock().await;
4254 producers.unassigned.take()
4255 };
4256 let producer = match unassigned {
4257 Some(producer) => producer,
4258 None => {
4259 let db_store_dir = crate::sync::test_helpers::test_store_dir();
4260 let db = crate::sync::test_helpers::open_test_db(db_store_dir.clone());
4261 let observer = {
4262 let producers = self.producers.lock().await;
4263 producers
4264 .by_name
4265 .values()
4266 .next()
4267 .ok_or_else(|| {
4268 TestError::invariant(
4269 "test Store has no active device observer".to_string(),
4270 )
4271 })?
4272 .clone()
4273 };
4274 TestDevice::activate_joined(
4275 observer,
4276 coven_database::StoreDatabase::new(&db),
4277 db_store_dir,
4278 &self.signer,
4279 "2026-07-16T00:00:00Z",
4280 std::sync::Arc::new(
4281 self.storage
4282 .connection_for_test_identity(self.signer.clone()),
4283 ),
4284 )
4285 .await?
4286 }
4287 };
4288 let mut producers = self.producers.lock().await;
4289 if producers
4290 .by_name
4291 .insert(name.to_string(), producer)
4292 .is_some()
4293 {
4294 return Err(TestError::invariant(format!(
4295 "test producer {name:?} was registered twice"
4296 )));
4297 }
4298 Ok(())
4299 }
4300
4301 pub async fn open_into(
4302 &self,
4303 db: &Database,
4304 store_dir: StoreDir,
4305 ) -> Result<TestDevice, crate::sync::store::StoreInitializationError> {
4306 TestDevice::open_with_database(
4307 coven_database::StoreDatabase::new(db),
4308 store_dir,
4309 std::sync::Arc::new(
4310 self.storage
4311 .connection_for_test_identity(self.signer.clone()),
4312 ),
4313 &self.root,
4314 &self.signer,
4315 )
4316 .await
4317 }
4318
4319 pub async fn open_into_store_database(
4320 &self,
4321 database: &coven_database::StoreDatabase,
4322 store_dir: StoreDir,
4323 ) -> Result<TestDevice, crate::sync::store::StoreInitializationError> {
4324 TestDevice::open_with_database(
4325 database.clone(),
4326 store_dir,
4327 std::sync::Arc::new(
4328 self.storage
4329 .connection_for_test_identity(self.signer.clone()),
4330 ),
4331 &self.root,
4332 &self.signer,
4333 )
4334 .await
4335 }
4336
4337 pub async fn publish_pending(
4338 &self,
4339 db: &Database,
4340 store_dir: &StoreDir,
4341 ) -> Result<bool, TestError> {
4342 self.publish_pending_store_database(&coven_database::StoreDatabase::new(db), store_dir)
4343 .await
4344 }
4345
4346 pub async fn publish_pending_store_database(
4347 &self,
4348 database: &coven_database::StoreDatabase,
4349 store_dir: &StoreDir,
4350 ) -> Result<bool, TestError> {
4351 let device = self
4352 .bind_store_device(database, store_dir.clone(), &self.signer)
4353 .await?;
4354 device.publish_pending_store_database().await
4355 }
4356
4357 #[cfg(test)]
4358 pub async fn cross_principal_device_for_test(
4359 &self,
4360 identity: &UserKeypair,
4361 peer_account_id: &str,
4362 ) -> Result<CrossPrincipalTestDevice, TestError> {
4363 let provider_binding =
4364 coven_storage::CloudSyncObjectStorage::provider_binding(&*self.storage).await?;
4365 let coven_protocol::objects::StoreProviderBinding::Dropbox { namespace_id } =
4366 &provider_binding.store
4367 else {
4368 return Err(TestError::invariant(
4369 "cross-principal test Store is not Dropbox".to_string(),
4370 ));
4371 };
4372 let peer_binding = coven_protocol::objects::ResolvedProviderBinding {
4373 store: provider_binding.store.clone(),
4374 device: coven_protocol::objects::ProviderDeviceBinding {
4375 principal: coven_protocol::objects::ProviderPrincipalId::Dropbox {
4376 account_id: peer_account_id.to_string(),
4377 },
4378 },
4379 };
4380 let peer_home: std::sync::Arc<dyn coven_storage::ExactCloudHome> = std::sync::Arc::new(
4381 self.home
4382 .as_ref()
4383 .clone()
4384 .with_provider_binding(peer_binding),
4385 );
4386 Ok(CrossPrincipalTestDevice {
4387 storage: std::sync::Arc::new(
4388 self.storage
4389 .connection_for_test_identity_and_home(identity.clone(), peer_home),
4390 ),
4391 access_administrator: TestDropboxAccessAdministrator {
4392 namespace_id: namespace_id.clone(),
4393 },
4394 })
4395 }
4396
4397 #[cfg(test)]
4406 pub fn install_cross_principal_device<'a>(
4407 &'a self,
4408 joining_store_dir: StoreDir,
4409 synced_tables: Vec<coven_protocol::synced_schema::SyncedTable>,
4410 migrations: Vec<coven_database::Migration>,
4411 binary_schema_version: u32,
4412 identity: &'a UserKeypair,
4413 peer_account_id: &'a str,
4414 published_at: &'a str,
4415 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<JoinedTestStore, TestError>> + 'a>>
4416 {
4417 Box::pin(async move {
4418 let open_synced_tables = synced_tables.clone();
4419 let observer = self.founder.clone();
4420 observer.ensure_device_join_snapshot_for_test().await?;
4423 let peer = self
4424 .cross_principal_device_for_test(identity, peer_account_id)
4425 .await?;
4426 let pending_dir = tempfile::tempdir()?;
4427 let pending = crate::sync::store::DeviceJoinJournalDatabase::open_for_test(
4428 pending_dir.path().join("pending-device-join.sqlite"),
4429 )?;
4430 let offer = observer.begin_device_join(&pubkey_hex(identity)).await?;
4431 let mut pending_join = peer
4432 .open_pending_device_join(&pending, identity, offer.clone())
4433 .await?;
4434 let access_request = pending_join.prepare_provider_access_request().await?;
4435 let approval = peer
4436 .authorize_device_provider_access(&observer, access_request)
4437 .await?;
4438 if !matches!(
4439 approval.admission,
4440 coven_protocol::store_commit::device_join_exchange::DeviceProviderAdmission::CrossPrincipal { .. }
4441 ) {
4442 return Err(
4443 TestError::invariant(
4444 "distinct provider principals produced same-principal admission",
4445 ),
4446 );
4447 }
4448 let registration_request = pending_join.prepare_registration_request(approval).await?;
4449 let provisional = observer
4450 .accept_device_registration_request(registration_request)
4451 .await?;
4452 let provider_ready = observer
4453 .publish_device_provider_challenge(provisional)
4454 .await?;
4455 drop(pending_join);
4456 let joined_device_id = offer.attempt_id.to_string();
4457 let restoring = peer
4458 .install_store_snapshot(
4459 &joining_store_dir,
4460 &offer.store_root,
4461 &observer.membership().await?,
4462 identity,
4463 offer.attempt_id.to_string(),
4464 binary_schema_version,
4465 synced_tables,
4466 &migrations,
4467 )
4468 .await?;
4469 let mut joining = restoring.begin_device_join(&pending, offer).await?;
4470 let readiness = joining
4471 .bootstrap(provider_ready, published_at, None)
4472 .await?;
4473 if !matches!(
4474 readiness.provider,
4475 coven_protocol::store_commit::device_join_exchange::DeviceProviderReadiness::CrossPrincipal(_)
4476 ) {
4477 return Err(
4478 TestError::invariant(
4479 "distinct provider principals produced same-principal readiness",
4480 ),
4481 );
4482 }
4483 let completion = observer
4484 .complete_device_provider_admission(readiness)
4485 .await?;
4486 if !matches!(
4487 completion,
4488 coven_protocol::store_commit::device_join_exchange::DeviceProviderAdmissionCompletion::CrossPrincipal { .. }
4489 ) {
4490 return Err(
4491 TestError::invariant(
4492 "distinct provider principals produced same-principal completion",
4493 ),
4494 );
4495 }
4496 let activation = observer.finalize_device_join(completion).await?;
4497 joining.complete(activation).await?;
4498 drop(joining);
4499 let database = open_joined_test_store(
4500 &joining_store_dir,
4501 joined_device_id,
4502 open_synced_tables,
4503 &migrations,
4504 )?;
4505 Ok(JoinedTestStore {
4506 database: coven_database::StoreDatabase::new(&database),
4507 })
4508 })
4509 }
4510
4511 #[cfg(test)]
4512 pub async fn push_circle_snapshots(
4513 &self,
4514 db: &Database,
4515 store_dir: StoreDir,
4516 temp_dir: std::path::PathBuf,
4517 schema_version: u32,
4518 created_at: &str,
4519 store_routing: &coven_keys::encryption::EncryptionService,
4520 ) -> Result<coven_protocol::store_commit::CircleSnapshotMeta, crate::sync::store::SnapshotError>
4521 {
4522 self.bind_device(db, store_dir, &self.signer)
4523 .await
4524 .map_err(|error| crate::sync::store::SnapshotError::PublicationStore(Box::new(error)))?
4525 .authorize_writer()
4526 .await
4527 .map_err(crate::sync::store::SnapshotError::from)?
4528 .circles()
4529 .snapshots()
4530 .author_one_circle_snapshot_for_test(
4531 temp_dir,
4532 schema_version,
4533 created_at,
4534 store_routing,
4535 )
4536 .await
4537 }
4538
4539 #[cfg(test)]
4540 pub async fn load_circle_snapshot_metas(
4541 &self,
4542 db: &Database,
4543 store_dir: StoreDir,
4544 circle_id: coven_protocol::circle::CircleId,
4545 access: &coven_protocol::circle_activation::CircleEpochAccess,
4546 ) -> Result<
4547 Vec<coven_protocol::store_commit::CircleSnapshotMeta>,
4548 crate::sync::store::SnapshotError,
4549 > {
4550 self.bind_device(db, store_dir, &self.signer)
4551 .await
4552 .map_err(|error| crate::sync::store::SnapshotError::PublicationStore(Box::new(error)))?
4553 .authorize_writer()
4554 .await
4555 .map_err(crate::sync::store::SnapshotError::from)?
4556 .circles()
4557 .snapshots()
4558 .load_circle_snapshot_metas_for_test(circle_id, access)
4559 .await
4560 }
4561
4562 #[cfg(test)]
4563 pub async fn verify_standalone_circle_snapshot_image(
4564 &self,
4565 db: &Database,
4566 store_dir: StoreDir,
4567 circle_id: coven_protocol::circle::CircleId,
4568 access: &coven_protocol::circle_activation::CircleEpochAccess,
4569 store_routing: &coven_keys::encryption::EncryptionService,
4570 ) -> Result<(), crate::sync::store::SnapshotError> {
4571 self.bind_device(db, store_dir, &self.signer)
4572 .await
4573 .map_err(|error| crate::sync::store::SnapshotError::PublicationStore(Box::new(error)))?
4574 .authorize_writer()
4575 .await
4576 .map_err(crate::sync::store::SnapshotError::from)?
4577 .circles()
4578 .snapshots()
4579 .verify_standalone_circle_snapshot_image_for_test(circle_id, access, store_routing)
4580 .await
4581 }
4582
4583 #[cfg(test)]
4584 pub async fn circle_snapshot_is_stable(
4585 &self,
4586 db: &Database,
4587 store_dir: StoreDir,
4588 circle_id: coven_protocol::circle::CircleId,
4589 snapshot_cut: &coven_protocol::store_commit::CommitFrontier,
4590 ) -> Result<bool, crate::sync::store::SnapshotError> {
4591 self.bind_device(db, store_dir, &self.signer)
4592 .await
4593 .map_err(|error| crate::sync::store::SnapshotError::PublicationStore(Box::new(error)))?
4594 .authorize_writer()
4595 .await
4596 .map_err(crate::sync::store::SnapshotError::from)?
4597 .circles()
4598 .snapshots()
4599 .circle_snapshot_is_stable(circle_id, snapshot_cut)
4600 .await
4601 }
4602
4603 #[cfg(test)]
4604 pub async fn load_circle_acknowledgement(
4605 &self,
4606 db: &Database,
4607 store_dir: StoreDir,
4608 reference: &coven_protocol::store_commit::CircleAckRef,
4609 ) -> Result<coven_protocol::store_commit::CircleAck, crate::sync::store::StoreAckError> {
4610 self.bind_device(db, store_dir, &self.signer)
4611 .await
4612 .map_err(crate::sync::store::StoreAckError::Outbound)?
4613 .load_circle_acknowledgement_for_test(reference)
4614 .await
4615 }
4616
4617 #[cfg(test)]
4618 pub async fn read_circle_snapshot_image(
4619 &self,
4620 selected: &coven_protocol::store_commit::CircleSnapshotMeta,
4621 access: &coven_protocol::circle_activation::CircleEpochAccess,
4622 ) -> Result<Vec<u8>, coven_protocol::objects::StorageError> {
4623 let context = access.protocol_context(
4624 self.root.store_root_hash,
4625 coven_protocol::objects::ProtocolObjectDomain::CircleSnapshotImage,
4626 );
4627 self.storage
4628 .read_protocol_object(
4629 &context,
4630 &selected.bootstrap.image.object,
4631 &coven_protocol::store_commit::circle_snapshot_image_semantic_prefix(
4632 selected.circle_id,
4633 &selected.author_registration.device_id.to_string(),
4634 selected.bootstrap.image.image_hash,
4635 ),
4636 )
4637 .await
4638 }
4639
4640 #[cfg(test)]
4641 pub async fn circle_snapshot_meta_is_unreadable(
4642 &self,
4643 circle_id: coven_protocol::circle::CircleId,
4644 encryption: coven_keys::encryption::EncryptionService,
4645 ) -> bool {
4646 let context = coven_protocol::objects::ProtocolObjectContext::circle(
4647 self.root.store_root_hash,
4648 coven_protocol::objects::ProtocolObjectDomain::CircleSnapshotMeta,
4649 encryption,
4650 );
4651 let prefix = coven_protocol::store_commit::circle_snapshot_slot_prefix(
4652 circle_id,
4653 &self.founder.device_id(),
4654 0,
4655 );
4656 let slot = coven_protocol::objects::ObjectSlot::logical(format!("{prefix}.json"))
4657 .expect("valid generation-zero Circle snapshot slot");
4658 self.storage
4659 .read_protocol_slot(&context, &slot, &prefix)
4660 .await
4661 .is_err()
4662 }
4663
4664 #[cfg(test)]
4665 pub fn store_root_hash(&self) -> coven_protocol::store_commit::ObjectHash {
4666 self.root.store_root_hash
4667 }
4668
4669 #[cfg(test)]
4670 pub async fn publish_changeset(
4671 &self,
4672 name: &str,
4673 sequence: u64,
4674 changeset: &[u8],
4675 schema_version: u32,
4676 ) -> Result<coven_protocol::store_commit::StoreBatchCommitRef, TestError> {
4677 self.ensure_producer_registered(name).await?;
4678 let device = {
4679 let producers = self.producers.lock().await;
4680 producers
4681 .by_name
4682 .get(name)
4683 .expect("registered test producer exists")
4684 .clone()
4685 };
4686 device
4687 .publish_changeset_for_test(sequence, changeset.to_vec(), schema_version)
4688 .await
4689 }
4690
4691 #[cfg(test)]
4692 pub async fn publish_founder_changeset(
4693 &self,
4694 changeset: Vec<u8>,
4695 previous_sequence: u64,
4696 ) -> Result<coven_protocol::store_commit::StoreBatchCommitRef, TestError> {
4697 self.founder
4698 .publish_changeset_after_for_test(changeset, previous_sequence)
4699 .await
4700 }
4701}
4702
4703#[cfg(test)]
4706pub fn plaintext_cipher() -> std::sync::RwLock<coven_storage::CloudCipher> {
4707 std::sync::RwLock::new(coven_storage::CloudCipher::Plaintext)
4708}
4709
4710#[cfg(any(test, feature = "test-utils"))]
4712#[derive(Debug, Clone, Copy, PartialEq, Eq)]
4713pub enum ProtocolRead {
4714 Object,
4715 Slot,
4716 PreparedSlot,
4717 Listing,
4722}
4723
4724#[cfg(any(test, feature = "test-utils"))]
4725pub enum ProviderObjectExistsInterception {
4726 Proceed,
4727 DeleteAndReportAbsent,
4728}
4729
4730#[cfg(any(test, feature = "test-utils"))]
4738#[async_trait::async_trait]
4739pub trait StorageInterceptor: Send + Sync {
4740 async fn before_protocol_create(
4741 &self,
4742 _prepared: &coven_protocol::objects::PreparedExactObject,
4743 ) -> Result<(), coven_protocol::objects::StorageError> {
4744 Ok(())
4745 }
4746
4747 async fn before_protocol_read(
4748 &self,
4749 _read: ProtocolRead,
4750 _semantic_prefix: &str,
4751 ) -> Result<(), coven_protocol::objects::StorageError> {
4752 Ok(())
4753 }
4754
4755 fn filter_protocol_slots(
4756 &self,
4757 _listing_prefix: &str,
4758 _slots: &mut Vec<coven_protocol::objects::ObjectSlot>,
4759 ) {
4760 }
4761
4762 async fn before_blob_allocate(&self) -> Result<(), coven_protocol::objects::StorageError> {
4763 Ok(())
4764 }
4765
4766 async fn before_blob_prepare(&self) -> Result<(), coven_protocol::objects::StorageError> {
4767 Ok(())
4768 }
4769
4770 async fn before_blob_create(
4771 &self,
4772 _blob: &coven_protocol::blob::locator::StoredBlobRef,
4773 ) -> Result<(), coven_protocol::objects::StorageError> {
4774 Ok(())
4775 }
4776
4777 async fn before_blob_stage(&self) -> Result<(), coven_protocol::objects::StorageError> {
4778 Ok(())
4779 }
4780
4781 async fn before_provider_object_read(
4782 &self,
4783 _key: &str,
4784 ) -> Result<(), coven_protocol::objects::StorageError> {
4785 Ok(())
4786 }
4787
4788 async fn before_provider_object_write(
4789 &self,
4790 _key: &str,
4791 ) -> Result<(), coven_protocol::objects::StorageError> {
4792 Ok(())
4793 }
4794
4795 async fn before_provider_object_exists(
4796 &self,
4797 _key: &str,
4798 ) -> Result<ProviderObjectExistsInterception, coven_protocol::objects::StorageError> {
4799 Ok(ProviderObjectExistsInterception::Proceed)
4800 }
4801
4802 async fn before_provider_object_delete(
4803 &self,
4804 _key: &str,
4805 ) -> Result<(), coven_protocol::objects::StorageError> {
4806 Ok(())
4807 }
4808}
4809
4810#[cfg(test)]
4811#[async_trait::async_trait]
4812impl<T> StorageInterceptor for std::sync::Arc<T>
4813where
4814 T: StorageInterceptor + ?Sized,
4815{
4816 async fn before_protocol_create(
4817 &self,
4818 prepared: &coven_protocol::objects::PreparedExactObject,
4819 ) -> Result<(), coven_protocol::objects::StorageError> {
4820 (**self).before_protocol_create(prepared).await
4821 }
4822
4823 async fn before_protocol_read(
4824 &self,
4825 read: ProtocolRead,
4826 semantic_prefix: &str,
4827 ) -> Result<(), coven_protocol::objects::StorageError> {
4828 (**self).before_protocol_read(read, semantic_prefix).await
4829 }
4830
4831 fn filter_protocol_slots(
4832 &self,
4833 listing_prefix: &str,
4834 slots: &mut Vec<coven_protocol::objects::ObjectSlot>,
4835 ) {
4836 (**self).filter_protocol_slots(listing_prefix, slots);
4837 }
4838
4839 async fn before_blob_allocate(&self) -> Result<(), coven_protocol::objects::StorageError> {
4840 (**self).before_blob_allocate().await
4841 }
4842
4843 async fn before_blob_prepare(&self) -> Result<(), coven_protocol::objects::StorageError> {
4844 (**self).before_blob_prepare().await
4845 }
4846
4847 async fn before_blob_create(
4848 &self,
4849 blob: &coven_protocol::blob::locator::StoredBlobRef,
4850 ) -> Result<(), coven_protocol::objects::StorageError> {
4851 (**self).before_blob_create(blob).await
4852 }
4853
4854 async fn before_blob_stage(&self) -> Result<(), coven_protocol::objects::StorageError> {
4855 (**self).before_blob_stage().await
4856 }
4857
4858 async fn before_provider_object_read(
4859 &self,
4860 key: &str,
4861 ) -> Result<(), coven_protocol::objects::StorageError> {
4862 (**self).before_provider_object_read(key).await
4863 }
4864
4865 async fn before_provider_object_write(
4866 &self,
4867 key: &str,
4868 ) -> Result<(), coven_protocol::objects::StorageError> {
4869 (**self).before_provider_object_write(key).await
4870 }
4871
4872 async fn before_provider_object_exists(
4873 &self,
4874 key: &str,
4875 ) -> Result<ProviderObjectExistsInterception, coven_protocol::objects::StorageError> {
4876 (**self).before_provider_object_exists(key).await
4877 }
4878
4879 async fn before_provider_object_delete(
4880 &self,
4881 key: &str,
4882 ) -> Result<(), coven_protocol::objects::StorageError> {
4883 (**self).before_provider_object_delete(key).await
4884 }
4885}
4886
4887#[cfg(any(test, feature = "test-utils"))]
4890pub struct InterceptedStorage<S, I: StorageInterceptor>
4891where
4892 S: std::ops::Deref,
4893{
4894 inner: S,
4895 interceptor: I,
4896}
4897
4898#[cfg(any(test, feature = "test-utils"))]
4899impl<S, I> coven_storage::CloudSyncCipherStateAccess for InterceptedStorage<S, I>
4900where
4901 S: std::ops::Deref + Send + Sync,
4902 S::Target: coven_storage::CloudSyncCipherStateAccess,
4903 I: StorageInterceptor,
4904{
4905 fn is_plaintext(&self) -> bool {
4906 self.inner.is_plaintext()
4907 }
4908
4909 fn suffix(&self) -> &'static str {
4910 self.inner.suffix()
4911 }
4912
4913 fn current_generation(&self) -> Option<u64> {
4914 self.inner.current_generation()
4915 }
4916
4917 fn current_fingerprint(&self) -> Option<String> {
4918 self.inner.current_fingerprint()
4919 }
4920
4921 fn open(
4922 &self,
4923 stored: Vec<u8>,
4924 aad_context: &[u8],
4925 ) -> Result<Vec<u8>, coven_keys::encryption::EncryptionError> {
4926 self.inner.open(stored, aad_context)
4927 }
4928
4929 fn seal(&self, plaintext: Vec<u8>, aad_context: &[u8]) -> Vec<u8> {
4930 self.inner.seal(plaintext, aad_context)
4931 }
4932
4933 fn open_sealed_blob_for_test(
4934 &self,
4935 stored: &[u8],
4936 aad_context: &[u8],
4937 ) -> Result<
4938 (coven_keys::encryption::KeyFingerprint, Vec<u8>),
4939 coven_keys::encryption::EncryptionError,
4940 > {
4941 self.inner.open_sealed_blob_for_test(stored, aad_context)
4942 }
4943
4944 fn merged_keyring(
4945 &self,
4946 new_encryption: &coven_keys::encryption::EncryptionService,
4947 ) -> Result<coven_storage::CloudKeyringMerge, coven_keys::encryption::EncryptionError> {
4948 self.inner.merged_keyring(new_encryption)
4949 }
4950
4951 fn merge_key_rotation(
4952 &self,
4953 new_encryption: &coven_keys::encryption::EncryptionService,
4954 custody: &dyn coven_keys::keys::MasterKeyCustody,
4955 ) -> Result<Option<String>, coven_keys::keys::KeyError> {
4956 self.inner.merge_key_rotation(new_encryption, custody)
4957 }
4958}
4959
4960#[cfg(any(test, feature = "test-utils"))]
4961impl<S, I> coven_storage::CloudSyncRotationStateAccess for InterceptedStorage<S, I>
4962where
4963 S: std::ops::Deref + Send + Sync,
4964 S::Target: coven_storage::CloudSyncRotationStateAccess,
4965 I: StorageInterceptor,
4966{
4967 fn mark_candidate(
4968 &self,
4969 generation: u64,
4970 mutation: coven_protocol::store_commit::ObjectHash,
4971 ) -> Result<(), coven_storage::RotationStateError> {
4972 self.inner.mark_candidate(generation, mutation)
4973 }
4974
4975 fn mark_committed_mutation(
4976 &self,
4977 generation: u64,
4978 mutation: coven_protocol::store_commit::ObjectHash,
4979 ) -> Result<(), coven_storage::RotationStateError> {
4980 self.inner.mark_committed_mutation(generation, mutation)
4981 }
4982
4983 fn gate(&self) -> Option<coven_protocol::objects::RotationGate> {
4984 self.inner.gate()
4985 }
4986
4987 fn install_durable_gate(&self, gate: Option<coven_protocol::objects::RotationGate>) {
4988 self.inner.install_durable_gate(gate);
4989 }
4990
4991 fn check(
4992 &self,
4993 live_generation: Option<u64>,
4994 ) -> Result<(), coven_protocol::objects::RotationPending> {
4995 self.inner.check(live_generation)
4996 }
4997}
4998
4999#[cfg(any(test, feature = "test-utils"))]
5000impl<S, I> crate::sync::cycle::CloudSyncCycleConnection for InterceptedStorage<S, I>
5001where
5002 S: std::ops::Deref + Send + Sync,
5003 S::Target: crate::sync::cycle::CloudSyncCycleConnection,
5004 I: StorageInterceptor,
5005{
5006}
5007
5008#[cfg(any(test, feature = "test-utils"))]
5009impl<S, I: StorageInterceptor> InterceptedStorage<S, I>
5010where
5011 S: std::ops::Deref,
5012{
5013 pub fn new(inner: S, interceptor: I) -> Self {
5014 Self { inner, interceptor }
5015 }
5016
5017 pub fn interceptor(&self) -> &I {
5018 &self.interceptor
5019 }
5020}
5021
5022#[cfg(any(test, feature = "test-utils"))]
5023#[async_trait::async_trait]
5024impl<S, I> coven_storage::CloudSyncObjectStorage for InterceptedStorage<S, I>
5025where
5026 S: std::ops::Deref + Send + Sync,
5027 S::Target: coven_storage::CloudSyncObjectStorage + coven_storage::CloudSyncCipherStateAccess,
5028 I: StorageInterceptor,
5029{
5030 fn blob_path_scheme(&self) -> coven_storage::BlobPathScheme {
5031 self.inner.blob_path_scheme()
5032 }
5033
5034 async fn probe_provider(&self) -> Result<(), coven_protocol::objects::StorageError> {
5035 self.inner.probe_provider().await
5036 }
5037
5038 fn provider_requests(
5039 &self,
5040 ) -> Option<std::sync::Arc<dyn coven_foundation::stage_timing::ProviderRequests>> {
5041 self.inner.provider_requests()
5042 }
5043
5044 async fn set_member_access(
5045 &self,
5046 state: coven_storage::cloud::CloudAccessState,
5047 ) -> Result<coven_storage::cloud::CloudAccessOutcome, coven_protocol::objects::StorageError>
5048 {
5049 self.inner.set_member_access(state).await
5050 }
5051
5052 async fn read_blob_tombstone(
5053 &self,
5054 stored: &coven_protocol::blob::locator::StoredBlobRef,
5055 ) -> Result<Option<Vec<u8>>, coven_protocol::objects::StorageError> {
5056 let key = coven_storage::blob_tombstone_key(
5057 stored,
5058 coven_storage::CloudSyncCipherStateAccess::suffix(&*self.inner),
5059 );
5060 self.interceptor.before_provider_object_read(&key).await?;
5061 self.inner.read_blob_tombstone(stored).await
5062 }
5063
5064 async fn write_blob_tombstone(
5065 &self,
5066 stored: &coven_protocol::blob::locator::StoredBlobRef,
5067 plaintext: Vec<u8>,
5068 ) -> Result<(), coven_protocol::objects::StorageError> {
5069 let key = coven_storage::blob_tombstone_key(
5070 stored,
5071 coven_storage::CloudSyncCipherStateAccess::suffix(&*self.inner),
5072 );
5073 self.interceptor.before_provider_object_write(&key).await?;
5074 self.inner.write_blob_tombstone(stored, plaintext).await
5075 }
5076
5077 async fn list_blob_tombstones(
5078 &self,
5079 ) -> Result<Vec<coven_storage::ListedBlobTombstone>, coven_protocol::objects::StorageError>
5080 {
5081 self.inner.list_blob_tombstones().await
5082 }
5083
5084 async fn blob_tombstone_exists(
5085 &self,
5086 stored: &coven_protocol::blob::locator::StoredBlobRef,
5087 ) -> Result<bool, coven_protocol::objects::StorageError> {
5088 let key = coven_storage::blob_tombstone_key(
5089 stored,
5090 coven_storage::CloudSyncCipherStateAccess::suffix(&*self.inner),
5091 );
5092 match self.interceptor.before_provider_object_exists(&key).await? {
5093 ProviderObjectExistsInterception::Proceed => {
5094 self.inner.blob_tombstone_exists(stored).await
5095 }
5096 ProviderObjectExistsInterception::DeleteAndReportAbsent => {
5097 self.inner.delete_blob_tombstone(stored).await?;
5098 Ok(false)
5099 }
5100 }
5101 }
5102
5103 async fn delete_blob_tombstone(
5104 &self,
5105 stored: &coven_protocol::blob::locator::StoredBlobRef,
5106 ) -> Result<(), coven_protocol::objects::StorageError> {
5107 let key = coven_storage::blob_tombstone_key(
5108 stored,
5109 coven_storage::CloudSyncCipherStateAccess::suffix(&*self.inner),
5110 );
5111 self.interceptor.before_provider_object_delete(&key).await?;
5112 self.inner.delete_blob_tombstone(stored).await
5113 }
5114
5115 async fn read_provider_bytes_for_test(
5116 &self,
5117 key: &str,
5118 ) -> Result<Vec<u8>, coven_protocol::objects::StorageError> {
5119 self.interceptor.before_provider_object_read(key).await?;
5120 self.inner.read_provider_bytes_for_test(key).await
5121 }
5122
5123 async fn write_provider_bytes_for_test(
5124 &self,
5125 key: &str,
5126 bytes: Vec<u8>,
5127 ) -> Result<(), coven_protocol::objects::StorageError> {
5128 self.interceptor.before_provider_object_write(key).await?;
5129 self.inner.write_provider_bytes_for_test(key, bytes).await
5130 }
5131
5132 async fn list_provider_keys_for_test(
5133 &self,
5134 prefix: &str,
5135 ) -> Result<Vec<String>, coven_protocol::objects::StorageError> {
5136 self.inner.list_provider_keys_for_test(prefix).await
5137 }
5138
5139 async fn provider_key_exists_for_test(
5140 &self,
5141 key: &str,
5142 ) -> Result<bool, coven_protocol::objects::StorageError> {
5143 match self.interceptor.before_provider_object_exists(key).await? {
5144 ProviderObjectExistsInterception::Proceed => {
5145 self.inner.provider_key_exists_for_test(key).await
5146 }
5147 ProviderObjectExistsInterception::DeleteAndReportAbsent => {
5148 Err(coven_protocol::objects::StorageError::InvalidContent(
5149 "raw test-key interception cannot delete through the production API"
5150 .to_string(),
5151 ))
5152 }
5153 }
5154 }
5155
5156 async fn reserve_cross_principal_response_slot(
5157 &self,
5158 probe_id: coven_protocol::provider::ProviderProbeId,
5159 ) -> Result<coven_protocol::objects::ObjectSlot, coven_protocol::provider::ProviderProbeError>
5160 {
5161 self.inner
5162 .reserve_cross_principal_response_slot(probe_id)
5163 .await
5164 }
5165
5166 async fn prepare_cross_principal_challenge(
5167 &self,
5168 publication_journal: &dyn coven_protocol::provider::DeviceJoinChallengePublicationJournal,
5169 probe_id: coven_protocol::provider::ProviderProbeId,
5170 store: &coven_protocol::StoreProviderBinding,
5171 context: &coven_protocol::provider::CrossPrincipalChallengeContext,
5172 administrator_signer: &dyn coven_keys::keys::DeviceSigningAuthority,
5173 ) -> Result<
5174 coven_protocol::provider::CrossPrincipalProbeChallenge,
5175 coven_protocol::provider::ProviderProbeError,
5176 > {
5177 self.inner
5178 .prepare_cross_principal_challenge(
5179 publication_journal,
5180 probe_id,
5181 store,
5182 context,
5183 administrator_signer,
5184 )
5185 .await
5186 }
5187
5188 async fn settle_cross_principal_challenge(
5189 &self,
5190 publication_journal: &dyn coven_protocol::provider::DeviceJoinChallengePublicationJournal,
5191 authorization: &coven_protocol::provider::DeviceJoinChallengePublicationAuthorization,
5192 challenge: &coven_protocol::provider::CrossPrincipalProbeChallenge,
5193 context: &coven_protocol::provider::CrossPrincipalChallengeContext,
5194 store: &coven_protocol::StoreProviderBinding,
5195 ) -> Result<
5196 coven_protocol::provider::CrossPrincipalProbeChallenge,
5197 coven_protocol::provider::ProviderProbeError,
5198 > {
5199 self.inner
5200 .settle_cross_principal_challenge(
5201 publication_journal,
5202 authorization,
5203 challenge,
5204 context,
5205 store,
5206 )
5207 .await
5208 }
5209
5210 async fn create_cross_principal_response(
5211 &self,
5212 challenge: &coven_protocol::provider::CrossPrincipalProbeChallenge,
5213 context: &coven_protocol::provider::CrossPrincipalResponseContext,
5214 store: &coven_protocol::StoreProviderBinding,
5215 administrator_signing_pubkey: &str,
5216 peer_signer: &coven_keys::keys::UserKeypair,
5217 ) -> Result<
5218 coven_protocol::provider::CrossPrincipalProbeResponse,
5219 coven_protocol::provider::ProviderProbeError,
5220 > {
5221 self.inner
5222 .create_cross_principal_response(
5223 challenge,
5224 context,
5225 store,
5226 administrator_signing_pubkey,
5227 peer_signer,
5228 )
5229 .await
5230 }
5231
5232 async fn complete_cross_principal_probe(
5233 &self,
5234 journal: &dyn coven_protocol::provider::ProviderProbeJournal,
5235 challenge: &coven_protocol::provider::CrossPrincipalProbeChallenge,
5236 response: &coven_protocol::provider::CrossPrincipalProbeResponse,
5237 context: &coven_protocol::provider::CrossPrincipalResponseContext,
5238 store: &coven_protocol::StoreProviderBinding,
5239 administrator_signer: &dyn coven_keys::keys::DeviceSigningAuthority,
5240 peer_signing_pubkey: &str,
5241 ) -> Result<
5242 coven_protocol::provider::CrossPrincipalProbeReceipt,
5243 coven_protocol::provider::ProviderProbeError,
5244 > {
5245 self.inner
5246 .complete_cross_principal_probe(
5247 journal,
5248 challenge,
5249 response,
5250 context,
5251 store,
5252 administrator_signer,
5253 peer_signing_pubkey,
5254 )
5255 .await
5256 }
5257
5258 async fn probe_exact_slots(
5259 &self,
5260 journal: &dyn coven_protocol::provider::ProviderProbeJournal,
5261 probe_id: coven_protocol::provider::ProviderProbeId,
5262 binding: &coven_protocol::objects::ResolvedProviderBinding,
5263 ) -> Result<
5264 coven_protocol::provider::ExactSlotProbeReceipt,
5265 coven_protocol::provider::ProviderProbeError,
5266 > {
5267 self.inner
5268 .probe_exact_slots(journal, probe_id, binding)
5269 .await
5270 }
5271
5272 async fn observe_exact_slot(
5273 &self,
5274 slot: &coven_protocol::objects::ObjectSlot,
5275 ) -> Result<
5276 Option<coven_protocol::objects::ExactObjectRef>,
5277 coven_protocol::objects::StorageError,
5278 > {
5279 self.inner.observe_exact_slot(slot).await
5280 }
5281
5282 async fn delete_exact_slot_and_verify_absent(
5283 &self,
5284 slot: &coven_protocol::objects::ObjectSlot,
5285 ) -> Result<(), coven_protocol::objects::StorageError> {
5286 self.inner.delete_exact_slot_and_verify_absent(slot).await
5287 }
5288
5289 fn store_blob_key_fingerprint(
5290 &self,
5291 ) -> Result<Option<coven_keys::encryption::KeyFingerprint>, coven_protocol::objects::StorageError>
5292 {
5293 self.inner.store_blob_key_fingerprint()
5294 }
5295
5296 fn create_store_key_confirmation(
5297 &self,
5298 creation_id: coven_protocol::store_commit::StoreCreationId,
5299 ) -> Result<
5300 coven_protocol::store_commit::StoreKeyConfirmation,
5301 coven_protocol::objects::StorageError,
5302 > {
5303 self.inner.create_store_key_confirmation(creation_id)
5304 }
5305
5306 fn verify_store_key_confirmation(
5307 &self,
5308 creation_id: coven_protocol::store_commit::StoreCreationId,
5309 confirmation: &coven_protocol::store_commit::StoreKeyConfirmation,
5310 ) -> Result<(), coven_protocol::objects::StorageError> {
5311 self.inner
5312 .verify_store_key_confirmation(creation_id, confirmation)
5313 }
5314
5315 async fn seal_store_blob_to_spool(
5316 &self,
5317 locator: &coven_protocol::blob::locator::BlobLocator,
5318 authority: &coven_protocol::objects::BlobWriteAuthority<'_>,
5319 plaintext_file: &std::path::Path,
5320 spool: coven_foundation::local_file::AtomicStagedFile,
5321 progress: coven_storage::cloud::PreparationProgress,
5322 ) -> Result<coven_protocol::objects::BlobSpoolWrite, coven_protocol::objects::StorageError>
5323 {
5324 self.inner
5325 .seal_store_blob_to_spool(locator, authority, plaintext_file, spool, progress)
5326 .await
5327 }
5328
5329 async fn stage_verified_store_blob_plaintext(
5330 &self,
5331 blob: &coven_protocol::blob::locator::StoredBlobRef,
5332 stage: coven_foundation::local_file::AtomicStagedFile,
5333 progress: coven_storage::cloud::DownloadProgress,
5334 ) -> Result<coven_foundation::local_file::AtomicStagedFile, coven_protocol::objects::StorageError>
5335 {
5336 self.interceptor.before_blob_stage().await?;
5337 self.inner
5338 .stage_verified_store_blob_plaintext(blob, stage, progress)
5339 .await
5340 }
5341
5342 async fn open_store_blob_range_reader(
5343 &self,
5344 blob: &coven_protocol::blob::locator::StoredBlobRef,
5345 ) -> Result<coven_storage::BlobRangeReader, coven_protocol::objects::StorageError> {
5346 self.inner.open_store_blob_range_reader(blob).await
5347 }
5348
5349 async fn provider_binding(
5350 &self,
5351 ) -> Result<
5352 coven_protocol::objects::ResolvedProviderBinding,
5353 coven_protocol::objects::StorageError,
5354 > {
5355 self.inner.provider_binding().await
5356 }
5357
5358 async fn allocate_protocol_slot(
5359 &self,
5360 context: &coven_protocol::objects::ProtocolObjectContext,
5361 semantic_prefix: &str,
5362 extension: &str,
5363 ) -> Result<coven_protocol::objects::ObjectSlot, coven_protocol::objects::StorageError> {
5364 self.inner
5365 .allocate_protocol_slot(context, semantic_prefix, extension)
5366 .await
5367 }
5368
5369 fn prepare_protocol_object(
5370 &self,
5371 context: &coven_protocol::objects::ProtocolObjectContext,
5372 slot: coven_protocol::objects::ObjectSlot,
5373 semantic_prefix: &str,
5374 data: Vec<u8>,
5375 ) -> Result<coven_protocol::objects::PreparedExactObject, coven_protocol::objects::StorageError>
5376 {
5377 self.inner
5378 .prepare_protocol_object(context, slot, semantic_prefix, data)
5379 }
5380
5381 async fn open_prepared_protocol_object(
5382 &self,
5383 context: &coven_protocol::objects::ProtocolObjectContext,
5384 prepared: &coven_protocol::objects::PreparedExactObject,
5385 semantic_prefix: &str,
5386 ) -> Result<Vec<u8>, coven_protocol::objects::StorageError> {
5387 self.inner
5388 .open_prepared_protocol_object(context, prepared, semantic_prefix)
5389 .await
5390 }
5391
5392 async fn create_protocol_object(
5393 &self,
5394 prepared: &coven_protocol::objects::PreparedExactObject,
5395 ) -> Result<(), coven_protocol::objects::StorageError> {
5396 self.interceptor.before_protocol_create(prepared).await?;
5397 self.inner.create_protocol_object(prepared).await
5398 }
5399
5400 async fn create_versioned_protocol_record(
5401 &self,
5402 context: &coven_protocol::objects::ProtocolObjectContext,
5403 prepared: &coven_protocol::objects::PreparedExactObject,
5404 semantic_prefix: &str,
5405 expected: &[u8],
5406 ) -> Result<coven_storage::CloudObjectVersion, coven_protocol::objects::StorageError> {
5407 self.inner
5408 .create_versioned_protocol_record(context, prepared, semantic_prefix, expected)
5409 .await
5410 }
5411
5412 async fn read_protocol_object(
5413 &self,
5414 context: &coven_protocol::objects::ProtocolObjectContext,
5415 object: &coven_protocol::objects::ExactObjectRef,
5416 semantic_prefix: &str,
5417 ) -> Result<Vec<u8>, coven_protocol::objects::StorageError> {
5418 self.interceptor
5419 .before_protocol_read(ProtocolRead::Object, semantic_prefix)
5420 .await?;
5421 self.inner
5422 .read_protocol_object(context, object, semantic_prefix)
5423 .await
5424 }
5425
5426 async fn read_protocol_object_with_progress(
5427 &self,
5428 context: &coven_protocol::objects::ProtocolObjectContext,
5429 object: &coven_protocol::objects::ExactObjectRef,
5430 semantic_prefix: &str,
5431 progress: coven_storage::cloud::DownloadProgress,
5432 ) -> Result<Vec<u8>, coven_protocol::objects::StorageError> {
5433 self.interceptor
5434 .before_protocol_read(ProtocolRead::Object, semantic_prefix)
5435 .await?;
5436 self.inner
5437 .read_protocol_object_with_progress(context, object, semantic_prefix, progress)
5438 .await
5439 }
5440
5441 async fn read_versioned_protocol_record(
5442 &self,
5443 context: &coven_protocol::objects::ProtocolObjectContext,
5444 slot: &coven_protocol::objects::ObjectSlot,
5445 semantic_prefix: &str,
5446 ) -> Result<(Vec<u8>, coven_storage::CloudObjectVersion), coven_protocol::objects::StorageError>
5447 {
5448 self.interceptor
5449 .before_provider_object_read(slot.logical_key())
5450 .await?;
5451 self.inner
5452 .read_versioned_protocol_record(context, slot, semantic_prefix)
5453 .await
5454 }
5455
5456 async fn replace_protocol_record_if_version(
5457 &self,
5458 context: &coven_protocol::objects::ProtocolObjectContext,
5459 slot: &coven_protocol::objects::ObjectSlot,
5460 semantic_prefix: &str,
5461 expected: &coven_storage::CloudObjectVersion,
5462 data: Vec<u8>,
5463 ) -> Result<coven_storage::cloud::ConditionalWriteOutcome, coven_protocol::objects::StorageError>
5464 {
5465 self.interceptor
5466 .before_provider_object_write(slot.logical_key())
5467 .await?;
5468 self.inner
5469 .replace_protocol_record_if_version(context, slot, semantic_prefix, expected, data)
5470 .await
5471 }
5472
5473 async fn list_protocol_slots(
5474 &self,
5475 context: &coven_protocol::objects::ProtocolObjectContext,
5476 listing_prefix: &str,
5477 ) -> Result<Vec<coven_protocol::objects::ObjectSlot>, coven_protocol::objects::StorageError>
5478 {
5479 self.interceptor
5480 .before_protocol_read(ProtocolRead::Listing, listing_prefix)
5481 .await?;
5482 let mut slots = self
5483 .inner
5484 .list_protocol_slots(context, listing_prefix)
5485 .await?;
5486 self.interceptor
5487 .filter_protocol_slots(listing_prefix, &mut slots);
5488 Ok(slots)
5489 }
5490
5491 async fn read_protocol_slot(
5492 &self,
5493 context: &coven_protocol::objects::ProtocolObjectContext,
5494 slot: &coven_protocol::objects::ObjectSlot,
5495 semantic_prefix: &str,
5496 ) -> Result<
5497 (Vec<u8>, coven_protocol::objects::ExactObjectRef),
5498 coven_protocol::objects::StorageError,
5499 > {
5500 self.interceptor
5501 .before_protocol_read(ProtocolRead::Slot, semantic_prefix)
5502 .await?;
5503 self.inner
5504 .read_protocol_slot(context, slot, semantic_prefix)
5505 .await
5506 }
5507
5508 async fn read_prepared_protocol_slot(
5509 &self,
5510 context: &coven_protocol::objects::ProtocolObjectContext,
5511 slot: &coven_protocol::objects::ObjectSlot,
5512 semantic_prefix: &str,
5513 ) -> Result<
5514 (Vec<u8>, coven_protocol::objects::PreparedExactObject),
5515 coven_protocol::objects::StorageError,
5516 > {
5517 self.interceptor
5518 .before_protocol_read(ProtocolRead::PreparedSlot, semantic_prefix)
5519 .await?;
5520 self.inner
5521 .read_prepared_protocol_slot(context, slot, semantic_prefix)
5522 .await
5523 }
5524
5525 async fn delete_protocol_object(
5526 &self,
5527 object: &coven_protocol::objects::ExactObjectRef,
5528 ) -> Result<(), coven_protocol::objects::StorageError> {
5529 self.inner.delete_protocol_object(object).await
5530 }
5531
5532 async fn allocate_blob_slot(
5533 &self,
5534 locator: &coven_protocol::blob::locator::BlobLocator,
5535 authority: &coven_protocol::objects::BlobWriteAuthority<'_>,
5536 ) -> Result<coven_protocol::objects::ObjectSlot, coven_protocol::objects::StorageError> {
5537 self.interceptor.before_blob_allocate().await?;
5538 self.inner.allocate_blob_slot(locator, authority).await
5539 }
5540
5541 async fn seal_blob_to_spool(
5542 &self,
5543 locator: &coven_protocol::blob::locator::BlobLocator,
5544 authority: &coven_protocol::objects::BlobWriteAuthority<'_>,
5545 protection: coven_protocol::objects::BlobSpoolProtection,
5546 plaintext_file: &std::path::Path,
5547 spool: coven_foundation::local_file::AtomicStagedFile,
5548 progress: coven_storage::cloud::PreparationProgress,
5549 ) -> Result<coven_protocol::objects::BlobSpoolWrite, coven_protocol::objects::StorageError>
5550 {
5551 self.inner
5552 .seal_blob_to_spool(
5553 locator,
5554 authority,
5555 protection,
5556 plaintext_file,
5557 spool,
5558 progress,
5559 )
5560 .await
5561 }
5562
5563 async fn prepare_blob_object(
5564 &self,
5565 locator: &coven_protocol::blob::locator::BlobLocator,
5566 authority: &coven_protocol::objects::BlobWriteAuthority<'_>,
5567 slot: coven_protocol::objects::ObjectSlot,
5568 stored_file: &std::path::Path,
5569 ) -> Result<coven_protocol::blob::locator::StoredBlobRef, coven_protocol::objects::StorageError>
5570 {
5571 self.interceptor.before_blob_prepare().await?;
5572 self.inner
5573 .prepare_blob_object(locator, authority, slot, stored_file)
5574 .await
5575 }
5576
5577 async fn create_blob_object_from_file(
5578 &self,
5579 blob: &coven_protocol::blob::locator::StoredBlobRef,
5580 authority: &coven_protocol::objects::BlobWriteAuthority<'_>,
5581 stored_file: &std::path::Path,
5582 control: &coven_storage::cloud::UploadControl,
5583 ) -> Result<(), coven_protocol::objects::StorageError> {
5584 self.interceptor.before_blob_create(blob).await?;
5585 self.inner
5586 .create_blob_object_from_file(blob, authority, stored_file, control)
5587 .await
5588 }
5589
5590 async fn verify_blob_object(
5591 &self,
5592 blob: &coven_protocol::blob::locator::StoredBlobRef,
5593 ) -> Result<(), coven_protocol::objects::StorageError> {
5594 self.inner.verify_blob_object(blob).await
5595 }
5596
5597 async fn stage_verified_blob_plaintext(
5598 &self,
5599 blob: &coven_protocol::blob::locator::StoredBlobRef,
5600 protection: coven_protocol::objects::BlobSpoolProtection,
5601 stage: coven_foundation::local_file::AtomicStagedFile,
5602 progress: coven_storage::cloud::DownloadProgress,
5603 ) -> Result<coven_foundation::local_file::AtomicStagedFile, coven_protocol::objects::StorageError>
5604 {
5605 self.interceptor.before_blob_stage().await?;
5606 self.inner
5607 .stage_verified_blob_plaintext(blob, protection, stage, progress)
5608 .await
5609 }
5610
5611 async fn open_blob_range_reader(
5612 &self,
5613 blob: &coven_protocol::blob::locator::StoredBlobRef,
5614 protection: coven_protocol::objects::BlobSpoolProtection,
5615 ) -> Result<coven_storage::BlobRangeReader, coven_protocol::objects::StorageError> {
5616 self.inner.open_blob_range_reader(blob, protection).await
5617 }
5618
5619 async fn delete_blob_object(
5620 &self,
5621 blob: &coven_protocol::blob::locator::StoredBlobRef,
5622 ) -> Result<(), coven_protocol::objects::StorageError> {
5623 self.inner.delete_blob_object(blob).await
5624 }
5625}