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