Skip to main content

coven_replication/sync/
test_helpers.rs

1/// Shared test helpers for sync module tests.
2///
3/// These drive a real [`Database`] over an in-memory connection carrying the
4/// synthetic test schema, so tests exercise the engine through the same path
5/// production does.
6use std::collections::HashMap;
7use std::sync::{Arc, Mutex};
8
9use coven_database::Database;
10use coven_foundation::store_dir::StoreDir;
11use coven_keys::encryption::MasterKeyring;
12use coven_keys::keys::{KeyError, MasterKeyCustody, UserKeypair};
13#[cfg(test)]
14use coven_protocol::store_commit::ObjectHash;
15use coven_storage::CloudSyncObjectStorage;
16
17/// The synthetic store's schema and `Database` constructors, which the database
18/// layer owns and its own tests open directly.
19pub 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/// In-memory [`MasterKeyCustody`] for tests, with a switch to force `persist`
37/// to fail. The switch models a device whose keyring is momentarily
38/// unwritable, so a test can drive a key adoption into its failure path and then
39/// clear the switch to prove the retry converges. Stores the serialized form
40/// (like the real `Keyring` preset), so `stored_key` reflects exactly what a
41/// caller wrote.
42#[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    /// Make the next and every subsequent `persist` fail until cleared.
61    pub fn fail_writes(&self) {
62        self.fail.store(true, std::sync::atomic::Ordering::SeqCst);
63    }
64
65    /// Let `persist` succeed again.
66    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
99/// Copy the file-backed payloads one store directory holds into another.
100///
101/// A store is a directory, not a file: rows name payload files beside the
102/// database, so a test that copies the database with `VACUUM INTO` and opens the
103/// copy has to bring those files along, exactly as a device carries its whole
104/// store directory rather than one file out of it.
105pub 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
128/// Hex-encoded ed25519 public key, as membership entries and the wrapped-key
129/// store identify a member.
130pub fn pubkey_hex(kp: &UserKeypair) -> String {
131    coven_keys::keys::public_key_hex(kp)
132}
133
134/// Ed25519 identity derived from exact test-owned seed bytes.
135pub 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
141/// Grants a Dropbox shared-folder membership to whichever peer account asks —
142/// the provider-side step a cross-principal admission needs before the joining
143/// device can write to the store's namespace.
144pub 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    /// Download and install the newest Store snapshot this joining device may
195    /// install, into `store_dir`, through its own provider access.
196    ///
197    /// The step production runs before it asks for the history published after
198    /// that snapshot. Skipping it leaves the joining device asking for the
199    /// closure back to genesis, which means a package per commit — and a store
200    /// that reclaims has deleted the ones its snapshot restates.
201    #[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    /// Every provider operation this Store's devices have asked for, counted at
275    /// the same boundary a shipped home counts them.
276    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/// Why a test pull did not produce a result. Keeps the three steps a test pull
422/// runs — opening the store, authorizing the writer, running the cycle — apart,
423/// so a test asserting on one of them cannot pass on another.
424#[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        /// One per device, for the device's life — the shape a sync loop has,
521        /// so repeated cycles here cost what repeated cycles cost there.
522        settled: std::sync::Arc<crate::sync::store::SettledCycle>,
523    }
524
525    impl TestDevice {
526        /// Replay this device's retained history and count `table`'s rows.
527        ///
528        /// The device performs the replay with the database it owns rather than
529        /// handing that database out, so a caller checking a replay never gets
530        /// a handle it could write through.
531        /// Answer a read-only query against this device's Store, for a test
532        /// that never holds the database — a device installed by a join or a
533        /// restore is opened by that install, not by its caller.
534        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        /// Join a device the way production does: install the owner's newest
711        /// snapshot, then carry only the history published after it.
712        ///
713        /// [`activate_joined`](Self::activate_joined) pulls the whole history
714        /// into an empty database instead, which leaves the joining device on a
715        /// genesis replay baseline holding — and pinning for replay — every
716        /// commit back to the beginning of the store. A store that reclaims has
717        /// deleted the packages its snapshot restates, so a device joined that
718        /// way cannot pull past the first reclaim it meets.
719        #[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            // The install owns the database it created, so the join has to be
796            // finished and dropped before a device opens it — the same order
797            // production has, where the join and the running device are
798            // separate processes over one file.
799            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        /// Publish one more generation over the current frontier and acknowledge
818        /// it, the way the cadence would.
819        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(&registration);
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(&registration);
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        /// The writer authority every Circle helper below has to take before it
2185        /// can name its operation, with the one error shape they all report.
2186        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        /// Publish an acknowledgement, then run the independent replay
2679        /// baseline stage and report what it retired.
2680        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        /// Publish the statement and leave baseline retirement to its own
2712        /// cycle stage.
2713        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        /// Stand on the snapshot this device has acknowledged, the way the
2736        /// cycle does, and report what it did or why it did nothing.
2737        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        /// Publish an acknowledgement and run the baseline stage that follows
2747        /// it in a cycle.
2748        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        /// Stage this device's acknowledgement of what it has materialized, and
2806        /// fail if there was nothing new to say — the tests that use this are
2807        /// testing what an acknowledgement does, so one has to be staged.
2808        /// [`Self::stage_current_acknowledgement_if_new`] is for the tests about
2809        /// whether one is staged at all.
2810        #[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        /// Pull, then run the eager cache fill the sync loop runs behind its
2942        /// cycles. A pull records what its rows bind and downloads none of it,
2943        /// so this is what makes an eager blob's bytes local.
2944        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/// The Store a completed join installed, for a test that only asks it
2989/// questions.
2990///
2991/// A joining device never opens a database of its own — the snapshot install is
2992/// what creates it — so a fixture cannot hand one back without handing out a
2993/// raw database. It hands back this instead: the questions the join is asserted
2994/// on, and nothing to write through.
2995#[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/// Open the Store a join installed, once the join has closed its own handle.
3018#[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/// Open the device a join installed, over the database the install created.
3038///
3039/// A joining device never opens a database of its own: the snapshot install is
3040/// what creates the file, so the join has to finish and close before anything
3041/// else opens it. This is that second open, and it is the only handle a test
3042/// gets — the fixtures hand back a device, not a database.
3043#[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    /// Plants a typed tombstone through the Store's exact cloud layout while
3160    /// bypassing the signing drain, so deletion tests can exercise rejected
3161    /// signatures and Store identities.
3162    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    /// Pull, then run the eager cache fill the sync loop runs behind its cycles.
3354    ///
3355    /// A pull records what its rows bind and downloads none of it, so a test
3356    /// that wants an eager blob's bytes on disk has to do what the loop does.
3357    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    /// A store whose home keeps blobs **browsable**: stored in the clear under
3514    /// readable paths. The counterpart of [`Self::create`], whose home is opaque
3515    /// (sealed under the store key, hashed paths). The pair is fixed per home,
3516    /// so a test that needs the browsable verification story needs this store.
3517    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        // Counted the way a shipped home is counted, at the one boundary every
3587        // provider call crosses, so a test can assert a settled cycle's budget
3588        // in the same unit the cycle log reports it in.
3589        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    /// Provider operations asked for so far. The unit the cycle log reports in,
3624    /// so a budget written here is the budget read there.
3625    /// Delay every whole-object read and record how many overlap, so a test can
3626    /// assert on the schedule of reads rather than on wall-clock time.
3627    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    /// Publish one object from the bytes and reference that identify it, for
3673    /// tests holding a candidate that carries references rather than uploads.
3674    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(&registration.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    /// Join a device through the production shape: snapshot install first,
4304    /// then only the history published after it. Hands back the database the
4305    /// install created, because the joining device's database is the image.
4306    #[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    /// Run a cross-principal device join end to end and hand back the database
4687    /// the joining device ends up with.
4688    ///
4689    /// The joining device installs the owner's newest snapshot first and then
4690    /// carries only the history published after it, which is the shape
4691    /// production has: a joiner that carried the closure back to genesis would
4692    /// need every package ever written, including the ones reclamation deletes
4693    /// once every device has acknowledged the snapshot restating them.
4694    #[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            // The joining device installs a snapshot, so the Store has to have
4710            // published one — the same precondition production has.
4711            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/// A plaintext cloud cipher — the default for tests that are not exercising
4993/// sealing.
4994#[cfg(test)]
4995pub fn plaintext_cipher() -> std::sync::RwLock<coven_storage::CloudCipher> {
4996    std::sync::RwLock::new(coven_storage::CloudCipher::Plaintext)
4997}
4998
4999/// Which protocol read an interceptor hook is running ahead of.
5000#[cfg(any(test, feature = "test-utils"))]
5001#[derive(Debug, Clone, Copy, PartialEq, Eq)]
5002pub enum ProtocolRead {
5003    Object,
5004    Slot,
5005    PreparedSlot,
5006    /// Naming the slots under a prefix, which fetches no object's bytes. Apart
5007    /// from `Slot` because a reader that lists a prefix and then reads what it
5008    /// found makes one of these and N of those, and a test counting reads or
5009    /// failing the Nth one means the reads.
5010    Listing,
5011}
5012
5013#[cfg(any(test, feature = "test-utils"))]
5014pub enum ProviderObjectExistsInterception {
5015    Proceed,
5016    DeleteAndReportAbsent,
5017}
5018
5019/// Test-side observation of a [`CloudSyncObjectStorage`] call.
5020///
5021/// Every hook runs before the wrapped storage does the work, and returning `Err`
5022/// fails the call without reaching it. All hooks default to doing nothing, so an
5023/// interceptor states only the operations its test is about — which is the point:
5024/// a test that intercepts two reads should not also have to restate the sixteen
5025/// operations it does not care about.
5026#[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/// A [`CloudSyncObjectStorage`] that forwards every call to `inner`, giving `interceptor`
5162/// its chance first.
5163#[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}