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}