1use super::*;
2
3#[derive(Clone)]
7pub struct Database {
8 connection: DatabaseConnection,
9}
10
11fn store_dir_of(path: &Path) -> coven_foundation::store_dir::StoreDir {
17 if path == Path::new(":memory:") {
18 return coven_foundation::store_dir::StoreDir::new_ephemeral(
19 std::env::temp_dir().join(format!("coven-in-memory-store-{}", uuid::Uuid::new_v4())),
20 );
21 }
22 let parent = path
23 .parent()
24 .filter(|parent| !parent.as_os_str().is_empty())
25 .unwrap_or_else(|| Path::new("."));
26 coven_foundation::store_dir::StoreDir::new(parent)
27}
28
29impl Database {
30 pub(crate) fn from_core(core: DatabaseCore, thread_name: &str) -> Result<Self, DbError> {
31 Ok(Self {
32 connection: DatabaseConnection::start(core, thread_name)?,
33 })
34 }
35
36 pub(crate) async fn call_database<F, R>(&self, operation: F) -> Result<R, DbError>
37 where
38 F: for<'session> FnOnce(
39 &mut crate::database_session::DatabaseSession<'session>,
40 ) -> Result<R, DbError>
41 + Send
42 + 'static,
43 R: Send + 'static,
44 {
45 self.connection.call_database(operation).await
46 }
47
48 pub(crate) async fn call_store<F, R>(&self, operation: F) -> Result<R, DbError>
49 where
50 F: for<'session> FnOnce(&mut crate::store::StoreSession<'session>) -> Result<R, DbError>
51 + Send
52 + 'static,
53 R: Send + 'static,
54 {
55 self.connection.call_store(operation).await
56 }
57
58 pub(crate) async fn read_store<F, R, E>(&self, read: F) -> Result<Result<R, E>, DbError>
59 where
60 F: for<'connection> FnOnce(crate::store::SqlReadContext<'connection>) -> Result<R, E>
61 + Send
62 + 'static,
63 R: Send + 'static,
64 E: Send + 'static,
65 {
66 self.connection.read_store(read).await
67 }
68
69 pub(crate) fn store_schema_version(&self) -> u32 {
70 self.connection.store_schema_version()
71 }
72
73 pub(crate) fn store_sync_routing_hash(&self) -> ObjectHash {
74 self.connection.store_sync_routing_hash()
75 }
76
77 pub(crate) fn store_has_synced_tables(&self) -> bool {
78 self.connection.store_has_synced_tables()
79 }
80
81 pub(crate) fn store_blob_transition_root(&self, table_name: &str) -> BlobTransitionRoot {
82 self.connection.store_blob_transition_root(table_name)
83 }
84
85 pub(crate) fn store_transfer_limits(&self) -> coven_protocol::blob::TransferLimits {
86 self.connection.store_transfer_limits()
87 }
88
89 pub(crate) fn set_store_transfer_limits(&self, limits: coven_protocol::blob::TransferLimits) {
90 self.connection.set_store_transfer_limits(limits)
91 }
92
93 pub(crate) fn store_blob_tombstone_grace(&self) -> chrono::Duration {
94 self.connection.store_blob_tombstone_grace()
95 }
96
97 pub(crate) fn store_has_scoped_graph(&self) -> bool {
98 self.connection.store_has_scoped_graph()
99 }
100
101 pub(crate) fn store_stamp(&self) -> String {
102 self.connection.store_stamp()
103 }
104
105 pub(crate) fn store_hlc_high_water(&self) -> String {
106 self.connection.store_hlc_high_water()
107 }
108
109 pub(crate) fn store_blob_ref_from_change(
110 &self,
111 change: &coven_foundation::changeset::RowChange,
112 ) -> Result<Option<coven_protocol::blob::BlobRef>, BlobDeclError> {
113 self.connection.store_blob_ref_from_change(change)
114 }
115
116 pub(crate) fn validate_store_local_blob_cleanup_changes(
117 &self,
118 old_changes: &[coven_foundation::changeset::RowChange],
119 new_changes: &[coven_foundation::changeset::RowChange],
120 ) -> Result<(), BlobDeclError> {
121 self.connection
122 .validate_store_local_blob_cleanup_changes(old_changes, new_changes)
123 }
124
125 pub(crate) fn store_receive_wall_ms(&self) -> u64 {
126 self.connection.store_receive_wall_ms()
127 }
128
129 #[cfg(any(test, feature = "test-utils"))]
130 pub(crate) fn assert_owns_payload_directory_for_test(
131 &self,
132 store_dir: &coven_foundation::store_dir::StoreDir,
133 ) {
134 self.connection
135 .assert_owns_payload_directory_for_test(store_dir);
136 }
137
138 pub(crate) fn new_store_id(&self) -> String {
139 self.connection.new_store_id()
140 }
141
142 pub(crate) fn notify_store_write_status(&self, write_id: WriteId, status: WriteStatus) {
143 self.connection.notify_store_write_status(write_id, status);
144 }
145
146 pub(crate) fn subscribe_store_write_status(
147 &self,
148 write_id: WriteId,
149 current: WriteStatus,
150 ) -> tokio::sync::watch::Receiver<WriteStatus> {
151 self.connection
152 .subscribe_store_write_status(write_id, current)
153 }
154
155 pub(crate) fn subscribe_committed_changes(
156 &self,
157 ) -> tokio::sync::broadcast::Receiver<Arc<crate::CommittedChanges>> {
158 self.connection.subscribe_committed_changes()
159 }
160
161 pub(crate) async fn membership_load_permit(&self) -> crate::store::MembershipLoadPermit {
162 self.connection.membership_load_permit().await
163 }
164
165 pub(crate) async fn membership_mutation_permit(
166 &self,
167 ) -> crate::store::MembershipMutationPermit {
168 self.connection.membership_mutation_permit().await
169 }
170
171 pub(crate) async fn store_creation_permit(&self) -> crate::store::StoreCreationPermit {
172 self.connection.store_creation_permit().await
173 }
174
175 pub(crate) async fn device_exclusion_permit(&self) -> crate::store::DeviceExclusionPermit {
176 self.connection.device_exclusion_permit().await
177 }
178
179 pub(crate) async fn author_own_store_stream(&self) -> crate::store::OwnStreamAuthorship {
180 self.connection.author_own_store_stream().await
181 }
182
183 pub(crate) async fn blob_upload_drain_permit(&self) -> crate::store::BlobUploadDrainPermit {
184 self.connection.blob_upload_drain_permit().await
185 }
186
187 pub(crate) async fn snapshot_publication_permit(
188 &self,
189 ) -> crate::store::SnapshotPublicationPermit {
190 self.connection.snapshot_publication_permit().await
191 }
192
193 pub(crate) async fn local_blob_cleanup_permit(&self) -> crate::store::LocalBlobCleanupPermit {
194 self.connection.local_blob_cleanup_permit().await
195 }
196
197 pub(crate) async fn apply_local_blob_cleanup_intent(
198 &self,
199 intent: &crate::local_blob_cleanup_intents::LocalBlobCleanupIntent,
200 ) -> Result<(), DbError> {
201 self.connection
202 .apply_local_blob_cleanup_intent(intent)
203 .await
204 }
205
206 pub(crate) async fn stage_host_write_blobs<E>(
207 &self,
208 blobs: Vec<crate::store::NewBlob>,
209 ) -> Result<crate::store::StagedBlobBatch, crate::HostWriteError<E>> {
210 self.connection.stage_host_write_blobs(blobs).await
211 }
212
213 pub(crate) async fn sync_store_parent_dir(
214 &self,
215 path: &Path,
216 ) -> Result<(), coven_foundation::atomic_file::FileError> {
217 self.connection.sync_store_parent_dir(path).await
218 }
219
220 #[cfg(any(test, feature = "test-utils"))]
221 pub(crate) async fn reach_store_test_point(&self, point: DatabaseTestPoint) {
222 self.connection.reach_store_test_point(point).await;
223 }
224
225 pub fn open(
234 path: &Path,
235 synced_tables: Vec<SyncedTable>,
236 blob_tombstone_grace: chrono::Duration,
237 transfer_limits: coven_protocol::blob::TransferLimits,
238 device_id: String,
239 clock: coven_foundation::clock::ClockRef,
240 coven_migration_policy: CovenMigrationPolicy,
241 migrations: &[Migration],
242 ) -> Result<Database, OpenError> {
243 let hlc = Hlc::try_new(device_id, clock).map_err(|e| DbError::context("device_id", e))?;
244 Self::open_with_hlc_and_coven_metadata(
245 path,
246 synced_tables,
247 blob_tombstone_grace,
248 transfer_limits,
249 Arc::new(hlc),
250 coven_migration_policy,
251 migrations,
252 CovenMetadataOpen::Detect,
253 )
254 }
255
256 pub fn open_initialized_store(
257 path: &Path,
258 install: &VerifiedSnapshotBootstrapInstall,
259 synced_tables: Vec<SyncedTable>,
260 blob_tombstone_grace: chrono::Duration,
261 transfer_limits: coven_protocol::blob::TransferLimits,
262 device_id: String,
263 clock: coven_foundation::clock::ClockRef,
264 coven_migration_policy: CovenMigrationPolicy,
265 migrations: &[Migration],
266 ) -> Result<Database, OpenError> {
267 let hlc = Hlc::try_new(device_id, clock).map_err(|e| DbError::context("device_id", e))?;
268 Self::open_with_hlc_and_coven_metadata(
269 path,
270 synced_tables,
271 blob_tombstone_grace,
272 transfer_limits,
273 Arc::new(hlc),
274 coven_migration_policy,
275 migrations,
276 CovenMetadataOpen::VerifiedSnapshot(install),
277 )
278 }
279
280 fn open_with_hlc_and_coven_metadata(
281 path: &Path,
282 synced_tables: Vec<SyncedTable>,
283 blob_tombstone_grace: chrono::Duration,
284 transfer_limits: coven_protocol::blob::TransferLimits,
285 hlc: Arc<Hlc>,
286 coven_migration_policy: CovenMigrationPolicy,
287 migrations: &[Migration],
288 metadata_open: CovenMetadataOpen<'_>,
289 ) -> Result<Database, OpenError> {
290 let store_dir = store_dir_of(path);
291 Self::open_with_hlc_and_coven_metadata_in_store_dir(
292 path,
293 store_dir,
294 crate::connection_io::ConnectionDurability::Full,
295 synced_tables,
296 blob_tombstone_grace,
297 transfer_limits,
298 hlc,
299 coven_migration_policy,
300 migrations,
301 metadata_open,
302 )
303 }
304
305 fn open_with_hlc_and_coven_metadata_in_store_dir(
306 path: &Path,
307 store_dir: coven_foundation::store_dir::StoreDir,
308 connection_durability: crate::connection_io::ConnectionDurability,
309 synced_tables: Vec<SyncedTable>,
310 blob_tombstone_grace: chrono::Duration,
311 transfer_limits: coven_protocol::blob::TransferLimits,
312 hlc: Arc<Hlc>,
313 coven_migration_policy: CovenMigrationPolicy,
314 migrations: &[Migration],
315 metadata_open: CovenMetadataOpen<'_>,
316 ) -> Result<Database, OpenError> {
317 let core = DatabaseCore::open(
318 path,
319 store_dir,
320 connection_durability,
321 synced_tables,
322 blob_tombstone_grace,
323 transfer_limits,
324 hlc,
325 coven_migration_policy,
326 migrations,
327 metadata_open,
328 )?;
329
330 Self::from_core(core, "coven-db").map_err(OpenError::from)
331 }
332
333 #[cfg(any(test, feature = "test-utils"))]
334 pub fn open_in_store_dir_for_test(
335 path: &Path,
336 store_dir: coven_foundation::store_dir::StoreDir,
337 synced_tables: Vec<SyncedTable>,
338 blob_tombstone_grace: chrono::Duration,
339 transfer_limits: coven_protocol::blob::TransferLimits,
340 device_id: String,
341 clock: coven_foundation::clock::ClockRef,
342 coven_migration_policy: CovenMigrationPolicy,
343 migrations: &[Migration],
344 ) -> Result<Database, OpenError> {
345 let hlc = Hlc::try_new(device_id, clock).map_err(|e| DbError::context("device_id", e))?;
346 Self::open_with_hlc_in_store_dir_for_test(
347 path,
348 store_dir,
349 synced_tables,
350 blob_tombstone_grace,
351 transfer_limits,
352 Arc::new(hlc),
353 coven_migration_policy,
354 migrations,
355 )
356 }
357
358 #[cfg(any(test, feature = "test-utils"))]
359 pub fn open_with_hlc_in_store_dir_for_test(
360 path: &Path,
361 store_dir: coven_foundation::store_dir::StoreDir,
362 synced_tables: Vec<SyncedTable>,
363 blob_tombstone_grace: chrono::Duration,
364 transfer_limits: coven_protocol::blob::TransferLimits,
365 hlc: Arc<Hlc>,
366 coven_migration_policy: CovenMigrationPolicy,
367 migrations: &[Migration],
368 ) -> Result<Database, OpenError> {
369 Self::open_with_hlc_and_coven_metadata_in_store_dir(
370 path,
371 store_dir,
372 crate::connection_io::ConnectionDurability::Disabled,
373 synced_tables,
374 blob_tombstone_grace,
375 transfer_limits,
376 hlc,
377 coven_migration_policy,
378 migrations,
379 CovenMetadataOpen::Detect,
380 )
381 }
382
383 pub fn open_read_only(
399 path: &Path,
400 synced_tables: Vec<SyncedTable>,
401 blob_tombstone_grace: chrono::Duration,
402 transfer_limits: coven_protocol::blob::TransferLimits,
403 device_id: String,
404 clock: coven_foundation::clock::ClockRef,
405 migrations: &[Migration],
406 ) -> Result<Database, OpenError> {
407 let hlc = Hlc::try_new(device_id, clock).map_err(|e| DbError::context("device_id", e))?;
408 let store_dir = store_dir_of(path);
409 let core = DatabaseCore::open_read_only(
410 path,
411 store_dir,
412 synced_tables,
413 blob_tombstone_grace,
414 transfer_limits,
415 Arc::new(hlc),
416 migrations,
417 )?;
418 Self::from_core(core, "coven-db-ro").map_err(OpenError::from)
419 }
420
421 #[cfg(any(test, feature = "test-utils"))]
422 #[doc(hidden)]
423 pub fn arm_test_pause(
424 &self,
425 point: DatabaseTestPoint,
426 ) -> (Arc<tokio::sync::Notify>, Arc<tokio::sync::Notify>) {
427 self.connection.arm_test_pause(point)
428 }
429
430 #[cfg(any(test, feature = "test-utils"))]
431 #[doc(hidden)]
432 pub fn observe_test_points(&self) -> tokio::sync::mpsc::UnboundedReceiver<DatabaseTestPoint> {
433 self.connection.observe_test_points()
434 }
435
436 #[cfg(any(test, feature = "test-utils"))]
437 #[doc(hidden)]
438 pub fn fail_next_merge_materialization_at(&self, point: MergeMaterializationFailurePoint) {
439 self.connection.fail_next_merge_materialization_at(point);
440 }
441
442 #[cfg(any(test, feature = "test-utils"))]
449 pub fn open_with_hlc(
450 path: &Path,
451 synced_tables: Vec<SyncedTable>,
452 blob_tombstone_grace: chrono::Duration,
453 transfer_limits: coven_protocol::blob::TransferLimits,
454 hlc: Arc<Hlc>,
455 coven_migration_policy: CovenMigrationPolicy,
456 migrations: &[Migration],
457 ) -> Result<Database, OpenError> {
458 Self::open_with_hlc_and_coven_metadata(
459 path,
460 synced_tables,
461 blob_tombstone_grace,
462 transfer_limits,
463 hlc,
464 coven_migration_policy,
465 migrations,
466 CovenMetadataOpen::Detect,
467 )
468 }
469
470 #[cfg(any(test, feature = "test-utils"))]
471 pub fn schema_version(&self) -> u32 {
472 self.connection.store_schema_version()
473 }
474
475 #[cfg(any(test, feature = "test-utils"))]
476 pub fn sync_routing_hash(&self) -> ObjectHash {
477 self.connection.store_sync_routing_hash()
478 }
479
480 #[cfg(any(test, feature = "test-utils"))]
485 pub fn receive_wall_ms(&self) -> u64 {
486 self.connection.store_receive_wall_ms()
487 }
488}
489
490#[cfg(test)]
491mod tests {
492 use super::*;
493
494 #[test]
495 fn a_relative_database_uses_the_working_directory_as_its_store_directory() {
496 assert_eq!(
497 store_dir_of(Path::new("store.sqlite")).as_ref(),
498 Path::new(".")
499 );
500 }
501}