Skip to main content

coven_replication/sync/
test_owner_graph.rs

1use std::sync::Arc;
2
3use crate::blob::transition::{ConnectedBlobTransitions, LocalBlobTransitions};
4use crate::sync::store::blob::{
5    CurrentRemoteBlobSource, LocalStoreBlobAccess, RemoteStoreBlobAccess, StoreBlobCache,
6};
7use coven_database::StoreDatabase;
8use coven_foundation::store_dir::StoreDir;
9use coven_storage::CloudSyncObjectStorage;
10
11#[derive(Clone)]
12pub struct TestOwnerGraph {
13    database: StoreDatabase,
14    store_dir: StoreDir,
15    local_access: LocalStoreBlobAccess,
16    local_transitions: LocalBlobTransitions,
17}
18
19fn blob_owner(database: StoreDatabase, store_dir: StoreDir) -> LocalStoreBlobAccess {
20    let cache = StoreBlobCache::new(database.clone(), store_dir.clone());
21    LocalStoreBlobAccess::new(database, store_dir, cache)
22}
23
24impl TestOwnerGraph {
25    pub fn new(database: StoreDatabase, store_dir: StoreDir) -> Self {
26        database.assert_owns_payload_directory_for_test(&store_dir);
27        let local_access = blob_owner(database.clone(), store_dir.clone());
28        let local_transitions = LocalBlobTransitions::new(database.clone(), store_dir.clone());
29        Self {
30            database,
31            store_dir,
32            local_access,
33            local_transitions,
34        }
35    }
36
37    /// Insert a Local release: a gated-off note plus a blob-bearing photo with an
38    /// external source file registered for it. Returns the external source path.
39    pub async fn seed_local_release(
40        &self,
41        user_dir: &std::path::Path,
42        note_id: &str,
43        photo_id: &str,
44        cloud_path: &str,
45        bytes: &[u8],
46    ) -> std::path::PathBuf {
47        self.database
48            .seed_local_release_rows_for_test(None, note_id, photo_id, cloud_path, bytes)
49            .await;
50        std::fs::create_dir_all(user_dir).expect("create external blob fixture directory");
51        let source = user_dir.join(format!("{photo_id}.jpg"));
52        std::fs::write(&source, bytes).expect("write external blob fixture");
53        self.database
54            .register_external_blob_for_test("note_photos", photo_id, &source)
55            .await;
56        source
57    }
58
59    /// Insert a Remote release: a gated-on note plus a photo whose blob is already
60    /// in cloud storage at the readable path the plaintext scheme derives.
61    #[allow(clippy::too_many_arguments)]
62    pub async fn seed_remote_release(
63        &self,
64        store: &crate::sync::test_helpers::TestStore,
65        routing_encryption: Option<&coven_keys::encryption::EncryptionService>,
66        note_id: &str,
67        photo_id: &str,
68        cloud_path: &str,
69        bytes: &[u8],
70    ) {
71        self.database
72            .seed_local_release_rows_for_test(
73                routing_encryption.cloned(),
74                note_id,
75                photo_id,
76                cloud_path,
77                bytes,
78            )
79            .await;
80        StoreDir::store_local_blob(&self.store_dir, "fixture_sources", photo_id, bytes)
81            .await
82            .expect("write exact remote fixture source");
83        let source = self
84            .store_dir
85            .local_blob_path("fixture_sources", photo_id)
86            .expect("build exact remote fixture source path");
87        self.database
88            .register_external_blob_for_test("note_photos", photo_id, &source)
89            .await;
90        store
91            .open_into_store_database(&self.database, self.store_dir.clone())
92            .await
93            .expect("open exact test Store");
94        self.make_remote("notes", note_id, "Notes Root", false)
95            .await
96            .expect("queue exact remote fixture upload");
97        let outcome = store
98            .drain_uploads(
99                &self.database,
100                &self.store_dir,
101                &coven_foundation::clock::SystemClock,
102                routing_encryption,
103                None,
104            )
105            .await
106            .expect("create exact remote fixture blob");
107        assert_eq!(outcome.uploaded(), 1);
108        assert!(outcome.yielded_for_publish());
109        assert!(
110            store
111                .publish_pending_store_database(&self.database, &self.store_dir)
112                .await
113                .expect("publish exact remote fixture"),
114            "remote fixture publishes its Store write",
115        );
116    }
117
118    pub async fn stage_pending_upload_for_test(
119        &self,
120        source_dir: &std::path::Path,
121        blob_id: &str,
122        bytes: &[u8],
123        created_at: &str,
124    ) {
125        let source = source_dir.join(blob_id);
126        coven_foundation::local_file::AtomicStagedFile::write_for_test(&source, bytes)
127            .await
128            .expect("write upload source");
129        let reference = self
130            .database
131            .row_blob_ref("note_photos", blob_id)
132            .await
133            .expect("load exact Local row blob reference");
134        self.database
135            .enqueue_blob_upload_for_test(
136                "notes",
137                &format!("note-{blob_id}"),
138                &reference,
139                &source,
140                created_at,
141            )
142            .await
143            .expect("enqueue exact Local row upload");
144    }
145
146    pub async fn drain_published_blob_drop_intents(
147        &self,
148        through_sequence: u64,
149    ) -> Result<(), crate::sync::test_helpers::TestError> {
150        Ok(self
151            .local_access
152            .drain_published_blob_drop_intents(through_sequence)
153            .await?)
154    }
155
156    pub async fn make_remote(
157        &self,
158        root_table: &str,
159        root_id: &str,
160        root_label: &str,
161        pin: bool,
162    ) -> Result<(), crate::blob::transition::MakeRemoteError> {
163        let refs = self
164            .database
165            .row_blob_refs_for_root(root_table, root_id)
166            .await?;
167        self.local_transitions
168            .make_remote(root_table, root_id, root_label, pin, refs)
169            .await
170    }
171
172    #[allow(clippy::too_many_arguments)]
173    pub async fn make_local(
174        &self,
175        storage: Arc<dyn CloudSyncObjectStorage>,
176        routing_encryption: Option<coven_keys::encryption::EncryptionService>,
177        observer: Option<Arc<dyn coven_protocol::blob::BlobTransitionObserver>>,
178        root_table: &str,
179        root_id: &str,
180        dest: &std::collections::HashMap<String, std::path::PathBuf>,
181        cancel: &tokio::sync::watch::Receiver<bool>,
182    ) -> Result<(), crate::blob::transition::MakeLocalError> {
183        self.connected_blob_transitions(storage, routing_encryption, observer)
184            .make_local(root_table, root_id, dest, cancel)
185            .await
186    }
187
188    fn remote_blob_access(
189        &self,
190        storage: Arc<dyn CloudSyncObjectStorage>,
191    ) -> RemoteStoreBlobAccess {
192        RemoteStoreBlobAccess::new(
193            self.local_access.clone(),
194            CurrentRemoteBlobSource::current(self.database.clone(), storage),
195        )
196    }
197
198    /// Run the eager cache fill the sync loop runs behind its cycles, which is
199    /// what makes an eager blob local now that a pull downloads nothing.
200    pub async fn fill_eager_cache(
201        &self,
202        storage: Arc<dyn CloudSyncObjectStorage>,
203    ) -> Result<(), std::sync::Arc<crate::sync::EagerCacheFillError>> {
204        let (_cancel, cancel_rx) = tokio::sync::watch::channel(false);
205        let (status, _status_rx) =
206            tokio::sync::watch::channel(crate::sync::EagerCacheFillStatus::Scanning);
207        crate::sync::store::blob::eager_cache::run(
208            &self.database,
209            &self.remote_blob_access(storage),
210            cancel_rx,
211            &status,
212        )
213        .await
214    }
215
216    pub async fn read_blob(
217        &self,
218        storage: Option<Arc<dyn CloudSyncObjectStorage>>,
219        reference: &coven_protocol::blob::RowBlobRef,
220    ) -> Result<Vec<u8>, crate::sync::BlobCacheError> {
221        match storage {
222            Some(storage) => self.remote_blob_access(storage).read(reference).await,
223            None => self.local_access.read(reference).await,
224        }
225    }
226
227    pub async fn open_blob_stream(
228        &self,
229        storage: Option<Arc<dyn CloudSyncObjectStorage>>,
230        reference: &coven_protocol::blob::RowBlobRef,
231    ) -> Result<crate::sync::BlobStream, crate::sync::BlobCacheError> {
232        match storage {
233            Some(storage) => {
234                self.remote_blob_access(storage)
235                    .open_stream(reference)
236                    .await
237            }
238            None => self.local_access.open_stream(reference).await,
239        }
240    }
241
242    pub async fn read_blob_range(
243        &self,
244        storage: Option<Arc<dyn CloudSyncObjectStorage>>,
245        reference: &coven_protocol::blob::RowBlobRef,
246        offset: u64,
247        len: u64,
248    ) -> Result<Vec<u8>, crate::sync::BlobCacheError> {
249        self.open_blob_stream(storage, reference)
250            .await?
251            .read_at(offset, len)
252            .await
253    }
254
255    pub async fn materialize_blob(
256        &self,
257        storage: Option<Arc<dyn CloudSyncObjectStorage>>,
258        reference: &coven_protocol::blob::RowBlobRef,
259    ) -> Result<(), crate::sync::BlobCacheError> {
260        match storage {
261            Some(storage) => {
262                self.remote_blob_access(storage)
263                    .materialize(reference)
264                    .await
265            }
266            None => self.local_access.materialize(reference).await,
267        }
268    }
269
270    pub async fn pin_blobs(
271        &self,
272        storage: Option<Arc<dyn CloudSyncObjectStorage>>,
273        references: &[coven_protocol::blob::RowBlobRef],
274    ) -> Result<(), crate::sync::BlobCacheError> {
275        match storage {
276            Some(storage) => self.remote_blob_access(storage).pin(references).await,
277            None => self.local_access.pin(references).await,
278        }
279    }
280
281    fn connected_blob_transitions(
282        &self,
283        storage: Arc<dyn CloudSyncObjectStorage>,
284        routing_encryption: Option<coven_keys::encryption::EncryptionService>,
285        observer: Option<Arc<dyn coven_protocol::blob::BlobTransitionObserver>>,
286    ) -> ConnectedBlobTransitions {
287        ConnectedBlobTransitions::new(
288            self.local_transitions.clone(),
289            Arc::new(crate::sync::store::blob::RemoteStoreBlobAccess::new(
290                self.local_access.clone(),
291                crate::sync::store::blob::CurrentRemoteBlobSource::current(
292                    self.database.clone(),
293                    storage,
294                ),
295            )),
296            routing_encryption,
297            observer,
298        )
299    }
300
301    pub async fn run_sync_cycle(
302        &self,
303        storage: impl Into<std::sync::Arc<coven_storage::CloudSyncConnection>>,
304        identity: coven_keys::keys::UserKeypair,
305    ) -> Result<crate::sync::cycle::SyncCycleResult, crate::sync::test_helpers::TestError> {
306        let expected_store_root = self.database.local_store_root_ref().await?.ok_or_else(|| {
307            crate::sync::test_helpers::TestError::invariant(
308                "cycle fixture database has no exact Store root",
309            )
310        })?;
311        let components = Box::pin(crate::sync::cycle::PreparedSyncComponents::prepare(
312            self.database.clone(),
313            self.store_dir.clone(),
314            storage,
315            identity,
316            crate::sync::cycle::StoreInitialization::OpenStore {
317                expected_store_root,
318            },
319            None,
320            std::sync::Arc::new(crate::sync::test_helpers::TestCustody::default()),
321        ))
322        .await?;
323        let components = Box::pin(components.initialize(None)).await?;
324        Ok(components
325            .run_cycle(&coven_foundation::clock::SystemClock, None)
326            .await?)
327    }
328}
329
330#[cfg(test)]
331mod tests {
332    use super::*;
333
334    #[test]
335    #[should_panic(expected = "payload directory does not belong to this database")]
336    fn owner_graph_rejects_a_database_payload_directory_mismatch() {
337        let database_store_dir = crate::sync::test_helpers::test_store_dir();
338        let database = crate::sync::test_helpers::open_test_db(database_store_dir);
339        let unrelated_store_dir = crate::sync::test_helpers::test_store_dir();
340
341        TestOwnerGraph::new(StoreDatabase::new(&database), unrelated_store_dir);
342    }
343}