Skip to main content

khive_storage/
sql.rs

1//! SQL access capability traits.
2
3use std::any::Any;
4use std::future::Future;
5use std::pin::Pin;
6
7use async_trait::async_trait;
8
9use crate::types::{PageRequest, SqlRow, SqlStatement, SqlValue, StorageResult};
10
11/// A boxed future, borrowing from the `&mut dyn SqlWriter` an
12/// [`AtomicUnitOp`] is called with (see [`SqlAccess::atomic_unit`]).
13pub type BoxFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
14
15/// A caller-supplied unit of work to run as ONE atomic operation via
16/// [`SqlAccess::atomic_unit`] (ADR-067 Component A, Fork C slice 2).
17///
18/// `op` receives a live `&mut dyn SqlWriter` already inside an open write
19/// transaction — it must issue DML only (no bare `BEGIN`/`COMMIT`/
20/// `ROLLBACK`; the caller-visible transaction boundary is owned entirely by
21/// `atomic_unit`, exactly like the existing `execute_batch` contract) — and
22/// returns its result type-erased via `Box<dyn Any + Send>` so this trait
23/// method stays object-safe (no method-level generics on a trait used as
24/// `dyn SqlAccess`). Callers downcast the returned box back to their own
25/// concrete outcome type.
26pub type AtomicUnitOp = Box<
27    dyn for<'w> FnOnce(&'w mut dyn SqlWriter) -> BoxFuture<'w, StorageResult<Box<dyn Any + Send>>>
28        + Send,
29>;
30
31/// Read-capable SQL connection.
32#[async_trait]
33pub trait SqlReader: Send + 'static {
34    /// Execute `statement` and return the first row, or `None` if the result set is empty.
35    /// Implementations must not convert later matching rows into owned values.
36    /// Implementations may stop stepping the statement early, so statements with
37    /// side effects (e.g. DML with `RETURNING`) must not be issued through this
38    /// method; use [`SqlWriter::execute`] or [`SqlWriter::execute_batch`] for writes.
39    async fn query_row(&mut self, statement: SqlStatement) -> StorageResult<Option<SqlRow>>;
40    /// Execute `statement` and return all rows.
41    ///
42    /// This compatibility primitive has no result-size bound. Callers that do
43    /// not already constrain their SQL should use [`Self::query_page`].
44    async fn query_all(&mut self, statement: SqlStatement) -> StorageResult<Vec<SqlRow>>;
45    /// Execute `statement` and return only the requested offset page.
46    ///
47    /// The default preserves source compatibility for alternate backends by
48    /// slicing [`Self::query_all`]. Backends should override this method when
49    /// they can stop row conversion at `page.limit`; khive-db's SQLite bridge
50    /// does so. On that path, materialization is bounded by the caller-supplied
51    /// `page.limit`, and callers own choosing a sane limit; the database engine
52    /// may still do work for the query plan and offset.
53    ///
54    /// Like [`Self::query_all`], the default has no result-size bound: it
55    /// materializes every matching row before slicing. Callers issuing
56    /// unconstrained SQL must not rely on the default to keep memory
57    /// proportional to `page.limit` — that bound holds only on backends that
58    /// override this method.
59    ///
60    /// Implementations may stop stepping the statement early, so statements with
61    /// side effects (e.g. DML with `RETURNING`) must not be issued through this
62    /// method; use [`SqlWriter::execute`] or [`SqlWriter::execute_batch`] for writes.
63    async fn query_page(
64        &mut self,
65        statement: SqlStatement,
66        page: PageRequest,
67    ) -> StorageResult<Vec<SqlRow>> {
68        let rows = self.query_all(statement).await?;
69        let offset = usize::try_from(page.offset).unwrap_or(usize::MAX);
70        let limit = usize::try_from(page.limit).unwrap_or(usize::MAX);
71        Ok(rows.into_iter().skip(offset).take(limit).collect())
72    }
73    /// Execute `statement` and return the first column of the first row as a scalar.
74    async fn query_scalar(&mut self, statement: SqlStatement) -> StorageResult<Option<SqlValue>>;
75    /// Run `EXPLAIN QUERY PLAN` for `statement` and return the plan rows.
76    async fn explain(&mut self, statement: SqlStatement) -> StorageResult<Vec<SqlRow>>;
77}
78
79/// Write-capable SQL connection (extends `SqlReader`).
80#[async_trait]
81pub trait SqlWriter: SqlReader + Send + 'static {
82    /// Execute a single DML statement and return the number of rows affected.
83    ///
84    /// Transaction-control rejection is an `execute_batch` contract; this
85    /// primitive remains available to internal transaction owners such as
86    /// `atomic_unit` for their `BEGIN`/`COMMIT`/`ROLLBACK` calls.
87    async fn execute(&mut self, statement: SqlStatement) -> StorageResult<u64>;
88    /// Execute multiple DML statements and return the total rows affected.
89    async fn execute_batch(&mut self, statements: Vec<SqlStatement>) -> StorageResult<u64>;
90    /// Execute a raw SQL script (no parameters; used for migrations).
91    ///
92    /// This script boundary is internal/migration-only and deliberately does
93    /// not inherit `execute_batch`'s transaction-control rejection.
94    async fn execute_script(&mut self, script: String) -> StorageResult<()>;
95
96    /// Execute a raw SQL script that MUST run outside any open transaction
97    /// (ADR-067 Component A, Fork C slice 2) — e.g.
98    /// `VACUUM`, which SQLite rejects if issued inside `BEGIN`/`COMMIT`.
99    ///
100    /// Default implementation delegates to [`Self::execute_script`]: every
101    /// writer implementation in this codebase except khive-db's
102    /// write-queue-routed `SqliteWriter` already runs `execute_script`
103    /// transaction-free (a plain connection call, or already inside a
104    /// caller-managed transaction where a top-level statement would be
105    /// invalid regardless of which method is called). `SqliteWriter`
106    /// overrides this to route around its writer task's per-request `BEGIN
107    /// IMMEDIATE` specifically for this call, while still serializing
108    /// through the single writer owner.
109    ///
110    /// This is an internal maintenance boundary (for example `VACUUM` or a
111    /// checkpoint script), not an `execute_batch` transaction-control guard.
112    async fn execute_script_top_level(&mut self, script: String) -> StorageResult<()> {
113        self.execute_script(script).await
114    }
115}
116
117/// Base SQL access capability.
118#[async_trait]
119pub trait SqlAccess: Send + Sync + 'static {
120    /// Canonical filesystem identity for this SQL database, when file-backed.
121    ///
122    /// Cross-resource operations use this only to derive advisory coordination
123    /// files outside SQLite. Only genuinely pathless, process-private
124    /// implementations may return `None`; every file-backed implementation
125    /// must expose its canonical path so cross-process exclusion cannot
126    /// silently degrade. The method performs no I/O.
127    fn database_path(&self) -> Option<std::path::PathBuf> {
128        None
129    }
130
131    /// Acquire a read-only connection from the pool.
132    async fn reader(&self) -> StorageResult<Box<dyn SqlReader>>;
133    /// Acquire a read-write connection from the pool.
134    async fn writer(&self) -> StorageResult<Box<dyn SqlWriter>>;
135
136    /// Run `op` as ONE atomic unit of work (ADR-067 Component A, Fork C
137    /// slice 2).
138    ///
139    /// Where a single-writer task is active (file-backed pool,
140    /// `KHIVE_WRITE_QUEUE=1`), `op` runs inside that task's one write
141    /// transaction for this request — no separate connection is opened, so
142    /// this call cannot compete with the writer task for SQLite's write
143    /// lock. Where no writer task applies (flag off, no runtime, or an
144    /// in-memory pool), `op` runs under a manual
145    /// `BEGIN IMMEDIATE`/`COMMIT`/`ROLLBACK` on a writer handle exactly like
146    /// calling [`Self::writer`] and driving the statements by hand — the
147    /// pre-ADR-067 behavior, preserved byte-for-byte on this path.
148    ///
149    /// **The atomic-unit suspend-free invariant (normative for every
150    /// caller):** `op`'s future must complete on its **first poll** — it may
151    /// issue only synchronous DML against the `&mut dyn SqlWriter` it is
152    /// handed and must never reach a real suspension point (no embedding
153    /// computation, no ANN warming, no service or channel `await`, no
154    /// network round-trip). Synchronous external work is forbidden too: no
155    /// filesystem/process/network I/O, sleeps, blocking waits, or unbounded CPU
156    /// may run between the transaction owner's `BEGIN IMMEDIATE` and
157    /// `COMMIT`. First-poll enforcement alone cannot detect those operations;
158    /// ADR-091's write-transaction audit table is the review-time guard, and a
159    /// new owner or caller must update that table before merge. On the
160    /// single-writer path this is enforced at runtime: the writer task drives
161    /// `op` through a single-poll driver and
162    /// returns a typed error the instant the future is `Pending`, so a
163    /// violation fails loudly rather than corrupting state. On the flag-off
164    /// path (no writer task active) a suspending `op` would currently
165    /// *succeed* — that path drives `op` as an ordinary `.await` under a
166    /// manual transaction — so the invariant is a correctness contract this
167    /// trait asks every caller to uphold, not something the type system or
168    /// every code path enforces. Callers must not rely on the flag-off
169    /// path's tolerance; behavior must be identical (synchronous DML only)
170    /// regardless of whether the single-writer flag is on.
171    async fn atomic_unit(&self, op: AtomicUnitOp) -> StorageResult<Box<dyn Any + Send>>;
172}