eidetica/instance/mod.rs
1//!
2//! Provides the main database structures (`Instance` and `Database`).
3//!
4//! `Instance` manages multiple `Database` instances and interacts with the storage `Database`.
5//! `Database` represents a single, independent history of data entries, analogous to a table or branch.
6
7use std::{
8 collections::HashMap,
9 future::Future,
10 path::PathBuf,
11 pin::Pin,
12 sync::{
13 Arc, Mutex, Weak,
14 atomic::{AtomicU64, Ordering},
15 },
16};
17
18use crate::{
19 Clock, Database, Entry, Result, SystemClock,
20 auth::crypto::{PrivateKey, PublicKey},
21 backend::{BackendImpl, InstanceMetadata, InstanceSecrets, VerificationStatus},
22 entry::ID,
23 snapshot::Snapshot,
24 sync::Sync,
25 user::User,
26};
27#[cfg(all(unix, feature = "service"))]
28use crate::{auth::SigKey, service::client::RemoteConnection};
29use handle_trait::Handle;
30
31pub mod backend;
32pub mod errors;
33pub mod new_user;
34pub mod settings_merge;
35pub mod url;
36
37#[cfg(test)]
38mod tests;
39
40// Re-export main types for easier access
41#[cfg(all(unix, feature = "service"))]
42use backend::RemoteBackend;
43use backend::{Backend, LocalBackend};
44pub use errors::InstanceError;
45pub use new_user::NewUser;
46
47/// Indicates whether an entry write originated locally or from a remote source (e.g., sync).
48///
49/// This distinction allows different callbacks to be triggered based on the write source,
50/// enabling behaviors like "only trigger sync for local writes" or "only update UI for remote writes".
51///
52/// `#[non_exhaustive]` covers Rust source compatibility only: downstream code
53/// that matches on this enum must include a wildcard arm, so a new variant
54/// does not break their builds. It does **not** cover wire compatibility —
55/// `WriteSource` is serialized into service-protocol frames
56/// ([`Notification::DatabaseWrite`](crate::service::protocol::Notification::DatabaseWrite)
57/// carries a `source`), and a peer on an older protocol version fails to
58/// deserialize a variant it does not know. Adding a variant is therefore a
59/// protocol change, not a backward-compatible addition: bump
60/// [`crate::service::protocol::PROTOCOL_VERSION`] (and
61/// [`crate::sync::protocol::PROTOCOL_VERSION`], which guards the sync wire)
62/// and treat old peers as incompatible. Always include a wildcard arm when
63/// matching.
64#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
65#[non_exhaustive]
66pub enum WriteSource {
67 /// Write originated from a local transaction commit
68 Local,
69 /// Write originated from a remote source (e.g., sync, replication)
70 Remote,
71}
72
73/// A cursor-advance notification delivered to a write callback.
74///
75/// A `WriteEvent` carries *no entry payloads* — only the cursor brackets
76/// (`previous_tips` → `post_tips`) and the [`WriteSource`]. Callbacks that
77/// only care that *something* changed (cache invalidation, UI wake-ups)
78/// can act on the event directly without touching the wire or the DAG.
79/// Callbacks that need to enumerate or fetch the new entries call
80/// [`Database::ids_added`](crate::Database::ids_added) with the brackets:
81///
82/// ```rust,no_run
83/// # use eidetica::{instance::WriteEvent, Database, Result};
84/// # async fn example(event: &WriteEvent, db: &Database) -> Result<()> {
85/// // Enumerate the IDs added between the two cursors
86/// let new_ids = db.ids_added(event.previous_tips(), event.post_tips()).await?;
87/// for id in new_ids {
88/// // … fetch entry bodies via db.get_entry(id) if needed …
89/// }
90/// # Ok(()) }
91/// ```
92///
93/// Cursor semantics: `previous_tips` is this callback's frontier *before*
94/// this fire — the user-supplied initial tips on the first fire, then the
95/// preceding fire's `post_tips` on each subsequent fire. The cursor
96/// advances to `post_tips` synchronously *before* the user closure is
97/// awaited, so the next fire is guaranteed to bracket against the latest
98/// observed frontier even if the closure is slow.
99///
100/// Triggered only by settled-state (Verified) writes, but the cursors are
101/// raw DAG frontiers — the bracket can span `Unverified`/`Failed` entries,
102/// which [`Database::ids_added`](crate::Database::ids_added) will
103/// enumerate. See
104/// [`Notification::DatabaseWrite`](crate::service::protocol::Notification::DatabaseWrite)
105/// rustdoc for the full verification contract.
106#[derive(Debug, Clone)]
107pub struct WriteEvent {
108 /// The database state this callback was last delivered at — its
109 /// cursor before this fire, as a canonical [`Snapshot`]. Subsequent
110 /// fires for the same callback will have `previous_tips = this fire's
111 /// post_tips`.
112 previous_tips: Snapshot,
113 /// The database state after this write, as a canonical [`Snapshot`].
114 /// Equal to this callback's cursor *after* the fire. Useful for
115 /// "what's the frontier I'm now caught up to" without an extra read.
116 post_tips: Snapshot,
117 /// Whether this write originated locally or from a remote sync.
118 source: WriteSource,
119}
120
121impl WriteEvent {
122 /// Get the database state at this callback's cursor *before* this fire.
123 ///
124 /// The first fire on a freshly-registered callback returns the
125 /// initial snapshot passed at registration time. Subsequent fires
126 /// return the previous fire's `post_tips`.
127 pub fn previous_tips(&self) -> &Snapshot {
128 &self.previous_tips
129 }
130
131 /// Get the database state at this callback's cursor *after* this fire.
132 ///
133 /// The cursor advances to this value before the callback is awaited,
134 /// so the next fire on the same callback will have
135 /// `previous_tips() == this fire's post_tips()`.
136 pub fn post_tips(&self) -> &Snapshot {
137 &self.post_tips
138 }
139
140 /// The source of this write (local commit or remote sync).
141 pub fn source(&self) -> WriteSource {
142 self.source
143 }
144}
145
146/// Boxed future returned by the internal async callback dispatcher.
147/// The future a callback returns is `'static` — it may borrow the `&WriteEvent`
148/// / `&Database` only for the synchronous prefix of the call, never past the
149/// returned future. Both registration paths already require `Fut: 'static`, so
150/// this is not a new constraint on callers; it lets `spawn_write_callbacks`
151/// invoke the callback *synchronously in cursor-advance order* (running any
152/// synchronous side effect — e.g. the service subscription's frame send — in
153/// canonical order under the tree lock) and then spawn only the returned future.
154pub(crate) type AsyncWriteCallbackFuture =
155 Pin<Box<dyn Future<Output = Result<()>> + Send + 'static>>;
156
157/// Internal async callback function type. The user-facing callback contract
158/// is documented on [`Database::on_write`](crate::Database::on_write).
159pub(crate) type AsyncWriteCallbackFn = Arc<
160 dyn for<'a> Fn(&'a WriteEvent, &'a Database) -> AsyncWriteCallbackFuture
161 + Send
162 + std::marker::Sync,
163>;
164
165/// Opaque identifier for a registered callback. Stable for the life of the registration.
166#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
167pub(crate) struct CallbackId(u64);
168
169/// One per-database callback registration plus the cursor that tracks the
170/// frontier this specific callback has observed.
171///
172/// Each fire reads the cursor, builds a [`WriteEvent`] with the cursor
173/// as `previous_tips`, advances the cursor to the post-write tips, and
174/// then invokes the callback. The cursor mutex is held synchronously
175/// during the read/advance — never across the user callback's `.await`.
176///
177/// Stored in `Vec<Arc<PerDbCallbackEntry>>` on each tree so
178/// `fire_write_callbacks` can snapshot Arcs under the registry mutex
179/// and run the dispatches outside it.
180pub(crate) struct PerDbCallbackEntry {
181 pub(crate) id: CallbackId,
182 /// Cursor — the post-write [`Snapshot`] of the most recent event this
183 /// callback has been delivered, or the user-provided initial snapshot
184 /// from the registration call before any event has fired.
185 pub(crate) last_tips: std::sync::Mutex<Snapshot>,
186 pub(crate) callback: AsyncWriteCallbackFn,
187}
188
189/// Type alias for the per-database callback list on a tree.
190type PerDbCallbackVec = Vec<Arc<PerDbCallbackEntry>>;
191
192/// Type alias for the global write callback list. Globals fire for every
193/// write on every tree (used by sync today). No cursor — globals are
194/// tree-agnostic and don't have a meaningful per-tree frontier to track,
195/// so they continue to use whatever `previous_tips` the caller passes in.
196type GlobalCallbackVec = Vec<(CallbackId, AsyncWriteCallbackFn)>;
197
198/// Handle to a registered write callback. **Drop to unregister.**
199///
200/// Returned by [`Database::on_write`](crate::Database::on_write). While this
201/// value is alive the callback fires on writes; dropping it removes the
202/// registration. Use [`detach`](Self::detach) to keep the callback registered
203/// for the life of the [`Instance`] when you don't want to manage the lifetime
204/// yourself.
205///
206/// Holds a weak reference to the [`Instance`], so a `WriteCallback` will not
207/// keep the Instance alive on its own.
208#[must_use = "dropping a WriteCallback unregisters it; call .detach() to keep the callback registered"]
209pub struct WriteCallback {
210 instance: WeakInstance,
211 tree_id: ID,
212 id: CallbackId,
213 detached: bool,
214}
215
216impl WriteCallback {
217 pub(crate) fn new_per_database(instance: WeakInstance, tree_id: ID, id: CallbackId) -> Self {
218 Self {
219 instance,
220 tree_id,
221 id,
222 detached: false,
223 }
224 }
225
226 /// Consume the handle without unregistering. The callback remains active
227 /// for the life of the [`Instance`].
228 ///
229 /// Implementation note: this sets a flag rather than calling `mem::forget`
230 /// so that field destructors (the `WeakInstance`'s weak count, the
231 /// `tree_id`'s heap allocation) still run — only our `Drop` impl is
232 /// short-circuited.
233 pub fn detach(mut self) {
234 self.detached = true;
235 }
236}
237
238impl std::fmt::Debug for WriteCallback {
239 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
240 f.debug_struct("WriteCallback")
241 .field("id", &self.id)
242 .field("tree_id", &self.tree_id)
243 .field("detached", &self.detached)
244 .finish()
245 }
246}
247
248impl Drop for WriteCallback {
249 fn drop(&mut self) {
250 if self.detached {
251 return;
252 }
253 if let Some(instance) = self.instance.upgrade() {
254 // On a connected instance, dropping the last local callback
255 // for this tree transitions the wire-side subscription to
256 // `Idle`. Daemon-side stays subscribed through a grace
257 // window so a quick re-registration is a no-op; the sweep
258 // task in the connection unsubscribes after the window
259 // elapses.
260 //
261 // The local-instance build doesn't read the `was_last`
262 // signal; bind it under the `cfg` so the inactive build
263 // doesn't carry a dead variable.
264 #[cfg(all(unix, feature = "service"))]
265 {
266 let was_last = instance.remove_write_callback(&self.tree_id, self.id);
267 if was_last && let Some(conn) = instance.remote_connection() {
268 conn.transition_to_idle(&self.tree_id);
269 }
270 }
271 #[cfg(not(all(unix, feature = "service")))]
272 {
273 let _ = instance.remove_write_callback(&self.tree_id, self.id);
274 }
275 }
276 }
277}
278
279/// Internal state for Instance
280///
281/// This structure holds the actual implementation data for Instance.
282/// Instance itself is just a cheap-to-clone handle wrapping Arc<InstanceInternal>.
283pub(crate) struct InstanceInternal {
284 /// The database storage backend
285 backend: Arc<dyn Backend>,
286 /// Time provider for timestamps
287 clock: Arc<dyn Clock>,
288 /// Synchronization module for this database instance
289 /// TODO: Overengineered, Sync can be created by default but disabled
290 sync: std::sync::OnceLock<Arc<Sync>>,
291 /// Public instance metadata (device identity, system database IDs)
292 metadata: InstanceMetadata,
293 /// Private instance secrets (None for remote instances without key access)
294 secrets: Option<InstanceSecrets>,
295 /// JSON snapshot file path for an in-memory backend constructed via
296 /// `memory:///path.json` (or set explicitly through
297 /// [`Instance::snapshot_to_path`]). [`Instance::flush`] and the
298 /// [`Drop`] safety net write through this. `None` on any non-snapshot
299 /// backend.
300 ///
301 /// The mutex serves double duty: it guards the path slot itself (so
302 /// `set_snapshot_path` doesn't race with readers) AND serializes the
303 /// actual write so concurrent callers from `flush` / `snapshot_to_path`
304 /// / `Drop` don't race on the shared `<path>.tmp` staging file in
305 /// [`InMemory::save_to_file`]. Held across sync I/O only — never
306 /// across an `.await`. Poison-tolerant: a panic mid-write leaves the
307 /// on-disk snapshot unchanged but must not strand the [`Instance`].
308 snapshot_path: Mutex<Option<PathBuf>>,
309 /// Per-database callbacks keyed by tree_id. Each entry carries its own
310 /// cursor (`last_tips`) so fires can build a callback-specific
311 /// `previous_tips` regardless of when the callback registered or when
312 /// the most recent fire actually advanced its frontier. Consumers
313 /// branch on [`WriteEvent::source`] if they only care about one
314 /// source.
315 write_callbacks: Mutex<HashMap<ID, PerDbCallbackVec>>,
316 /// Global callbacks fired for every write across every database.
317 /// Tree-agnostic — no per-callback cursor.
318 global_write_callbacks: Mutex<GlobalCallbackVec>,
319 /// Monotonic id source for [`CallbackId`].
320 next_callback_id: AtomicU64,
321 /// Per-tree async locks serializing the
322 /// `snapshot` → backend write → callback dispatch sequence so
323 /// `WriteEvent::previous_tips` is consistent for concurrent writers
324 /// to the same tree.
325 tree_locks: Mutex<HashMap<ID, Arc<tokio::sync::Mutex<()>>>>,
326}
327
328impl std::fmt::Debug for InstanceInternal {
329 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
330 f.debug_struct("InstanceInternal")
331 .field("backend", &"<BackendDB>")
332 .field("clock", &self.clock)
333 .field("sync", &self.sync)
334 .field("metadata", &self.metadata)
335 .field("secrets", &self.secrets.is_some())
336 .field(
337 "write_callbacks",
338 &format!(
339 "<{} per-db callbacks>",
340 self.write_callbacks
341 .lock()
342 .unwrap_or_else(|p| p.into_inner())
343 .len()
344 ),
345 )
346 .field(
347 "global_write_callbacks",
348 &format!(
349 "<{} global callbacks>",
350 self.global_write_callbacks
351 .lock()
352 .unwrap_or_else(|p| p.into_inner())
353 .len()
354 ),
355 )
356 .field(
357 "next_callback_id",
358 &self.next_callback_id.load(Ordering::Relaxed),
359 )
360 .finish()
361 }
362}
363
364impl InstanceInternal {
365 /// Synchronously write a JSON snapshot of the underlying backend to `path`.
366 ///
367 /// Returns [`InstanceError::SnapshotNotSupported`] for any backend other
368 /// than the local in-memory backend. Shared by [`Instance::snapshot_to_path`],
369 /// [`Instance::flush`], and the [`Drop`] fallback so the three can't drift.
370 ///
371 /// **Caller must hold the [`snapshot_path`](Self::snapshot_path) mutex.**
372 /// That lock serializes the write — without it, concurrent callers
373 /// would race on the shared `<path>.tmp` staging file in
374 /// [`InMemory::save_to_file`]. The critical section is fully sync; no
375 /// `.await` happens while the lock is held.
376 fn save_snapshot_locked(&self, path: &std::path::Path) -> Result<()> {
377 use crate::backend::database::InMemory;
378 let engine = self
379 .backend
380 .local_engine()
381 .ok_or(InstanceError::SnapshotNotSupported)?;
382 let in_memory = engine
383 .as_any()
384 .downcast_ref::<InMemory>()
385 .ok_or(InstanceError::SnapshotNotSupported)?;
386 in_memory.save_to_file(path)
387 }
388}
389
390/// Best-effort snapshot save on the *last* `InstanceInternal` drop.
391///
392/// Fires when the `Arc<InstanceInternal>` reaches refcount 0 and a snapshot
393/// path is armed (i.e. the `Instance` was constructed via a
394/// `memory:///path.json` URL). [`Instance::flush`] does **not** clear the
395/// snapshot path, so Drop fires even after a successful `flush()` — the
396/// write is idempotent (same atomic tmp+rename), so the worst case is one
397/// extra write of unchanged JSON.
398///
399/// **Errors are logged via `tracing::error!`, not surfaced** — `Drop` can't
400/// return a `Result` and panicking would be worse than logging. Apps that
401/// care about snapshot durability should call [`Instance::flush`] at
402/// well-defined checkpoints and inspect its `Result`; Drop is a safety net,
403/// not the primary persistence path. If `flush()` failed with a permanent
404/// error (e.g. nonexistent parent directory), Drop will fail the same way
405/// and emit a second log line — accept this redundancy as the cost of a
406/// best-effort fallback.
407///
408/// **Blocking I/O warning:** the snapshot write is synchronous
409/// (`std::fs::write` + `rename`). If the `Instance` is dropped on a tokio
410/// worker thread, this blocks that worker for the duration of the write —
411/// negligible for small snapshots, but pathological for very large ones.
412/// Prefer `flush().await` (which still blocks briefly, but does so under
413/// explicit caller control).
414impl Drop for InstanceInternal {
415 fn drop(&mut self) {
416 // Drop runs at Arc refcount 0, so no other handle can race here —
417 // any in-flight `flush()` future holds `&self` and thus an Arc
418 // clone, which would have prevented Drop from firing. We still
419 // acquire the lock for the write so the locking discipline in
420 // `save_snapshot_locked`'s doc-comment holds uniformly. The lock
421 // is uncontended at this point.
422 let mut guard = match self.snapshot_path.lock() {
423 Ok(g) => g,
424 Err(p) => p.into_inner(),
425 };
426 let Some(path) = guard.take() else { return };
427
428 if let Err(e) = self.save_snapshot_locked(&path) {
429 tracing::error!(
430 snapshot_path = %path.display(),
431 error = %e,
432 "Drop: snapshot save failed. Call `Instance::flush().await` at \
433 checkpoints to inspect the error via Result; Drop is a safety net only.",
434 );
435 }
436 }
437}
438/// Database implementation on top of the storage backend.
439///
440/// Instance manages infrastructure only:
441/// - Backend storage and device identity
442/// - System databases (_users, _databases, _sync)
443/// - User account management (create, login, list)
444///
445/// All database creation and key operations happen through User after login.
446///
447/// Instance is a cheap-to-clone handle around `Arc<InstanceInternal>`.
448///
449/// ## Example
450///
451/// ```
452/// # use eidetica::{Instance, NewUser, crdt::Doc};
453/// # #[tokio::main]
454/// # async fn main() -> eidetica::Result<()> {
455/// // Bootstrap a fresh instance with an initial admin user. The first user
456/// // created on an instance is automatically granted Admin on the system
457/// // databases.
458/// let (instance, maybe_user) = Instance::connect_or_create(
459/// "memory://",
460/// NewUser::passwordless("alice"),
461/// ).await?;
462/// let mut user = maybe_user.expect("memory:// is always fresh");
463///
464/// // Use User API for operations
465/// let mut settings = Doc::new();
466/// settings.set("name", "my_database");
467/// let default_key = user.get_default_key()?;
468/// let db = user.create_database(settings, &default_key).await?;
469/// # Ok(())
470/// # }
471/// ```
472#[derive(Clone, Debug, Handle)]
473pub struct Instance {
474 inner: Arc<InstanceInternal>,
475}
476
477/// Weak reference to an Instance.
478///
479/// This is a weak handle that does not prevent the Instance from being dropped.
480/// Dependent objects (Database, Sync, BackgroundSync) hold weak references to avoid
481/// circular reference cycles that would leak memory.
482///
483/// Use `upgrade()` to convert to a strong `Instance` reference.
484#[derive(Clone, Debug, Handle)]
485pub struct WeakInstance {
486 inner: Weak<InstanceInternal>,
487}
488
489impl Instance {
490 /// Open a connection to an eidetica instance described by a connection URL.
491 ///
492 /// Strict load: returns [`InstanceError::NotInitialized`] when the URL
493 /// points at an embedded backend (`sqlite://`, `postgres://`, `memory://`)
494 /// that has no eidetica metadata yet. Use
495 /// [`Instance::connect_or_create`] to bootstrap an embedded backend on
496 /// first run.
497 ///
498 /// Supported URL schemes:
499 /// - `sqlite://./app.db` — embedded sqlite backend; URL is passed through
500 /// to `sqlx::sqlite`, so any sqlx-accepted form works
501 /// (`?mode=rwc&journal_mode=WAL` etc.).
502 /// - `postgres://user:pwd@host/db` — embedded postgres backend; URL is
503 /// passed through to `sqlx::postgres`.
504 /// - `unix:///run/eidetica/sock` — thin client to a running daemon.
505 /// - `memory://` — empty in-memory backend. Strict load against an
506 /// empty in-memory backend always errors `NotInitialized`; use
507 /// `connect_or_create` for a fresh in-memory instance.
508 /// - `memory:///path/to/snap.json` — in-memory backend with a JSON
509 /// snapshot file (load-on-start; snapshot writes via
510 /// [`Instance::flush`] / [`Instance::snapshot_to_path`] / Drop fallback).
511 ///
512 /// See [`crate::instance::url`] for the full URL grammar.
513 ///
514 /// # Example
515 ///
516 /// `connect()` only succeeds against an already-initialised backend.
517 /// The two-phase pattern below bootstraps once, then re-opens with
518 /// the strict load:
519 ///
520 /// ```
521 /// # #[tokio::main]
522 /// # async fn main() -> eidetica::Result<()> {
523 /// use eidetica::{Instance, NewUser};
524 ///
525 /// let temp = tempfile::tempdir()?;
526 /// let snapshot = temp.path().join("app.json");
527 /// let url = format!("memory://{}", snapshot.display());
528 ///
529 /// // First run: bootstrap and flush a snapshot to disk.
530 /// {
531 /// let (instance, maybe_user) =
532 /// Instance::connect_or_create(&url, NewUser::passwordless("alice")).await?;
533 /// let _user = maybe_user.expect("fresh bootstrap on first run");
534 /// instance.flush()?;
535 /// }
536 ///
537 /// // Later: strict connect against the persisted snapshot.
538 /// let instance = Instance::connect(&url).await?;
539 /// let _user = instance.login_user("alice", None).await?;
540 /// # Ok(())
541 /// # }
542 /// ```
543 ///
544 /// Calling `connect()` on a backend with no eidetica metadata returns
545 /// [`InstanceError::NotInitialized`]; reach for
546 /// [`Instance::connect_or_create`] when first-run bootstrap is part
547 /// of the expected lifecycle.
548 pub async fn connect(url: impl AsRef<str>) -> Result<Self> {
549 Self::connect_impl(url.as_ref(), Arc::new(SystemClock)).await
550 }
551
552 /// Open or initialise an eidetica instance described by a connection URL.
553 ///
554 /// On the load arm: identical to [`Instance::connect`]; `initial` is
555 /// silently ignored and the second tuple element is `None`.
556 ///
557 /// On the bootstrap arm: initialises the backend at the URL with the
558 /// supplied [`NewUser`] as the first admin and returns
559 /// `(Instance, Some(User))`. Only embedded backends
560 /// (sqlite/postgres/memory) ever take the bootstrap arm — `unix://` URLs
561 /// degrade to `connect` (the daemon owns its own initialisation), so the
562 /// returned `Option<User>` is always `None` for `unix://`.
563 ///
564 /// # Example
565 /// ```
566 /// # use eidetica::{Instance, NewUser};
567 /// # #[tokio::main]
568 /// # async fn main() -> eidetica::Result<()> {
569 /// let (instance, maybe_user) = Instance::connect_or_create(
570 /// "memory://",
571 /// NewUser::passwordless("alice"),
572 /// ).await?;
573 /// let mut user = match maybe_user {
574 /// Some(u) => u,
575 /// None => instance.login_user("alice", None).await?,
576 /// };
577 /// # let _ = user.get_default_key()?;
578 /// # Ok(())
579 /// # }
580 /// ```
581 pub async fn connect_or_create(
582 url: impl AsRef<str>,
583 initial: NewUser,
584 ) -> Result<(Self, Option<User>)> {
585 Self::connect_or_create_impl(url.as_ref(), initial, Arc::new(SystemClock)).await
586 }
587
588 /// Escape hatch: open or initialise an eidetica instance against a
589 /// pre-built [`BackendImpl`] (sqlite, postgres, in-memory, or custom).
590 ///
591 /// Same load-or-bootstrap semantics as [`Instance::connect_or_create`]
592 /// but skips URL parsing. Useful for tests, embedded apps that want to
593 /// configure the backend's pool/runtime manually, or backends not yet
594 /// exposed via a URL scheme.
595 pub async fn connect_or_create_backend(
596 backend: Box<dyn BackendImpl>,
597 initial: NewUser,
598 ) -> Result<(Self, Option<User>)> {
599 Self::connect_or_create_backend_impl(
600 Arc::from(backend),
601 initial,
602 Arc::new(SystemClock),
603 None,
604 )
605 .await
606 }
607
608 /// Strict-load escape hatch: open an eidetica instance against a
609 /// pre-built [`BackendImpl`] that's already been initialised. Mirrors
610 /// [`Instance::connect`]'s strict semantics for the URL-less case.
611 ///
612 /// Errors with [`InstanceError::NotInitialized`] if the backend has no
613 /// instance metadata; use [`Instance::connect_or_create_backend`] when you want
614 /// to bootstrap on an empty backend.
615 pub async fn open_backend(backend: Box<dyn BackendImpl>) -> Result<Self> {
616 Self::open_impl(backend, Arc::new(SystemClock)).await
617 }
618
619 /// Test variant of [`Instance::open_backend`] with an injectable clock.
620 ///
621 /// Gated behind the `testing` feature and not for production use: an
622 /// instance whose clock a caller controls can be made to write entries with
623 /// arbitrary timestamps.
624 #[cfg(any(test, feature = "testing"))]
625 pub async fn open_backend_with_clock(
626 backend: Box<dyn BackendImpl>,
627 clock: Arc<dyn Clock>,
628 ) -> Result<Self> {
629 Self::open_impl(backend, clock).await
630 }
631
632 /// Strict-create escape hatch: initialise an eidetica instance on a
633 /// fresh pre-built [`BackendImpl`] and bootstrap an initial admin user.
634 ///
635 /// Errors with [`InstanceError::InstanceAlreadyExists`] if the backend
636 /// is already initialised; use [`Instance::connect_or_create_backend`] when the
637 /// caller doesn't want to choose between load and create up front.
638 pub async fn create_backend(
639 backend: Box<dyn BackendImpl>,
640 initial: NewUser,
641 ) -> Result<(Self, User)> {
642 Self::create_backend_impl(backend, initial, Arc::new(SystemClock)).await
643 }
644
645 /// Test variant of [`Instance::create_backend`] with an injectable clock.
646 ///
647 /// Arg order: backend, clock, initial — clock goes in the middle so
648 /// migrating from the prior `create_with_clock` is a pure rename.
649 ///
650 /// Gated behind the `testing` feature and not for production use: an
651 /// instance whose clock a caller controls can be made to write entries with
652 /// arbitrary timestamps.
653 #[cfg(any(test, feature = "testing"))]
654 pub async fn create_backend_with_clock(
655 backend: Box<dyn BackendImpl>,
656 clock: Arc<dyn Clock>,
657 initial: NewUser,
658 ) -> Result<(Self, User)> {
659 Self::create_backend_impl(backend, initial, clock).await
660 }
661
662 async fn create_backend_impl(
663 backend: Box<dyn BackendImpl>,
664 initial: NewUser,
665 clock: Arc<dyn Clock>,
666 ) -> Result<(Self, User)> {
667 let backend: Arc<dyn BackendImpl> = Arc::from(backend);
668 if backend.get_instance_metadata().await?.is_some() {
669 return Err(InstanceError::InstanceAlreadyExists.into());
670 }
671 Self::create_internal(backend, clock, initial).await
672 }
673
674 // Clock injection is exposed only through the pre-built-backend
675 // variants ([`open_backend_with_clock`] and [`create_backend_with_clock`]).
676 // The URL-based `connect_*` constructors deliberately have no
677 // `_with_clock` siblings: every existing test that needs deterministic
678 // timestamps already builds an `InMemory` backend directly, so a URL-
679 // shaped clock entry point would be dead weight.
680
681 // ============ Internal URL dispatchers ============
682
683 async fn connect_impl(url: &str, clock: Arc<dyn Clock>) -> Result<Self> {
684 let parsed = url::parse(url)?;
685 match parsed {
686 url::ConnectionUrl::Sqlite { url } => Self::connect_sqlite(&url, clock).await,
687 url::ConnectionUrl::Postgres { url } => Self::connect_postgres(&url, clock).await,
688 url::ConnectionUrl::Unix { socket_path } => {
689 Self::connect_unix_socket(socket_path, clock).await
690 }
691 url::ConnectionUrl::Memory { snapshot_path } => {
692 Self::connect_memory(snapshot_path, clock).await
693 }
694 }
695 }
696
697 async fn connect_or_create_impl(
698 url: &str,
699 initial: NewUser,
700 clock: Arc<dyn Clock>,
701 ) -> Result<(Self, Option<User>)> {
702 let parsed = url::parse(url)?;
703 match parsed {
704 url::ConnectionUrl::Sqlite { url } => {
705 let backend = open_sqlite_backend(&url).await?;
706 Self::connect_or_create_backend_impl(Arc::from(backend), initial, clock, None).await
707 }
708 url::ConnectionUrl::Postgres { url } => {
709 let backend = open_postgres_backend(&url).await?;
710 Self::connect_or_create_backend_impl(Arc::from(backend), initial, clock, None).await
711 }
712 url::ConnectionUrl::Unix { socket_path } => {
713 // Daemons own their own initialisation. connect_or_create
714 // against `unix://` degrades to a plain connect; `initial`
715 // is unused on this arm. Log it so the silent drop is
716 // discoverable when debugging "why didn't my initial user
717 // get created?" on a remote URL.
718 tracing::debug!(
719 socket_path = %socket_path.display(),
720 username = %initial.username,
721 "connect_or_create against `unix://` is degrading to `connect`; \
722 `initial` is ignored — daemons own their own initialisation. \
723 Run `eidetica daemon init` to bootstrap a daemon-side instance."
724 );
725 let instance = Self::connect_unix_socket(socket_path, clock).await?;
726 Ok((instance, None))
727 }
728 url::ConnectionUrl::Memory { snapshot_path } => {
729 use crate::backend::database::InMemory;
730 // Build the backend from a single `try_load_from_file` call.
731 // `Ok(None)` means the file didn't exist at read time (the
732 // bootstrap-friendly "first run" case → empty backend).
733 // `Ok(Some(loaded))` means the file existed and parsed; if
734 // it carries no instance metadata it's foreign data that
735 // happened to satisfy the `SerializableDatabase` shape, and
736 // we refuse to bootstrap over it (the next snapshot would
737 // silently overwrite the caller's file). Doing the existence
738 // test and the read in one call removes the TOCTOU window a
739 // separate `path.exists()` check would open.
740 let backend: Box<dyn BackendImpl> = match snapshot_path.as_deref() {
741 None => Box::new(InMemory::new()),
742 Some(path) => {
743 let loaded = InMemory::try_load_from_file(path).await.map_err(|e| {
744 InstanceError::InvalidSnapshot {
745 path: path.to_path_buf(),
746 reason: e.to_string(),
747 }
748 })?;
749 match loaded {
750 None => Box::new(InMemory::new()),
751 Some(loaded) => {
752 let boxed: Box<dyn BackendImpl> = Box::new(loaded);
753 if boxed.get_instance_metadata().await?.is_none() {
754 return Err(InstanceError::InvalidSnapshot {
755 path: path.to_path_buf(),
756 reason: "snapshot file exists but contains no instance \
757 metadata; refusing to bootstrap on top of foreign \
758 data. Delete or move the file to create a fresh \
759 instance at this path."
760 .into(),
761 }
762 .into());
763 }
764 boxed
765 }
766 }
767 }
768 };
769 Self::connect_or_create_backend_impl(
770 Arc::from(backend),
771 initial,
772 clock,
773 snapshot_path,
774 )
775 .await
776 }
777 }
778 }
779
780 /// Internal: load-or-bootstrap against a pre-built backend, optionally
781 /// remembering a snapshot path so Drop / flush can write to it.
782 async fn connect_or_create_backend_impl(
783 backend: Arc<dyn BackendImpl>,
784 initial: NewUser,
785 clock: Arc<dyn Clock>,
786 snapshot_path: Option<PathBuf>,
787 ) -> Result<(Self, Option<User>)> {
788 if let Some(metadata) = backend.get_instance_metadata().await? {
789 let instance = Self::open_impl_arc_with_metadata(backend, clock, metadata).await?;
790 instance.set_snapshot_path(snapshot_path);
791 Ok((instance, None))
792 } else {
793 let (instance, user) = Self::create_internal(backend, clock, initial).await?;
794 instance.set_snapshot_path(snapshot_path);
795 Ok((instance, Some(user)))
796 }
797 }
798
799 // ============ Backend connection helpers ============
800
801 #[cfg(all(unix, feature = "service"))]
802 async fn connect_unix_socket(socket_path: PathBuf, clock: Arc<dyn Clock>) -> Result<Self> {
803 let conn = crate::service::client::RemoteConnection::connect(&socket_path).await?;
804 // Keep a clone for the post-construction `attach_instance` call:
805 // the reader task spawned inside `RemoteConnection::connect` needs a
806 // `WeakInstance` to route push notifications into, and the Instance
807 // doesn't exist yet here.
808 let conn_for_attach = conn.clone();
809 let backend: Arc<dyn Backend> = Arc::new(RemoteBackend::new(conn, None));
810
811 // Load metadata from the remote backend
812 let metadata = backend
813 .get_instance_metadata()
814 .await?
815 .ok_or(InstanceError::DeviceKeyNotFound)?;
816
817 // No local secrets — keys are held server-side after login.
818 let inner = Arc::new(InstanceInternal {
819 backend,
820 clock,
821 sync: std::sync::OnceLock::new(),
822 metadata,
823 secrets: None,
824 snapshot_path: Mutex::new(None),
825 write_callbacks: Mutex::new(HashMap::new()),
826 global_write_callbacks: Mutex::new(Vec::new()),
827 next_callback_id: AtomicU64::new(0),
828 tree_locks: Mutex::new(HashMap::new()),
829 });
830 let instance = Self { inner };
831 // Hand the reader task a Weak reference so it can dispatch
832 // server-pushed `Notification::DatabaseWrite` frames into this
833 // Instance's callback registry. Must run *after* the Instance is
834 // constructed; until then the reader task drops notifications
835 // (none can arrive because the client only subscribes lazily on
836 // the first `Database::on_write` call).
837 conn_for_attach.attach_instance(instance.downgrade());
838 Ok(instance)
839 }
840
841 #[cfg(not(all(unix, feature = "service")))]
842 async fn connect_unix_socket(_socket_path: PathBuf, _clock: Arc<dyn Clock>) -> Result<Self> {
843 Err(InstanceError::BackendUnavailable {
844 scheme: "unix",
845 missing_feature: "service",
846 }
847 .into())
848 }
849
850 async fn connect_sqlite(url: &str, clock: Arc<dyn Clock>) -> Result<Self> {
851 let backend = open_sqlite_backend(url).await?;
852 Self::open_impl(backend, clock).await
853 }
854
855 async fn connect_postgres(url: &str, clock: Arc<dyn Clock>) -> Result<Self> {
856 let backend = open_postgres_backend(url).await?;
857 Self::open_impl(backend, clock).await
858 }
859
860 async fn connect_memory(snapshot_path: Option<PathBuf>, clock: Arc<dyn Clock>) -> Result<Self> {
861 use crate::backend::database::InMemory;
862 // Strict load: a snapshot URL that points at a non-existent file
863 // cannot satisfy `connect`'s "must already be initialised" contract.
864 // `try_load_from_file` returns `Ok(None)` when the file doesn't
865 // exist, which we translate into a pointed `InvalidSnapshot`. Using
866 // the same call for the existence test and the read removes the
867 // TOCTOU window a separate `path.exists()` check would open: if the
868 // file vanishes mid-call, the underlying `read_to_string` surfaces
869 // `NotFound` and lands us in the same `None` arm.
870 let backend: Box<dyn BackendImpl> = match snapshot_path.as_deref() {
871 None => Box::new(InMemory::new()),
872 Some(path) => {
873 let loaded = InMemory::try_load_from_file(path).await.map_err(|e| {
874 InstanceError::InvalidSnapshot {
875 path: path.to_path_buf(),
876 reason: e.to_string(),
877 }
878 })?;
879 match loaded {
880 Some(loaded) => Box::new(loaded),
881 None => {
882 return Err(InstanceError::InvalidSnapshot {
883 path: path.to_path_buf(),
884 reason: "snapshot file does not exist; \
885 use `Instance::connect_or_create` to bootstrap a new instance \
886 at this path, or pass `memory://` for an ephemeral instance"
887 .into(),
888 }
889 .into());
890 }
891 }
892 }
893 };
894 let instance = Self::open_impl(backend, clock).await?;
895 instance.set_snapshot_path(snapshot_path);
896 Ok(instance)
897 }
898
899 /// Flush deferred persistence state to disk.
900 ///
901 /// For an `Instance` constructed via a `memory:///path.json` URL, this
902 /// writes the current backend state to the snapshot path (atomic on
903 /// POSIX — `<path>.tmp` then rename). For sqlite/postgres/unix
904 /// backends this is a no-op; those storage layers handle persistence
905 /// inline.
906 ///
907 /// Idempotent and reentrant — call it as often as you like at
908 /// well-defined checkpoints. The snapshot path stays armed, so the
909 /// [`Drop`] fallback continues to fire on the last handle as a safety
910 /// net. The `Instance` (and any clones) remain fully usable after
911 /// `flush()` returns; this is not a shutdown.
912 ///
913 /// If `flush()` fails (e.g. nonexistent parent directory), the error
914 /// surfaces in the `Result`. Drop will later try the same write and
915 /// fail the same way, logging via `tracing::error!`. The duplicate
916 /// signal is intentional — Drop must report what it sees.
917 ///
918 /// **Blocking I/O note:** the snapshot write is synchronous
919 /// (`std::fs::write` + `rename`) and runs inline on the caller. Hence
920 /// the sync signature — there is no `.await` inside. If you're calling
921 /// from a tokio task, this briefly blocks the runtime worker;
922 /// negligible for small snapshots.
923 pub fn flush(&self) -> Result<()> {
924 // Acquire the snapshot_path lock once and hold it across the
925 // write — the lock both gates the path slot and serializes the
926 // sync I/O so concurrent flushes don't clobber each other's
927 // staging tempfile. The path stays armed (we read, don't take)
928 // so subsequent flushes and the Drop safety net keep working.
929 let guard = self
930 .inner
931 .snapshot_path
932 .lock()
933 .unwrap_or_else(|p| p.into_inner());
934 if let Some(path) = guard.as_deref() {
935 self.inner.save_snapshot_locked(path)?;
936 }
937 Ok(())
938 }
939
940 /// Write a JSON snapshot of the in-memory backend to `path`.
941 ///
942 /// The write goes to `<path>.tmp` and then renames into place. On POSIX
943 /// the rename is atomic; on Windows it is not atomic when the
944 /// destination already exists. Returns
945 /// [`InstanceError::SnapshotNotSupported`] on any backend other than
946 /// the in-memory backend.
947 pub fn snapshot_to_path(&self, path: impl AsRef<std::path::Path>) -> Result<()> {
948 let _guard = self
949 .inner
950 .snapshot_path
951 .lock()
952 .unwrap_or_else(|p| p.into_inner());
953 self.inner.save_snapshot_locked(path.as_ref())
954 }
955
956 /// Stash the snapshot path on the InstanceInternal so Drop / close can
957 /// find it. Only meaningful for in-memory backends — no-op for others.
958 fn set_snapshot_path(&self, path: Option<PathBuf>) {
959 if path.is_none() {
960 return;
961 }
962 // Poison-tolerant: a panic in another holder must not strand the
963 // Instance — the snapshot path is a simple swappable Option.
964 let mut guard = self
965 .inner
966 .snapshot_path
967 .lock()
968 .unwrap_or_else(|p| p.into_inner());
969 *guard = path;
970 }
971
972 /// Internal load-only implementation that works with any clock.
973 async fn open_impl(backend: Box<dyn BackendImpl>, clock: Arc<dyn Clock>) -> Result<Self> {
974 let backend: Arc<dyn BackendImpl> = Arc::from(backend);
975
976 // Strict: require existing InstanceMetadata. Initialisation is the
977 // caller's responsibility (`connect_or_create` / `connect_or_create_backend`).
978 let metadata = backend
979 .get_instance_metadata()
980 .await?
981 .ok_or(InstanceError::NotInitialized)?;
982
983 // Load secrets (contains the private key)
984 let secrets = backend.get_instance_secrets().await?;
985
986 // If secrets are present, verify they match the metadata
987 if let Some(ref secrets) = secrets {
988 let derived_id = secrets.signing_key.public_key();
989 if derived_id != metadata.id {
990 return Err(InstanceError::DeviceKeyMismatch.into());
991 }
992 }
993
994 // Existing backend: load from metadata + secrets
995 let inner = Arc::new(InstanceInternal {
996 backend: Arc::new(LocalBackend::new(backend)),
997 clock,
998 sync: std::sync::OnceLock::new(),
999 metadata,
1000 secrets,
1001 snapshot_path: Mutex::new(None),
1002 write_callbacks: Mutex::new(HashMap::new()),
1003 global_write_callbacks: Mutex::new(Vec::new()),
1004 next_callback_id: AtomicU64::new(0),
1005 tree_locks: Mutex::new(HashMap::new()),
1006 });
1007 Ok(Self { inner })
1008 }
1009
1010 /// Load-only helper that accepts an already-arc'd backend and the
1011 /// already-fetched metadata. Used by `connect_or_create_backend_impl`,
1012 /// which has already inspected metadata to choose between the load
1013 /// and bootstrap arms — passing it through avoids a redundant
1014 /// `get_instance_metadata` round-trip.
1015 async fn open_impl_arc_with_metadata(
1016 backend: Arc<dyn BackendImpl>,
1017 clock: Arc<dyn Clock>,
1018 metadata: InstanceMetadata,
1019 ) -> Result<Self> {
1020 let secrets = backend.get_instance_secrets().await?;
1021 if let Some(ref secrets) = secrets {
1022 let derived_id = secrets.signing_key.public_key();
1023 if derived_id != metadata.id {
1024 return Err(InstanceError::DeviceKeyMismatch.into());
1025 }
1026 }
1027 let inner = Arc::new(InstanceInternal {
1028 backend: Arc::new(LocalBackend::new(backend)),
1029 clock,
1030 sync: std::sync::OnceLock::new(),
1031 metadata,
1032 secrets,
1033 snapshot_path: Mutex::new(None),
1034 write_callbacks: Mutex::new(HashMap::new()),
1035 global_write_callbacks: Mutex::new(Vec::new()),
1036 next_callback_id: AtomicU64::new(0),
1037 tree_locks: Mutex::new(HashMap::new()),
1038 });
1039 Ok(Self { inner })
1040 }
1041
1042 /// Internal create implementation. Returns the new `Instance` along with
1043 /// the just-bootstrapped initial `User`, materialised directly from the
1044 /// keys we generated (no redundant login round-trip).
1045 pub(crate) async fn create_internal(
1046 backend: Arc<dyn BackendImpl>,
1047 clock: Arc<dyn Clock>,
1048 initial: NewUser,
1049 ) -> Result<(Self, User)> {
1050 use crate::user::system_databases::{create_databases_tracking, create_users_database};
1051
1052 // 1. Generate device key
1053 let device_key = PrivateKey::generate();
1054 let device_id = device_key.public_key();
1055
1056 // 2. Create system databases with device_key passed directly
1057 // Create a temporary Instance for database creation (databases will store full IDs later)
1058 //
1059 // SAFETY: The temporary instance has empty users_db_id and databases_db_id placeholders.
1060 // This is safe because:
1061 // 1. We only use it to create new system databases via Database::create()
1062 // 2. Database::create() doesn't access the instance's system database IDs
1063 // 3. The system databases don't exist yet, so their IDs can't be referenced
1064 // 4. The temporary instance is only used during initial setup and discarded
1065 // 5. The real instance is constructed afterward with the correct database IDs
1066 let temp_instance = Self {
1067 inner: Arc::new(InstanceInternal {
1068 backend: Arc::new(LocalBackend::new(Arc::clone(&backend))),
1069 clock: Arc::clone(&clock),
1070 sync: std::sync::OnceLock::new(),
1071 metadata: InstanceMetadata {
1072 id: device_id.clone(),
1073 users_db: ID::default(), // Placeholder - system DBs don't exist yet
1074 databases_db: ID::default(), // Placeholder - system DBs don't exist yet
1075 sync_db: None,
1076 },
1077 secrets: Some(InstanceSecrets {
1078 signing_key: device_key.clone(),
1079 }),
1080 snapshot_path: Mutex::new(None),
1081 write_callbacks: Mutex::new(HashMap::new()),
1082 global_write_callbacks: Mutex::new(Vec::new()),
1083 next_callback_id: AtomicU64::new(0),
1084 tree_locks: Mutex::new(HashMap::new()),
1085 }),
1086 };
1087 let users_db = create_users_database(&temp_instance, &device_key).await?;
1088 let databases_db = create_databases_tracking(&temp_instance, &device_key).await?;
1089
1090 // 3. Save metadata and secrets (marks instance as initialized)
1091 // NB: Ordering matters. Secrets are stored first, then Metadata.
1092 // The presence of the Metadata indicates the instance is fully initialized.
1093 let secrets = InstanceSecrets {
1094 signing_key: device_key,
1095 };
1096 backend.set_instance_secrets(&secrets).await?;
1097
1098 let metadata = InstanceMetadata {
1099 id: device_id,
1100 users_db: users_db.root_id().clone(),
1101 databases_db: databases_db.root_id().clone(),
1102 sync_db: None,
1103 };
1104 backend.set_instance_metadata(&metadata).await?;
1105
1106 // 4. Build real instance
1107 let inner = Arc::new(InstanceInternal {
1108 backend: Arc::new(LocalBackend::new(backend)),
1109 clock,
1110 sync: std::sync::OnceLock::new(),
1111 metadata,
1112 secrets: Some(secrets),
1113 snapshot_path: Mutex::new(None),
1114 write_callbacks: Mutex::new(HashMap::new()),
1115 global_write_callbacks: Mutex::new(Vec::new()),
1116 next_callback_id: AtomicU64::new(0),
1117 tree_locks: Mutex::new(HashMap::new()),
1118 });
1119
1120 let instance = Self { inner };
1121
1122 // 5. Bootstrap the initial user. The first user created on an
1123 // instance is automatically promoted to Admin on the system
1124 // databases by `system_databases::create_user`'s
1125 // first-user-becomes-admin logic.
1126 let users_db = instance.users_db().await?;
1127 let (user_uuid, user_info, root_key) = crate::user::system_databases::create_user(
1128 &users_db,
1129 &instance,
1130 &initial.username,
1131 initial.password.as_deref(),
1132 )
1133 .await?;
1134
1135 // 6. Materialise the User session directly from the keys we just
1136 // generated — skips a redundant `login_user` round-trip that would
1137 // otherwise re-derive the encryption key from the password.
1138 let user = crate::user::system_databases::build_user_session(
1139 &instance,
1140 &user_uuid,
1141 &user_info,
1142 root_key,
1143 initial.password.as_deref(),
1144 )
1145 .await?;
1146
1147 Ok((instance, user))
1148 }
1149
1150 /// Get a reference to the backend seam.
1151 pub fn backend(&self) -> &Arc<dyn Backend> {
1152 &self.inner.backend
1153 }
1154
1155 /// The concrete in-process storage engine, or [`OperationNotSupported`] on
1156 /// a remote instance.
1157 ///
1158 /// Off-seam local-only operations (instance secrets, verification-status
1159 /// mutation, `all_roots`/`get_tree` raw dumps, scope-keyed cache) are
1160 /// performed through this accessor, so they are reachable only where a
1161 /// concrete local backend exists.
1162 ///
1163 /// [`OperationNotSupported`]: InstanceError::OperationNotSupported
1164 pub(crate) fn require_local_engine(&self) -> Result<Arc<dyn BackendImpl>> {
1165 self.inner.backend.local_engine().ok_or_else(|| {
1166 InstanceError::OperationNotSupported {
1167 operation: "local backend engine on remote instance".to_string(),
1168 }
1169 .into()
1170 })
1171 }
1172
1173 /// The remote connection backing this instance, if it was created via
1174 /// [`connect`](Self::connect). Returns `None` for local instances.
1175 ///
1176 /// Useful for constructing a [`Database`](crate::Database) that routes
1177 /// reads through the Database-level wire API while sharing the same
1178 /// connection and session as the instance's write path.
1179 #[cfg(all(unix, feature = "service"))]
1180 pub fn remote_connection(&self) -> Option<RemoteConnection> {
1181 self.inner.backend.remote_connection()
1182 }
1183
1184 /// Check if an entry exists in storage.
1185 pub async fn has_entry(&self, id: &ID) -> bool {
1186 self.inner.backend.get(id).await.is_ok()
1187 }
1188
1189 /// Check if a database is present locally.
1190 ///
1191 /// This differs from `has_entry` in that it checks for the active tracking
1192 /// of the database by the Instance. This method checks if we're tracking
1193 /// the database's tip state.
1194 pub async fn has_database(&self, root_id: &ID) -> bool {
1195 match self.inner.backend.snapshot(root_id).await {
1196 Ok(snap) => !snap.is_empty(),
1197 Err(_) => false,
1198 }
1199 }
1200
1201 /// Get a reference to the clock.
1202 ///
1203 /// The clock is used for timestamps in height calculations and peer tracking.
1204 pub(crate) fn clock(&self) -> &dyn Clock {
1205 &*self.inner.clock
1206 }
1207
1208 /// Get a cloned Arc of the clock.
1209 ///
1210 /// Used when passing the clock to components that need ownership (e.g., HeightCalculator).
1211 pub(crate) fn clock_arc(&self) -> Arc<dyn Clock> {
1212 self.inner.clock.clone()
1213 }
1214
1215 // === Backend pass-through methods (pub(crate) for internal use) ===
1216
1217 /// Get an entry from the backend
1218 pub(crate) async fn get(&self, id: &crate::entry::ID) -> Result<crate::entry::Entry> {
1219 self.inner.backend.get(id).await
1220 }
1221
1222 /// Put an entry into the backend. Always stored Unverified — see
1223 /// [`crate::backend::BackendImpl::put`].
1224 pub(crate) async fn put(&self, entry: crate::entry::Entry) -> Result<()> {
1225 self.inner.backend.put(entry).await
1226 }
1227
1228 /// Returns the current [`crate::Snapshot`] of `tree` — its DAG tips. See
1229 /// [`Database::snapshot`] for the public entry point.
1230 pub(crate) async fn snapshot(
1231 &self,
1232 tree: &crate::entry::ID,
1233 ) -> Result<crate::snapshot::Snapshot> {
1234 self.inner.backend.snapshot(tree).await
1235 }
1236
1237 // === System database accessors ===
1238
1239 /// Get the _users database
1240 ///
1241 /// This constructs a Database instance on-the-fly to avoid circular references.
1242 /// On a local instance the device signing key is attached so users-table
1243 /// writes (e.g., the local `create_user` path) can sign. On a remote
1244 /// instance the device key lives on the daemon side and isn't available
1245 /// locally, so no key is attached — the returned handle is read-only.
1246 /// Write paths on a remote instance must instead go through
1247 /// [`Instance::users_db_for_session`], which attaches the caller's
1248 /// session signing key (e.g. admin's key on the `InstanceAdmin`
1249 /// `create_user` path) and routes through `Database::open_remote`.
1250 pub(crate) async fn users_db(&self) -> Result<Database> {
1251 let db = Database::open(self, &self.inner.metadata.users_db).await?;
1252 #[cfg(all(unix, feature = "service"))]
1253 if self.remote_connection().is_some() {
1254 return Ok(db);
1255 }
1256 Ok(db.with_key(self.signing_key()?.clone()))
1257 }
1258
1259 /// Open the _users system database with a specific signing key (not the device
1260 /// key). Used by the admin-session paths
1261 /// ([`InstanceAdmin`](crate::user::InstanceAdmin), `User::admin_check`) on
1262 /// remote instances where the device key is unavailable.
1263 pub(crate) async fn users_db_for_session(&self, signing_key: &PrivateKey) -> Result<Database> {
1264 self.open_system_db_for_session(&self.inner.metadata.users_db, signing_key)
1265 .await
1266 }
1267
1268 /// Open a system database for an authenticated session.
1269 ///
1270 /// On a remote instance this routes every read through the connection's
1271 /// Database wire protocol ([`Database::open_remote`], a per-handle
1272 /// `RemoteBackend`), gated by the session key's identity — the plain
1273 /// [`Database::open`] path instead clones the instance's session backend,
1274 /// so on a connected instance its reads carry the connection's login
1275 /// identity. On a local instance it opens against the local backend as
1276 /// before. The signing key is attached for writes.
1277 pub(crate) async fn open_system_db_for_session(
1278 &self,
1279 root_id: &ID,
1280 signing_key: &PrivateKey,
1281 ) -> Result<Database> {
1282 #[cfg(all(unix, feature = "service"))]
1283 if let Some(conn) = self.remote_connection() {
1284 // The daemon gates per-tree reads against the acting pubkey
1285 // from the request's identity hint, and the hint here is
1286 // `signing_key.public_key()` (the caller's chosen identity for
1287 // this DB). The hint must be in the connection's session keyset
1288 // — register it now so subsequent reads through the returned
1289 // `RemoteBackend` are accepted.
1290 conn.register_session_key(signing_key).await?;
1291 let identity = SigKey::from_pubkey(&signing_key.public_key());
1292 return Ok(Database::open_remote(self, conn, root_id, identity)
1293 .await?
1294 .with_key(signing_key.clone()));
1295 }
1296 Ok(Database::open(self, root_id)
1297 .await?
1298 .with_key(signing_key.clone()))
1299 }
1300
1301 /// Get the _databases tracking database
1302 ///
1303 /// Parallel to `users_db()` — opens the instance's database-registry
1304 /// system DB with the device signing key attached. Used by the
1305 /// instance-admin bootstrap path (`system_databases::create_user`) to add
1306 /// the first user's pubkey as `Admin(0)` on the registry, so subsequent
1307 /// admin-gated instance ops (e.g., `SetInstanceMetadata`) can authorize
1308 /// against the user's key instead of the device key.
1309 pub(crate) async fn databases_db(&self) -> Result<Database> {
1310 Ok(Database::open(self, &self.inner.metadata.databases_db)
1311 .await?
1312 .with_key(self.signing_key()?.clone()))
1313 }
1314
1315 /// Root id of the `_databases` system DB.
1316 ///
1317 /// The service daemon uses this to gate admin-only ops
1318 /// (e.g., `SetInstanceMetadata`) against `_databases.auth_settings`:
1319 /// an instance admin is a user with `Admin` on `_databases`.
1320 pub(crate) fn databases_db_id(&self) -> &ID {
1321 &self.inner.metadata.databases_db
1322 }
1323
1324 /// Root id of the `_users` system DB.
1325 ///
1326 /// Parallel to `databases_db_id()`. Lets an instance admin open `_users`
1327 /// keyed by their own signing key (rather than the device key that
1328 /// `users_db()` attaches), so admin-gated edits to `_users.auth_settings`
1329 /// resolve against the admin's identity.
1330 pub(crate) fn users_db_id(&self) -> &ID {
1331 &self.inner.metadata.users_db
1332 }
1333
1334 // === User Management ===
1335
1336 /// Login a user with flexible password handling.
1337 ///
1338 /// Returns a User session object that provides access to user operations.
1339 /// For password-protected users, provide the password. For passwordless users, pass None.
1340 ///
1341 /// # Arguments
1342 /// * `user_id` - User identifier (username)
1343 /// * `password` - Optional password. None for passwordless users.
1344 ///
1345 /// # Returns
1346 /// A Result containing the User session
1347 pub async fn login_user(&self, user_id: &str, password: Option<&str>) -> Result<User> {
1348 // On a remote instance, the `TrustedLogin*` handshake authenticates the
1349 // socket connection AND ships back the user's full `UserInfo` plus the
1350 // decrypted root signing key. Build the `User` session from those —
1351 // the per-tree gate means a freshly-logged-in user with no permissions
1352 // on `_users` couldn't re-read it over the wire anyway, so we don't try.
1353 #[cfg(all(unix, feature = "service"))]
1354 if let Some(conn) = self.remote_connection() {
1355 let (user_uuid, user_info, signing_key) = conn.trusted_login(user_id, password).await?;
1356 return crate::user::system_databases::build_user_session(
1357 self,
1358 &user_uuid,
1359 &user_info,
1360 signing_key,
1361 password,
1362 )
1363 .await;
1364 }
1365
1366 use crate::user::system_databases::login_user;
1367 let users_db = self.users_db().await?;
1368 login_user(&users_db, self, user_id, password).await
1369 }
1370
1371 // === User-Sync Integration ===
1372
1373 // === Device Identity Management ===
1374 //
1375 // The Instance's public identity is stored in InstanceMetadata, and the private
1376 // signing key is stored in InstanceSecrets. Both are cached in memory.
1377
1378 /// Get the device signing key.
1379 ///
1380 /// # Internal Use Only
1381 ///
1382 /// This method provides direct access to the instance's cryptographic identity
1383 /// and is intended for internal operations that require the device key (sync,
1384 /// system database creation, authentication validation, etc.).
1385 ///
1386 /// These operations should only be performed by the server/instance administrator,
1387 /// but we don't verify that yet. Future versions may add admin permission checks.
1388 ///
1389 /// Similar to `Database::open` (without a key), this is a controlled escape hatch
1390 /// for internal library operations. Use with care - prefer User API for normal operations.
1391 ///
1392 /// Returns an error if this is a remote Instance that does not have access to the
1393 /// device key (e.g., connected via RPC where secrets are never transmitted).
1394 #[cfg(not(any(test, feature = "testing")))]
1395 pub(crate) fn signing_key(&self) -> Result<&PrivateKey> {
1396 self.inner
1397 .secrets
1398 .as_ref()
1399 .map(|s| &s.signing_key)
1400 .ok_or_else(|| InstanceError::DeviceKeyNotFound.into())
1401 }
1402
1403 /// Test-only: Get the device signing key.
1404 ///
1405 /// This is exposed for testing purposes only. In production, use the User API.
1406 ///
1407 /// Returns an error if this is a remote Instance that does not have access to the
1408 /// device key.
1409 #[cfg(any(test, feature = "testing"))]
1410 pub fn signing_key(&self) -> Result<&PrivateKey> {
1411 self.inner
1412 .secrets
1413 .as_ref()
1414 .map(|s| &s.signing_key)
1415 .ok_or_else(|| InstanceError::DeviceKeyNotFound.into())
1416 }
1417
1418 /// Get the instance identity (public key).
1419 ///
1420 /// # Returns
1421 /// The instance's public key identity.
1422 pub fn id(&self) -> PublicKey {
1423 self.inner.metadata.id.clone()
1424 }
1425
1426 // === Synchronization Management ===
1427 //
1428 // These methods provide access to the Sync module for managing synchronization
1429 // settings and state for this database instance.
1430
1431 /// Initializes the Sync module for this instance.
1432 ///
1433 /// Enables synchronization operations for this instance. This method is idempotent;
1434 /// calling it multiple times has no effect.
1435 ///
1436 /// # Errors
1437 /// Returns an error if the sync settings database cannot be created or if device key
1438 /// generation/storage fails.
1439 pub async fn enable_sync(&self) -> Result<()> {
1440 // Check if there is an existing Sync already loaded
1441 if self.inner.sync.get().is_some() {
1442 return Ok(());
1443 }
1444
1445 // A remote Instance must not run sync client-side: building a Sync
1446 // here would spin up a background sync engine that drives RPCs against
1447 // the daemon's backend — duplicating (and racing) the daemon's own
1448 // sync. Sync is owned by the process that owns the Instance.
1449 //
1450 // Return `Ok(())` so callers on a connected instance get the same
1451 // no-op success they would on a local instance where sync is already
1452 // running daemon-side. Long-term this should become an admin-gated
1453 // operation that lets a client ask the daemon to enable its sync
1454 // subsystem; until that ships, the client-side `enable_sync` is
1455 // intentionally a silent no-op because the daemon either already
1456 // has sync running or it doesn't, and the client can't change that.
1457 //
1458 // TODO(service): expose an admin-gated `enable_sync` on
1459 // `InstanceAdmin` so a client can enable sync remotely.
1460 #[cfg(all(unix, feature = "service"))]
1461 if self.remote_connection().is_some() {
1462 return Ok(());
1463 }
1464
1465 // Check InstanceMetadata for existing sync_db
1466 let metadata = self
1467 .backend()
1468 .get_instance_metadata()
1469 .await?
1470 .ok_or(InstanceError::DeviceKeyNotFound)?; // Metadata must exist if instance is initialized
1471
1472 let sync = if let Some(ref sync_db) = metadata.sync_db {
1473 // Load existing sync tree
1474 Sync::load(self.clone(), sync_db).await?
1475 } else {
1476 // Create new sync tree
1477 let sync = Sync::new(self.clone()).await?;
1478
1479 // Save sync_db to metadata
1480 let mut new_metadata = metadata;
1481 new_metadata.sync_db = Some(sync.sync_tree_root_id().clone());
1482 self.backend().set_instance_metadata(&new_metadata).await?;
1483
1484 sync
1485 };
1486
1487 let sync_arc = Arc::new(sync);
1488
1489 // Initialize the sync engine (no transports registered yet)
1490 // Users should call register_transport() to add transports
1491 sync_arc.start_background_sync()?;
1492
1493 // Sync wants to observe writes across *every* tree, including
1494 // trees created after this point — there's no fixed tree set to
1495 // register per-db callbacks against, so this is one of the few
1496 // legitimate uses of `register_global_write_callback`. Idempotent
1497 // because `enable_sync` early-returns at the top if sync is
1498 // already initialized (`self.inner.sync.get().is_some()`); without
1499 // that guard this would register a duplicate hook every call.
1500 let sync_for_callback = Arc::clone(&sync_arc);
1501 self.register_global_write_callback(move |event, database| {
1502 let sync = Arc::clone(&sync_for_callback);
1503 let event = event.clone();
1504 let database = database.clone();
1505 async move {
1506 if event.source() == WriteSource::Local
1507 && sync.is_reconciliation_source(database.root_id()).await?
1508 {
1509 sync.reconcile_user_settings().await?;
1510 }
1511
1512 if event.source() == WriteSource::Local {
1513 sync.on_local_write(&event, &database).await
1514 } else {
1515 Ok(())
1516 }
1517 }
1518 });
1519
1520 let _ = self.inner.sync.set(sync_arc);
1521 Ok(())
1522 }
1523
1524 /// Get a reference to the Sync module.
1525 ///
1526 /// Returns a cheap-to-clone Arc handle to the Sync module. The Sync module
1527 /// uses interior mutability (AtomicBool and OnceLock) so &self methods are sufficient.
1528 ///
1529 /// # Returns
1530 /// An `Option` containing an `Arc<Sync>` if the Sync module is initialized.
1531 pub fn sync(&self) -> Option<Arc<Sync>> {
1532 self.inner.sync.get().map(Arc::clone)
1533 }
1534
1535 /// Flush all pending sync operations.
1536 ///
1537 /// This is a convenience method that processes all queued entries and
1538 /// retries any failed sends. If sync is not enabled, returns Ok(()).
1539 ///
1540 /// This is useful to force pending syncs to complete, e.g. on program shutdown.
1541 ///
1542 /// # Returns
1543 /// `Ok(())` if sync is not enabled or all operations completed successfully,
1544 /// or an error if sends failed.
1545 pub async fn flush_sync(&self) -> Result<()> {
1546 if let Some(sync) = self.sync() {
1547 sync.flush().await
1548 } else {
1549 Ok(())
1550 }
1551 }
1552
1553 // === Entry Write Coordination ===
1554 //
1555 // All entry writes go through Instance::put_entry() which handles backend storage
1556 // and callback dispatch. This centralizes write coordination and ensures hooks fire.
1557
1558 /// Register a per-database callback. Fires for writes to `tree_id` on
1559 /// this Instance.
1560 ///
1561 /// `initial_tips` seeds the callback's cursor — the first
1562 /// [`WriteEvent`] this callback receives will have `previous_tips`
1563 /// equal to `initial_tips`, and the cursor advances on each
1564 /// subsequent fire to that fire's post-write snapshot. Callers that
1565 /// want "tell me about everything after the point I just read at"
1566 /// pass the snapshot they just read; callers that want "tell me about
1567 /// everything from this empty cursor forward" can pass
1568 /// [`Snapshot::EMPTY`] (the first fire's `previous_tips` will be
1569 /// empty, and the subscriber walks the DAG to discover the gap).
1570 ///
1571 /// Returns the [`CallbackId`] of the registration. Callers wrap
1572 /// this in a [`WriteCallback`] handle (see
1573 /// [`Database::on_write_at_tips`]) to manage lifetime.
1574 pub(crate) fn register_write_callback<F, Fut>(
1575 &self,
1576 tree_id: ID,
1577 initial_tips: Snapshot,
1578 callback: F,
1579 ) -> CallbackId
1580 where
1581 F: for<'a> Fn(&'a WriteEvent, &'a Database) -> Fut + Send + std::marker::Sync + 'static,
1582 Fut: Future<Output = Result<()>> + Send + 'static,
1583 {
1584 let id = CallbackId(self.inner.next_callback_id.fetch_add(1, Ordering::Relaxed));
1585 let cb: AsyncWriteCallbackFn = Arc::new(move |event: &WriteEvent, database: &Database| {
1586 let fut = callback(event, database);
1587 Box::pin(fut) as AsyncWriteCallbackFuture
1588 });
1589 let entry = Arc::new(PerDbCallbackEntry {
1590 id,
1591 last_tips: std::sync::Mutex::new(initial_tips),
1592 callback: cb,
1593 });
1594 self.inner
1595 .write_callbacks
1596 .lock()
1597 .unwrap_or_else(|p| p.into_inner())
1598 .entry(tree_id)
1599 .or_default()
1600 .push(entry);
1601 id
1602 }
1603
1604 /// Register a non-removable callback fired for **every** write on **every**
1605 /// database for the life of the Instance.
1606 ///
1607 /// This is purpose-built for hooks that need to observe writes across
1608 /// the entire Instance — including writes to trees created *after* the
1609 /// hook is registered. The only legitimate use is something that
1610 /// genuinely doesn't know its target tree set up front: today, just
1611 /// sync (which wants to react to every local write so it can propagate
1612 /// to peers, and registers its hook once during `enable_sync`).
1613 ///
1614 /// **Not the right primitive for connection-scoped fan-out.** If you
1615 /// know up front which trees a consumer cares about (e.g. service
1616 /// clients subscribing per-tree via `DatabaseOp::SubscribeWrites`),
1617 /// register per-database callbacks with [`Self::register_write_callback`]
1618 /// instead. Per-db callbacks have a removal path
1619 /// ([`Self::remove_write_callback`]) which lets you tear them down on
1620 /// disconnect; this API does not.
1621 ///
1622 /// Callers branch on [`WriteEvent::source`] inside the closure if they
1623 /// only care about one source. Caller is responsible for idempotency
1624 /// (registering N times produces N firings); the canonical pattern is
1625 /// to guard the registration site with a `OnceLock` so it cannot run
1626 /// twice on the same Instance.
1627 pub(crate) fn register_global_write_callback<F, Fut>(&self, callback: F)
1628 where
1629 F: for<'a> Fn(&'a WriteEvent, &'a Database) -> Fut + Send + std::marker::Sync + 'static,
1630 Fut: Future<Output = Result<()>> + Send + 'static,
1631 {
1632 let id = CallbackId(self.inner.next_callback_id.fetch_add(1, Ordering::Relaxed));
1633 let cb: AsyncWriteCallbackFn = Arc::new(move |event: &WriteEvent, database: &Database| {
1634 let fut = callback(event, database);
1635 Box::pin(fut) as AsyncWriteCallbackFuture
1636 });
1637 self.inner
1638 .global_write_callbacks
1639 .lock()
1640 .unwrap_or_else(|p| p.into_inner())
1641 .push((id, cb));
1642 }
1643
1644 /// Remove a per-database callback by id. Returns `true` iff the
1645 /// removal emptied the per-tree callback list (i.e. this was the
1646 /// last live callback for `tree_id` on this Instance). No-op if
1647 /// the id isn't registered.
1648 ///
1649 /// Callers use the `true` return to drive lifecycle hooks on the
1650 /// connection's subscription state: dropping the last local
1651 /// callback for a tree on a connected instance is the trigger to
1652 /// transition the wire subscription to `Idle` (see
1653 /// [`crate::service::client::RemoteConnection::transition_to_idle`]).
1654 pub(crate) fn remove_write_callback(&self, tree_id: &ID, id: CallbackId) -> bool {
1655 let mut callbacks = self
1656 .inner
1657 .write_callbacks
1658 .lock()
1659 .unwrap_or_else(|p| p.into_inner());
1660 if let Some(vec) = callbacks.get_mut(tree_id) {
1661 let before = vec.len();
1662 vec.retain(|entry| entry.id != id);
1663 let removed = vec.len() < before;
1664 if vec.is_empty() {
1665 callbacks.remove(tree_id);
1666 return removed;
1667 }
1668 }
1669 false
1670 }
1671
1672 /// Whether any per-database write callback is currently registered for
1673 /// `tree_id`.
1674 ///
1675 /// Used by [`crate::service::client::RemoteConnection::transition_to_idle`]
1676 /// to re-check the registry while holding the subscription-state lock,
1677 /// so a registration racing the last callback's drop can't leave a live
1678 /// callback stranded on an `Idle` wire subscription.
1679 #[cfg(all(unix, feature = "service"))]
1680 pub(crate) fn has_write_callbacks(&self, tree_id: &ID) -> bool {
1681 self.inner
1682 .write_callbacks
1683 .lock()
1684 .unwrap_or_else(|p| p.into_inner())
1685 .contains_key(tree_id)
1686 }
1687
1688 /// Acquire (or create) the per-tree async lock that serializes the
1689 /// `snapshot` → backend write → callback dispatch sequence.
1690 ///
1691 /// Without this, two concurrent writers to the same tree both snapshot
1692 /// `previous_tips` before either writes, so the second callback's
1693 /// `previous_tips` would not reflect the first write — breaking the
1694 /// "diff against current tips" contract documented on [`WriteEvent`].
1695 pub(crate) fn tree_lock(&self, tree_id: &ID) -> Arc<tokio::sync::Mutex<()>> {
1696 let mut locks = self
1697 .inner
1698 .tree_locks
1699 .lock()
1700 .unwrap_or_else(|p| p.into_inner());
1701 Arc::clone(
1702 locks
1703 .entry(tree_id.clone())
1704 .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(()))),
1705 )
1706 }
1707
1708 /// Write an entry to the backend and dispatch callbacks.
1709 ///
1710 /// This is the central coordination point for all entry writes in the system.
1711 /// All writes must go through this method to ensure:
1712 /// - Entries are persisted to the backend
1713 /// - Appropriate callbacks are triggered based on write source
1714 /// - Hooks have full context (entry, database, instance)
1715 ///
1716 /// Serialized per-tree against [`Self::put_remote_entries`] so
1717 /// [`WriteEvent::previous_tips`] is consistent.
1718 ///
1719 /// # Arguments
1720 /// * `tree_id` - The root ID of the database being written to
1721 /// * `verification` - Authentication verification status of the entry
1722 /// * `entry` - The entry to write
1723 /// * `source` - Whether this is a local or remote write
1724 ///
1725 /// # Returns
1726 /// A Result indicating success or failure
1727 pub async fn put_entry(
1728 &self,
1729 tree_id: &ID,
1730 verification: crate::backend::VerificationStatus,
1731 entry: Entry,
1732 source: WriteSource,
1733 ) -> Result<()> {
1734 let lock = self.tree_lock(tree_id);
1735 let guard = lock.lock_owned().await;
1736 self.put_entry_under_tree_lock(guard, tree_id, verification, entry, source)
1737 .await
1738 }
1739
1740 /// Write an entry while the caller holds this tree's write lock.
1741 pub(crate) async fn put_entry_under_tree_lock(
1742 &self,
1743 guard: tokio::sync::OwnedMutexGuard<()>,
1744 tree_id: &ID,
1745 verification: crate::backend::VerificationStatus,
1746 entry: Entry,
1747 source: WriteSource,
1748 ) -> Result<()> {
1749 // 1. Capture tips before the write so callbacks know what changed.
1750 //
1751 // On a connected (remote) instance, the daemon owns the canonical
1752 // DAG and the client's local backend has nothing to read; reading
1753 // tips here would also gate against the *connection's* login pubkey
1754 // (not the per-DB acting identity from the Database handle), which
1755 // breaks the legitimate "create-a-tree-with-a-non-login-key" flow.
1756 // The client also skips firing callbacks locally on this path (see
1757 // step 3 below) — the daemon round-trips a `Notification::DatabaseWrite`
1758 // back with its own canonical `previous_tips` and we fire from
1759 // there instead.
1760 #[cfg(all(unix, feature = "service"))]
1761 let is_connected = self.remote_connection().is_some();
1762 #[cfg(not(all(unix, feature = "service")))]
1763 let is_connected = false;
1764
1765 let previous_tips = if is_connected {
1766 Snapshot::EMPTY
1767 } else {
1768 self.snapshot(tree_id).await?
1769 };
1770
1771 // 2. Persist to backend storage (and notify server for remote backends)
1772 self.backend()
1773 .write_entry(verification, entry.clone(), source)
1774 .await?;
1775
1776 // 3. Build event and fire callbacks — but only on a local
1777 // instance, and only for entries that arrive `Verified`.
1778 //
1779 // **Connected instance**: the daemon is the sole publisher of
1780 // write events. It fires its own callback registry when it
1781 // stores the entry, then pushes a `Notification::DatabaseWrite`
1782 // back to every subscribed client (including this one). Firing
1783 // here too would double-deliver. See `Database::on_write` for
1784 // the timing contract.
1785 //
1786 // **Unverified path**: skipped on purpose. Only a Verified write
1787 // *triggers* a fire. An entry that arrives `Unverified` (over the
1788 // wire as a `SubmitSignedEntry` body, or via sync as a remote
1789 // batch) is ingested silently here; the subsequent local
1790 // verification pass (the caller's responsibility to schedule)
1791 // decides whether it ever becomes a fire-eligible Verified entry,
1792 // and if so fires from there.
1793 //
1794 // This gates the *trigger*, not the bracket: the cursors below are
1795 // raw backend snapshots, so an Unverified or Failed entry that is
1796 // a raw tip still falls inside a subsequent event's bracket and
1797 // `ids_added` will enumerate it. Narrowing that to the Verified
1798 // frontier needs an incremental frontier first — `verified_frontier`
1799 // is an O(N) walk, too expensive on the per-commit path.
1800 // coding: raw-frontier cursors; switch to the Verified frontier
1801 // once an incremental one exists.
1802 let joins = if !is_connected && verification == VerificationStatus::Verified {
1803 // Compute the post-write snapshot for cursor advance. Cheap:
1804 // just re-read the backend snapshot post-put. Each per-callback
1805 // cursor advances to this value.
1806 // Propagate rather than defaulting: `Snapshot::EMPTY` is a
1807 // *meaningful* cursor ("no initial state"), not a neutral
1808 // fallback. Masking an error here would fire this event with
1809 // `post = EMPTY` (so `ids_added` reports nothing and sync's
1810 // global hook queues nothing — a silently unsynced commit) and
1811 // leave every per-callback cursor at EMPTY, replaying the full
1812 // history to every subscriber on the next event.
1813 let post_tips = self.snapshot(tree_id).await?;
1814 // Commit cursor advances + spawn under the tree lock (ordering),
1815 // but drain after dropping the guard (below) — a callback that
1816 // reads tips can re-enter `verify()` → `tree_lock` and would
1817 // deadlock against a still-held `_guard`. See
1818 // `Instance::spawn_write_callbacks`.
1819 Some(
1820 self.spawn_write_callbacks(tree_id, &previous_tips, &post_tips, source)
1821 .await,
1822 )
1823 } else {
1824 None
1825 };
1826
1827 // Release the per-tree lock before awaiting user callbacks.
1828 drop(guard);
1829
1830 if let Some(mut joins) = joins {
1831 while joins.join_next().await.is_some() {}
1832 }
1833
1834 Ok(())
1835 }
1836
1837 /// Store a batch of remotely-received entries and fire callbacks once.
1838 ///
1839 /// This is the correct way to ingest entries from sync. All entries are
1840 /// persisted first, then callbacks fire exactly once with the full batch
1841 /// and the tips from before ingestion. This ensures:
1842 ///
1843 /// - The database is fully consistent when callbacks execute
1844 /// - Callbacks fire once per sync exchange, not once per entry
1845 /// - `previous_tips` lets consumers reconstruct exactly what changed
1846 ///
1847 /// Entries that fail to store are logged and skipped — remaining entries
1848 /// are still stored and callbacks still fire for whatever was persisted.
1849 /// Returns the number of entries that were successfully persisted.
1850 ///
1851 /// Serialized per-tree against [`Self::put_entry`] and other concurrent
1852 /// `put_remote_entries` calls so `previous_tips` is consistent across
1853 /// writers.
1854 ///
1855 /// Entries are stored as [`VerificationStatus::Unverified`] without
1856 /// exception: they arrive from outside this node's local validation pass,
1857 /// so this node has not verified them and a peer cannot assert that it
1858 /// did. A later local re-verification pass may promote them.
1859 ///
1860 /// # Arguments
1861 /// * `tree_id` - The root ID of the database receiving the batch
1862 /// * `entries` - The entries to ingest
1863 pub(crate) async fn put_remote_entries(
1864 &self,
1865 tree_id: &ID,
1866 entries: Vec<Entry>,
1867 ) -> Result<usize> {
1868 if entries.is_empty() {
1869 return Ok(0);
1870 }
1871
1872 // Store the batch under the tree lock; release before calling
1873 // `verify`, which acquires its own lock for the pass + fire.
1874 let stored_count = {
1875 let lock = self.tree_lock(tree_id);
1876 let _guard = lock.lock().await;
1877 let mut stored = 0usize;
1878 for entry in entries {
1879 match self.backend().put(entry.clone()).await {
1880 Ok(_) => stored += 1,
1881 Err(e) => tracing::error!(
1882 tree_id = %tree_id,
1883 entry_id = %entry.id(),
1884 "Failed to store remote entry: {}", e
1885 ),
1886 }
1887 }
1888 stored
1889 };
1890
1891 // Run verify inline. `Database::verify` walks the Unverified
1892 // region in O(K), promotes whatever can be settled, and fires
1893 // one batched `Verified` event for the promotions. Sync-ingest
1894 // subscribers see the promotion without needing to schedule
1895 // their own verify pass.
1896 if stored_count > 0 {
1897 Database::open(self, tree_id).await?.verify().await?;
1898 }
1899
1900 Ok(stored_count)
1901 }
1902
1903 /// Demote `entry_id` to [`VerificationStatus::Unverified`] and
1904 /// cascade the demotion to every `Verified` descendant in
1905 /// `tree_id`'s DAG.
1906 ///
1907 /// **Why the cascade.** The `Verified` set on a tree is
1908 /// prefix-closed: an entry is `Verified` only if every one of its
1909 /// ancestors is. Demoting an ancestor without also demoting its
1910 /// `Verified` descendants breaks that invariant, which means
1911 /// `Database::verify`'s targeted walk-from-tips cannot find the
1912 /// demoted entry — it's hidden behind a still-`Verified` descendant
1913 /// and would be stranded. The cascade restores the invariant.
1914 ///
1915 /// **Scope.** Today's only callers are tests (and the v0
1916 /// re-verification scenarios they exercise). Production code does
1917 /// not demote `Verified` → `Unverified` at all under the current
1918 /// verify implementation. If a future demotion path appears (e.g.
1919 /// retroactive settings-change-driven invalidation) it should route
1920 /// through this method.
1921 ///
1922 /// O(N) per call: walks `get_tree` to build a children index. Cheap
1923 /// for the test sizes this targets; not a hot path.
1924 ///
1925 /// Exposed publicly under `cfg(test)` and the `testing` feature so
1926 /// integration tests in the `it` crate can use it; otherwise
1927 /// `pub(crate)`-equivalent.
1928 #[cfg(any(test, feature = "testing"))]
1929 pub async fn demote_to_unverified(&self, tree_id: &ID, entry_id: &ID) -> Result<()> {
1930 self.demote_to_unverified_impl(tree_id, entry_id).await
1931 }
1932
1933 #[cfg(not(any(test, feature = "testing")))]
1934 pub(crate) async fn demote_to_unverified(&self, tree_id: &ID, entry_id: &ID) -> Result<()> {
1935 self.demote_to_unverified_impl(tree_id, entry_id).await
1936 }
1937
1938 async fn demote_to_unverified_impl(&self, tree_id: &ID, entry_id: &ID) -> Result<()> {
1939 use std::collections::{HashMap, HashSet, VecDeque};
1940 let backend = self.require_local_engine()?;
1941 let entries = backend.get_tree(tree_id).await?;
1942
1943 // Build children index from each entry's parents.
1944 let mut children: HashMap<ID, Vec<ID>> = HashMap::new();
1945 for entry in &entries {
1946 for p in entry.parents().unwrap_or_default() {
1947 children.entry(p).or_default().push(entry.id());
1948 }
1949 }
1950
1951 // BFS from the target. The target itself always gets demoted
1952 // (caller's intent); descendants get demoted only if they are
1953 // currently `Verified`. `Failed` or `Unverified` descendants
1954 // are left as-is — `Failed` is terminal, `Unverified` is
1955 // already at the target state.
1956 let mut queue: VecDeque<ID> = VecDeque::new();
1957 queue.push_back(entry_id.clone());
1958 let mut visited: HashSet<ID> = HashSet::new();
1959 while let Some(id) = queue.pop_front() {
1960 if !visited.insert(id.clone()) {
1961 continue;
1962 }
1963 let status = match backend.get_verification_status(&id).await {
1964 Ok(s) => s,
1965 Err(e) if e.is_not_found() => continue,
1966 Err(e) => return Err(e),
1967 };
1968 let is_target = id == *entry_id;
1969 if is_target || status == VerificationStatus::Verified {
1970 backend
1971 .update_verification_status(&id, VerificationStatus::Unverified)
1972 .await?;
1973 }
1974 if let Some(kids) = children.get(&id) {
1975 for kid in kids {
1976 queue.push_back(kid.clone());
1977 }
1978 }
1979 }
1980 Ok(())
1981 }
1982
1983 /// Dispatch callbacks for a write event.
1984 ///
1985 /// Per-database callbacks for `tree_id` get **per-callback events**:
1986 /// each callback's `previous_tips` is read from its own cursor and
1987 /// the cursor advances to `post_tips` synchronously around the
1988 /// fire. The cursor mutex is released before the user callback is
1989 /// awaited, so a slow callback does not stall other callbacks'
1990 /// cursor reads on a concurrent fire.
1991 ///
1992 /// **Per-callback dispatch is concurrent.** Cursor advancement
1993 /// happens synchronously in arrival order under each callback's
1994 /// own mutex, then every callback's closure is spawned on its own
1995 /// tokio task. A slow callback for one subscriber doesn't stall
1996 /// other subscribers' callbacks for the same event. The dispatcher
1997 /// awaits every spawned task before returning, so the per-tree
1998 /// dispatch worker's "this notification is finished" point still
1999 /// serialises against the next event on the same tree — the
2000 /// inter-event ordering contract documented on
2001 /// [`Database::on_write`](crate::Database::on_write) is preserved.
2002 ///
2003 /// Global callbacks fire with a single shared event whose
2004 /// `previous_tips` is the caller-supplied `previous_tips`
2005 /// argument — globals don't track per-tree cursors and continue
2006 /// to receive the pre-write tips view. Globals also dispatch
2007 /// concurrently across subscribers.
2008 ///
2009 /// `pub(crate)` so the service module's reader task can drive this
2010 /// directly when a `Notification::DatabaseWrite` arrives from the
2011 /// daemon — that path is the *sole* publisher on a connected
2012 /// instance.
2013 ///
2014 /// Convenience wrapper over [`Self::spawn_write_callbacks`]: spawns the
2015 /// dispatches and awaits them all. Use this from callers that do **not**
2016 /// hold a [`tree_lock`](Self::tree_lock) across the fire. Callers that
2017 /// hold the lock (the local `put_entry` path, `Database::verify`) must
2018 /// instead call `spawn_write_callbacks` under the lock, drop the guard,
2019 /// then drain the returned `JoinSet` — see that method's contract.
2020 pub(crate) async fn fire_write_callbacks(
2021 &self,
2022 tree_id: &ID,
2023 previous_tips: &Snapshot,
2024 post_tips: &Snapshot,
2025 source: WriteSource,
2026 ) {
2027 let mut joins = self
2028 .spawn_write_callbacks(tree_id, previous_tips, post_tips, source)
2029 .await;
2030 while joins.join_next().await.is_some() {}
2031 }
2032
2033 /// Phase 1 of callback dispatch: advance every callback's cursor and
2034 /// spawn its closure, returning the in-flight [`JoinSet`] **without
2035 /// awaiting it**.
2036 ///
2037 /// Splitting the dispatch in two lets a caller that holds a
2038 /// [`tree_lock`](Self::tree_lock) commit the cursor advances *under*
2039 /// the lock — which is what preserves event ordering against a
2040 /// concurrent writer — and then release the lock *before* awaiting the
2041 /// user closures.
2042 ///
2043 /// **Why the lock must be dropped before draining.** Awaiting a user
2044 /// callback while holding `tree_id`'s lock risks a reentrant deadlock: a
2045 /// callback that reads tips (`Database::snapshot` and friends) can trip
2046 /// the access-time auto-verify hook, which calls `Database::verify`,
2047 /// which acquires the very same `tree_lock`. The awaiting caller still
2048 /// holds it, and `verify` is waiting on the callback it spawned →
2049 /// circular wait.
2050 ///
2051 /// Everything this method does up to (and including) the callback
2052 /// *invocation* is synchronous — the cursor read/advance under each
2053 /// callback's own `std::Mutex`, and calling the callback to obtain its
2054 /// `'static` future, which runs only the closure's synchronous prefix (the
2055 /// service subscription's non-blocking `frame_tx.send`; for a plain
2056 /// `async move { … }` callback the prefix is empty). That prefix must not
2057 /// itself block on the tree lock, which no in-tree callback does. Only the
2058 /// returned future's `.await` can re-enter the lock, and that is what the
2059 /// caller drains lock-free. Running the invocation under the lock — rather
2060 /// than deferring it into the spawned task — is deliberate: it makes each
2061 /// callback's synchronous side effect commit in canonical cursor order, so
2062 /// concurrent same-tree writers can't reorder a subscriber's notification
2063 /// stream.
2064 ///
2065 /// [`JoinSet`]: tokio::task::JoinSet
2066 pub(crate) async fn spawn_write_callbacks(
2067 &self,
2068 tree_id: &ID,
2069 previous_tips: &Snapshot,
2070 post_tips: &Snapshot,
2071 source: WriteSource,
2072 ) -> tokio::task::JoinSet<()> {
2073 let per_db_callbacks = self
2074 .inner
2075 .write_callbacks
2076 .lock()
2077 .unwrap_or_else(|p| p.into_inner())
2078 .get(tree_id)
2079 .cloned();
2080
2081 let global_callbacks = self
2082 .inner
2083 .global_write_callbacks
2084 .lock()
2085 .unwrap_or_else(|p| p.into_inner())
2086 .clone();
2087
2088 let has_callbacks = per_db_callbacks.is_some() || !global_callbacks.is_empty();
2089 if !has_callbacks {
2090 return tokio::task::JoinSet::new();
2091 }
2092
2093 // Callback dispatch is triggered by a write the Instance just accepted,
2094 // so the root is already present. Build the keyless handle without a
2095 // backend read: on a connected Instance that read would be re-gated as
2096 // the client session and could suppress the authorization-change event
2097 // that revokes that very session.
2098 let database =
2099 Database::from_parts(tree_id.clone(), self.downgrade(), self.backend().clone());
2100
2101 // Single JoinSet across per-db + global callbacks. Two things happen
2102 // synchronously, in arrival order, before any task is spawned — both
2103 // under the caller's tree lock:
2104 //
2105 // 1. Each callback's cursor read+advance (under its own std::Mutex),
2106 // so `previous_tips` brackets commit deterministically.
2107 // 2. The callback *invocation* itself. A callback returns a `'static`
2108 // future, so calling it runs the closure's synchronous prefix now,
2109 // in canonical order — e.g. the service subscription's
2110 // `frame_tx.send` (a non-blocking push whose future is a no-op).
2111 // Same-tree events therefore reach each subscriber's channel in
2112 // order even under concurrent writers; only the returned future's
2113 // async remainder runs concurrently on the drained set.
2114 let mut joins = tokio::task::JoinSet::new();
2115
2116 if let Some(callbacks) = per_db_callbacks {
2117 for entry in callbacks {
2118 let cb_previous = {
2119 let mut guard = entry
2120 .last_tips
2121 .lock()
2122 .unwrap_or_else(|poisoned| poisoned.into_inner());
2123 std::mem::replace(&mut *guard, post_tips.clone())
2124 };
2125 let event = WriteEvent {
2126 previous_tips: cb_previous,
2127 post_tips: post_tips.clone(),
2128 source,
2129 };
2130 // Invoke synchronously in cursor order; spawn only the tail. A
2131 // callback's ordering-critical side effect must live in this
2132 // synchronous prefix (before its future's first await) — that is
2133 // what keeps it under the tree lock and in canonical order. See
2134 // the `frame_tx.send` invariant in `service::server`'s
2135 // `SubscribeWrites` handler.
2136 let fut = (entry.callback)(&event, &database);
2137 let tree_id_for_cb = tree_id.clone();
2138 let cb_id = entry.id;
2139 joins.spawn(async move {
2140 if let Err(e) = fut.await {
2141 tracing::error!(
2142 tree_id = %tree_id_for_cb,
2143 source = ?source,
2144 callback_id = ?cb_id,
2145 "Per-database callback failed: {}", e
2146 );
2147 }
2148 });
2149 }
2150 }
2151
2152 // Globals fire with the shared pre-write tips (no per-callback cursor).
2153 // Invoked synchronously here too, in registration order, for the same
2154 // in-order-side-effect guarantee; only the async remainder is spawned.
2155 for (id, callback) in global_callbacks {
2156 let event = WriteEvent {
2157 previous_tips: previous_tips.clone(),
2158 post_tips: post_tips.clone(),
2159 source,
2160 };
2161 let fut = callback(&event, &database);
2162 let tree_id_for_cb = tree_id.clone();
2163 joins.spawn(async move {
2164 if let Err(e) = fut.await {
2165 tracing::error!(
2166 tree_id = %tree_id_for_cb,
2167 source = ?source,
2168 callback_id = ?id,
2169 "Global callback failed: {}", e
2170 );
2171 }
2172 });
2173 }
2174
2175 joins
2176 }
2177
2178 /// Downgrade to a weak reference.
2179 ///
2180 /// Creates a weak reference that does not prevent the Instance from being dropped.
2181 /// This is useful for preventing circular reference cycles in dependent objects.
2182 ///
2183 /// # Returns
2184 /// A `WeakInstance` that can be upgraded back to a strong reference.
2185 pub fn downgrade(&self) -> WeakInstance {
2186 WeakInstance {
2187 inner: Arc::downgrade(&self.inner),
2188 }
2189 }
2190}
2191
2192impl WeakInstance {
2193 /// Upgrade to a strong reference.
2194 ///
2195 /// Attempts to upgrade this weak reference to a strong `Instance` reference.
2196 /// Returns `None` if the Instance has already been dropped.
2197 ///
2198 /// # Returns
2199 /// `Some(Instance)` if the Instance still exists, `None` otherwise.
2200 ///
2201 /// # Example
2202 /// ```
2203 /// # use eidetica::{Instance, NewUser};
2204 /// # #[tokio::main]
2205 /// # async fn main() -> eidetica::Result<()> {
2206 /// let (instance, maybe_user) = Instance::connect_or_create(
2207 /// "memory://",
2208 /// NewUser::passwordless("alice"),
2209 /// ).await?;
2210 /// let user = maybe_user.expect("memory:// is always fresh");
2211 /// let weak = instance.downgrade();
2212 ///
2213 /// // Upgrade works while instance exists
2214 /// assert!(weak.upgrade().is_some());
2215 ///
2216 /// // User holds its own strong handle to the Instance — drop it too so
2217 /// // the weak upgrade can fail.
2218 /// drop(user);
2219 /// drop(instance);
2220 /// // Upgrade fails after instance is dropped
2221 /// assert!(weak.upgrade().is_none());
2222 /// # Ok(())
2223 /// # }
2224 /// ```
2225 pub fn upgrade(&self) -> Option<Instance> {
2226 self.inner.upgrade().map(|inner| Instance { inner })
2227 }
2228}
2229
2230// ============ URL-dispatch backend constructors ============
2231
2232#[cfg(feature = "sqlite")]
2233async fn open_sqlite_backend(url: &str) -> Result<Box<dyn BackendImpl>> {
2234 let backend = crate::backend::database::Sqlite::connect(url).await?;
2235 Ok(Box::new(backend))
2236}
2237
2238#[cfg(not(feature = "sqlite"))]
2239async fn open_sqlite_backend(_url: &str) -> Result<Box<dyn BackendImpl>> {
2240 Err(InstanceError::BackendUnavailable {
2241 scheme: "sqlite",
2242 missing_feature: "sqlite",
2243 }
2244 .into())
2245}
2246
2247#[cfg(feature = "postgres")]
2248async fn open_postgres_backend(url: &str) -> Result<Box<dyn BackendImpl>> {
2249 let backend = crate::backend::database::Postgres::connect(url).await?;
2250 Ok(Box::new(backend))
2251}
2252
2253#[cfg(not(feature = "postgres"))]
2254async fn open_postgres_backend(_url: &str) -> Result<Box<dyn BackendImpl>> {
2255 Err(InstanceError::BackendUnavailable {
2256 scheme: "postgres",
2257 missing_feature: "postgres",
2258 }
2259 .into())
2260}