Skip to main content

coven_database/
database_runtime.rs

1use super::*;
2
3/// A cloneable handle to one owned database. The connection capability retains
4/// both the worker and its matching database context; this handle has no second
5/// path to either.
6#[derive(Clone)]
7pub struct Database {
8    connection: DatabaseConnection,
9}
10
11/// Bind a database to the store directory that owns its payload files.
12///
13/// Production databases live under their store directory. In-memory databases
14/// have no parent path, so the database opening boundary creates their
15/// process-local directory once and passes that dependency into the core.
16fn 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    /// Open and own the connection at `path`.
226    ///
227    /// Runs the host migration ladder and validates its final sync-routing
228    /// contract in one transaction. A fresh database creates Coven metadata in
229    /// that transaction; an initialized database commits only when the final
230    /// contract exactly matches its pinned bytes. Then seeds the register clock
231    /// from on-disk rows. The `_updated_at` stamper remains inside the database
232    /// boundary and is used by every synced-row write.
233    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    /// Open the store at `path` read-only for a same-store secondary reader
384    /// (e.g. a separate process reading while another holds the writer open).
385    ///
386    /// Distinct from [`Database::open`] in three ways, all so the reader never
387    /// mutates shared state a concurrent writer owns: the connection is
388    /// `SQLITE_OPEN_READONLY`; no migration ladder or bookkeeping DDL runs (it
389    /// opens against the schema the writer left, and refuses one newer than this
390    /// binary knows — the writer's `SchemaTooNew` policy); and it returns no
391    /// stamper, because a reader mints no `_updated_at`. Reads are safe across
392    /// processes because a read-only connection can coexist with the writer and
393    /// observes commits after each read transaction ends.
394    ///
395    /// The caller takes no store open-lock for a read-only open: the exclusive
396    /// advisory lock guards against a second *writer*, and a read-only connection
397    /// cannot write, so multiple readers can coexist with one writer.
398    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    /// Open with a caller-supplied register clock instead of a fresh
443    /// system-wall-clock one. Lets a test inject an [`Hlc`] over a controlled
444    /// wall clock to exercise the skew/restart-seeding guarantees, sharing the
445    /// production open path (migration, seed, session) so the test drives the
446    /// real unit.
447    ///
448    #[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    /// The receiver's current wall-clock millis, read from this database's
481    /// register clock. The pull reads it once and passes it down to bound an
482    /// incoming `_updated_at`'s physical component (a grossly-future stamp must not
483    /// win last-writer-wins or ratchet the clock).
484    #[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}