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