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}