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}