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 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 #[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 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}