Development Documentation (main branch) - For stable release docs, see docs.rs/eidetica
Skip to main content

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}