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 may issue DML and synchronous read assertions (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/// Closed set of parameter-free maintenance statements executed outside the
32/// writer's per-request transaction wrapper.
33///
34/// Every variant renders a static SQL literal. Commands that require a caller
35/// value, such as a destination path, need a separate parameter-bound primitive;
36/// this set deliberately has no raw SQL variant or string conversion.
37#[derive(Clone, Copy, Debug, PartialEq, Eq)]
38pub enum TopLevelMaintenance {
39    /// Reclaim unused database pages with `VACUUM`.
40    Vacuum,
41    /// Checkpoint the write-ahead log and truncate it to zero bytes.
42    WalCheckpointTruncate,
43}
44
45impl TopLevelMaintenance {
46    /// Render the reviewed, parameter-free SQL literal for this operation.
47    pub const fn as_sql(self) -> &'static str {
48        match self {
49            Self::Vacuum => "VACUUM;",
50            Self::WalCheckpointTruncate => "PRAGMA wal_checkpoint(TRUNCATE);",
51        }
52    }
53}
54
55/// Read-capable SQL connection.
56#[async_trait]
57pub trait SqlReader: Send + 'static {
58    /// Execute `statement` and return the first row, or `None` if the result set is empty.
59    /// Implementations must not convert later matching rows into owned values.
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_row(&mut self, statement: SqlStatement) -> StorageResult<Option<SqlRow>>;
64    /// Execute `statement` and return all rows.
65    ///
66    /// This compatibility primitive has no result-size bound. Callers that do
67    /// not already constrain their SQL should use [`Self::query_page`].
68    async fn query_all(&mut self, statement: SqlStatement) -> StorageResult<Vec<SqlRow>>;
69    /// Execute `statement` and return only the requested offset page.
70    ///
71    /// The default preserves source compatibility for alternate backends by
72    /// slicing [`Self::query_all`]. Backends should override this method when
73    /// they can stop row conversion at `page.limit`; khive-db's SQLite bridge
74    /// does so. On that path, materialization is bounded by the caller-supplied
75    /// `page.limit`, and callers own choosing a sane limit; the database engine
76    /// may still do work for the query plan and offset.
77    ///
78    /// Like [`Self::query_all`], the default has no result-size bound: it
79    /// materializes every matching row before slicing. Callers issuing
80    /// unconstrained SQL must not rely on the default to keep memory
81    /// proportional to `page.limit` — that bound holds only on backends that
82    /// override this method.
83    ///
84    /// Implementations may stop stepping the statement early, so statements with
85    /// side effects (e.g. DML with `RETURNING`) must not be issued through this
86    /// method; use [`SqlWriter::execute`] or [`SqlWriter::execute_batch`] for writes.
87    async fn query_page(
88        &mut self,
89        statement: SqlStatement,
90        page: PageRequest,
91    ) -> StorageResult<Vec<SqlRow>> {
92        let rows = self.query_all(statement).await?;
93        let offset = usize::try_from(page.offset).unwrap_or(usize::MAX);
94        let limit = usize::try_from(page.limit).unwrap_or(usize::MAX);
95        Ok(rows.into_iter().skip(offset).take(limit).collect())
96    }
97    /// Execute `statement` and return the first column of the first row as a scalar.
98    async fn query_scalar(&mut self, statement: SqlStatement) -> StorageResult<Option<SqlValue>>;
99    /// Run `EXPLAIN QUERY PLAN` for `statement` and return the plan rows.
100    async fn explain(&mut self, statement: SqlStatement) -> StorageResult<Vec<SqlRow>>;
101}
102
103/// Write-capable SQL connection (extends `SqlReader`).
104#[async_trait]
105pub trait SqlWriter: SqlReader + Send + 'static {
106    /// Execute a single DML statement and return the number of rows affected.
107    ///
108    /// Transaction-control rejection is an `execute_batch` contract; this
109    /// primitive remains available to internal transaction owners such as
110    /// `atomic_unit` for their `BEGIN`/`COMMIT`/`ROLLBACK` calls.
111    async fn execute(&mut self, statement: SqlStatement) -> StorageResult<u64>;
112    /// Execute multiple DML statements and return the total rows affected.
113    async fn execute_batch(&mut self, statements: Vec<SqlStatement>) -> StorageResult<u64>;
114    /// Execute a raw SQL script (no parameters; used for migrations).
115    ///
116    /// This script boundary is internal/migration-only and deliberately does
117    /// not inherit `execute_batch`'s transaction-control rejection.
118    async fn execute_script(&mut self, script: String) -> StorageResult<()>;
119
120    /// Execute one closed-set maintenance operation that MUST run outside any
121    /// open transaction (ADR-067 Component A, Fork C slice 2). `VACUUM`, for
122    /// example, is rejected by SQLite inside `BEGIN`/`COMMIT`.
123    ///
124    /// Default implementation delegates to [`Self::execute_script`]: every
125    /// writer implementation in this codebase except khive-db's
126    /// write-queue-routed `SqliteWriter` already runs `execute_script`
127    /// transaction-free (a plain connection call, or already inside a
128    /// caller-managed transaction where a top-level statement would be
129    /// invalid regardless of which method is called). `SqliteWriter`
130    /// overrides this to route around its writer task's per-request `BEGIN
131    /// IMMEDIATE` specifically for this call, while still serializing
132    /// through the single writer owner.
133    ///
134    /// SQL is rendered only from [`TopLevelMaintenance`]'s static literals;
135    /// this API does not accept caller-supplied or formatted SQL. It is an
136    /// internal maintenance boundary, not an `execute_batch` transaction-control
137    /// guard. The separate migration-only [`Self::execute_script`] is unchanged.
138    ///
139    /// ```no_run
140    /// use khive_storage::{SqlWriter, StorageResult, TopLevelMaintenance};
141    ///
142    /// async fn maintain(writer: &mut dyn SqlWriter) -> StorageResult<()> {
143    ///     writer.execute_script_top_level(TopLevelMaintenance::Vacuum).await?;
144    ///     writer.execute_script_top_level(TopLevelMaintenance::WalCheckpointTruncate).await
145    /// }
146    /// ```
147    ///
148    /// A runtime string cannot cross this boundary, even if its current contents
149    /// happen to be a supported maintenance statement:
150    ///
151    /// ```compile_fail,E0308
152    /// use khive_storage::{SqlWriter, StorageResult};
153    ///
154    /// async fn reject_script(writer: &mut dyn SqlWriter, script: String) -> StorageResult<()> {
155    ///     writer.execute_script_top_level(script).await
156    /// }
157    /// ```
158    async fn execute_script_top_level(
159        &mut self,
160        maintenance: TopLevelMaintenance,
161    ) -> StorageResult<()> {
162        self.execute_script(maintenance.as_sql().to_owned()).await
163    }
164}
165
166/// Base SQL access capability.
167#[async_trait]
168pub trait SqlAccess: Send + Sync + 'static {
169    /// Canonical filesystem identity for this SQL database, when file-backed.
170    ///
171    /// Cross-resource operations use this only to derive advisory coordination
172    /// files outside SQLite. Only genuinely pathless, process-private
173    /// implementations may return `None`; every file-backed implementation
174    /// must expose its canonical path so cross-process exclusion cannot
175    /// silently degrade. The method performs no I/O.
176    fn database_path(&self) -> Option<std::path::PathBuf> {
177        None
178    }
179
180    /// Acquire a read-only connection from the pool.
181    async fn reader(&self) -> StorageResult<Box<dyn SqlReader>>;
182    /// Acquire a read-write connection from the pool.
183    async fn writer(&self) -> StorageResult<Box<dyn SqlWriter>>;
184
185    /// Run `op` as ONE atomic unit of work (ADR-067 Component A, Fork C
186    /// slice 2).
187    ///
188    /// Where a single-writer task is active (file-backed pool,
189    /// `KHIVE_WRITE_QUEUE=1`), `op` runs inside that task's one write
190    /// transaction for this request — no separate connection is opened, so
191    /// this call cannot compete with the writer task for SQLite's write
192    /// lock. For an in-memory pool, one pool connection guard is retained
193    /// throughout the transaction. Where no writer task applies to a
194    /// file-backed pool, `op` runs under a manual
195    /// `BEGIN IMMEDIATE`/`COMMIT`/`ROLLBACK` on a writer handle exactly like
196    /// calling [`Self::writer`] and driving the statements by hand — the
197    /// pre-ADR-067 behavior, preserved byte-for-byte on this path.
198    ///
199    /// **The atomic-unit suspend-free invariant (normative for every
200    /// caller):** `op`'s future must complete on its **first poll** — it may
201    /// issue only synchronous DML against the `&mut dyn SqlWriter` it is
202    /// handed and must never reach a real suspension point (no embedding
203    /// computation, no ANN warming, no service or channel `await`, no
204    /// network round-trip). Synchronous external work is forbidden too: no
205    /// filesystem/process/network I/O, sleeps, blocking waits, or unbounded CPU
206    /// may run between the transaction owner's `BEGIN IMMEDIATE` and
207    /// `COMMIT`. First-poll enforcement alone cannot detect those operations;
208    /// ADR-091's write-transaction audit table is the review-time guard, and a
209    /// new owner or caller must update that table before merge. On the
210    /// single-writer and in-memory paths this is enforced at runtime: the backend drives
211    /// `op` through a single-poll driver and
212    /// returns a typed error the instant the future is `Pending`, so a
213    /// violation fails loudly rather than corrupting state. On the file-backed
214    /// flag-off path (no writer task active) a suspending `op` would currently
215    /// *succeed* — that path drives `op` as an ordinary `.await` under a
216    /// manual transaction — so the invariant is a correctness contract this
217    /// trait asks every caller to uphold, not something the type system or
218    /// every code path enforces. Callers must not rely on the flag-off
219    /// path's tolerance; behavior must be identical (synchronous DML only)
220    /// regardless of whether the single-writer flag is on.
221    async fn atomic_unit(&self, op: AtomicUnitOp) -> StorageResult<Box<dyn Any + Send>>;
222}