Skip to main content

khive_db/
sql_bridge.rs

1//! SqlAccess bridge: connects `ConnectionPool` to `khive_storage::SqlAccess`.
2//!
3//! Two modes:
4//! - **File-backed**: Opens standalone connections per reader/writer handle under pool-wide caps.
5//!   Cross-statement atomicity goes through `atomic_unit`, which drives a single
6//!   registered raw transaction span rather than a caller-held per-tx connection.
7//! - **Memory**: Uses pool-backed approach (acquire pool connection per-query inside `spawn_blocking`).
8
9use std::any::Any;
10use std::sync::Arc;
11use std::time::Instant;
12
13use async_trait::async_trait;
14
15use khive_storage::error::StorageError;
16use khive_storage::types::{PageRequest, SqlColumn, SqlRow, SqlStatement, SqlValue};
17use khive_storage::{AtomicUnitOp, StorageCapability};
18use tokio::sync::{OwnedSemaphorePermit, Semaphore};
19
20use crate::error::SqliteError;
21use crate::pool::ConnectionPool;
22
23// =============================================================================
24// Shared helpers
25// =============================================================================
26
27/// Convert a rusqlite `Row` into an owned `SqlRow`.
28fn row_to_sql_row(row: &rusqlite::Row<'_>, col_count: usize, col_names: &[String]) -> SqlRow {
29    #[cfg(test)]
30    ROW_CONVERSIONS.with(|count| count.set(count.get() + 1));
31
32    let mut columns = Vec::with_capacity(col_count);
33    for i in 0..col_count {
34        let value = match row.get_ref(i) {
35            Ok(rusqlite::types::ValueRef::Null) => SqlValue::Null,
36            Ok(rusqlite::types::ValueRef::Integer(v)) => SqlValue::Integer(v),
37            Ok(rusqlite::types::ValueRef::Real(v)) => SqlValue::Float(v),
38            Ok(rusqlite::types::ValueRef::Text(bytes)) => {
39                SqlValue::Text(String::from_utf8_lossy(bytes).into_owned())
40            }
41            Ok(rusqlite::types::ValueRef::Blob(bytes)) => SqlValue::Blob(bytes.to_vec()),
42            Err(_) => SqlValue::Null,
43        };
44        columns.push(SqlColumn {
45            name: col_names.get(i).cloned().unwrap_or_default(),
46            value,
47        });
48    }
49    SqlRow { columns }
50}
51
52#[cfg(test)]
53thread_local! {
54    static ROW_CONVERSIONS: std::cell::Cell<usize> = const { std::cell::Cell::new(0) };
55}
56
57/// Bind `SqlValue` parameters to a rusqlite statement.
58///
59/// `pub(crate)` (ADR-099 B3 r6 structural cut): reused by the pure
60/// `*_statement` builders in `stores::{entity,note,graph,text,vectors}` so
61/// that every store's async execution path and the ADR-099 `--atomic`
62/// prepare path bind params identically — one implementation, not two.
63pub(crate) fn bind_params(
64    stmt: &mut rusqlite::Statement<'_>,
65    params: &[SqlValue],
66) -> Result<(), rusqlite::Error> {
67    for (i, param) in params.iter().enumerate() {
68        let idx = i + 1; // rusqlite uses 1-based indexing
69        match param {
70            SqlValue::Null => stmt.raw_bind_parameter(idx, rusqlite::types::Null)?,
71            SqlValue::Bool(v) => stmt.raw_bind_parameter(idx, *v as i64)?,
72            SqlValue::Integer(v) => stmt.raw_bind_parameter(idx, *v)?,
73            SqlValue::Float(v) => stmt.raw_bind_parameter(idx, *v)?,
74            SqlValue::Text(v) => stmt.raw_bind_parameter(idx, v.as_str())?,
75            SqlValue::Blob(v) => stmt.raw_bind_parameter(idx, v.as_slice())?,
76            SqlValue::Json(v) => {
77                let s = serde_json::to_string(v).unwrap_or_default();
78                stmt.raw_bind_parameter(idx, s.as_str())?;
79            }
80            SqlValue::Uuid(v) => stmt.raw_bind_parameter(idx, v.to_string().as_str())?,
81            SqlValue::Timestamp(v) => {
82                stmt.raw_bind_parameter(idx, v.timestamp_micros())?;
83            }
84        }
85    }
86    Ok(())
87}
88
89/// Prepare exactly one [`SqlStatement`] SQL string.
90///
91/// rusqlite's `Connection::prepare` checks SQLite's returned tail and returns
92/// `rusqlite::Error::MultipleStatement` when the tail contains another
93/// executable statement (tail comments remain valid). Keeping this wrapper at
94/// the bridge boundary makes the single-statement `SqlStatement` contract
95/// explicit for queries and writes alike.
96fn prepare_sql_statement<'conn>(
97    conn: &'conn rusqlite::Connection,
98    sql: &str,
99) -> Result<rusqlite::Statement<'conn>, rusqlite::Error> {
100    conn.prepare(sql)
101}
102
103/// Prepare one [`SqlStatement`] through rusqlite's per-connection LRU cache
104/// while retaining the same single-statement tail validation as
105/// [`prepare_sql_statement`].
106fn prepare_cached_sql_statement<'conn>(
107    conn: &'conn rusqlite::Connection,
108    sql: &str,
109) -> Result<rusqlite::CachedStatement<'conn>, rusqlite::Error> {
110    conn.prepare_cached(sql)
111}
112
113/// A batch statement prepared once before execution. Ordinary SQL errors are
114/// retried in the execution phase so they still exercise the owning
115/// transaction's rollback path and can observe schema changes made by an
116/// earlier statement in the same batch; only `MultipleStatement` aborts
117/// preflight.
118enum PreparedBatchStatement<'conn> {
119    Ready(rusqlite::Statement<'conn>),
120    PrepareAtExecution,
121}
122
123/// Prepare each batch statement exactly once while rejecting an executable
124/// tail before any statement runs.
125fn prepare_batch_statements<'conn>(
126    conn: &'conn rusqlite::Connection,
127    statements: &[SqlStatement],
128) -> Result<Vec<PreparedBatchStatement<'conn>>, rusqlite::Error> {
129    let mut prepared = Vec::with_capacity(statements.len());
130    for statement in statements {
131        match prepare_sql_statement(conn, &statement.sql) {
132            Ok(statement) => prepared.push(PreparedBatchStatement::Ready(statement)),
133            Err(error @ rusqlite::Error::MultipleStatement) => return Err(error),
134            Err(_) => prepared.push(PreparedBatchStatement::PrepareAtExecution),
135        }
136    }
137    Ok(prepared)
138}
139
140/// Bind and execute handles returned by [`prepare_batch_statements`].
141fn execute_prepared_batch<'conn>(
142    conn: &'conn rusqlite::Connection,
143    prepared: Vec<PreparedBatchStatement<'conn>>,
144    statements: &[SqlStatement],
145) -> Result<u64, rusqlite::Error> {
146    debug_assert_eq!(prepared.len(), statements.len());
147    let mut total = 0u64;
148    for (prepared, statement) in prepared.into_iter().zip(statements) {
149        let mut prepared = match prepared {
150            PreparedBatchStatement::Ready(prepared) => prepared,
151            PreparedBatchStatement::PrepareAtExecution => {
152                prepare_sql_statement(conn, &statement.sql)?
153            }
154        };
155        bind_params(&mut prepared, &statement.params)?;
156        total += prepared.raw_execute()? as u64;
157    }
158    Ok(total)
159}
160
161/// SQL statement heads that are transaction control. `execute_batch` owns the
162/// `BEGIN`/`COMMIT` boundary for the whole batch (the standalone path wraps
163/// the list in its own `BEGIN IMMEDIATE`, and the queue-backed path runs
164/// inside the writer task's per-request transaction), so a caller-supplied
165/// statement that itself starts, ends, or branches a transaction can commit
166/// or roll back early and break the batch's all-or-nothing contract. Cached
167/// read-only handles use the same lexical classification to drive their
168/// separately admitted single-level read-transaction state machine below.
169/// `START` is classified as the alternate transaction-opening spelling so
170/// callers get a typed boundary error; `END` is SQLite's `COMMIT` spelling.
171const TRANSACTION_CONTROL_KEYWORDS: [&str; 7] = [
172    "BEGIN",
173    "START",
174    "COMMIT",
175    "END",
176    "ROLLBACK",
177    "SAVEPOINT",
178    "RELEASE",
179];
180
181/// Skip the same leading whitespace, UTF-8 BOMs, empty statements (`;`), and
182/// line/block comments SQLite accepts before an executable statement.
183fn skip_sqlite_empty_prefix(mut rest: &[u8]) -> &[u8] {
184    loop {
185        let mut idx = 0;
186        while idx < rest.len() && rest[idx].is_ascii_whitespace() {
187            idx += 1;
188        }
189        rest = &rest[idx..];
190        if let Some(tail) = rest.strip_prefix(b"\xEF\xBB\xBF") {
191            rest = tail;
192            continue;
193        }
194        if let Some(tail) = rest.strip_prefix(b";") {
195            rest = tail;
196            continue;
197        }
198        if let Some(tail) = rest.strip_prefix(b"--") {
199            let mut idx = 0;
200            while idx < tail.len() && tail[idx] != b'\n' {
201                idx += 1;
202            }
203            rest = if idx < tail.len() {
204                &tail[idx + 1..]
205            } else {
206                &[]
207            };
208            continue;
209        }
210        if let Some(tail) = rest.strip_prefix(b"/*") {
211            let mut idx = 0;
212            while idx + 1 < tail.len() && !(tail[idx] == b'*' && tail[idx + 1] == b'/') {
213                idx += 1;
214            }
215            rest = if idx + 1 < tail.len() {
216                &tail[idx + 2..]
217            } else {
218                &[]
219            };
220            continue;
221        }
222        break;
223    }
224    rest
225}
226
227/// Return one ASCII SQL token after SQLite whitespace/comments/BOM trivia.
228/// Empty-statement separators are deliberately not trivia here: callers use
229/// this only after the executable statement head has already been consumed.
230fn next_sqlite_token(mut rest: &[u8]) -> Option<(&[u8], &[u8])> {
231    loop {
232        let mut idx = 0;
233        while idx < rest.len() && rest[idx].is_ascii_whitespace() {
234            idx += 1;
235        }
236        rest = &rest[idx..];
237        if let Some(tail) = rest.strip_prefix(b"\xEF\xBB\xBF") {
238            rest = tail;
239            continue;
240        }
241        if let Some(tail) = rest.strip_prefix(b"--") {
242            let mut idx = 0;
243            while idx < tail.len() && tail[idx] != b'\n' {
244                idx += 1;
245            }
246            rest = if idx < tail.len() {
247                &tail[idx + 1..]
248            } else {
249                &[]
250            };
251            continue;
252        }
253        if let Some(tail) = rest.strip_prefix(b"/*") {
254            let mut idx = 0;
255            while idx + 1 < tail.len() && !(tail[idx] == b'*' && tail[idx + 1] == b'/') {
256                idx += 1;
257            }
258            rest = if idx + 1 < tail.len() {
259                &tail[idx + 2..]
260            } else {
261                &[]
262            };
263            continue;
264        }
265        break;
266    }
267
268    let len = rest
269        .iter()
270        .take_while(|byte| byte.is_ascii_alphanumeric() || **byte == b'_')
271        .count();
272    (len != 0).then_some((&rest[..len], &rest[len..]))
273}
274
275/// Return a transaction-control keyword and the bytes following it, if any.
276/// Matching is case-insensitive and requires a word boundary, so an identifier
277/// that merely starts with `begin` or `commit` never matches.
278fn transaction_control_parts(sql: &str) -> Option<(&'static str, &[u8])> {
279    let rest = skip_sqlite_empty_prefix(sql.as_bytes());
280    TRANSACTION_CONTROL_KEYWORDS
281        .iter()
282        .copied()
283        .find_map(|keyword| {
284            let kw = keyword.as_bytes();
285            if rest.len() < kw.len() || !rest[..kw.len()].eq_ignore_ascii_case(kw) {
286                return None;
287            }
288            let boundary = match rest.get(kw.len()) {
289                Some(next) => !(next.is_ascii_alphanumeric() || *next == b'_'),
290                None => true,
291            };
292            boundary.then_some((keyword, &rest[kw.len()..]))
293        })
294}
295
296/// Return the transaction-control keyword heading `sql`, if any.
297fn transaction_control_head(sql: &str) -> Option<&'static str> {
298    transaction_control_parts(sql).map(|(keyword, _)| keyword)
299}
300
301#[derive(Clone, Copy, Debug, PartialEq, Eq)]
302enum CachedReadTransactionControl {
303    /// `BEGIN`, `BEGIN TRANSACTION`, or their explicitly `DEFERRED` form.
304    BeginDeferred,
305    /// A transaction-ending control statement and its diagnostic keyword.
306    Finish(&'static str),
307    /// Transaction control that cannot be represented by the single-level
308    /// admitted read-transaction state machine.
309    Unsupported(&'static str),
310}
311
312/// Classify transaction control for a cached read-only connection.
313///
314/// A cached reader may own exactly one top-level deferred read transaction.
315/// Immediate/exclusive starts could reserve write-side locks, while nested
316/// savepoint/rollback-to controls would require a second lifecycle level, so
317/// both remain rejected. Batch/write paths continue using
318/// [`transaction_control_head`] and reject every variant without exception.
319fn cached_read_transaction_control(sql: &str) -> Option<CachedReadTransactionControl> {
320    let (keyword, tail) = transaction_control_parts(sql)?;
321    match keyword {
322        "BEGIN" => {
323            // Accept exactly `BEGIN`, `BEGIN TRANSACTION`, `BEGIN DEFERRED`,
324            // or `BEGIN DEFERRED TRANSACTION`, with no trailing tokens.
325            // SQLite's grammar also admits `BEGIN TRANSACTION <name>` (the
326            // name parses as an identifier and is ignored), so a mode keyword
327            // in that trailing position — `BEGIN TRANSACTION IMMEDIATE` —
328            // still parses, and classifying it by its first token alone would
329            // launder what reads as a write-reserving start into a deferred
330            // one. Every trailing token is therefore Unsupported — and the
331            // check cannot stop at `next_sqlite_token` returning `None`,
332            // because that tokenizer returns `None` for any non-identifier
333            // byte, not only end-of-input: a quoted or bracketed tail
334            // (`BEGIN TRANSACTION "IMMEDIATE"`, `[IMMEDIATE]`) would fall
335            // out of the loop and read as the end of an accepted form. After
336            // the accepted keywords, the remainder must reduce to nothing
337            // under the same trivia/empty-statement skipping SQLite applies
338            // (whitespace, comments, `;`), or the statement is Unsupported.
339            let mut rest = tail;
340            let mut saw_deferred = false;
341            let mut saw_transaction = false;
342            while let Some((token, next)) = next_sqlite_token(rest) {
343                if !saw_deferred && !saw_transaction && token.eq_ignore_ascii_case(b"DEFERRED") {
344                    saw_deferred = true;
345                } else if !saw_transaction && token.eq_ignore_ascii_case(b"TRANSACTION") {
346                    saw_transaction = true;
347                } else {
348                    return Some(CachedReadTransactionControl::Unsupported(keyword));
349                }
350                rest = next;
351            }
352            if !skip_sqlite_empty_prefix(rest).is_empty() {
353                return Some(CachedReadTransactionControl::Unsupported(keyword));
354            }
355            Some(CachedReadTransactionControl::BeginDeferred)
356        }
357        "COMMIT" | "END" => Some(CachedReadTransactionControl::Finish(keyword)),
358        "ROLLBACK" => {
359            let first = next_sqlite_token(tail);
360            let rollback_target = match first {
361                Some((token, rest)) if token.eq_ignore_ascii_case(b"TRANSACTION") => {
362                    next_sqlite_token(rest).map(|(token, _)| token)
363                }
364                Some((token, _)) => Some(token),
365                None => None,
366            };
367            if rollback_target.is_some_and(|token| token.eq_ignore_ascii_case(b"TO")) {
368                Some(CachedReadTransactionControl::Unsupported(keyword))
369            } else {
370                Some(CachedReadTransactionControl::Finish(keyword))
371            }
372        }
373        _ => Some(CachedReadTransactionControl::Unsupported(keyword)),
374    }
375}
376
377/// Reject transaction-control statements in `statements` with a typed
378/// [`StorageError::InvalidInput`] BEFORE anything executes, preserving the
379/// batch's all-or-nothing contract (see [`TRANSACTION_CONTROL_KEYWORDS`]).
380fn reject_transaction_control_statements(
381    statements: &[SqlStatement],
382    operation: &'static str,
383) -> khive_storage::types::StorageResult<()> {
384    for (index, statement) in statements.iter().enumerate() {
385        if let Some(keyword) = transaction_control_head(&statement.sql) {
386            return Err(StorageError::InvalidInput {
387                capability: StorageCapability::Sql,
388                operation: operation.into(),
389                message: format!(
390                    "statement at index {index} is transaction control ({keyword}); \
391                     execute_batch owns the BEGIN/COMMIT boundary for the whole \
392                     batch — remove transaction-control statements from the batch"
393                ),
394            });
395        }
396    }
397    Ok(())
398}
399
400/// One standalone-`execute_batch` failure, paired with the reason the handle
401/// was poisoned (dropped instead of restored), if it was.
402struct BatchFailure {
403    error: rusqlite::Error,
404    poison_reason: Option<BatchPoisonReason>,
405}
406
407#[derive(Clone, Copy, Debug, PartialEq, Eq)]
408enum BatchHandleDisposition {
409    Retain,
410    Poison,
411}
412
413/// Execute one standalone batch while every prepared statement remains scoped
414/// to the borrowed connection. The owned [`StandaloneHandle`] stays outside
415/// this helper, so it can be restored or dropped only after all statement
416/// borrows have ended.
417fn execute_standalone_batch(
418    conn: &rusqlite::Connection,
419    statements: &[SqlStatement],
420    origin: khive_storage::tx_registry::TxOrigin,
421) -> (BatchHandleDisposition, Result<u64, BatchFailure>) {
422    let prepared = match prepare_batch_statements(conn, statements) {
423        Ok(prepared) => prepared,
424        Err(error) => {
425            return (
426                BatchHandleDisposition::Retain,
427                Err(BatchFailure {
428                    error,
429                    poison_reason: None,
430                }),
431            );
432        }
433    };
434    if let Err(begin_error) = conn.execute_batch("BEGIN IMMEDIATE") {
435        // Busy/locked is transient contention (another writer held SQLite's
436        // write lock past `busy_timeout`); the connection itself is untouched,
437        // so the handle remains reusable. Any other failure leaves transaction
438        // state suspect and poisons the handle.
439        drop(prepared);
440        let (disposition, poison_reason) = if crate::timeout_sink::is_busy_or_locked(&begin_error) {
441            (BatchHandleDisposition::Retain, None)
442        } else {
443            tracing::warn!(
444                %begin_error,
445                "execute_batch: BEGIN IMMEDIATE failed non-transiently; \
446                 poisoning the standalone connection — the handle is \
447                 dropped and must be re-acquired"
448            );
449            (
450                BatchHandleDisposition::Poison,
451                Some(BatchPoisonReason::BeginFailed),
452            )
453        };
454        return (
455            disposition,
456            Err(BatchFailure {
457                error: begin_error,
458                poison_reason,
459            }),
460        );
461    }
462
463    // Registered only after BEGIN succeeds, and retained through COMMIT or
464    // ROLLBACK so the registry never reports a transaction as finished early.
465    let _tx_handle =
466        khive_storage::tx_registry::register_scoped(Some("execute_batch".to_string()), origin);
467    let result = (|| -> Result<u64, rusqlite::Error> {
468        let total = execute_prepared_batch(conn, prepared, statements)?;
469        conn.execute_batch("COMMIT")?;
470        Ok(total)
471    })();
472
473    let mut disposition = BatchHandleDisposition::Retain;
474    let mut poison_reason = None;
475    if let Err(error) = &result {
476        if let Err(rollback_error) = conn.execute_batch("ROLLBACK") {
477            // A failed ROLLBACK leaves the connection in an unknown
478            // transaction state. Preserve the original statement error while
479            // making the poison cause explicit to the caller.
480            tracing::warn!(
481                %error,
482                %rollback_error,
483                "execute_batch: ROLLBACK after statement failure failed; \
484                 poisoning the standalone connection — the handle is \
485                 dropped and must be re-acquired"
486            );
487            disposition = BatchHandleDisposition::Poison;
488            poison_reason = Some(BatchPoisonReason::RollbackFailed(rollback_error));
489        }
490    }
491
492    (
493        disposition,
494        result.map_err(|error| BatchFailure {
495            error,
496            poison_reason,
497        }),
498    )
499}
500
501#[derive(Debug)]
502enum BatchPoisonReason {
503    BeginFailed,
504    RollbackFailed(rusqlite::Error),
505}
506
507impl std::fmt::Display for BatchPoisonReason {
508    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
509        match self {
510            Self::BeginFailed => f.write_str(
511                "BEGIN IMMEDIATE failed non-transiently; connection transaction state is suspect",
512            ),
513            Self::RollbackFailed(error) => {
514                write!(f, "ROLLBACK after statement failure failed: {error}")
515            }
516        }
517    }
518}
519
520/// A `rusqlite::Error` whose display carries the poison context, so a
521/// poisoned handle is visible to the caller in the returned error instead of
522/// being discoverable only through later calls' generic "connection already
523/// consumed" failures.
524#[derive(Debug)]
525struct PoisonedBatchError {
526    original: rusqlite::Error,
527    poison_reason: BatchPoisonReason,
528}
529
530impl std::fmt::Display for PoisonedBatchError {
531    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
532        write!(
533            f,
534            "{}; original error: {}",
535            self.poison_reason, self.original
536        )
537    }
538}
539
540impl std::error::Error for PoisonedBatchError {
541    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
542        Some(&self.original)
543    }
544}
545
546fn prepare_bound_statement<'conn>(
547    conn: &'conn rusqlite::Connection,
548    statement: &SqlStatement,
549) -> Result<rusqlite::Statement<'conn>, rusqlite::Error> {
550    let mut stmt = prepare_sql_statement(conn, &statement.sql)?;
551    bind_params(&mut stmt, &statement.params)?;
552    Ok(stmt)
553}
554
555fn execute_prepared_query(
556    mut stmt: rusqlite::Statement<'_>,
557) -> Result<Vec<SqlRow>, rusqlite::Error> {
558    let col_count = stmt.column_count();
559    let col_names: Vec<String> = (0..col_count)
560        .map(|i| stmt.column_name(i).unwrap_or("").to_string())
561        .collect();
562
563    let mut rows = Vec::new();
564    let mut raw_rows = stmt.raw_query();
565    while let Some(row) = raw_rows.next()? {
566        rows.push(row_to_sql_row(row, col_count, &col_names));
567    }
568    Ok(rows)
569}
570
571fn execute_prepared_query_row(
572    mut stmt: rusqlite::Statement<'_>,
573) -> Result<Option<SqlRow>, rusqlite::Error> {
574    let col_count = stmt.column_count();
575    let col_names: Vec<String> = (0..col_count)
576        .map(|i| stmt.column_name(i).unwrap_or("").to_string())
577        .collect();
578
579    let mut raw_rows = stmt.raw_query();
580    Ok(raw_rows
581        .next()?
582        .map(|row| row_to_sql_row(row, col_count, &col_names)))
583}
584
585fn execute_prepared_query_page(
586    mut stmt: rusqlite::Statement<'_>,
587    page: &PageRequest,
588) -> Result<Vec<SqlRow>, rusqlite::Error> {
589    // A zero-limit page still prepares and binds the statement, so invalid
590    // SQL fails identically across every limit; it skips the row cursor
591    // entirely and returns no rows.
592    if page.limit == 0 {
593        return Ok(Vec::new());
594    }
595
596    let col_count = stmt.column_count();
597    let col_names: Vec<String> = (0..col_count)
598        .map(|i| stmt.column_name(i).unwrap_or("").to_string())
599        .collect();
600
601    let mut rows = Vec::new();
602    let mut offset = page.offset;
603    let mut remaining = u64::from(page.limit);
604    let mut raw_rows = stmt.raw_query();
605    // The bound covers owned Rust rows only — this function advances past
606    // `offset`, owns at most the caller-supplied `page.limit` rows, and drops
607    // the statement cursor immediately afterward (ADR-005's bounded-
608    // materialization amendment). Callers own choosing a sane limit.
609    // Engine work is the query plan's own cost, not O(offset + limit):
610    // SQLite still produces and discards `offset` rows, and an unindexed
611    // ORDER BY can force a full sort of the result set before the first row
612    // is stepped. Callers deep-paging a large result set should prefer
613    // keyset pagination over growing offsets.
614    while remaining > 0 {
615        let Some(row) = raw_rows.next()? else {
616            break;
617        };
618        if offset > 0 {
619            offset -= 1;
620            continue;
621        }
622        rows.push(row_to_sql_row(row, col_count, &col_names));
623        remaining -= 1;
624    }
625    Ok(rows)
626}
627
628/// Execute a query on a `rusqlite::Connection` and return owned rows.
629fn execute_query(
630    conn: &rusqlite::Connection,
631    statement: &SqlStatement,
632) -> Result<Vec<SqlRow>, rusqlite::Error> {
633    execute_prepared_query(prepare_bound_statement(conn, statement)?)
634}
635
636fn execute_query_row(
637    conn: &rusqlite::Connection,
638    statement: &SqlStatement,
639) -> Result<Option<SqlRow>, rusqlite::Error> {
640    execute_prepared_query_row(prepare_bound_statement(conn, statement)?)
641}
642
643fn execute_query_page(
644    conn: &rusqlite::Connection,
645    statement: &SqlStatement,
646    page: &PageRequest,
647) -> Result<Vec<SqlRow>, rusqlite::Error> {
648    execute_prepared_query_page(prepare_bound_statement(conn, statement)?, page)
649}
650
651/// SQLite's prepared-statement classifier is authoritative for the safety
652/// boundary. Row-producing DML (`UPDATE ... RETURNING`) and transaction
653/// control may be called through `SqlReader`, but they are admitted writes and
654/// must never register an interrupt target.
655fn statement_is_cancellable_read(stmt: &rusqlite::Statement<'_>, sql: &str) -> bool {
656    stmt.readonly() && transaction_control_head(sql).is_none()
657}
658
659fn execute_query_interruptibly(
660    scope: &crate::read_cancellation::InterruptibleReadScope,
661    conn: &rusqlite::Connection,
662    statement: &SqlStatement,
663    operation: &'static str,
664    rollback_interrupted_transaction: bool,
665    interruptible: bool,
666) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
667    let stmt = prepare_bound_statement(conn, statement)
668        .map_err(|error| map_rusqlite_err(error, operation))?;
669    if interruptible && statement_is_cancellable_read(&stmt, &statement.sql) {
670        scope.run_with_interrupted_cleanup(
671            conn,
672            move || {
673                execute_prepared_query(stmt).map_err(|error| map_rusqlite_err(error, operation))
674            },
675            || {
676                rollback_interrupted_read_transaction(
677                    conn,
678                    operation,
679                    rollback_interrupted_transaction,
680                )
681            },
682        )
683    } else {
684        scope.mark_write_committed()?;
685        execute_prepared_query(stmt).map_err(|error| map_rusqlite_err(error, operation))
686    }
687}
688
689fn execute_query_row_interruptibly(
690    scope: &crate::read_cancellation::InterruptibleReadScope,
691    conn: &rusqlite::Connection,
692    statement: &SqlStatement,
693    operation: &'static str,
694    rollback_interrupted_transaction: bool,
695    interruptible: bool,
696) -> khive_storage::types::StorageResult<Option<SqlRow>> {
697    let stmt = prepare_bound_statement(conn, statement)
698        .map_err(|error| map_rusqlite_err(error, operation))?;
699    if interruptible && statement_is_cancellable_read(&stmt, &statement.sql) {
700        scope.run_with_interrupted_cleanup(
701            conn,
702            move || {
703                execute_prepared_query_row(stmt).map_err(|error| map_rusqlite_err(error, operation))
704            },
705            || {
706                rollback_interrupted_read_transaction(
707                    conn,
708                    operation,
709                    rollback_interrupted_transaction,
710                )
711            },
712        )
713    } else {
714        scope.mark_write_committed()?;
715        execute_prepared_query_row(stmt).map_err(|error| map_rusqlite_err(error, operation))
716    }
717}
718
719fn execute_query_page_interruptibly(
720    scope: &crate::read_cancellation::InterruptibleReadScope,
721    conn: &rusqlite::Connection,
722    statement: &SqlStatement,
723    page: &PageRequest,
724    operation: &'static str,
725    rollback_interrupted_transaction: bool,
726    interruptible: bool,
727) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
728    let stmt = prepare_bound_statement(conn, statement)
729        .map_err(|error| map_rusqlite_err(error, operation))?;
730    if interruptible && statement_is_cancellable_read(&stmt, &statement.sql) {
731        scope.run_with_interrupted_cleanup(
732            conn,
733            move || {
734                execute_prepared_query_page(stmt, page)
735                    .map_err(|error| map_rusqlite_err(error, operation))
736            },
737            || {
738                rollback_interrupted_read_transaction(
739                    conn,
740                    operation,
741                    rollback_interrupted_transaction,
742                )
743            },
744        )
745    } else {
746        scope.mark_write_committed()?;
747        execute_prepared_query_page(stmt, page).map_err(|error| map_rusqlite_err(error, operation))
748    }
749}
750
751fn rollback_interrupted_read_transaction(
752    conn: &rusqlite::Connection,
753    operation: &'static str,
754    enabled: bool,
755) -> khive_storage::types::StorageResult<()> {
756    if !enabled || conn.is_autocommit() {
757        return Ok(());
758    }
759    conn.execute_batch("ROLLBACK")
760        .map_err(|error| map_rusqlite_err(error, operation))?;
761    if conn.is_autocommit() {
762        Ok(())
763    } else {
764        Err(StorageError::Transaction {
765            operation: operation.into(),
766            message: "interrupted read transaction rollback did not restore autocommit".into(),
767        })
768    }
769}
770
771/// Map a rusqlite error to `StorageError`.
772fn map_rusqlite_err(e: rusqlite::Error, op: &'static str) -> StorageError {
773    StorageError::driver(StorageCapability::Sql, op, e)
774}
775
776/// How an elapsed handle-slot deadline is classified. ADR-005 pins the
777/// raw-SQL reader admission paths (`sql_bridge.reader_open` and
778/// `sql_bridge.reader_operation`) to `StorageError::Timeout`; every other
779/// handle slot reports the typed `AdmissionTimeout`.
780#[derive(Clone, Copy)]
781enum SlotTimeoutClass {
782    Admission,
783    ReaderContract,
784}
785
786async fn acquire_handle_slot(
787    slots: Arc<Semaphore>,
788    timeout: std::time::Duration,
789    operation: &'static str,
790    class: SlotTimeoutClass,
791) -> Result<OwnedSemaphorePermit, StorageError> {
792    tokio::time::timeout(timeout, slots.acquire_owned())
793        .await
794        .map_err(|_| match class {
795            SlotTimeoutClass::Admission => StorageError::AdmissionTimeout {
796                operation: operation.into(),
797                timeout_ms: u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX),
798            },
799            SlotTimeoutClass::ReaderContract => StorageError::Timeout {
800                operation: operation.into(),
801            },
802        })?
803        .map_err(|error| StorageError::Pool {
804            operation: operation.into(),
805            message: error.to_string(),
806        })
807}
808
809// =============================================================================
810// Standalone connection readers/writers (file-backed databases)
811// =============================================================================
812
813fn open_standalone_reader(pool: &ConnectionPool) -> Result<rusqlite::Connection, StorageError> {
814    pool.open_standalone_reader()
815        .map_err(|error| StorageError::driver(StorageCapability::Sql, "open_reader", error))
816}
817
818fn open_standalone_writer(pool: &ConnectionPool) -> Result<rusqlite::Connection, StorageError> {
819    let config = pool.config();
820    let conn = pool
821        .open_standalone_writer()
822        .map_err(|e| StorageError::driver(StorageCapability::Sql, "open_writer", e))?;
823
824    conn.busy_timeout(config.busy_timeout)
825        .map_err(|e| map_rusqlite_err(e, "open_writer"))?;
826    conn.pragma_update(None, "cache_size", "-65536")
827        .map_err(|e| map_rusqlite_err(e, "open_writer"))?;
828    conn.pragma_update(None, "mmap_size", "1073741824")
829        .map_err(|e| map_rusqlite_err(e, "open_writer"))?;
830
831    Ok(conn)
832}
833
834/// Lift a standalone open onto the blocking thread pool while carrying the
835/// connection-cap permit into the same closure.
836///
837/// If the awaiting future is cancelled, the detached blocking closure owns
838/// both the connection result and the permit until it finishes, so the cap
839/// cannot be released while an open is still running.
840async fn open_standalone_on_blocking<F>(
841    pool: Arc<ConnectionPool>,
842    slot: OwnedSemaphorePermit,
843    operation: &'static str,
844    open: F,
845) -> khive_storage::types::StorageResult<(rusqlite::Connection, OwnedSemaphorePermit)>
846where
847    F: FnOnce(&ConnectionPool) -> Result<rusqlite::Connection, StorageError> + Send + 'static,
848{
849    tokio::task::spawn_blocking(move || open(&pool).map(|conn| (conn, slot)))
850        .await
851        .map_err(|e| StorageError::driver(StorageCapability::Sql, operation, e))?
852}
853
854/// [`open_standalone_reader`] lifted onto the blocking thread pool.
855///
856/// Opening a SQLite connection is filesystem I/O (open the file, read the
857/// database header) followed by pragmas executed through SQLite. No
858/// database lock is acquired at open itself — locks are taken on the first
859/// statement — but filesystem latency is unbounded, and this module already
860/// runs every other rusqlite call under `spawn_blocking`, so the open gets
861/// the same treatment instead of blocking an async worker thread. The reader
862/// permit is supplied to the helper and returned with the connection.
863async fn open_standalone_reader_on_blocking(
864    pool: Arc<ConnectionPool>,
865    slot: OwnedSemaphorePermit,
866) -> khive_storage::types::StorageResult<(rusqlite::Connection, OwnedSemaphorePermit)> {
867    open_standalone_on_blocking(pool, slot, "open_reader", open_standalone_reader).await
868}
869
870/// [`open_standalone_writer`] lifted onto the blocking thread pool; see
871/// [`open_standalone_reader_on_blocking`] for the blocking rationale.
872async fn open_standalone_writer_on_blocking(
873    pool: Arc<ConnectionPool>,
874    slot: OwnedSemaphorePermit,
875) -> khive_storage::types::StorageResult<(rusqlite::Connection, OwnedSemaphorePermit)> {
876    open_standalone_on_blocking(pool, slot, "open_writer", open_standalone_writer).await
877}
878
879// =============================================================================
880// File-backed: SqliteReader (standalone connection)
881// =============================================================================
882
883const CACHED_READ_TRANSACTION_LABEL: &str = "sql_bridge_cached_read_transaction";
884
885/// Admission and observability guards for one explicit cached-reader
886/// transaction. Both guards are installed only after SQLite accepts `BEGIN`
887/// and are retained together until SQLite reports autocommit again or the
888/// owning connection is closed.
889///
890/// Field order is deliberate: after [`StandaloneHandle::conn`] closes, the
891/// reader permit is returned before the registry evidence disappears. There
892/// is therefore no interval in which SQLite can still own the snapshot while
893/// the transaction is absent from `tx_registry`.
894struct CachedReadTransaction {
895    _slot: OwnedSemaphorePermit,
896    _tx_handle: khive_storage::tx_registry::TxHandle,
897    /// When this explicit `BEGIN` was admitted. Read on every subsequent
898    /// reuse of the owning cached-reader handle (#1846): a transaction whose
899    /// age has crossed `read_tx_max_age` is rolled back instead of being
900    /// extended by another call, bounding how long any one reader can pin
901    /// the WAL snapshot regardless of how many further requests it makes.
902    opened_at: Instant,
903}
904
905struct StandaloneHandle {
906    conn: rusqlite::Connection,
907    /// Present only for a standalone read-write handle, whose one-permit
908    /// connection budget remains handle-scoped. Read-only connections are
909    /// cached without a retained permit; each ordinary read operation acquires
910    /// the shared permit, while an explicit read transaction retains that
911    /// operation's permit in `read_transaction_slot` until its terminal call.
912    _retained_slot: Option<OwnedSemaphorePermit>,
913    /// Present only while a cached read-only connection owns one explicit
914    /// multi-call read transaction. Field order is load-bearing: Rust drops
915    /// `conn` before these guards, so cancellation or handle drop closes the
916    /// SQLite transaction before returning reader admission or deregistering
917    /// the transaction span.
918    read_transaction_slot: Option<CachedReadTransaction>,
919}
920
921impl StandaloneHandle {
922    /// Whether this is an idle-cacheable read-only connection rather than a
923    /// writer connection covered by the handle-scoped writer permit.
924    fn is_cached_reader(&self) -> bool {
925        self._retained_slot.is_none()
926    }
927
928    fn has_read_transaction(&self) -> bool {
929        self.read_transaction_slot.is_some()
930    }
931}
932
933struct SqliteReader {
934    handle: Option<StandaloneHandle>,
935    pool: Arc<ConnectionPool>,
936}
937
938async fn open_cached_reader_handle(
939    pool: Arc<ConnectionPool>,
940) -> khive_storage::types::StorageResult<StandaloneHandle> {
941    let open_slot = crate::await_request_read_phase(
942        "sql_bridge.reader_open",
943        acquire_handle_slot(
944            pool.sql_bridge_reader_slots(),
945            pool.config().checkout_timeout,
946            "sql_bridge.reader_open",
947            SlotTimeoutClass::ReaderContract,
948        ),
949    )
950    .await??;
951    let (conn, open_slot) = crate::await_request_read_phase(
952        "sql_bridge.reader_open",
953        open_standalone_reader_on_blocking(pool, open_slot),
954    )
955    .await??;
956    drop(open_slot);
957    Ok(StandaloneHandle {
958        conn,
959        _retained_slot: None,
960        read_transaction_slot: None,
961    })
962}
963
964/// Run one file-backed read while coupling its connection state to the active
965/// reader permit.
966///
967/// An ordinary cached read is admitted for one call and must return to
968/// autocommit before its permit is released. A successful top-level deferred
969/// `BEGIN` instead moves that operation permit into the handle; subsequent
970/// reads reuse it until `COMMIT`/`END`/`ROLLBACK` restores autocommit. The
971/// connection is declared before that retained permit, so dropping or
972/// cancelling the handle closes SQLite first and releases admission second.
973/// Standalone read-write handles are exempt from cached-reader transaction
974/// state because their writer permit remains held across legitimate manual
975/// transactions. Their reader-supertrait calls still take ordinary active-read
976/// admission; once the connection is outside autocommit, that acquisition and
977/// the SELECT are completion-preserving rather than request-cancellable.
978async fn execute_standalone_read<R, F>(
979    handle: &mut Option<StandaloneHandle>,
980    pool: Arc<ConnectionPool>,
981    operation: &'static str,
982    transaction_control: Option<CachedReadTransactionControl>,
983    read: F,
984) -> khive_storage::types::StorageResult<R>
985where
986    R: Send + 'static,
987    F: FnOnce(
988            &crate::read_cancellation::InterruptibleReadScope,
989            &rusqlite::Connection,
990            bool,
991            bool,
992        ) -> khive_storage::types::StorageResult<R>
993        + Send
994        + 'static,
995{
996    if handle.is_none() {
997        return Err(StorageError::Pool {
998            operation: operation.into(),
999            message: "connection already consumed".into(),
1000        });
1001    }
1002    let active_read_transaction = handle
1003        .as_ref()
1004        .is_some_and(|handle| handle.is_cached_reader() && handle.has_read_transaction());
1005    let completion_preserving_writer_transaction = handle
1006        .as_ref()
1007        .is_some_and(|handle| !handle.is_cached_reader() && !handle.conn.is_autocommit());
1008    let mut operation_slot = if active_read_transaction {
1009        None
1010    } else if completion_preserving_writer_transaction {
1011        // A read inside an admitted write transaction still counts against the
1012        // reader budget, but request cancellation cannot skip that admission
1013        // and strand the transaction between statements.
1014        Some(
1015            acquire_handle_slot(
1016                pool.sql_bridge_reader_slots(),
1017                pool.config().checkout_timeout,
1018                "sql_bridge.reader_operation",
1019                SlotTimeoutClass::ReaderContract,
1020            )
1021            .await?,
1022        )
1023    } else {
1024        Some(
1025            crate::await_request_read_phase(
1026                "sql_bridge.reader_operation",
1027                acquire_handle_slot(
1028                    pool.sql_bridge_reader_slots(),
1029                    pool.config().checkout_timeout,
1030                    "sql_bridge.reader_operation",
1031                    SlotTimeoutClass::ReaderContract,
1032                ),
1033            )
1034            .await??,
1035        )
1036    };
1037    let Some(owned_handle) = handle.take() else {
1038        return Err(StorageError::Pool {
1039            operation: operation.into(),
1040            message: "connection already consumed".into(),
1041        });
1042    };
1043    let origin = pool.origin();
1044    let read_tx_max_age = pool.config().read_tx_max_age;
1045    let (owned_handle, result) = crate::read_cancellation::run_interruptible_read(
1046        StorageCapability::Sql,
1047        operation,
1048        move |scope| {
1049            let mut owned_handle = owned_handle;
1050            let cached_reader = owned_handle.is_cached_reader();
1051            let entered_with_transaction = owned_handle.has_read_transaction();
1052            let entered_autocommit = owned_handle.conn.is_autocommit();
1053            let mut restore_handle = true;
1054            let mut result = if cached_reader && entered_with_transaction && entered_autocommit {
1055                // A connection cannot pin a snapshot in autocommit. Repair the
1056                // admission state before returning the invariant failure.
1057                drop(owned_handle.read_transaction_slot.take());
1058                Err(StorageError::InvalidInput {
1059                    capability: StorageCapability::Sql,
1060                    operation: operation.into(),
1061                    message: "cached read-only handle retained transaction admission after SQLite \
1062                          had already returned to autocommit; the stale permit was released"
1063                        .into(),
1064                })
1065            } else if cached_reader && !entered_with_transaction && !entered_autocommit {
1066                Err(StorageError::InvalidInput {
1067                    capability: StorageCapability::Sql,
1068                    operation: operation.into(),
1069                    message: "cached read-only handle entered the operation outside autocommit; \
1070                          its transaction was rolled back before releasing the reader permit"
1071                        .into(),
1072                })
1073            } else if cached_reader
1074                && entered_with_transaction
1075                && owned_handle
1076                    .read_transaction_slot
1077                    .as_ref()
1078                    .is_some_and(|tx| tx.opened_at.elapsed() >= read_tx_max_age)
1079            {
1080                // #1846: this handle's admitted read transaction has pinned a
1081                // WAL snapshot for at least `read_tx_max_age` — reject the
1082                // continuation and roll it back instead of extending the pin
1083                // for another call, regardless of what the caller asked for.
1084                crate::checkpoint::note_read_tx_max_age_eviction();
1085                match owned_handle.conn.execute_batch("ROLLBACK") {
1086                    Ok(()) if owned_handle.conn.is_autocommit() => {
1087                        drop(owned_handle.read_transaction_slot.take());
1088                        Err(StorageError::ReadTransactionAgeEvicted {
1089                            operation: operation.into(),
1090                            max_age_secs: read_tx_max_age.as_secs(),
1091                        })
1092                    }
1093                    Ok(()) => {
1094                        restore_handle = false;
1095                        Err(StorageError::ReadTransactionAgeEvictionCleanupFailed {
1096                            operation: operation.into(),
1097                            max_age_secs: read_tx_max_age.as_secs(),
1098                            message: "rollback did not restore autocommit".into(),
1099                        })
1100                    }
1101                    Err(error) => {
1102                        restore_handle = false;
1103                        Err(StorageError::ReadTransactionAgeEvictionCleanupFailed {
1104                            operation: operation.into(),
1105                            max_age_secs: read_tx_max_age.as_secs(),
1106                            message: format!("rollback failed: {error}"),
1107                        })
1108                    }
1109                }
1110            } else if cached_reader && entered_with_transaction {
1111                match transaction_control {
1112                    None | Some(CachedReadTransactionControl::Finish(_)) => {
1113                        read(scope, &owned_handle.conn, true, true)
1114                    }
1115                    Some(CachedReadTransactionControl::BeginDeferred) => {
1116                        Err(StorageError::InvalidInput {
1117                            capability: StorageCapability::Sql,
1118                            operation: operation.into(),
1119                            message: "cached read-only handle already owns an admitted read \
1120                                  transaction; nested BEGIN is not supported"
1121                                .into(),
1122                        })
1123                    }
1124                    Some(CachedReadTransactionControl::Unsupported(keyword)) => {
1125                        Err(StorageError::InvalidInput {
1126                            capability: StorageCapability::Sql,
1127                            operation: operation.into(),
1128                            message: format!(
1129                                "cached read-only transaction does not support nested or \
1130                             write-locking transaction control ({keyword})"
1131                            ),
1132                        })
1133                    }
1134                }
1135            } else if cached_reader {
1136                match transaction_control {
1137                    None | Some(CachedReadTransactionControl::BeginDeferred) => {
1138                        read(scope, &owned_handle.conn, false, true)
1139                    }
1140                    Some(CachedReadTransactionControl::Finish(keyword))
1141                    | Some(CachedReadTransactionControl::Unsupported(keyword)) => {
1142                        Err(StorageError::InvalidInput {
1143                            capability: StorageCapability::Sql,
1144                            operation: operation.into(),
1145                            message: format!(
1146                                "cached read-only handle has no admitted transaction for \
1147                             transaction control ({keyword})"
1148                            ),
1149                        })
1150                    }
1151                }
1152            } else {
1153                read(scope, &owned_handle.conn, false, entered_autocommit)
1154            };
1155
1156            if scope.cleanup_failed() {
1157                // A connection-global callback that could not be removed may
1158                // fire for an unrelated future borrower. Closing this handle
1159                // is the only safe recovery; its transaction, if any, ends
1160                // before reader admission is released below.
1161                restore_handle = false;
1162            }
1163
1164            // An interrupted explicit read transaction must never be restored to
1165            // the cached handle: it may still own a WAL snapshot and SQLite's
1166            // interrupted flag applies to the transaction as a whole. Roll back
1167            // before releasing its retained admission; if rollback cannot prove
1168            // autocommit, discard the connection.
1169            if cached_reader
1170                && matches!(result, Err(StorageError::Timeout { .. }))
1171                && !owned_handle.conn.is_autocommit()
1172            {
1173                match owned_handle.conn.execute_batch("ROLLBACK") {
1174                    Ok(()) if owned_handle.conn.is_autocommit() => {
1175                        drop(owned_handle.read_transaction_slot.take());
1176                    }
1177                    Ok(()) => {
1178                        restore_handle = false;
1179                        result = Err(StorageError::Transaction {
1180                            operation: operation.into(),
1181                            message:
1182                                "interrupted read transaction rollback did not restore autocommit; \
1183                                  the connection was discarded"
1184                                    .into(),
1185                        });
1186                    }
1187                    Err(error) => {
1188                        restore_handle = false;
1189                        result = Err(StorageError::Transaction {
1190                            operation: operation.into(),
1191                            message: format!(
1192                                "failed to roll back interrupted read transaction ({error}); \
1193                             the connection was discarded"
1194                            ),
1195                        });
1196                    }
1197                }
1198            }
1199
1200            if cached_reader && entered_with_transaction {
1201                if owned_handle.conn.is_autocommit() {
1202                    // SQLite has ended the snapshot; release only after observing
1203                    // that terminal state. This also fails closed if an ordinary
1204                    // statement unexpectedly ended the transaction.
1205                    drop(owned_handle.read_transaction_slot.take());
1206                    if result.is_ok()
1207                        && !matches!(
1208                            transaction_control,
1209                            Some(CachedReadTransactionControl::Finish(_))
1210                        )
1211                    {
1212                        result = Err(StorageError::InvalidInput {
1213                            capability: StorageCapability::Sql,
1214                            operation: operation.into(),
1215                            message: "cached read-only operation unexpectedly ended its admitted \
1216                                  transaction; reader admission was released after autocommit"
1217                                .into(),
1218                        });
1219                    }
1220                } else if result.is_ok()
1221                    && matches!(
1222                        transaction_control,
1223                        Some(CachedReadTransactionControl::Finish(_))
1224                    )
1225                {
1226                    result = Err(StorageError::InvalidInput {
1227                        capability: StorageCapability::Sql,
1228                        operation: operation.into(),
1229                        message: "transaction-ending control completed but the cached reader \
1230                              remained outside autocommit; its reader permit remains retained"
1231                            .into(),
1232                    });
1233                }
1234            } else if cached_reader
1235                && entered_autocommit
1236                && matches!(
1237                    transaction_control,
1238                    Some(CachedReadTransactionControl::BeginDeferred)
1239                )
1240                && result.is_ok()
1241            {
1242                if owned_handle.conn.is_autocommit() {
1243                    result = Err(StorageError::InvalidInput {
1244                        capability: StorageCapability::Sql,
1245                        operation: operation.into(),
1246                        message: "deferred BEGIN completed without opening a read transaction"
1247                            .into(),
1248                    });
1249                } else {
1250                    match operation_slot.take() {
1251                        Some(slot) => {
1252                            let tx_handle = khive_storage::tx_registry::register_scoped(
1253                                Some(CACHED_READ_TRANSACTION_LABEL.to_string()),
1254                                origin.clone(),
1255                            );
1256                            owned_handle.read_transaction_slot = Some(CachedReadTransaction {
1257                                _slot: slot,
1258                                _tx_handle: tx_handle,
1259                                opened_at: Instant::now(),
1260                            });
1261                        }
1262                        None => {
1263                            result = Err(StorageError::Pool {
1264                                operation: operation.into(),
1265                                message: "successful cached-reader BEGIN had no operation permit; \
1266                                      its transaction was rolled back before returning"
1267                                    .into(),
1268                            });
1269                        }
1270                    }
1271                }
1272            }
1273
1274            // Any non-autocommit state without its retained admission is stale or
1275            // was opened by a statement the transaction classifier did not admit.
1276            if cached_reader
1277                && owned_handle.read_transaction_slot.is_none()
1278                && !owned_handle.conn.is_autocommit()
1279            {
1280                match owned_handle.conn.execute_batch("ROLLBACK") {
1281                    Ok(()) if owned_handle.conn.is_autocommit() => {
1282                        if result.is_ok() {
1283                            result = Err(StorageError::InvalidInput {
1284                                capability: StorageCapability::Sql,
1285                                operation: operation.into(),
1286                                message: "cached read-only operation left the connection outside \
1287                                      autocommit; its transaction was rolled back before \
1288                                      releasing the reader permit"
1289                                    .into(),
1290                            });
1291                        }
1292                    }
1293                    Ok(()) => {
1294                        restore_handle = false;
1295                        result = Err(StorageError::Transaction {
1296                            operation: operation.into(),
1297                            message: "ROLLBACK completed but the cached reader remained outside \
1298                                  autocommit; the connection was discarded before releasing \
1299                                  the reader permit"
1300                                .into(),
1301                        });
1302                    }
1303                    Err(error) => {
1304                        restore_handle = false;
1305                        result = Err(StorageError::Transaction {
1306                            operation: operation.into(),
1307                            message: format!(
1308                            "failed to roll back a cached reader outside autocommit ({error}); \
1309                             the connection was discarded before releasing the reader permit"
1310                        ),
1311                        });
1312                    }
1313                }
1314            }
1315
1316            let owned_handle = if restore_handle {
1317                Some(owned_handle)
1318            } else {
1319                // Closing the poisoned connection ends any remaining transaction.
1320                // This must precede the active-reader permit release below.
1321                drop(owned_handle);
1322                None
1323            };
1324            // For ordinary reads and rejected controls this is the operation
1325            // permit. A successful BEGIN moved it into `owned_handle`; poisoned
1326            // handles were closed above before this remaining permit is released.
1327            drop(operation_slot);
1328            Ok((owned_handle, result))
1329        },
1330    )
1331    .await?;
1332    *handle = owned_handle;
1333    result
1334}
1335
1336#[async_trait]
1337impl khive_storage::SqlReader for SqliteReader {
1338    async fn query_row(
1339        &mut self,
1340        statement: SqlStatement,
1341    ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
1342        let transaction_control = cached_read_transaction_control(&statement.sql);
1343        execute_standalone_read(
1344            &mut self.handle,
1345            Arc::clone(&self.pool),
1346            "query_row",
1347            transaction_control,
1348            move |scope, conn, rollback, interruptible| {
1349                execute_query_row_interruptibly(
1350                    scope,
1351                    conn,
1352                    &statement,
1353                    "query_row",
1354                    rollback,
1355                    interruptible,
1356                )
1357            },
1358        )
1359        .await
1360    }
1361
1362    async fn query_all(
1363        &mut self,
1364        statement: SqlStatement,
1365    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1366        let transaction_control = cached_read_transaction_control(&statement.sql);
1367        execute_standalone_read(
1368            &mut self.handle,
1369            Arc::clone(&self.pool),
1370            "query_all",
1371            transaction_control,
1372            move |scope, conn, rollback, interruptible| {
1373                execute_query_interruptibly(
1374                    scope,
1375                    conn,
1376                    &statement,
1377                    "query_all",
1378                    rollback,
1379                    interruptible,
1380                )
1381            },
1382        )
1383        .await
1384    }
1385
1386    async fn query_page(
1387        &mut self,
1388        statement: SqlStatement,
1389        page: PageRequest,
1390    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1391        let transaction_control = cached_read_transaction_control(&statement.sql);
1392        execute_standalone_read(
1393            &mut self.handle,
1394            Arc::clone(&self.pool),
1395            "query_page",
1396            transaction_control,
1397            move |scope, conn, rollback, interruptible| {
1398                execute_query_page_interruptibly(
1399                    scope,
1400                    conn,
1401                    &statement,
1402                    &page,
1403                    "query_page",
1404                    rollback,
1405                    interruptible,
1406                )
1407            },
1408        )
1409        .await
1410    }
1411
1412    async fn query_scalar(
1413        &mut self,
1414        statement: SqlStatement,
1415    ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
1416        let row = self.query_row(statement).await?;
1417        Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
1418    }
1419
1420    async fn explain(
1421        &mut self,
1422        statement: SqlStatement,
1423    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1424        let explain_stmt = SqlStatement {
1425            sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
1426            params: statement.params,
1427            label: statement.label,
1428        };
1429        self.query_all(explain_stmt).await
1430    }
1431}
1432
1433// =============================================================================
1434// File-backed: SqliteWriter (standalone connection)
1435// =============================================================================
1436
1437struct SqliteWriter {
1438    /// `None` at construction when a `WriterTaskHandle` was obtained (ADR-136
1439    /// D1 gate 1: queue-first `writer()` skips the standalone open in that
1440    /// case). Lazily opened by [`SqliteWriter::ensure_conn`] on first read
1441    /// (`query_row`/`query_all`/`query_page`) — production callers do read
1442    /// through a `writer()` handle (e.g. `khive-pack-comm` and
1443    /// `khive-pack-gtd`), so the `SqlReader` supertrait's lazy-open path is
1444    /// live, not just a capability formality. An eagerly opened read-write
1445    /// connection retains its one-permit writer budget for the handle's
1446    /// lifetime. A lazily opened read-only connection is cached without a
1447    /// permit; every ordinary query acquires one for exactly the blocking
1448    /// operation. An explicit deferred read transaction retains its opening
1449    /// permit across queries until terminal control. In either state,
1450    /// cancellation cannot release admission before SQLite finishes/closes.
1451    handle: Option<StandaloneHandle>,
1452    /// ADR-067 Component A: when the write queue is enabled, `execute_batch`
1453    /// routes the whole caller-supplied statement list through the
1454    /// single-writer task instead of opening its own `BEGIN IMMEDIATE` on
1455    /// the standalone connection. `None` when the flag is off or no writer
1456    /// task is available
1457    /// (best-effort — degrades to the standalone-connection path below).
1458    writer_task: Option<crate::writer_task::WriterTaskHandle>,
1459    /// The origin (ADR-091 backend-scoped attribution) of the pool this
1460    /// standalone connection was opened against.
1461    origin: khive_storage::tx_registry::TxOrigin,
1462    /// This connection's pool's writer-timeout sink identity (`db_label`),
1463    /// captured at construction so the standalone-path busy/locked mapping
1464    /// below doesn't need a `&ConnectionPool` reference to report against.
1465    db: String,
1466    /// Needed only for [`Self::ensure_conn`]'s lazy standalone-connection
1467    /// open when `handle` was skipped at construction (`writer_task`
1468    /// present).
1469    pool: Arc<ConnectionPool>,
1470}
1471
1472impl SqliteWriter {
1473    /// Return the open handle if present, else lazily open a standalone
1474    /// **read-only** handle now. See the `handle` field doc comment for when
1475    /// this lazy path is reached — it is only reached from the `SqlReader`
1476    /// methods (`query_row`/`query_all`/`query_page`), never from a
1477    /// `SqlWriter` method: every `SqlWriter` method on this type either
1478    /// routes through `writer_task` (when present, the same condition that
1479    /// causes `handle` to start `None`) or uses the writer connection opened
1480    /// eagerly at construction (when `writer_task` is absent). Opening a
1481    /// read-only connection here (ADR-136 D1 gate 3 amendment) closes the
1482    /// gap where a caller holding a queue-backed writer handle could issue
1483    /// an `INSERT ... RETURNING` (or any other DML) through `query_row` /
1484    /// `query_all` / `query_page` and have it execute on an untracked
1485    /// read-write connection, outside the `WriterTask` — SQLite rejects DML
1486    /// against a read-only connection instead. The lazy open acquires a
1487    /// pool-wide READER permit rather than the one-permit writer budget:
1488    /// the connection cannot write, and charging it against the writer
1489    /// budget would let a queue-backed handle's reads block standalone
1490    /// writers. The reader permit covers the open itself, then each ordinary
1491    /// query independently or one admitted explicit read-transaction lifetime.
1492    /// Associated function rather than a `&self` method so the caller can
1493    /// pass a cloned pool handle and keep no borrow of `SqliteWriter` live
1494    /// across the permit await — `&SqliteWriter` is not `Sync` (the held
1495    /// `rusqlite::Connection` is not), and the `SqlReader` futures must
1496    /// stay `Send`.
1497    async fn ensure_conn(
1498        pool: Arc<ConnectionPool>,
1499    ) -> khive_storage::types::StorageResult<StandaloneHandle> {
1500        open_cached_reader_handle(pool).await
1501    }
1502}
1503
1504#[async_trait]
1505impl khive_storage::SqlReader for SqliteWriter {
1506    async fn query_row(
1507        &mut self,
1508        statement: SqlStatement,
1509    ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
1510        if self.handle.is_none() && self.writer_task.is_some() {
1511            self.handle = Some(Self::ensure_conn(Arc::clone(&self.pool)).await?);
1512        }
1513        let transaction_control = cached_read_transaction_control(&statement.sql);
1514        execute_standalone_read(
1515            &mut self.handle,
1516            Arc::clone(&self.pool),
1517            "writer.query_row",
1518            transaction_control,
1519            move |scope, conn, rollback, interruptible| {
1520                execute_query_row_interruptibly(
1521                    scope,
1522                    conn,
1523                    &statement,
1524                    "writer.query_row",
1525                    rollback,
1526                    interruptible,
1527                )
1528            },
1529        )
1530        .await
1531    }
1532
1533    async fn query_all(
1534        &mut self,
1535        statement: SqlStatement,
1536    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1537        if self.handle.is_none() && self.writer_task.is_some() {
1538            self.handle = Some(Self::ensure_conn(Arc::clone(&self.pool)).await?);
1539        }
1540        let transaction_control = cached_read_transaction_control(&statement.sql);
1541        execute_standalone_read(
1542            &mut self.handle,
1543            Arc::clone(&self.pool),
1544            "writer.query_all",
1545            transaction_control,
1546            move |scope, conn, rollback, interruptible| {
1547                execute_query_interruptibly(
1548                    scope,
1549                    conn,
1550                    &statement,
1551                    "writer.query_all",
1552                    rollback,
1553                    interruptible,
1554                )
1555            },
1556        )
1557        .await
1558    }
1559
1560    async fn query_page(
1561        &mut self,
1562        statement: SqlStatement,
1563        page: PageRequest,
1564    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1565        if self.handle.is_none() && self.writer_task.is_some() {
1566            self.handle = Some(Self::ensure_conn(Arc::clone(&self.pool)).await?);
1567        }
1568        let transaction_control = cached_read_transaction_control(&statement.sql);
1569        execute_standalone_read(
1570            &mut self.handle,
1571            Arc::clone(&self.pool),
1572            "writer.query_page",
1573            transaction_control,
1574            move |scope, conn, rollback, interruptible| {
1575                execute_query_page_interruptibly(
1576                    scope,
1577                    conn,
1578                    &statement,
1579                    &page,
1580                    "writer.query_page",
1581                    rollback,
1582                    interruptible,
1583                )
1584            },
1585        )
1586        .await
1587    }
1588
1589    async fn query_scalar(
1590        &mut self,
1591        statement: SqlStatement,
1592    ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
1593        let row = khive_storage::SqlReader::query_row(self, statement).await?;
1594        Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
1595    }
1596
1597    async fn explain(
1598        &mut self,
1599        statement: SqlStatement,
1600    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1601        let explain_stmt = SqlStatement {
1602            sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
1603            params: statement.params,
1604            label: statement.label,
1605        };
1606        khive_storage::SqlReader::query_all(self, explain_stmt).await
1607    }
1608}
1609
1610#[async_trait]
1611impl khive_storage::SqlWriter for SqliteWriter {
1612    async fn execute(
1613        &mut self,
1614        statement: SqlStatement,
1615    ) -> khive_storage::types::StorageResult<u64> {
1616        // ADR-067 Component A (Fork C slice 2): a single statement is
1617        // self-contained, just like `execute_batch`'s full statement list —
1618        // transaction-control rejection remains an `execute_batch` contract;
1619        // this primitive is also used by internal atomic transaction owners.
1620        // route it through the writer task when available. `self.handle` is
1621        // left untouched so a subsequent `execute`/`execute_script` call on
1622        // this same handle still works over the standalone connection.
1623        if let Some(writer_task) = self.writer_task.clone() {
1624            return writer_task
1625                .send_bounded(move |conn| {
1626                    let mut stmt = prepare_cached_sql_statement(conn, &statement.sql)
1627                        .map_err(|e| map_rusqlite_err(e, "execute"))?;
1628                    bind_params(&mut stmt, &statement.params)
1629                        .map_err(|e| map_rusqlite_err(e, "execute"))?;
1630                    let affected = stmt
1631                        .raw_execute()
1632                        .map_err(|e| map_rusqlite_err(e, "execute"))?;
1633                    Ok(affected as u64)
1634                })
1635                .await;
1636        }
1637
1638        let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
1639            operation: "execute".into(),
1640            message: "connection already consumed".into(),
1641        })?;
1642        let (handle, result) = tokio::task::spawn_blocking(move || {
1643            let res = (|| -> Result<usize, rusqlite::Error> {
1644                let mut stmt = prepare_cached_sql_statement(&handle.conn, &statement.sql)?;
1645                bind_params(&mut stmt, &statement.params)?;
1646                stmt.raw_execute()
1647            })();
1648            (handle, res)
1649        })
1650        .await
1651        .map_err(|e| StorageError::driver(StorageCapability::Sql, "execute", e))?;
1652        self.handle = Some(handle);
1653        let affected = result.map_err(|e| {
1654            crate::timeout_sink::maybe_emit_busy(
1655                &self.db,
1656                crate::timeout_sink::Site::StandaloneSqlBridge,
1657                &e,
1658            );
1659            map_rusqlite_err(e, "execute")
1660        })?;
1661        Ok(affected as u64)
1662    }
1663
1664    async fn execute_batch(
1665        &mut self,
1666        statements: Vec<SqlStatement>,
1667    ) -> khive_storage::types::StorageResult<u64> {
1668        // ADR-067 Component A: this call is self-contained (the full statement
1669        // list is supplied up front and the whole thing commits or rolls back
1670        // as one unit) — unlike `writer()`'s live incrementally-driven handle,
1671        // it maps cleanly onto a single `WriteRequest`. Route it through the
1672        // writer task when available; `self.handle` is left untouched so a
1673        // subsequent `execute`/`execute_script` call on this same handle still
1674        // works over the standalone connection (that dispatch is unmigrated —
1675        // see `SqlBridge::writer()`).
1676        //
1677        // Both paths reject transaction-control statements BEFORE executing
1678        // anything: the queue-backed branch runs inside the writer task's own
1679        // `BEGIN IMMEDIATE` (a caller `COMMIT` there would close the task's
1680        // transaction and terminate the writer task), and the standalone
1681        // branch below wraps the list in its own `BEGIN IMMEDIATE` (a caller
1682        // `COMMIT` would commit early and break all-or-nothing).
1683        reject_transaction_control_statements(&statements, "execute_batch")?;
1684        if let Some(writer_task) = self.writer_task.clone() {
1685            return writer_task
1686                .send_bounded(move |conn| {
1687                    let prepared = prepare_batch_statements(conn, &statements)
1688                        .map_err(|e| map_rusqlite_err(e, "execute_batch"))?;
1689                    execute_prepared_batch(conn, prepared, &statements)
1690                        .map_err(|e| map_rusqlite_err(e, "execute_batch"))
1691                })
1692                .await;
1693        }
1694
1695        let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
1696            operation: "execute_batch".into(),
1697            message: "connection already consumed".into(),
1698        })?;
1699        let origin = self.origin.clone();
1700        let (handle, result) = tokio::task::spawn_blocking(move || {
1701            let (disposition, result) = execute_standalone_batch(&handle.conn, &statements, origin);
1702            let retained = match disposition {
1703                BatchHandleDisposition::Retain => Some(handle),
1704                BatchHandleDisposition::Poison => None,
1705            };
1706            (retained, result)
1707        })
1708        .await
1709        .map_err(|e| StorageError::driver(StorageCapability::Sql, "execute_batch", e))?;
1710        self.handle = handle;
1711        result.map_err(|failure| {
1712            crate::timeout_sink::maybe_emit_busy(
1713                &self.db,
1714                crate::timeout_sink::Site::StandaloneSqlBridge,
1715                &failure.error,
1716            );
1717            match failure.poison_reason {
1718                Some(poison_reason) => StorageError::driver(
1719                    StorageCapability::Sql,
1720                    "execute_batch",
1721                    PoisonedBatchError {
1722                        original: failure.error,
1723                        poison_reason,
1724                    },
1725                ),
1726                None => map_rusqlite_err(failure.error, "execute_batch"),
1727            }
1728        })
1729    }
1730
1731    async fn execute_script(&mut self, script: String) -> khive_storage::types::StorageResult<()> {
1732        // ADR-067 Component A (Fork C slice 2): the script text is
1733        // self-contained (supplied up front, runs as one unit), just like
1734        // `execute_batch` — route it through the writer task when
1735        // available. `self.handle` is left untouched so a subsequent
1736        // `execute`/`execute_script` call on this same handle still works
1737        // over the standalone connection. Callers must supply a DML-only
1738        // script (no bare `BEGIN`/`COMMIT`/`ROLLBACK`) on the flag-on path,
1739        // since it runs inside the writer task's own transaction — same
1740        // Boundary: transaction-control rejection is an `execute_batch`
1741        // contract; this raw script path is internal/migration-only. The
1742        // queue-backed branch still requires a DML-only script because it
1743        // runs inside the writer task's transaction.
1744        if let Some(writer_task) = self.writer_task.clone() {
1745            return writer_task
1746                .send_bounded(move |conn| {
1747                    conn.execute_batch(&script)
1748                        .map_err(|e| map_rusqlite_err(e, "execute_script"))
1749                })
1750                .await;
1751        }
1752
1753        let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
1754            operation: "execute_script".into(),
1755            message: "connection already consumed".into(),
1756        })?;
1757        let (handle, result) = tokio::task::spawn_blocking(move || {
1758            let res = handle.conn.execute_batch(&script);
1759            (handle, res)
1760        })
1761        .await
1762        .map_err(|e| StorageError::driver(StorageCapability::Sql, "execute_script", e))?;
1763        self.handle = Some(handle);
1764        result.map_err(|e| {
1765            crate::timeout_sink::maybe_emit_busy(
1766                &self.db,
1767                crate::timeout_sink::Site::StandaloneSqlBridge,
1768                &e,
1769            );
1770            map_rusqlite_err(e, "execute_script")
1771        })
1772    }
1773
1774    async fn execute_script_top_level(
1775        &mut self,
1776        script: String,
1777    ) -> khive_storage::types::StorageResult<()> {
1778        // Boundary: this internal maintenance/migration path deliberately
1779        // bypasses the `execute_batch` transaction-control rejection.
1780        // ADR-067 Component A: unlike
1781        // `execute_script`, this must NOT run inside the writer task's
1782        // per-request `BEGIN IMMEDIATE` — statements such as VACUUM are
1783        // rejected by SQLite inside any open transaction. Route through
1784        // `WriterTaskHandle::send_top_level`, which still serializes this
1785        // call through the single writer owner but skips the transaction
1786        // wrap entirely.
1787        if let Some(writer_task) = self.writer_task.clone() {
1788            return writer_task
1789                .send_top_level_bounded(move |conn| {
1790                    conn.execute_batch(&script)
1791                        .map_err(|e| map_rusqlite_err(e, "execute_script_top_level"))
1792                })
1793                .await;
1794        }
1795
1796        // Flag off / no writer task: identical to `execute_script`'s own
1797        // flag-off path — a bare `execute_batch` on the standalone
1798        // connection, already transaction-free.
1799        let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
1800            operation: "execute_script_top_level".into(),
1801            message: "connection already consumed".into(),
1802        })?;
1803        let (handle, result) = tokio::task::spawn_blocking(move || {
1804            let res = handle.conn.execute_batch(&script);
1805            (handle, res)
1806        })
1807        .await
1808        .map_err(|e| StorageError::driver(StorageCapability::Sql, "execute_script_top_level", e))?;
1809        self.handle = Some(handle);
1810        result.map_err(|e| {
1811            crate::timeout_sink::maybe_emit_busy(
1812                &self.db,
1813                crate::timeout_sink::Site::StandaloneSqlBridge,
1814                &e,
1815            );
1816            map_rusqlite_err(e, "execute_script_top_level")
1817        })
1818    }
1819}
1820
1821// =============================================================================
1822// Pool-backed reader/writer (in-memory databases)
1823// =============================================================================
1824
1825async fn run_pool_reader_query<T, F>(
1826    pool: Arc<ConnectionPool>,
1827    operation: &'static str,
1828    query: F,
1829) -> khive_storage::types::StorageResult<T>
1830where
1831    T: Send + 'static,
1832    F: FnOnce(
1833            &crate::read_cancellation::InterruptibleReadScope,
1834            &rusqlite::Connection,
1835        ) -> khive_storage::types::StorageResult<T>
1836        + Send
1837        + 'static,
1838{
1839    crate::read_cancellation::run_interruptible_read(
1840        StorageCapability::Sql,
1841        operation,
1842        move |scope| {
1843            // Checkout tri-state (cancelled -> Timeout, admission expiry ->
1844            // retryable AdmissionTimeout, other -> Driver) lives in ONE place:
1845            // `ConnectionPool::resolve_reader_checkout`.
1846            let mut guard = pool.resolve_reader_checkout(
1847                StorageCapability::Sql,
1848                operation,
1849                pool.reader_until(|| scope.should_stop()),
1850            )?;
1851            scope.with_pooled_reader(&mut guard, |conn| query(scope, conn))
1852        },
1853    )
1854    .await
1855}
1856
1857async fn run_pool_writer_query<T, F>(
1858    pool: Arc<ConnectionPool>,
1859    operation: &'static str,
1860    query: F,
1861) -> khive_storage::types::StorageResult<T>
1862where
1863    T: Send + 'static,
1864    F: FnOnce(
1865            &crate::read_cancellation::InterruptibleReadScope,
1866            &rusqlite::Connection,
1867            bool,
1868        ) -> khive_storage::types::StorageResult<T>
1869        + Send
1870        + 'static,
1871{
1872    crate::read_cancellation::run_interruptible_read(
1873        StorageCapability::Sql,
1874        operation,
1875        move |scope| {
1876            let guard = pool.try_writer().map_err(|error: SqliteError| {
1877                StorageError::driver(StorageCapability::Sql, operation, error)
1878            })?;
1879            scope.with_pooled_writer(&pool, &guard, |conn| {
1880                let interruptible = conn.is_autocommit();
1881                query(scope, conn, interruptible)
1882            })
1883        },
1884    )
1885    .await
1886}
1887
1888struct PoolBackedReader {
1889    pool: Arc<ConnectionPool>,
1890}
1891
1892#[async_trait]
1893impl khive_storage::SqlReader for PoolBackedReader {
1894    async fn query_row(
1895        &mut self,
1896        statement: SqlStatement,
1897    ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
1898        let pool = Arc::clone(&self.pool);
1899        run_pool_reader_query(pool, "pool_reader.query_row", move |scope, conn| {
1900            execute_query_row_interruptibly(
1901                scope,
1902                conn,
1903                &statement,
1904                "pool_reader.query_row",
1905                false,
1906                true,
1907            )
1908        })
1909        .await
1910    }
1911
1912    async fn query_all(
1913        &mut self,
1914        statement: SqlStatement,
1915    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1916        let pool = Arc::clone(&self.pool);
1917        run_pool_reader_query(pool, "pool_reader.query_all", move |scope, conn| {
1918            execute_query_interruptibly(
1919                scope,
1920                conn,
1921                &statement,
1922                "pool_reader.query_all",
1923                false,
1924                true,
1925            )
1926        })
1927        .await
1928    }
1929
1930    async fn query_page(
1931        &mut self,
1932        statement: SqlStatement,
1933        page: PageRequest,
1934    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1935        let pool = Arc::clone(&self.pool);
1936        run_pool_reader_query(pool, "pool_reader.query_page", move |scope, conn| {
1937            execute_query_page_interruptibly(
1938                scope,
1939                conn,
1940                &statement,
1941                &page,
1942                "pool_reader.query_page",
1943                false,
1944                true,
1945            )
1946        })
1947        .await
1948    }
1949
1950    async fn query_scalar(
1951        &mut self,
1952        statement: SqlStatement,
1953    ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
1954        let row = self.query_row(statement).await?;
1955        Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
1956    }
1957
1958    async fn explain(
1959        &mut self,
1960        statement: SqlStatement,
1961    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1962        let explain_stmt = SqlStatement {
1963            sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
1964            params: statement.params,
1965            label: statement.label,
1966        };
1967        self.query_all(explain_stmt).await
1968    }
1969}
1970
1971struct PoolBackedWriter {
1972    pool: Arc<ConnectionPool>,
1973}
1974
1975#[async_trait]
1976impl khive_storage::SqlReader for PoolBackedWriter {
1977    async fn query_row(
1978        &mut self,
1979        statement: SqlStatement,
1980    ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
1981        let pool = Arc::clone(&self.pool);
1982        run_pool_writer_query(
1983            pool,
1984            "pool_writer.query_row",
1985            move |scope, conn, interruptible| {
1986                execute_query_row_interruptibly(
1987                    scope,
1988                    conn,
1989                    &statement,
1990                    "pool_writer.query_row",
1991                    false,
1992                    interruptible,
1993                )
1994            },
1995        )
1996        .await
1997    }
1998
1999    async fn query_all(
2000        &mut self,
2001        statement: SqlStatement,
2002    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2003        let pool = Arc::clone(&self.pool);
2004        run_pool_writer_query(
2005            pool,
2006            "pool_writer.query_all",
2007            move |scope, conn, interruptible| {
2008                execute_query_interruptibly(
2009                    scope,
2010                    conn,
2011                    &statement,
2012                    "pool_writer.query_all",
2013                    false,
2014                    interruptible,
2015                )
2016            },
2017        )
2018        .await
2019    }
2020
2021    async fn query_page(
2022        &mut self,
2023        statement: SqlStatement,
2024        page: PageRequest,
2025    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2026        let pool = Arc::clone(&self.pool);
2027        run_pool_writer_query(
2028            pool,
2029            "pool_writer.query_page",
2030            move |scope, conn, interruptible| {
2031                execute_query_page_interruptibly(
2032                    scope,
2033                    conn,
2034                    &statement,
2035                    &page,
2036                    "pool_writer.query_page",
2037                    false,
2038                    interruptible,
2039                )
2040            },
2041        )
2042        .await
2043    }
2044
2045    async fn query_scalar(
2046        &mut self,
2047        statement: SqlStatement,
2048    ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
2049        let row = khive_storage::SqlReader::query_row(self, statement).await?;
2050        Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
2051    }
2052
2053    async fn explain(
2054        &mut self,
2055        statement: SqlStatement,
2056    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2057        let explain_stmt = SqlStatement {
2058            sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
2059            params: statement.params,
2060            label: statement.label,
2061        };
2062        khive_storage::SqlReader::query_all(self, explain_stmt).await
2063    }
2064}
2065
2066#[async_trait]
2067impl khive_storage::SqlWriter for PoolBackedWriter {
2068    async fn execute(
2069        &mut self,
2070        statement: SqlStatement,
2071    ) -> khive_storage::types::StorageResult<u64> {
2072        // Boundary: `execute_batch` owns transaction-control rejection;
2073        // this one-statement primitive is used by internal DML/transaction
2074        // owners and is still guarded by the SqlStatement single-statement
2075        // prepare contract.
2076        let pool = Arc::clone(&self.pool);
2077        tokio::task::spawn_blocking(move || {
2078            let guard = pool.try_writer().map_err(|e: SqliteError| {
2079                StorageError::driver(StorageCapability::Sql, "pool_writer.execute", e)
2080            })?;
2081            let mut stmt = prepare_cached_sql_statement(&guard, &statement.sql)
2082                .map_err(|e| map_rusqlite_err(e, "pool_writer.execute"))?;
2083            bind_params(&mut stmt, &statement.params)
2084                .map_err(|e| map_rusqlite_err(e, "pool_writer.execute"))?;
2085            let rows = stmt
2086                .raw_execute()
2087                .map_err(|e| map_rusqlite_err(e, "pool_writer.execute"))?;
2088            Ok(rows as u64)
2089        })
2090        .await
2091        .map_err(|e| StorageError::driver(StorageCapability::Sql, "pool_writer.execute", e))?
2092    }
2093
2094    async fn execute_batch(
2095        &mut self,
2096        statements: Vec<SqlStatement>,
2097    ) -> khive_storage::types::StorageResult<u64> {
2098        // Same all-or-nothing contract as the file-backed path: this batch
2099        // wraps its list in its own `BEGIN IMMEDIATE`, so reject caller
2100        // transaction-control statements before executing anything.
2101        reject_transaction_control_statements(&statements, "pool_writer.execute_batch")?;
2102        let pool = Arc::clone(&self.pool);
2103        tokio::task::spawn_blocking(move || {
2104            let guard = pool.try_writer().map_err(|e: SqliteError| {
2105                StorageError::driver(StorageCapability::Sql, "pool_writer.execute_batch", e)
2106            })?;
2107            let prepared = prepare_batch_statements(&guard, &statements)
2108                .map_err(|e| map_rusqlite_err(e, "pool_writer.execute_batch"))?;
2109            guard
2110                .execute_batch("BEGIN IMMEDIATE")
2111                .map_err(|e| map_rusqlite_err(e, "pool_writer.execute_batch"))?;
2112            let _tx_handle = khive_storage::tx_registry::register_scoped(
2113                Some("pool_writer.execute_batch".to_string()),
2114                pool.origin(),
2115            );
2116            let result = execute_prepared_batch(&guard, prepared, &statements)
2117                .map_err(|e| map_rusqlite_err(e, "pool_writer.execute_batch"));
2118            match result {
2119                Ok(total) => {
2120                    if let Err(e) = guard.execute_batch("COMMIT") {
2121                        let _ = guard.execute_batch("ROLLBACK");
2122                        Err(map_rusqlite_err(e, "pool_writer.execute_batch"))
2123                    } else {
2124                        Ok(total)
2125                    }
2126                }
2127                Err(e) => {
2128                    let _ = guard.execute_batch("ROLLBACK");
2129                    Err(e)
2130                }
2131            }
2132        })
2133        .await
2134        .map_err(|e| StorageError::driver(StorageCapability::Sql, "pool_writer.execute_batch", e))?
2135    }
2136
2137    async fn execute_script(&mut self, script: String) -> khive_storage::types::StorageResult<()> {
2138        // Boundary: raw scripts are internal/migration-only and do not inherit
2139        // `execute_batch`'s transaction-control rejection.
2140        let pool = Arc::clone(&self.pool);
2141        tokio::task::spawn_blocking(move || {
2142            let guard = pool.try_writer().map_err(|e: SqliteError| {
2143                StorageError::driver(StorageCapability::Sql, "pool_writer.execute_script", e)
2144            })?;
2145            guard
2146                .execute_batch(&script)
2147                .map_err(|e| map_rusqlite_err(e, "pool_writer.execute_script"))
2148        })
2149        .await
2150        .map_err(|e| {
2151            StorageError::driver(StorageCapability::Sql, "pool_writer.execute_script", e)
2152        })?
2153    }
2154}
2155
2156// =============================================================================
2157// atomic_unit (ADR-067 Component A, Fork C slice 2)
2158// =============================================================================
2159
2160/// A purely-synchronous `SqlReader`/`SqlWriter` over a borrowed connection,
2161/// used ONLY to drive an [`AtomicUnitOp`] on the flag-on path, where the
2162/// closure body runs inside the writer task's `spawn_blocking` (synchronous
2163/// `FnOnce(&rusqlite::Connection) -> ...`) rather than a real async context.
2164///
2165/// Every method here does plain, non-suspending rusqlite work — there is no
2166/// real `.await` point anywhere in this impl — so [`block_on_sync`] driving
2167/// the resulting future to completion with a single poll is sound, not a
2168/// hack: the future can never actually be `Pending`.
2169///
2170/// `SqlReader`/`SqlWriter` both carry a `'static` supertrait bound (they are
2171/// used as `Box<dyn ...>` elsewhere in this module), so this type cannot
2172/// hold a real `&'c Connection` borrow — it would tie `InlineWriter` to a
2173/// non-`'static` lifetime and, independently, `&Connection` is not `Send`
2174/// (`Connection` is `!Sync`), which the `#[async_trait]`-generated futures
2175/// require. A raw pointer sidesteps both: `*const Connection` is `Send` and
2176/// `'static` on its face, and the safety burden (the pointee outliving
2177/// every dereference) is upheld by construction — see `atomic_unit`, the
2178/// only call site: it builds an `InlineWriter` from `conn: &Connection`,
2179/// drives `op` to completion via `block_on_sync` synchronously, and drops
2180/// the `InlineWriter` before that borrow ends, all within one stack frame.
2181struct InlineWriter {
2182    conn: *const rusqlite::Connection,
2183}
2184
2185// SAFETY: `InlineWriter` is never actually shared across a real thread
2186// boundary — it is constructed, driven to completion synchronously via
2187// `block_on_sync`, and dropped within a single call frame inside the
2188// writer task's `spawn_blocking` closure (see `atomic_unit`). The `Send`
2189// bound `async_trait` imposes on the futures below is a static
2190// over-approximation for this restricted, single-threaded usage pattern.
2191unsafe impl Send for InlineWriter {}
2192
2193impl InlineWriter {
2194    /// SAFETY: valid for the lifetime of the enclosing synchronous scope in
2195    /// `atomic_unit` (see the struct doc comment above) — the pointee is
2196    /// never dereferenced after that scope ends.
2197    fn conn(&self) -> &rusqlite::Connection {
2198        unsafe { &*self.conn }
2199    }
2200}
2201
2202#[async_trait]
2203impl khive_storage::SqlReader for InlineWriter {
2204    async fn query_row(
2205        &mut self,
2206        statement: SqlStatement,
2207    ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
2208        execute_query_row(self.conn(), &statement)
2209            .map_err(|e| map_rusqlite_err(e, "inline.query_row"))
2210    }
2211
2212    async fn query_all(
2213        &mut self,
2214        statement: SqlStatement,
2215    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2216        execute_query(self.conn(), &statement).map_err(|e| map_rusqlite_err(e, "inline.query_all"))
2217    }
2218
2219    async fn query_page(
2220        &mut self,
2221        statement: SqlStatement,
2222        page: PageRequest,
2223    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2224        execute_query_page(self.conn(), &statement, &page)
2225            .map_err(|e| map_rusqlite_err(e, "inline.query_page"))
2226    }
2227
2228    async fn query_scalar(
2229        &mut self,
2230        statement: SqlStatement,
2231    ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
2232        let row = khive_storage::SqlReader::query_row(self, statement).await?;
2233        Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
2234    }
2235
2236    async fn explain(
2237        &mut self,
2238        statement: SqlStatement,
2239    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2240        let explain_stmt = SqlStatement {
2241            sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
2242            params: statement.params,
2243            label: statement.label,
2244        };
2245        khive_storage::SqlReader::query_all(self, explain_stmt).await
2246    }
2247}
2248
2249#[async_trait]
2250impl khive_storage::SqlWriter for InlineWriter {
2251    async fn execute(
2252        &mut self,
2253        statement: SqlStatement,
2254    ) -> khive_storage::types::StorageResult<u64> {
2255        // Boundary: `execute_batch` owns transaction-control rejection;
2256        // `atomic_unit` uses this one-statement primitive for its own boundary.
2257        let mut stmt = prepare_cached_sql_statement(self.conn(), &statement.sql)
2258            .map_err(|e| map_rusqlite_err(e, "inline.execute"))?;
2259        bind_params(&mut stmt, &statement.params)
2260            .map_err(|e| map_rusqlite_err(e, "inline.execute"))?;
2261        let affected = stmt
2262            .raw_execute()
2263            .map_err(|e| map_rusqlite_err(e, "inline.execute"))?;
2264        Ok(affected as u64)
2265    }
2266
2267    async fn execute_batch(
2268        &mut self,
2269        statements: Vec<SqlStatement>,
2270    ) -> khive_storage::types::StorageResult<u64> {
2271        // Runs inside the writer task's per-request `BEGIN IMMEDIATE`
2272        // (atomic_unit flag-on path), so a caller `COMMIT` would close the
2273        // task's transaction — reject transaction-control statements up
2274        // front, same contract as every other `execute_batch`.
2275        reject_transaction_control_statements(&statements, "inline.execute_batch")?;
2276        let prepared = prepare_batch_statements(self.conn(), &statements)
2277            .map_err(|e| map_rusqlite_err(e, "inline.execute_batch"))?;
2278        execute_prepared_batch(self.conn(), prepared, &statements)
2279            .map_err(|e| map_rusqlite_err(e, "inline.execute_batch"))
2280    }
2281
2282    async fn execute_script(&mut self, script: String) -> khive_storage::types::StorageResult<()> {
2283        // Boundary: this raw script path is internal maintenance only and is
2284        // outside the `execute_batch` transaction-control contract.
2285        self.conn()
2286            .execute_batch(&script)
2287            .map_err(|e| map_rusqlite_err(e, "inline.execute_script"))
2288    }
2289}
2290
2291/// Poll `fut` exactly once with a no-op waker and return its output.
2292///
2293/// Only sound for futures that never actually suspend — every caller in
2294/// this module drives an [`InlineWriter`], whose methods are pure
2295/// synchronous rusqlite calls with no real `.await` point.
2296///
2297/// ADR-067 Component A: this used to
2298/// `unreachable!()`-panic on `Poll::Pending`, and a panicking closure
2299/// running inside the writer task's `spawn_blocking` (see
2300/// `SqlBridge::atomic_unit`'s flag-on branch) would surface as a
2301/// `JoinError` in `run_writer_task`, which is treated as fatal — the writer
2302/// task exits and every subsequent `WriterTaskHandle::send` on this pool
2303/// fails for the rest of the process. A future `atomic_unit` caller whose
2304/// closure ever gains a real suspend point (this file's own contract
2305/// already forbids it, but the invariant is enforced by convention, not the
2306/// type system) would take down the writer task for the whole daemon.
2307/// Returning `Err` instead lets `Pending` flow through the SAME error path
2308/// as any other `atomic_unit` op failure: `WriteRequest::execute_and_reply`
2309/// treats it as an ordinary `Err`, issues `ROLLBACK` on the writer task's
2310/// held transaction, replies the error to the caller, and the writer task's
2311/// `spawn_blocking` closure returns normally (not via panic) — so the task
2312/// keeps draining subsequent requests instead of dying with the whole pool.
2313fn block_on_sync<F: std::future::Future>(fut: F) -> Result<F::Output, StorageError> {
2314    use std::task::{Context, Poll, RawWaker, RawWakerVTable, Waker};
2315
2316    fn no_op(_: *const ()) {}
2317    fn clone_waker(_: *const ()) -> RawWaker {
2318        RawWaker::new(std::ptr::null(), &VTABLE)
2319    }
2320    static VTABLE: RawWakerVTable = RawWakerVTable::new(clone_waker, no_op, no_op, no_op);
2321
2322    // SAFETY: every `RawWakerVTable` function is a no-op that never
2323    // dereferences the data pointer, so a null data pointer is sound.
2324    let raw_waker = RawWaker::new(std::ptr::null(), &VTABLE);
2325    let waker = unsafe { Waker::from_raw(raw_waker) };
2326    let mut cx = Context::from_waker(&waker);
2327
2328    let mut fut = std::pin::pin!(fut);
2329    match fut.as_mut().poll(&mut cx) {
2330        Poll::Ready(v) => Ok(v),
2331        Poll::Pending => {
2332            tracing::error!(
2333                "block_on_sync: atomic_unit future suspended on its first poll — \
2334                 the closure passed to SqlAccess::atomic_unit must be non-blocking \
2335                 (synchronous InlineWriter calls only, no real .await point)"
2336            );
2337            Err(StorageError::Internal(
2338                "atomic_unit future suspended — closure must be non-blocking".to_string(),
2339            ))
2340        }
2341    }
2342}
2343
2344/// Run `op` under a manual `BEGIN IMMEDIATE`/`COMMIT`/`ROLLBACK` on `writer`
2345/// — the pre-ADR-067 shape, used by [`SqlBridge::atomic_unit`] whenever no
2346/// writer task applies (flag off, no runtime, or an in-memory pool),
2347/// preserving that path byte-for-byte.
2348async fn run_manual_atomic_unit(
2349    writer: &mut dyn khive_storage::SqlWriter,
2350    op: AtomicUnitOp,
2351    origin: khive_storage::tx_registry::TxOrigin,
2352) -> khive_storage::types::StorageResult<Box<dyn Any + Send>> {
2353    fn tx_stmt(sql: &str, label: &str) -> SqlStatement {
2354        SqlStatement {
2355            sql: sql.to_string(),
2356            params: vec![],
2357            label: Some(label.to_string()),
2358        }
2359    }
2360    khive_storage::SqlWriter::execute(writer, tx_stmt("BEGIN IMMEDIATE", "begin")).await?;
2361    let _tx_handle =
2362        khive_storage::tx_registry::register_scoped(Some("atomic_unit".to_string()), origin);
2363
2364    let result = op(writer).await;
2365
2366    match result {
2367        Ok(value) => {
2368            match khive_storage::SqlWriter::execute(writer, tx_stmt("COMMIT", "commit")).await {
2369                Ok(_) => Ok(value),
2370                Err(e) => {
2371                    let _ =
2372                        khive_storage::SqlWriter::execute(writer, tx_stmt("ROLLBACK", "rollback"))
2373                            .await;
2374                    Err(e)
2375                }
2376            }
2377        }
2378        Err(e) => {
2379            let _ =
2380                khive_storage::SqlWriter::execute(writer, tx_stmt("ROLLBACK", "rollback")).await;
2381            Err(e)
2382        }
2383    }
2384}
2385
2386// =============================================================================
2387// SqlBridge: the SqlAccess implementor
2388// =============================================================================
2389
2390/// Bridges `ConnectionPool` to `khive_storage::SqlAccess`.
2391///
2392/// Dispatches based on whether the pool is file-backed or in-memory:
2393/// - File-backed: cached standalone reader connections with ordinary
2394///   operation-scoped admission and explicitly admitted multi-call read
2395///   transactions, plus standalone writer connections capped at one live
2396///   handle; atomic units drive a single registered raw transaction span.
2397/// - In-memory: pool-backed connections per query (single shared connection).
2398pub struct SqlBridge {
2399    pool: Arc<ConnectionPool>,
2400    is_file_backed: bool,
2401}
2402
2403impl SqlBridge {
2404    /// Create a new bridge wrapping the given pool.
2405    pub fn new(pool: Arc<ConnectionPool>, is_file_backed: bool) -> Self {
2406        Self {
2407            pool,
2408            is_file_backed,
2409        }
2410    }
2411}
2412
2413#[async_trait]
2414impl khive_storage::SqlAccess for SqlBridge {
2415    fn database_path(&self) -> Option<std::path::PathBuf> {
2416        self.pool.canonical_path().map(std::path::Path::to_path_buf)
2417    }
2418
2419    async fn reader(
2420        &self,
2421    ) -> khive_storage::types::StorageResult<Box<dyn khive_storage::SqlReader>> {
2422        if self.is_file_backed {
2423            Ok(Box::new(SqliteReader {
2424                handle: Some(open_cached_reader_handle(Arc::clone(&self.pool)).await?),
2425                pool: Arc::clone(&self.pool),
2426            }))
2427        } else {
2428            Ok(Box::new(PoolBackedReader {
2429                pool: Arc::clone(&self.pool),
2430            }))
2431        }
2432    }
2433
2434    async fn writer(
2435        &self,
2436    ) -> khive_storage::types::StorageResult<Box<dyn khive_storage::SqlWriter>> {
2437        if self.is_file_backed {
2438            if self.pool.config().read_only {
2439                return Err(StorageError::Pool {
2440                    operation: "writer".into(),
2441                    message: "backend is read-only".into(),
2442                });
2443            }
2444            let db = crate::timeout_sink::db_label(&self.pool);
2445            // ADR-136 D1 gate 1: queue-first. The handle lookup runs BEFORE
2446            // any standalone connection is opened, and a lookup failure is
2447            // propagated (never silently degraded) when strict routing is
2448            // on. Only the flag-off/degraded case still opens a standalone
2449            // connection.
2450            let writer_task = match self.pool.writer_task_handle() {
2451                Ok(handle) => handle,
2452                Err(e) => {
2453                    if self.pool.config().write_routing_strict {
2454                        return Err(e);
2455                    }
2456                    tracing::warn!(
2457                        error = %e,
2458                        "KHIVE_WRITE_ROUTING is not strict; writer() degrades to the \
2459                         standalone-connection path"
2460                    );
2461                    None
2462                }
2463            };
2464            if writer_task.is_none() && self.pool.config().write_routing_strict {
2465                return Err(StorageError::Pool {
2466                    operation: "writer".into(),
2467                    message: "KHIVE_WRITE_ROUTING=strict but no writer-task handle is \
2468                              available; refusing to fall back to a direct connection"
2469                        .into(),
2470                });
2471            }
2472            if writer_task.is_none() && self.pool.write_queue_active() {
2473                // The queue is enabled but this call didn't get a handle
2474                // (spawn/runtime degrade) — a direct-route violation in the
2475                // making once this writer's execute*/query* methods run.
2476                // In-memory pools are excluded: they never spawn a writer
2477                // task by documented design (explicit `Some(true)` degrades),
2478                // so a violation row there would be noise, not signal.
2479                crate::timeout_sink::emit_direct_route_violation(
2480                    &db,
2481                    crate::timeout_sink::Site::DirectRouteSqlBridgeWriter,
2482                );
2483            }
2484            // A standalone read-write connection is opened only when there is
2485            // no queue handle to route writes through — `SqliteWriter`'s
2486            // `SqlReader` methods (`query_row`/`query_all`/`query_page`)
2487            // lazily open a read-only one on first use in the handle-present
2488            // case; production callers do read through a `writer()` handle,
2489            // so this lazy path is live (see `SqliteWriter::ensure_conn`).
2490            // The standalone open acquires the pool-wide one-permit writer
2491            // budget first, and the permit travels in the handle for the
2492            // handle's whole lifetime — a queue-backed handle holds no
2493            // writer permit (its writes route through the writer task), so
2494            // this budget caps exactly the standalone read-write
2495            // connections.
2496            let handle = if writer_task.is_none() {
2497                let handle_slot = acquire_handle_slot(
2498                    self.pool.sql_bridge_writer_slots(),
2499                    self.pool.config().checkout_timeout,
2500                    "sql_bridge.writer_handle",
2501                    SlotTimeoutClass::Admission,
2502                )
2503                .await?;
2504                let (conn, handle_slot) =
2505                    open_standalone_writer_on_blocking(Arc::clone(&self.pool), handle_slot).await?;
2506                Some(StandaloneHandle {
2507                    conn,
2508                    _retained_slot: Some(handle_slot),
2509                    read_transaction_slot: None,
2510                })
2511            } else {
2512                None
2513            };
2514            Ok(Box::new(SqliteWriter {
2515                handle,
2516                writer_task,
2517                origin: self.pool.origin(),
2518                db,
2519                pool: Arc::clone(&self.pool),
2520            }))
2521        } else {
2522            Ok(Box::new(PoolBackedWriter {
2523                pool: Arc::clone(&self.pool),
2524            }))
2525        }
2526    }
2527
2528    /// Implements the trait's atomic-unit suspend-free invariant
2529    /// (`SqlAccess::atomic_unit`'s doc comment): on the flag-on branch below,
2530    /// `op` is driven through `block_on_sync` on an `InlineWriter` — a
2531    /// single-poll driver that returns `Err` the instant `op`'s future is
2532    /// `Pending` instead of ever actually suspending. `op` must therefore
2533    /// issue only synchronous DML; see `InlineWriter`'s and
2534    /// `block_on_sync`'s doc comments for the full mechanics and why this
2535    /// restriction is load-bearing (a suspended poll inside the writer
2536    /// task's `spawn_blocking` would otherwise block that task on external
2537    /// async work while holding the single write connection).
2538    async fn atomic_unit(
2539        &self,
2540        op: AtomicUnitOp,
2541    ) -> khive_storage::types::StorageResult<Box<dyn Any + Send>> {
2542        if self.is_file_backed {
2543            if self.pool.config().read_only {
2544                return Err(StorageError::Pool {
2545                    operation: "atomic_unit".into(),
2546                    message: "backend is read-only".into(),
2547                });
2548            }
2549            // Best-effort, same guard `writer()` uses: `Ok(None)` on flag-off;
2550            // `Err(WriterTaskNoRuntime)` propagates loud rather than silently
2551            // falling back to a competing connection from a sync caller. ADR-136
2552            // D1 gate 3: `Ok(None)` under strict routing is ALSO a fail-closed
2553            // error (queue was requested but unavailable), not just a degrade.
2554            let handle = self.pool.writer_task_handle()?;
2555            if handle.is_none() && self.pool.config().write_routing_strict {
2556                return Err(StorageError::Pool {
2557                    operation: "atomic_unit".into(),
2558                    message: "KHIVE_WRITE_ROUTING=strict but no writer-task handle is \
2559                              available; refusing to fall back to a direct connection"
2560                        .into(),
2561                });
2562            }
2563            if handle.is_none() && self.pool.write_queue_active() {
2564                crate::timeout_sink::emit_direct_route_violation(
2565                    &crate::timeout_sink::db_label(&self.pool),
2566                    crate::timeout_sink::Site::DirectRouteAtomicUnit,
2567                );
2568            }
2569            if let Some(writer_task) = handle {
2570                // Flag-on: ONE queued WriteRequest. `run_writer_task` already
2571                // has an open `BEGIN IMMEDIATE` on its dedicated connection
2572                // before this closure runs and issues `COMMIT`/`ROLLBACK`
2573                // after it returns — `op` must not (and, via `InlineWriter`,
2574                // does not) issue its own transaction control.
2575                return writer_task
2576                    .send_bounded(move |conn| {
2577                        let mut inline = InlineWriter {
2578                            conn: conn as *const rusqlite::Connection,
2579                        };
2580                        // Flatten: `block_on_sync` now returns `Result<F::Output,
2581                        // StorageError>` (outer = "did the future actually
2582                        // resolve on first poll", inner = the op's own
2583                        // `StorageResult`) instead of panicking on `Pending`
2584                        // (ADR-067 Component A). Either
2585                        // error flows through this closure's ordinary `Err`
2586                        // return, which `WriteRequest::execute_and_reply`
2587                        // already turns into a normal ROLLBACK + error reply —
2588                        // no panic, so the writer task survives.
2589                        match block_on_sync(op(&mut inline)) {
2590                            Ok(inner) => inner,
2591                            Err(e) => Err(e),
2592                        }
2593                    })
2594                    .await;
2595            }
2596            // Flag-off (or no writer task available): manual
2597            // BEGIN IMMEDIATE/COMMIT/ROLLBACK on a standalone writer —
2598            // byte-for-byte the pre-ADR-067 shape.
2599            //
2600            // Contract: this acquire waits on the pool-wide one-permit
2601            // writer-handle budget — the same permit a live `writer()` handle
2602            // holds for its lifetime — so it times out with
2603            // `StorageError::AdmissionTimeout` after `checkout_timeout` while a writer
2604            // handle is checked out (and a `writer()` call times out while
2605            // this unit runs). Callers must not hold a boxed writer handle
2606            // across an `atomic_unit()` call on the same pool; drop the
2607            // handle first. The `writer_task` branch above never touches this
2608            // budget.
2609            let handle_slot = acquire_handle_slot(
2610                self.pool.sql_bridge_writer_slots(),
2611                self.pool.config().checkout_timeout,
2612                "sql_bridge.atomic_unit_handle",
2613                SlotTimeoutClass::Admission,
2614            )
2615            .await?;
2616            let (conn, handle_slot) =
2617                open_standalone_writer_on_blocking(Arc::clone(&self.pool), handle_slot).await?;
2618            let mut writer = SqliteWriter {
2619                handle: Some(StandaloneHandle {
2620                    conn,
2621                    _retained_slot: Some(handle_slot),
2622                    read_transaction_slot: None,
2623                }),
2624                writer_task: None,
2625                origin: self.pool.origin(),
2626                db: crate::timeout_sink::db_label(&self.pool),
2627                pool: Arc::clone(&self.pool),
2628            };
2629            run_manual_atomic_unit(&mut writer, op, self.pool.origin()).await
2630        } else {
2631            // In-memory pools are exempt (not accept-loop reachable, per the
2632            // rework spec's "Out of scope") — preserve the existing
2633            // pool-backed manual-transaction behavior.
2634            let mut writer = PoolBackedWriter {
2635                pool: Arc::clone(&self.pool),
2636            };
2637            run_manual_atomic_unit(&mut writer, op, self.pool.origin()).await
2638        }
2639    }
2640}
2641
2642#[cfg(test)]
2643mod tests {
2644    use super::*;
2645    use crate::pool::PoolConfig;
2646    use khive_storage::types::{SqlStatement, SqlValue};
2647    use khive_storage::{SqlAccess as _, SqlReader as _};
2648
2649    fn database_tx_view(pool: &ConnectionPool) -> khive_storage::tx_registry::TxOriginFilter {
2650        match pool.origin() {
2651            khive_storage::tx_registry::TxOrigin::Database(identity) => {
2652                khive_storage::tx_registry::TxOriginFilter::Secondary(identity)
2653            }
2654            other => panic!("expected a file-backed database origin, got {other:?}"),
2655        }
2656    }
2657
2658    struct NotifyOnDrop(Arc<tokio::sync::Notify>);
2659
2660    impl Drop for NotifyOnDrop {
2661        fn drop(&mut self) {
2662            self.0.notify_one();
2663        }
2664    }
2665
2666    fn blocking_non_interrupting_progress_gate(
2667        conn: &rusqlite::Connection,
2668    ) -> (
2669        Arc<tokio::sync::Notify>,
2670        Arc<std::sync::Barrier>,
2671        Arc<tokio::sync::Notify>,
2672    ) {
2673        let entered = Arc::new(tokio::sync::Notify::new());
2674        let callback_entered = Arc::clone(&entered);
2675        let release = Arc::new(std::sync::Barrier::new(2));
2676        let callback_release = Arc::clone(&release);
2677        let completed = Arc::new(tokio::sync::Notify::new());
2678        let notify_on_drop = NotifyOnDrop(Arc::clone(&completed));
2679        let blocked_once = Arc::new(std::sync::atomic::AtomicBool::new(false));
2680        let callback_blocked_once = Arc::clone(&blocked_once);
2681        conn.progress_handler(
2682            1_000,
2683            Some(move || {
2684                let _keep_until_connection_drop = &notify_on_drop;
2685                if !callback_blocked_once.swap(true, std::sync::atomic::Ordering::SeqCst) {
2686                    callback_entered.notify_one();
2687                    callback_release.wait();
2688                    // The gate deliberately stalls but never asks SQLite to
2689                    // abort. It proves completion-preserving SQLite work ignores
2690                    // request-read cancellation and finishes normally.
2691                    return false;
2692                }
2693                false
2694            }),
2695        )
2696        .unwrap();
2697        (entered, release, completed)
2698    }
2699
2700    fn progress_gate_statement() -> SqlStatement {
2701        SqlStatement {
2702            sql: "WITH RECURSIVE rows(value) AS (\
2703                  SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 999\
2704                  ) SELECT SUM(value) FROM rows"
2705                .into(),
2706            params: vec![],
2707            label: None,
2708        }
2709    }
2710
2711    fn slow_insert_statement() -> SqlStatement {
2712        SqlStatement {
2713            sql: "INSERT INTO cancellation_write_probe(value) \
2714                  WITH RECURSIVE rows(value) AS (\
2715                  SELECT 1 UNION ALL SELECT value + 1 FROM rows WHERE value < 10000\
2716                  ) SELECT value FROM rows"
2717                .into(),
2718            params: vec![],
2719            label: Some("non-interruptible-write-probe".into()),
2720        }
2721    }
2722
2723    fn passive_checkpoint(conn: &rusqlite::Connection) -> (i64, i64, i64) {
2724        conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
2725            Ok((row.get(0)?, row.get(1)?, row.get(2)?))
2726        })
2727        .unwrap()
2728    }
2729
2730    fn deliberately_slow_read_statement() -> SqlStatement {
2731        SqlStatement {
2732            sql: "WITH RECURSIVE numbers(value) AS (\
2733                  SELECT 1 UNION ALL SELECT value + 1 FROM numbers WHERE value < 1000\
2734                  ) SELECT SUM(a.value * b.value * c.value) \
2735                  FROM numbers AS a CROSS JOIN numbers AS b CROSS JOIN numbers AS c"
2736                .into(),
2737            params: vec![],
2738            label: Some("read-cancellation-progress-probe".into()),
2739        }
2740    }
2741
2742    async fn wait_for_progress(probe: &std::sync::atomic::AtomicUsize) {
2743        tokio::time::timeout(std::time::Duration::from_secs(1), async {
2744            while probe.load(std::sync::atomic::Ordering::SeqCst) == 0 {
2745                tokio::task::yield_now().await;
2746            }
2747        })
2748        .await
2749        .expect("slow SQLite statement never reached its progress callback");
2750    }
2751
2752    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2753    async fn cancellation_before_reader_checkout_is_prompt_and_executes_no_statement() {
2754        let dir = tempfile::tempdir().unwrap();
2755        let config = PoolConfig {
2756            path: Some(dir.path().join("sql_bridge_cancel_before_checkout.db")),
2757            max_readers: 1,
2758            checkout_timeout: std::time::Duration::from_secs(5),
2759            ..PoolConfig::default()
2760        };
2761        let pool = Arc::new(ConnectionPool::new(config).unwrap());
2762        pool.writer()
2763            .unwrap()
2764            .conn()
2765            .execute_batch(
2766                "CREATE TABLE checkout_cancel_probe(value INTEGER NOT NULL); \
2767                 INSERT INTO checkout_cancel_probe VALUES (0);",
2768            )
2769            .unwrap();
2770        let held_reader = pool.reader().expect("hold the sole pooled reader");
2771        // Deliberately force the pool-backed bridge over this file pool: the
2772        // held guard and `PoolBackedReader::reader_until` then contend for the
2773        // exact same one-connection queue (not the standalone semaphore).
2774        let bridge = SqlBridge::new(Arc::clone(&pool), false);
2775        let mut waiting_reader = bridge.reader().await.unwrap();
2776        let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
2777        let waiting = tokio::spawn(crate::scope_request_read_cancellation(
2778            cancel_rx,
2779            async move {
2780                waiting_reader
2781                    .query_row(SqlStatement {
2782                        sql: "UPDATE checkout_cancel_probe SET value = value + 1 RETURNING value"
2783                            .into(),
2784                        params: vec![],
2785                        label: Some("must-not-run-after-cancelled-checkout".into()),
2786                    })
2787                    .await
2788            },
2789        ));
2790
2791        tokio::task::yield_now().await;
2792        cancel_tx.send(true).unwrap();
2793        let result = tokio::time::timeout(std::time::Duration::from_millis(100), waiting)
2794            .await
2795            .expect("cancelled reader checkout waited for the five-second pool timeout")
2796            .expect("checkout task panicked");
2797        // Cancellation is NOT an admission wait: it must stay the non-admission
2798        // Timeout, never the retryable AdmissionTimeout, so a cancelled request
2799        // does not signal clients to retry into a saturated pool.
2800        assert!(matches!(result, Err(StorageError::Timeout { .. })));
2801
2802        drop(held_reader);
2803        tokio::time::sleep(std::time::Duration::from_millis(25)).await;
2804        let value: i64 = pool
2805            .reader()
2806            .unwrap()
2807            .conn()
2808            .query_row("SELECT value FROM checkout_cancel_probe", [], |row| {
2809                row.get(0)
2810            })
2811            .unwrap();
2812        assert_eq!(
2813            value, 0,
2814            "a DML statement started after its pre-admission checkout was cancelled"
2815        );
2816        assert_eq!(
2817            pool.available_readers(),
2818            1,
2819            "reader checkout leaked a permit"
2820        );
2821    }
2822
2823    /// A pooled-reader checkout that exhausts `checkout_timeout` WITHOUT any
2824    /// cancellation is a genuine admission wait and must surface as the
2825    /// retryable AdmissionTimeout. Before the fix, `reader_until`'s
2826    /// pool-exhausted error was mapped to `StorageError::Driver`, so a
2827    /// saturated pooled read stayed a non-retryable driver failure and the new
2828    /// AdmissionTimeout branch (reachable only for `Ok(None)` cancellation) was
2829    /// dead for real timeouts.
2830    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2831    async fn pooled_reader_checkout_timeout_is_a_retryable_admission_timeout() {
2832        let dir = tempfile::tempdir().unwrap();
2833        let config = PoolConfig {
2834            path: Some(dir.path().join("sql_bridge_reader_admission_timeout.db")),
2835            max_readers: 1,
2836            checkout_timeout: std::time::Duration::from_millis(200),
2837            ..PoolConfig::default()
2838        };
2839        let pool = Arc::new(ConnectionPool::new(config).unwrap());
2840        pool.writer()
2841            .unwrap()
2842            .conn()
2843            .execute_batch(
2844                "CREATE TABLE reader_admission_probe(value INTEGER NOT NULL); \
2845                 INSERT INTO reader_admission_probe VALUES (0);",
2846            )
2847            .unwrap();
2848        // Hold the sole pooled reader so the contending checkout cannot succeed
2849        // and must run `checkout_timeout` to exhaustion — no cancellation.
2850        let held_reader = pool.reader().expect("hold the sole pooled reader");
2851
2852        let bridge = SqlBridge::new(Arc::clone(&pool), false);
2853        let mut contender = bridge.reader().await.unwrap();
2854        let blocked = contender
2855            .query_row(SqlStatement {
2856                sql: "SELECT value FROM reader_admission_probe".into(),
2857                params: vec![],
2858                label: Some("reader-admission-timeout-probe".into()),
2859            })
2860            .await;
2861        assert!(
2862            matches!(blocked, Err(StorageError::AdmissionTimeout { .. })),
2863            "an exhausted pooled-reader checkout must be a retryable AdmissionTimeout; got {blocked:?}"
2864        );
2865
2866        drop(held_reader);
2867        tokio::time::sleep(std::time::Duration::from_millis(25)).await;
2868        assert_eq!(
2869            pool.available_readers(),
2870            1,
2871            "reader checkout leaked a permit"
2872        );
2873    }
2874
2875    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2876    async fn abandoned_read_interrupts_sqlite_releases_permit_and_stops_work() {
2877        let dir = tempfile::tempdir().unwrap();
2878        let config = PoolConfig {
2879            path: Some(dir.path().join("sql_bridge_abandoned_read.db")),
2880            max_readers: 1,
2881            checkout_timeout: std::time::Duration::from_millis(500),
2882            ..PoolConfig::default()
2883        };
2884        let pool = Arc::new(ConnectionPool::new(config).unwrap());
2885        let bridge = SqlBridge::new(Arc::clone(&pool), true);
2886        let mut reader = SqliteReader {
2887            handle: Some(open_cached_reader_handle(Arc::clone(&pool)).await.unwrap()),
2888            pool: Arc::clone(&pool),
2889        };
2890        let mut contender = bridge.reader().await.unwrap();
2891        let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
2892        let progress_in_scope = Arc::clone(&progress);
2893
2894        let query = tokio::spawn(crate::scope_test_read_progress(
2895            progress_in_scope,
2896            async move { reader.query_all(deliberately_slow_read_statement()).await },
2897        ));
2898        wait_for_progress(progress.as_ref()).await;
2899        query.abort();
2900        assert!(matches!(query.await, Err(error) if error.is_cancelled()));
2901
2902        tokio::time::timeout(
2903            std::time::Duration::from_millis(500),
2904            contender.query_row(SqlStatement {
2905                sql: "SELECT 1".into(),
2906                params: vec![],
2907                label: None,
2908            }),
2909        )
2910        .await
2911        .expect("abandoned SQLite statement did not return the sole reader promptly")
2912        .expect("reader probe failed after cancellation");
2913
2914        let stopped_at = progress.load(std::sync::atomic::Ordering::SeqCst);
2915        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
2916        assert_eq!(
2917            progress.load(std::sync::atomic::Ordering::SeqCst),
2918            stopped_at,
2919            "SQLite progress kept advancing after the abandoned request returned its reader"
2920        );
2921    }
2922
2923    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2924    async fn request_deadline_interrupts_statement_without_outer_timeout() {
2925        let dir = tempfile::tempdir().unwrap();
2926        let config = PoolConfig {
2927            path: Some(dir.path().join("sql_bridge_request_deadline.db")),
2928            max_readers: 1,
2929            checkout_timeout: std::time::Duration::from_millis(500),
2930            ..PoolConfig::default()
2931        };
2932        let pool = Arc::new(ConnectionPool::new(config).unwrap());
2933        let bridge = SqlBridge::new(Arc::clone(&pool), true);
2934        let mut reader = bridge.reader().await.unwrap();
2935        let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
2936
2937        let result = crate::scope_test_read_progress(
2938            Arc::clone(&progress),
2939            crate::scope_request_read_deadline(std::time::Duration::from_millis(25), async move {
2940                reader.query_all(deliberately_slow_read_statement()).await
2941            }),
2942        )
2943        .await;
2944        assert!(
2945            matches!(result, Err(StorageError::Timeout { .. })),
2946            "deadline must surface as a typed timeout, got {result:?}"
2947        );
2948
2949        let stopped_at = progress.load(std::sync::atomic::Ordering::SeqCst);
2950        assert!(stopped_at > 0, "deadline test never exercised SQLite work");
2951        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
2952        assert_eq!(
2953            progress.load(std::sync::atomic::Ordering::SeqCst),
2954            stopped_at,
2955            "deadline returned while SQLite kept consuming work"
2956        );
2957    }
2958
2959    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2960    async fn progress_handler_cleanup_failure_discards_pooled_connection() {
2961        let dir = tempfile::tempdir().unwrap();
2962        let config = PoolConfig {
2963            path: Some(dir.path().join("sql_bridge_cleanup_failure.db")),
2964            max_readers: 1,
2965            ..PoolConfig::default()
2966        };
2967        let pool = Arc::new(ConnectionPool::new(config).unwrap());
2968        let pool_for_read = Arc::clone(&pool);
2969        let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
2970        let result = crate::read_cancellation::scope_test_read_cleanup_failure(
2971            crate::scope_test_read_progress(
2972                Arc::clone(&progress),
2973                crate::read_cancellation::run_interruptible_read(
2974                    StorageCapability::Sql,
2975                    "cleanup_failure_probe",
2976                    move |scope| {
2977                        let mut guard = pool_for_read.reader().map_err(|error| {
2978                            StorageError::driver(
2979                                StorageCapability::Sql,
2980                                "cleanup_failure_probe",
2981                                error,
2982                            )
2983                        })?;
2984                        scope.run_pooled_reader(&mut guard, |conn| {
2985                            conn.query_row("SELECT 1", [], |row| row.get::<_, i64>(0))
2986                                .map_err(|error| map_rusqlite_err(error, "cleanup_failure_probe"))
2987                        })
2988                    },
2989                ),
2990            ),
2991        )
2992        .await;
2993        assert!(
2994            matches!(result, Err(StorageError::Internal(ref message)) if message.contains("clear failure")),
2995            "injected cleanup failure must be surfaced; got {result:?}"
2996        );
2997        assert_eq!(
2998            pool.available_readers(),
2999            1,
3000            "discard must install a replacement"
3001        );
3002
3003        let calls_after_failed_read = progress.load(std::sync::atomic::Ordering::SeqCst);
3004        let guard = pool.reader().unwrap();
3005        let sum: i64 = guard
3006            .conn()
3007            .query_row(
3008                "WITH RECURSIVE n(x) AS (VALUES(0) UNION ALL SELECT x + 1 FROM n WHERE x < 10000) \
3009                 SELECT sum(x) FROM n",
3010                [],
3011                |row| row.get(0),
3012            )
3013            .unwrap();
3014        assert_eq!(sum, 50_005_000);
3015        assert_eq!(
3016            progress.load(std::sync::atomic::Ordering::SeqCst),
3017            calls_after_failed_read,
3018            "a connection whose handler could not be cleared was reused"
3019        );
3020    }
3021
3022    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3023    async fn raw_pooled_reader_quarantines_cleanup_failure_during_unwind() {
3024        let dir = tempfile::tempdir().unwrap();
3025        let pool = Arc::new(
3026            ConnectionPool::new(PoolConfig {
3027                path: Some(dir.path().join("raw_reader_unwind_cleanup.db")),
3028                max_readers: 1,
3029                ..PoolConfig::default()
3030            })
3031            .unwrap(),
3032        );
3033        let worker_pool = Arc::clone(&pool);
3034        let result = crate::read_cancellation::scope_test_read_cleanup_failure(
3035            crate::read_cancellation::run_interruptible_read(
3036                StorageCapability::Sql,
3037                "raw_reader_unwind_cleanup",
3038                move |scope| {
3039                    let mut guard = worker_pool.reader().map_err(|error| {
3040                        StorageError::driver(
3041                            StorageCapability::Sql,
3042                            "raw_reader_unwind_cleanup",
3043                            error,
3044                        )
3045                    })?;
3046                    scope.with_pooled_reader(&mut guard, |conn| {
3047                        scope.run(conn, || -> khive_storage::types::StorageResult<()> {
3048                            panic!("injected raw reader panic after progress registration")
3049                        })
3050                    })
3051                },
3052            ),
3053        )
3054        .await;
3055        assert!(
3056            result.is_err(),
3057            "blocking panic must surface as a join error"
3058        );
3059        assert_eq!(
3060            pool.available_readers(),
3061            pool.max_readers(),
3062            "unwind cleanup failure must close and replace the raw pooled reader"
3063        );
3064    }
3065
3066    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3067    async fn raw_pooled_writer_retires_cleanup_failure_during_unwind() {
3068        let pool = Arc::new(ConnectionPool::new(PoolConfig::default()).unwrap());
3069        let worker_pool = Arc::clone(&pool);
3070        let result = crate::read_cancellation::scope_test_read_cleanup_failure(
3071            crate::read_cancellation::run_interruptible_read(
3072                StorageCapability::Sql,
3073                "raw_writer_unwind_cleanup",
3074                move |scope| {
3075                    let guard = worker_pool.try_writer().map_err(|error| {
3076                        StorageError::driver(
3077                            StorageCapability::Sql,
3078                            "raw_writer_unwind_cleanup",
3079                            error,
3080                        )
3081                    })?;
3082                    scope.with_pooled_writer(&worker_pool, &guard, |conn| {
3083                        scope.run(conn, || -> khive_storage::types::StorageResult<()> {
3084                            panic!("injected raw writer panic after progress registration")
3085                        })
3086                    })
3087                },
3088            ),
3089        )
3090        .await;
3091        assert!(
3092            result.is_err(),
3093            "blocking panic must surface as a join error"
3094        );
3095        assert!(
3096            pool.try_writer().is_err(),
3097            "unwind cleanup failure must retire the raw pooled writer"
3098        );
3099    }
3100
3101    #[test]
3102    fn query_row_converts_only_the_first_matching_row() {
3103        let conn = rusqlite::Connection::open_in_memory().unwrap();
3104        let statement = SqlStatement {
3105            sql: "WITH RECURSIVE rows(value) AS (\
3106                  SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 99\
3107                  ) SELECT value FROM rows ORDER BY value"
3108                .into(),
3109            params: vec![],
3110            label: None,
3111        };
3112
3113        ROW_CONVERSIONS.with(|count| count.set(0));
3114        let row = execute_query_row(&conn, &statement).unwrap().unwrap();
3115
3116        assert!(matches!(row.get("value"), Some(SqlValue::Integer(0))));
3117        ROW_CONVERSIONS.with(|count| assert_eq!(count.get(), 1));
3118    }
3119
3120    #[test]
3121    fn query_page_bounds_owned_rows_before_full_materialization() {
3122        let conn = rusqlite::Connection::open_in_memory().unwrap();
3123        let statement = SqlStatement {
3124            sql: "WITH RECURSIVE rows(value) AS (\
3125                  SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 99\
3126                  ) SELECT value FROM rows ORDER BY value"
3127                .into(),
3128            params: vec![],
3129            label: None,
3130        };
3131
3132        ROW_CONVERSIONS.with(|count| count.set(0));
3133        let rows = execute_query_page(
3134            &conn,
3135            &statement,
3136            &PageRequest {
3137                offset: 40,
3138                limit: 3,
3139            },
3140        )
3141        .unwrap();
3142
3143        assert_eq!(rows.len(), 3);
3144        assert!(matches!(rows[0].get("value"), Some(SqlValue::Integer(40))));
3145        assert!(matches!(rows[2].get("value"), Some(SqlValue::Integer(42))));
3146        ROW_CONVERSIONS.with(|count| assert_eq!(count.get(), 3));
3147    }
3148
3149    #[test]
3150    fn query_page_zero_limit_converts_no_rows_but_still_validates_sql() {
3151        let conn = rusqlite::Connection::open_in_memory().unwrap();
3152        let statement = SqlStatement {
3153            sql: "WITH RECURSIVE rows(value) AS (\
3154                  SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 99\
3155                  ) SELECT value FROM rows ORDER BY value"
3156                .into(),
3157            params: vec![],
3158            label: None,
3159        };
3160
3161        ROW_CONVERSIONS.with(|count| count.set(0));
3162        let rows = execute_query_page(
3163            &conn,
3164            &statement,
3165            &PageRequest {
3166                offset: 0,
3167                limit: 0,
3168            },
3169        )
3170        .unwrap();
3171
3172        assert!(rows.is_empty());
3173        ROW_CONVERSIONS.with(|count| assert_eq!(count.get(), 0));
3174
3175        let invalid = SqlStatement {
3176            sql: "SELECT FROM WHERE".into(),
3177            params: vec![],
3178            label: None,
3179        };
3180        assert!(
3181            execute_query_page(
3182                &conn,
3183                &invalid,
3184                &PageRequest {
3185                    offset: 0,
3186                    limit: 0
3187                }
3188            )
3189            .is_err(),
3190            "a zero-limit page must still fail on invalid SQL at prepare time"
3191        );
3192    }
3193
3194    #[test]
3195    fn cached_writer_prepare_preserves_the_single_statement_boundary() {
3196        let conn = rusqlite::Connection::open_in_memory().unwrap();
3197        assert!(matches!(
3198            prepare_cached_sql_statement(&conn, "SELECT 1; SELECT 2"),
3199            Err(rusqlite::Error::MultipleStatement)
3200        ));
3201    }
3202
3203    #[tokio::test]
3204    async fn queue_backed_execute_reuses_the_persistent_connection_statement_cache() {
3205        use rusqlite::hooks::{AuthAction, AuthContext, Authorization};
3206        use std::sync::atomic::{AtomicUsize, Ordering};
3207
3208        let dir = tempfile::tempdir().unwrap();
3209        let pool = Arc::new(
3210            ConnectionPool::new(PoolConfig {
3211                path: Some(dir.path().join("sql_bridge_writer_cache.db")),
3212                write_queue_enabled: Some(true),
3213                write_routing_strict: true,
3214                ..PoolConfig::default()
3215            })
3216            .unwrap(),
3217        );
3218        pool.writer()
3219            .unwrap()
3220            .conn()
3221            .execute_batch(
3222                "CREATE TABLE writer_cache_test (id INTEGER PRIMARY KEY, value TEXT NOT NULL)",
3223            )
3224            .unwrap();
3225
3226        let writer_task = pool
3227            .writer_task_handle()
3228            .unwrap()
3229            .expect("file-backed queue-enabled pool must expose its writer task");
3230        let prepare_count = Arc::new(AtomicUsize::new(0));
3231        let hook_count = Arc::clone(&prepare_count);
3232        writer_task
3233            .send_top_level(move |conn| {
3234                conn.authorizer(Some(move |context: AuthContext<'_>| {
3235                    if matches!(
3236                        context.action,
3237                        AuthAction::Insert { table_name } if table_name == "writer_cache_test"
3238                    ) {
3239                        hook_count.fetch_add(1, Ordering::SeqCst);
3240                    }
3241                    Authorization::Allow
3242                }))
3243                .map_err(|error| map_rusqlite_err(error, "test.install_authorizer"))
3244            })
3245            .await
3246            .unwrap();
3247
3248        let bridge = SqlBridge::new(Arc::clone(&pool), true);
3249        let mut writer = bridge.writer().await.unwrap();
3250        for id in [1, 2] {
3251            khive_storage::SqlWriter::execute(
3252                &mut *writer,
3253                SqlStatement {
3254                    sql: "INSERT INTO writer_cache_test (id, value) VALUES (?1, ?2)".into(),
3255                    params: vec![SqlValue::Integer(id), SqlValue::Text(format!("value-{id}"))],
3256                    label: None,
3257                },
3258            )
3259            .await
3260            .unwrap();
3261        }
3262
3263        assert_eq!(
3264            prepare_count.load(Ordering::SeqCst),
3265            1,
3266            "the second identical execute on the writer task's persistent connection must reuse \
3267             the cached SQLite statement instead of compiling it again"
3268        );
3269        writer_task
3270            .send_top_level(|conn| {
3271                conn.authorizer(None::<fn(AuthContext<'_>) -> Authorization>)
3272                    .map_err(|error| map_rusqlite_err(error, "test.remove_authorizer"))
3273            })
3274            .await
3275            .unwrap();
3276    }
3277
3278    #[test]
3279    fn inline_execute_batch_prepares_each_statement_once() {
3280        use rusqlite::hooks::{AuthAction, AuthContext, Authorization};
3281        use std::sync::atomic::{AtomicUsize, Ordering};
3282
3283        let conn = rusqlite::Connection::open_in_memory().unwrap();
3284        conn.execute_batch(
3285            "CREATE TABLE single_prepare_test (id INTEGER PRIMARY KEY, value TEXT NOT NULL)",
3286        )
3287        .unwrap();
3288        let prepare_count = Arc::new(AtomicUsize::new(0));
3289        let hook_count = Arc::clone(&prepare_count);
3290        conn.authorizer(Some(move |context: AuthContext<'_>| {
3291            if matches!(
3292                context.action,
3293                AuthAction::Insert { table_name } if table_name == "single_prepare_test"
3294            ) {
3295                hook_count.fetch_add(1, Ordering::SeqCst);
3296            }
3297            Authorization::Allow
3298        }))
3299        .unwrap();
3300
3301        let mut writer = InlineWriter {
3302            conn: &conn as *const rusqlite::Connection,
3303        };
3304        let affected = block_on_sync(khive_storage::SqlWriter::execute_batch(
3305            &mut writer,
3306            vec![SqlStatement {
3307                sql: "INSERT INTO single_prepare_test (id, value) VALUES (?1, ?2)".into(),
3308                params: vec![SqlValue::Integer(1), SqlValue::Text("once".into())],
3309                label: None,
3310            }],
3311        ))
3312        .expect("InlineWriter operations must resolve on their first poll")
3313        .expect("valid batch must execute");
3314
3315        assert_eq!(affected, 1);
3316        assert_eq!(
3317            prepare_count.load(Ordering::SeqCst),
3318            1,
3319            "classification and execution must share one prepared statement handle"
3320        );
3321        conn.authorizer(None::<fn(AuthContext<'_>) -> Authorization>)
3322            .unwrap();
3323    }
3324
3325    #[test]
3326    fn inline_execute_batch_preserves_schema_dependencies_between_statements() {
3327        let conn = rusqlite::Connection::open_in_memory().unwrap();
3328        let mut writer = InlineWriter {
3329            conn: &conn as *const rusqlite::Connection,
3330        };
3331
3332        let affected = block_on_sync(khive_storage::SqlWriter::execute_batch(
3333            &mut writer,
3334            vec![
3335                SqlStatement {
3336                    sql: "CREATE TABLE dependent_prepare_test (id INTEGER PRIMARY KEY)".into(),
3337                    params: vec![],
3338                    label: None,
3339                },
3340                SqlStatement {
3341                    sql: "INSERT INTO dependent_prepare_test (id) VALUES (1)".into(),
3342                    params: vec![],
3343                    label: None,
3344                },
3345            ],
3346        ))
3347        .expect("InlineWriter operations must resolve on their first poll")
3348        .expect("a later statement must be prepared after its prerequisite schema change");
3349
3350        assert_eq!(affected, 1);
3351        let count: i64 = conn
3352            .query_row("SELECT COUNT(*) FROM dependent_prepare_test", [], |row| {
3353                row.get(0)
3354            })
3355            .unwrap();
3356        assert_eq!(count, 1);
3357    }
3358
3359    #[tokio::test]
3360    async fn pool_backed_query_page_beyond_result_set_returns_empty() {
3361        let config = PoolConfig {
3362            path: None,
3363            ..PoolConfig::default()
3364        };
3365        let pool = Arc::new(ConnectionPool::new(config).unwrap());
3366        {
3367            let writer = pool.writer().unwrap();
3368            writer
3369                .conn()
3370                .execute_batch(
3371                    "CREATE TABLE page_test (id INTEGER PRIMARY KEY, val TEXT NOT NULL);\
3372                     INSERT INTO page_test (id, val) VALUES (1, 'a'), (2, 'b'), (3, 'c');",
3373                )
3374                .unwrap();
3375        }
3376        let bridge = SqlBridge::new(Arc::clone(&pool), false);
3377
3378        let statement = || SqlStatement {
3379            sql: "SELECT val FROM page_test ORDER BY id".into(),
3380            params: vec![],
3381            label: None,
3382        };
3383
3384        let mut reader = bridge.reader().await.unwrap();
3385        let page = reader
3386            .query_page(
3387                statement(),
3388                PageRequest {
3389                    offset: 1,
3390                    limit: 2,
3391                },
3392            )
3393            .await
3394            .unwrap();
3395        assert_eq!(page.len(), 2);
3396        assert!(matches!(page[0].get("val"), Some(SqlValue::Text(v)) if v == "b"));
3397        assert!(matches!(page[1].get("val"), Some(SqlValue::Text(v)) if v == "c"));
3398
3399        let empty = reader
3400            .query_page(
3401                statement(),
3402                PageRequest {
3403                    offset: 99,
3404                    limit: 10,
3405                },
3406            )
3407            .await
3408            .unwrap();
3409        assert!(
3410            empty.is_empty(),
3411            "offset past the last row must return an empty page, got {empty:?}"
3412        );
3413        drop(reader);
3414
3415        let mut writer = bridge.writer().await.unwrap();
3416        let empty = writer
3417            .query_page(
3418                statement(),
3419                PageRequest {
3420                    offset: 99,
3421                    limit: 10,
3422                },
3423            )
3424            .await
3425            .unwrap();
3426        assert!(
3427            empty.is_empty(),
3428            "offset past the last row must return an empty page, got {empty:?}"
3429        );
3430    }
3431
3432    #[tokio::test]
3433    async fn file_bridge_scopes_reader_permits_to_operations_and_caps_writer_handles() {
3434        let dir = tempfile::tempdir().unwrap();
3435        let config = PoolConfig {
3436            path: Some(dir.path().join("sql_bridge_handle_cap.db")),
3437            write_queue_enabled: Some(false),
3438            max_readers: 2,
3439            checkout_timeout: std::time::Duration::from_millis(20),
3440            ..PoolConfig::default()
3441        };
3442        let pool = Arc::new(ConnectionPool::new(config).unwrap());
3443        let bridge = SqlBridge::new(Arc::clone(&pool), true);
3444        let second_bridge = SqlBridge::new(Arc::clone(&pool), true);
3445
3446        let mut retained_readers = Vec::new();
3447        for expected in 0..3 {
3448            let mut reader = second_bridge.reader().await.unwrap();
3449            let value = reader
3450                .query_scalar(SqlStatement {
3451                    sql: format!("SELECT {expected}"),
3452                    params: vec![],
3453                    label: None,
3454                })
3455                .await
3456                .unwrap();
3457            assert!(matches!(value, Some(SqlValue::Integer(value)) if value == expected));
3458            retained_readers.push(reader);
3459        }
3460        assert_eq!(retained_readers.len(), 3);
3461
3462        let mut additional_reader = bridge.reader().await.unwrap();
3463        let page = additional_reader
3464            .query_page(
3465                SqlStatement {
3466                    sql: "WITH RECURSIVE rows(value) AS (\
3467                          SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 9\
3468                          ) SELECT value FROM rows ORDER BY value"
3469                        .into(),
3470                    params: vec![],
3471                    label: None,
3472                },
3473                PageRequest {
3474                    offset: 7,
3475                    limit: 2,
3476                },
3477            )
3478            .await
3479            .unwrap();
3480        assert_eq!(page.len(), 2);
3481        assert!(matches!(page[0].get("value"), Some(SqlValue::Integer(7))));
3482        assert!(matches!(page[1].get("value"), Some(SqlValue::Integer(8))));
3483        drop((additional_reader, retained_readers));
3484
3485        let writer = bridge.writer().await.unwrap();
3486        let writer_error = match second_bridge.writer().await {
3487            Ok(_) => panic!("a second live writer handle exceeded the one-handle cap"),
3488            Err(error) => error,
3489        };
3490        assert!(matches!(
3491            writer_error,
3492            StorageError::AdmissionTimeout { ref operation, .. }
3493                if operation.as_ref() == "sql_bridge.writer_handle"
3494        ));
3495        drop(writer);
3496        let writer_after_release = bridge.writer().await.unwrap();
3497        drop(writer_after_release);
3498    }
3499
3500    #[tokio::test]
3501    #[serial_test::serial(tx_registry)]
3502    async fn cached_read_transaction_retains_one_permit_until_commit_or_rollback() {
3503        let dir = tempfile::tempdir().unwrap();
3504        let config = PoolConfig {
3505            path: Some(dir.path().join("sql_bridge_reader_tx_control.db")),
3506            write_queue_enabled: Some(true),
3507            max_readers: 1,
3508            checkout_timeout: std::time::Duration::from_millis(20),
3509            ..PoolConfig::default()
3510        };
3511        let pool = Arc::new(ConnectionPool::new(config).unwrap());
3512        let origin = pool.origin();
3513        let origin_view = database_tx_view(&pool);
3514        let unrelated_view = khive_storage::tx_registry::TxOriginFilter::Secondary(
3515            khive_storage::tx_registry::DbIdentity::new("unrelated-sql-bridge.db"),
3516        );
3517        let bridge = SqlBridge::new(Arc::clone(&pool), true);
3518        let mut reader = bridge.reader().await.unwrap();
3519        let mut contender = bridge.reader().await.unwrap();
3520
3521        assert!(
3522            khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
3523            "an idle cached reader must not register a transaction"
3524        );
3525
3526        reader
3527            .query_all(SqlStatement {
3528                sql: "BEGIN DEFERRED".into(),
3529                params: vec![],
3530                label: None,
3531            })
3532            .await
3533            .expect("BEGIN DEFERRED must open an admitted cached-reader snapshot");
3534        let opened = khive_storage::tx_registry::oldest_for(&origin_view)
3535            .expect("successful BEGIN must register the cached-reader transaction");
3536        assert_eq!(opened.label.as_deref(), Some(CACHED_READ_TRANSACTION_LABEL));
3537        assert_eq!(opened.origin, origin);
3538        assert!(
3539            khive_storage::tx_registry::oldest_for(&unrelated_view).is_none(),
3540            "the read transaction must be attributed only to its own backend"
3541        );
3542        assert_eq!(
3543            pool.sql_bridge_reader_slots().available_permits(),
3544            0,
3545            "the successful BEGIN must retain its operation permit"
3546        );
3547
3548        let value = reader
3549            .query_scalar(SqlStatement {
3550                sql: "SELECT 7".into(),
3551                params: vec![],
3552                label: None,
3553            })
3554            .await
3555            .expect("a query inside the admitted transaction must reuse its retained permit");
3556        assert!(matches!(value, Some(SqlValue::Integer(7))));
3557        assert_eq!(
3558            khive_storage::tx_registry::oldest_for(&origin_view)
3559                .expect("queries must retain the transaction registration")
3560                .id,
3561            opened.id,
3562            "queries inside the transaction must retain the original span"
3563        );
3564
3565        let blocked = contender
3566            .query_scalar(SqlStatement {
3567                sql: "SELECT 8".into(),
3568                params: vec![],
3569                label: None,
3570            })
3571            .await;
3572        assert!(
3573            matches!(
3574                &blocked,
3575                Err(StorageError::Timeout { operation })
3576                    if operation.as_ref() == "sql_bridge.reader_operation"
3577            ),
3578            "a second logical read must contend with the admitted transaction \
3579             and time out with the ADR-005 reader contract error; got {blocked:?}"
3580        );
3581
3582        reader
3583            .query_all(SqlStatement {
3584                sql: "COMMIT".into(),
3585                params: vec![],
3586                label: None,
3587            })
3588            .await
3589            .expect("COMMIT must close the admitted cached-reader snapshot");
3590        assert!(
3591            khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
3592            "COMMIT must deregister after SQLite returns to autocommit"
3593        );
3594        assert_eq!(
3595            pool.sql_bridge_reader_slots().available_permits(),
3596            1,
3597            "COMMIT may release the permit only after autocommit is restored"
3598        );
3599        let value = contender
3600            .query_scalar(SqlStatement {
3601                sql: "SELECT 8".into(),
3602                params: vec![],
3603                label: None,
3604            })
3605            .await
3606            .expect("the contender must run after COMMIT releases admission");
3607        assert!(matches!(value, Some(SqlValue::Integer(8))));
3608
3609        reader
3610            .query_all(SqlStatement {
3611                sql: "BEGIN TRANSACTION".into(),
3612                params: vec![],
3613                label: None,
3614            })
3615            .await
3616            .expect("plain deferred BEGIN TRANSACTION must also be admitted");
3617        let reopened = khive_storage::tx_registry::oldest_for(&origin_view)
3618            .expect("the second successful BEGIN must register a fresh span");
3619        assert_ne!(reopened.id, opened.id);
3620        assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 0);
3621        let nested = reader
3622            .query_all(SqlStatement {
3623                sql: "ROLLBACK TO stale_snapshot".into(),
3624                params: vec![],
3625                label: None,
3626            })
3627            .await;
3628        assert!(
3629            matches!(&nested, Err(StorageError::InvalidInput { .. })),
3630            "ROLLBACK TO requires unsupported nested state; got {nested:?}"
3631        );
3632        assert_eq!(
3633            pool.sql_bridge_reader_slots().available_permits(),
3634            0,
3635            "rejected nested control must not release the still-live transaction admission"
3636        );
3637        assert_eq!(
3638            khive_storage::tx_registry::oldest_for(&origin_view)
3639                .expect("ROLLBACK TO rejection must retain the live span")
3640                .id,
3641            reopened.id
3642        );
3643        let savepoint = reader
3644            .query_all(SqlStatement {
3645                sql: "SAVEPOINT nested_snapshot".into(),
3646                params: vec![],
3647                label: None,
3648            })
3649            .await;
3650        assert!(
3651            matches!(&savepoint, Err(StorageError::InvalidInput { .. })),
3652            "SAVEPOINT must be rejected inside the admitted transaction; got {savepoint:?}"
3653        );
3654        assert_eq!(
3655            khive_storage::tx_registry::oldest_for(&origin_view)
3656                .expect("SAVEPOINT rejection must retain the live span")
3657                .id,
3658            reopened.id
3659        );
3660        reader
3661            .query_all(SqlStatement {
3662                sql: "ROLLBACK".into(),
3663                params: vec![],
3664                label: None,
3665            })
3666            .await
3667            .expect("ROLLBACK must close the admitted cached-reader snapshot");
3668        assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
3669        assert!(
3670            khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
3671            "full ROLLBACK must deregister after SQLite returns to autocommit"
3672        );
3673    }
3674
3675    #[tokio::test]
3676    async fn failed_cached_reader_begin_does_not_register_a_transaction() {
3677        use rusqlite::hooks::{AuthAction, AuthContext, Authorization, TransactionOperation};
3678
3679        fn deny_begin(ctx: AuthContext<'_>) -> Authorization {
3680            match ctx.action {
3681                AuthAction::Transaction {
3682                    operation: TransactionOperation::Begin,
3683                } => Authorization::Deny,
3684                _ => Authorization::Allow,
3685            }
3686        }
3687
3688        let dir = tempfile::tempdir().unwrap();
3689        let config = PoolConfig {
3690            path: Some(dir.path().join("sql_bridge_reader_failed_begin.db")),
3691            max_readers: 1,
3692            ..PoolConfig::default()
3693        };
3694        let pool = Arc::new(ConnectionPool::new(config).unwrap());
3695        let origin_view = database_tx_view(&pool);
3696        let conn = open_standalone_reader(&pool).unwrap();
3697        conn.authorizer(Some(deny_begin)).unwrap();
3698        let mut reader = SqliteReader {
3699            handle: Some(StandaloneHandle {
3700                conn,
3701                _retained_slot: None,
3702                read_transaction_slot: None,
3703            }),
3704            pool: Arc::clone(&pool),
3705        };
3706
3707        let begin = reader
3708            .query_all(SqlStatement {
3709                sql: "BEGIN DEFERRED".into(),
3710                params: vec![],
3711                label: None,
3712            })
3713            .await;
3714        assert!(begin.is_err(), "the authorizer must reject BEGIN");
3715        assert!(
3716            khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
3717            "a failed BEGIN must never enter the transaction registry"
3718        );
3719        assert_eq!(
3720            pool.sql_bridge_reader_slots().available_permits(),
3721            1,
3722            "a failed BEGIN must return the operation permit"
3723        );
3724    }
3725
3726    #[tokio::test]
3727    #[serial_test::serial(tx_registry)]
3728    async fn failed_cached_reader_rollback_deregisters_only_when_connection_is_discarded() {
3729        use rusqlite::hooks::{AuthAction, AuthContext, Authorization, TransactionOperation};
3730
3731        fn deny_rollback(ctx: AuthContext<'_>) -> Authorization {
3732            match ctx.action {
3733                AuthAction::Transaction {
3734                    operation: TransactionOperation::Rollback,
3735                } => Authorization::Deny,
3736                _ => Authorization::Allow,
3737            }
3738        }
3739
3740        let dir = tempfile::tempdir().unwrap();
3741        let config = PoolConfig {
3742            path: Some(dir.path().join("sql_bridge_reader_failed_rollback.db")),
3743            max_readers: 1,
3744            ..PoolConfig::default()
3745        };
3746        let pool = Arc::new(ConnectionPool::new(config).unwrap());
3747        let origin_view = database_tx_view(&pool);
3748        let conn = open_standalone_reader(&pool).unwrap();
3749        let mut reader = SqliteReader {
3750            handle: Some(StandaloneHandle {
3751                conn,
3752                _retained_slot: None,
3753                read_transaction_slot: None,
3754            }),
3755            pool: Arc::clone(&pool),
3756        };
3757
3758        reader
3759            .query_all(SqlStatement {
3760                sql: "BEGIN DEFERRED".into(),
3761                params: vec![],
3762                label: None,
3763            })
3764            .await
3765            .expect("BEGIN must establish the registered transaction");
3766        let opened = khive_storage::tx_registry::oldest_for(&origin_view)
3767            .expect("the admitted transaction must be registered");
3768        reader
3769            .handle
3770            .as_ref()
3771            .expect("reader must retain its connection")
3772            .conn
3773            .authorizer(Some(deny_rollback))
3774            .unwrap();
3775
3776        let rollback = reader
3777            .query_all(SqlStatement {
3778                sql: "ROLLBACK".into(),
3779                params: vec![],
3780                label: None,
3781            })
3782            .await;
3783        assert!(rollback.is_err(), "the authorizer must reject ROLLBACK");
3784        assert_eq!(
3785            khive_storage::tx_registry::oldest_for(&origin_view)
3786                .expect("failed ROLLBACK must retain registry evidence")
3787                .id,
3788            opened.id
3789        );
3790        assert_eq!(
3791            pool.sql_bridge_reader_slots().available_permits(),
3792            0,
3793            "failed ROLLBACK must retain reader admission"
3794        );
3795
3796        drop(reader);
3797        assert!(
3798            khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
3799            "discarding the connection must not leak its registry entry"
3800        );
3801        assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
3802    }
3803
3804    #[tokio::test]
3805    #[serial_test::serial(tx_registry)]
3806    async fn cached_read_only_handles_reject_unsupported_transaction_control_without_consumption() {
3807        let dir = tempfile::tempdir().unwrap();
3808        let config = PoolConfig {
3809            path: Some(
3810                dir.path()
3811                    .join("sql_bridge_reader_unsupported_tx_control.db"),
3812            ),
3813            write_queue_enabled: Some(true),
3814            max_readers: 1,
3815            ..PoolConfig::default()
3816        };
3817        let pool = Arc::new(ConnectionPool::new(config).unwrap());
3818        let bridge = SqlBridge::new(Arc::clone(&pool), true);
3819
3820        let mut reader = bridge.reader().await.unwrap();
3821        for (sql, keyword) in [
3822            ("BEGIN IMMEDIATE", "BEGIN"),
3823            ("BEGIN EXCLUSIVE", "BEGIN"),
3824            // Trailing-mode spellings parse in SQLite as a NAMED deferred
3825            // transaction, but the mode keyword in name position reads as
3826            // lock intent; the classifier must refuse rather than launder
3827            // them into a deferred start the cached reader would then hold.
3828            ("BEGIN TRANSACTION IMMEDIATE", "BEGIN"),
3829            ("BEGIN TRANSACTION EXCLUSIVE", "BEGIN"),
3830            ("BEGIN DEFERRED TRANSACTION trailing", "BEGIN"),
3831            // Quoted/bracketed tails tokenize as no identifier at all, so a
3832            // classifier that stops at the tokenizer's `None` reads them as
3833            // an accepted form's end. They must be refused exactly like the
3834            // bare-word spellings.
3835            ("BEGIN TRANSACTION \"IMMEDIATE\"", "BEGIN"),
3836            ("BEGIN TRANSACTION [IMMEDIATE]", "BEGIN"),
3837            ("BEGIN TRANSACTION `IMMEDIATE`", "BEGIN"),
3838            ("BEGIN TRANSACTION 'IMMEDIATE'", "BEGIN"),
3839            ("BEGIN \"DEFERRED\"", "BEGIN"),
3840            ("BEGIN; COMMIT", "BEGIN"),
3841            ("START TRANSACTION", "START"),
3842            ("COMMIT", "COMMIT"),
3843        ] {
3844            let rejected = reader
3845                .query_all(SqlStatement {
3846                    sql: sql.into(),
3847                    params: vec![],
3848                    label: None,
3849                })
3850                .await;
3851            assert!(
3852                matches!(
3853                    &rejected,
3854                    Err(StorageError::InvalidInput { message, .. })
3855                        if message.contains(keyword)
3856                ),
3857                "unsupported cached-reader control {sql:?} must fail closed; got {rejected:?}"
3858            );
3859            assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
3860        }
3861
3862        let mut queue_backed_writer = bridge.writer().await.unwrap();
3863        let rejected = queue_backed_writer
3864            .query_all(SqlStatement {
3865                sql: "SAVEPOINT stale_snapshot".into(),
3866                params: vec![],
3867                label: None,
3868            })
3869            .await;
3870        assert!(
3871            matches!(
3872                &rejected,
3873                Err(StorageError::InvalidInput {
3874                    operation,
3875                    message,
3876                    ..
3877                }) if operation.as_ref() == "writer.query_all"
3878                    && message.contains("transaction control")
3879                    && message.contains("SAVEPOINT")
3880            ),
3881            "a queue-backed writer's cached read-only connection must reject transaction \
3882             control; got {rejected:?}"
3883        );
3884        assert_eq!(
3885            pool.sql_bridge_reader_slots().available_permits(),
3886            1,
3887            "queue-backed rejection must leave the operation permit available"
3888        );
3889        let value = queue_backed_writer
3890            .query_scalar(SqlStatement {
3891                sql: "SELECT 8".into(),
3892                params: vec![],
3893                label: None,
3894            })
3895            .await
3896            .expect("transaction-control rejection must not consume the queue-backed handle");
3897        assert!(matches!(value, Some(SqlValue::Integer(8))));
3898
3899        queue_backed_writer
3900            .query_all(SqlStatement {
3901                sql: "BEGIN DEFERRED".into(),
3902                params: vec![],
3903                label: None,
3904            })
3905            .await
3906            .expect("queue-backed cached reader must share explicit read-transaction admission");
3907        assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 0);
3908        queue_backed_writer
3909            .query_all(SqlStatement {
3910                sql: "END".into(),
3911                params: vec![],
3912                label: None,
3913            })
3914            .await
3915            .expect("END must release queue-backed cached-reader admission");
3916        assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
3917    }
3918
3919    #[tokio::test]
3920    #[serial_test::serial(tx_registry)]
3921    async fn dropping_cached_reader_transaction_closes_snapshot_before_releasing_permit() {
3922        let dir = tempfile::tempdir().unwrap();
3923        let config = PoolConfig {
3924            path: Some(dir.path().join("sql_bridge_reader_tx_drop.db")),
3925            max_readers: 1,
3926            checkout_timeout: std::time::Duration::from_millis(20),
3927            ..PoolConfig::default()
3928        };
3929        let pool = Arc::new(ConnectionPool::new(config).unwrap());
3930        let origin_view = database_tx_view(&pool);
3931        let bridge = SqlBridge::new(Arc::clone(&pool), true);
3932        let mut reader = bridge.reader().await.unwrap();
3933        let mut contender = bridge.reader().await.unwrap();
3934
3935        reader
3936            .query_all(SqlStatement {
3937                sql: "BEGIN DEFERRED".into(),
3938                params: vec![],
3939                label: None,
3940            })
3941            .await
3942            .expect("begin admitted transaction");
3943        reader
3944            .query_all(SqlStatement {
3945                sql: "SELECT * FROM sqlite_schema".into(),
3946                params: vec![],
3947                label: None,
3948            })
3949            .await
3950            .expect("materialize read snapshot");
3951        assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 0);
3952        assert!(
3953            khive_storage::tx_registry::oldest_for(&origin_view).is_some(),
3954            "the live snapshot must remain registered until handle drop"
3955        );
3956
3957        drop(reader);
3958        assert!(
3959            khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
3960            "handle drop must close SQLite before deregistering the snapshot"
3961        );
3962        assert_eq!(
3963            pool.sql_bridge_reader_slots().available_permits(),
3964            1,
3965            "dropping the handle must close its transaction before returning admission"
3966        );
3967        contender
3968            .query_all(SqlStatement {
3969                sql: "SELECT * FROM sqlite_schema".into(),
3970                params: vec![],
3971                label: None,
3972            })
3973            .await
3974            .expect("a new operation must run after the transactional handle drops");
3975    }
3976
3977    /// #1846 regression: a cached-reader explicit read transaction that is
3978    /// never explicitly finished (a stuck/leaked caller that keeps reusing
3979    /// the handle without COMMIT/ROLLBACK) would otherwise pin the WAL
3980    /// snapshot open for as long as the caller kept calling in. Without the
3981    /// age check this reddens: the second `query_all` would return the row
3982    /// materialized inside the still-open transaction instead of an error,
3983    /// and `tx_registry::oldest_for` would keep reporting the same span
3984    /// open past `read_tx_max_age`.
3985    #[tokio::test]
3986    #[serial_test::serial(tx_registry)]
3987    async fn expired_cached_reader_transaction_is_rolled_back_on_reuse() {
3988        let dir = tempfile::tempdir().unwrap();
3989        let config = PoolConfig {
3990            path: Some(dir.path().join("sql_bridge_reader_tx_max_age.db")),
3991            max_readers: 1,
3992            checkout_timeout: std::time::Duration::from_millis(20),
3993            read_tx_max_age: std::time::Duration::from_millis(20),
3994            ..PoolConfig::default()
3995        };
3996        let pool = Arc::new(ConnectionPool::new(config).unwrap());
3997        let origin_view = database_tx_view(&pool);
3998        let bridge = SqlBridge::new(Arc::clone(&pool), true);
3999        let mut reader = bridge.reader().await.unwrap();
4000
4001        reader
4002            .query_all(SqlStatement {
4003                sql: "BEGIN DEFERRED".into(),
4004                params: vec![],
4005                label: None,
4006            })
4007            .await
4008            .expect("begin admitted transaction");
4009        reader
4010            .query_all(SqlStatement {
4011                sql: "SELECT * FROM sqlite_schema".into(),
4012                params: vec![],
4013                label: None,
4014            })
4015            .await
4016            .expect("materialize read snapshot");
4017        assert!(
4018            khive_storage::tx_registry::oldest_for(&origin_view).is_some(),
4019            "the open transaction must be registered before it ages out"
4020        );
4021
4022        tokio::time::sleep(std::time::Duration::from_millis(40)).await;
4023
4024        let evictions_before = crate::checkpoint::read_tx_max_age_evictions();
4025        let error = reader
4026            .query_all(SqlStatement {
4027                sql: "SELECT * FROM sqlite_schema".into(),
4028                params: vec![],
4029                label: None,
4030            })
4031            .await
4032            .expect_err("reusing a transaction past read_tx_max_age must be refused");
4033        assert!(
4034            error.is_retryable(),
4035            "an evicted-transaction error must be retryable so the caller can open a fresh \
4036             snapshot: {error}"
4037        );
4038        match &error {
4039            StorageError::ReadTransactionAgeEvicted {
4040                operation,
4041                max_age_secs,
4042            } => {
4043                assert_eq!(operation.as_ref(), "query_all");
4044                assert_eq!(
4045                    *max_age_secs, 0,
4046                    "a 20ms read_tx_max_age truncates to 0 whole seconds"
4047                );
4048            }
4049            other => panic!(
4050                "a clean age-triggered rollback must surface the dedicated \
4051                 ReadTransactionAgeEvicted variant, not a generic classification: {other:?}"
4052            ),
4053        }
4054        assert_eq!(
4055            crate::checkpoint::read_tx_max_age_evictions(),
4056            evictions_before + 1,
4057            "the eviction must be counted in the #1846 diagnostics gauge"
4058        );
4059        assert!(
4060            khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
4061            "the expired transaction must be rolled back and deregistered rather than \
4062             continuing to pin the WAL snapshot"
4063        );
4064
4065        reader
4066            .query_all(SqlStatement {
4067                sql: "SELECT * FROM sqlite_schema".into(),
4068                params: vec![],
4069                label: None,
4070            })
4071            .await
4072            .expect("the handle must remain usable for a fresh autocommit read after eviction");
4073    }
4074
4075    /// #1846 follow-up: the age-eviction branch
4076    /// has a rollback-failure path distinct from
4077    /// `failed_cached_reader_rollback_deregisters_only_when_connection_is_discarded`
4078    /// above (which covers an explicit caller-issued `ROLLBACK`, not the
4079    /// age-triggered cleanup rollback). When SQLite denies the age-triggered
4080    /// `ROLLBACK`, the branch must still discard the poisoned connection,
4081    /// deregister the expired transaction span, and release the reader
4082    /// admission permit rather than leaking either.
4083    #[tokio::test]
4084    #[serial_test::serial(tx_registry)]
4085    async fn expired_cached_reader_transaction_rollback_denial_discards_connection_and_releases_admission(
4086    ) {
4087        use rusqlite::hooks::{AuthAction, AuthContext, Authorization, TransactionOperation};
4088
4089        fn deny_rollback(ctx: AuthContext<'_>) -> Authorization {
4090            match ctx.action {
4091                AuthAction::Transaction {
4092                    operation: TransactionOperation::Rollback,
4093                } => Authorization::Deny,
4094                _ => Authorization::Allow,
4095            }
4096        }
4097
4098        let dir = tempfile::tempdir().unwrap();
4099        let config = PoolConfig {
4100            path: Some(
4101                dir.path()
4102                    .join("sql_bridge_reader_tx_max_age_rollback_denied.db"),
4103            ),
4104            max_readers: 1,
4105            checkout_timeout: std::time::Duration::from_millis(20),
4106            read_tx_max_age: std::time::Duration::from_millis(20),
4107            ..PoolConfig::default()
4108        };
4109        let pool = Arc::new(ConnectionPool::new(config).unwrap());
4110        let origin_view = database_tx_view(&pool);
4111        let conn = open_standalone_reader(&pool).unwrap();
4112        let mut reader = SqliteReader {
4113            handle: Some(StandaloneHandle {
4114                conn,
4115                _retained_slot: None,
4116                read_transaction_slot: None,
4117            }),
4118            pool: Arc::clone(&pool),
4119        };
4120
4121        reader
4122            .query_all(SqlStatement {
4123                sql: "BEGIN DEFERRED".into(),
4124                params: vec![],
4125                label: None,
4126            })
4127            .await
4128            .expect("begin admitted transaction");
4129        reader
4130            .query_all(SqlStatement {
4131                sql: "SELECT * FROM sqlite_schema".into(),
4132                params: vec![],
4133                label: None,
4134            })
4135            .await
4136            .expect("materialize read snapshot");
4137        assert!(
4138            khive_storage::tx_registry::oldest_for(&origin_view).is_some(),
4139            "the open transaction must be registered before it ages out"
4140        );
4141
4142        reader
4143            .handle
4144            .as_ref()
4145            .expect("reader must retain its connection")
4146            .conn
4147            .authorizer(Some(deny_rollback))
4148            .unwrap();
4149
4150        tokio::time::sleep(std::time::Duration::from_millis(40)).await;
4151
4152        let evictions_before = crate::checkpoint::read_tx_max_age_evictions();
4153        let error = reader
4154            .query_all(SqlStatement {
4155                sql: "SELECT * FROM sqlite_schema".into(),
4156                params: vec![],
4157                label: None,
4158            })
4159            .await
4160            .expect_err("a denied rollback on an expired transaction must surface an error");
4161        assert!(
4162            error.is_retryable(),
4163            "even a failed cleanup rollback must remain classified retryable so callers open a \
4164             fresh handle: {error}"
4165        );
4166        match &error {
4167            StorageError::ReadTransactionAgeEvictionCleanupFailed {
4168                operation,
4169                max_age_secs,
4170                message,
4171            } => {
4172                assert_eq!(operation.as_ref(), "query_all");
4173                assert_eq!(
4174                    *max_age_secs, 0,
4175                    "a 20ms read_tx_max_age truncates to 0 whole seconds"
4176                );
4177                assert!(
4178                    message.contains("rollback failed"),
4179                    "the failure must be attributable to the denied ROLLBACK, not silent \
4180                     success: {message}"
4181                );
4182            }
4183            other => panic!(
4184                "a denied cleanup rollback must surface the dedicated \
4185                 ReadTransactionAgeEvictionCleanupFailed variant, not a generic Transaction \
4186                 error the caller cannot machine-detect: {other:?}"
4187            ),
4188        }
4189        assert_eq!(
4190            crate::checkpoint::read_tx_max_age_evictions(),
4191            evictions_before + 1,
4192            "the eviction attempt must still be counted even though cleanup failed"
4193        );
4194        assert!(
4195            khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
4196            "a denied rollback must discard the connection and deregister the expired \
4197             transaction span rather than leaking it"
4198        );
4199        assert_eq!(
4200            pool.sql_bridge_reader_slots().available_permits(),
4201            1,
4202            "discarding the poisoned connection must release the reader admission slot"
4203        );
4204
4205        let reuse = reader
4206            .query_all(SqlStatement {
4207                sql: "SELECT * FROM sqlite_schema".into(),
4208                params: vec![],
4209                label: None,
4210            })
4211            .await;
4212        let message = match reuse {
4213            Err(StorageError::Pool { message, .. }) => message,
4214            other => panic!(
4215                "reusing this discarded reader must fail loudly with 'connection already \
4216                 consumed' rather than silently reopening; got {other:?}"
4217            ),
4218        };
4219        assert!(
4220            message.contains("connection already consumed"),
4221            "expected the discarded reader's reuse error to name the pinned failure; got \
4222             {message:?}"
4223        );
4224    }
4225
4226    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4227    #[serial_test::serial(tx_registry)]
4228    async fn cancelled_cached_reader_transaction_releases_guards_after_connection_closes() {
4229        let dir = tempfile::tempdir().unwrap();
4230        let config = PoolConfig {
4231            path: Some(dir.path().join("sql_bridge_reader_tx_drop_cancel.db")),
4232            max_readers: 1,
4233            checkout_timeout: std::time::Duration::from_millis(50),
4234            ..PoolConfig::default()
4235        };
4236        let pool = Arc::new(ConnectionPool::new(config).unwrap());
4237        let origin_view = database_tx_view(&pool);
4238        let bridge = SqlBridge::new(Arc::clone(&pool), true);
4239        let mut reader = SqliteReader {
4240            handle: Some(open_cached_reader_handle(Arc::clone(&pool)).await.unwrap()),
4241            pool: Arc::clone(&pool),
4242        };
4243        let mut contender = bridge.reader().await.unwrap();
4244
4245        reader
4246            .query_all(SqlStatement {
4247                sql: "BEGIN DEFERRED".into(),
4248                params: vec![],
4249                label: None,
4250            })
4251            .await
4252            .expect("begin admitted transaction");
4253        assert!(
4254            khive_storage::tx_registry::oldest_for(&origin_view).is_some(),
4255            "the explicit transaction must be registered before cancellation"
4256        );
4257        assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 0);
4258
4259        let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
4260        let query = tokio::spawn(crate::scope_test_read_progress(
4261            Arc::clone(&progress),
4262            async move { reader.query_all(deliberately_slow_read_statement()).await },
4263        ));
4264        wait_for_progress(progress.as_ref()).await;
4265        query.abort();
4266        assert!(matches!(query.await, Err(error) if error.is_cancelled()));
4267        tokio::time::timeout(std::time::Duration::from_secs(1), async {
4268            while khive_storage::tx_registry::oldest_for(&origin_view).is_some()
4269                || pool.sql_bridge_reader_slots().available_permits() != 1
4270            {
4271                tokio::task::yield_now().await;
4272            }
4273        })
4274        .await
4275        .expect("connection cleanup leaked transaction evidence or reader admission");
4276        contender
4277            .query_all(SqlStatement {
4278                sql: "SELECT 1".into(),
4279                params: vec![],
4280                label: None,
4281            })
4282            .await
4283            .expect("admission must recover after the cancelled connection closes");
4284    }
4285
4286    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4287    #[serial_test::serial(tx_registry)]
4288    async fn cancelled_cached_reader_rolls_back_releases_wal_and_clears_handler() {
4289        let dir = tempfile::tempdir().unwrap();
4290        let config = PoolConfig {
4291            path: Some(dir.path().join("sql_bridge_reader_tx_request_cancel.db")),
4292            max_readers: 1,
4293            checkout_timeout: std::time::Duration::from_millis(500),
4294            ..PoolConfig::default()
4295        };
4296        let pool = Arc::new(ConnectionPool::new(config).unwrap());
4297        let bridge = SqlBridge::new(Arc::clone(&pool), true);
4298        let writer = open_standalone_writer(&pool).unwrap();
4299        writer
4300            .execute_batch(
4301                "CREATE TABLE snapshot_probe(id INTEGER PRIMARY KEY, value TEXT NOT NULL); \
4302                 INSERT INTO snapshot_probe(value) VALUES ('seed');",
4303            )
4304            .unwrap();
4305        let mut reader = bridge.reader().await.unwrap();
4306        let mut contender = bridge.reader().await.unwrap();
4307
4308        reader
4309            .query_all(SqlStatement {
4310                sql: "BEGIN DEFERRED".into(),
4311                params: vec![],
4312                label: None,
4313            })
4314            .await
4315            .expect("begin admitted transaction");
4316        reader
4317            .query_all(SqlStatement {
4318                sql: "SELECT * FROM snapshot_probe".into(),
4319                params: vec![],
4320                label: None,
4321            })
4322            .await
4323            .expect("materialize a real WAL snapshot");
4324        writer
4325            .execute_batch(
4326                "WITH RECURSIVE rows(value) AS (\
4327                 SELECT 1 UNION ALL SELECT value + 1 FROM rows WHERE value < 100\
4328                 ) INSERT INTO snapshot_probe(value) SELECT printf('row-%d', value) FROM rows;",
4329            )
4330            .unwrap();
4331        let (_, log_before, checkpointed_before) = passive_checkpoint(&writer);
4332        assert!(
4333            log_before > checkpointed_before,
4334            "the explicit reader snapshot must pin WAL frames before cancellation"
4335        );
4336
4337        let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
4338        let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
4339        let query = tokio::spawn(crate::scope_test_read_progress(
4340            Arc::clone(&progress),
4341            crate::scope_request_read_cancellation(cancel_rx, async move {
4342                let result = reader.query_all(deliberately_slow_read_statement()).await;
4343                (reader, result)
4344            }),
4345        ));
4346        wait_for_progress(progress.as_ref()).await;
4347        cancel_tx.send(true).unwrap();
4348        let (mut reader, result) = tokio::time::timeout(std::time::Duration::from_secs(1), query)
4349            .await
4350            .expect("interrupted explicit read transaction did not stop promptly")
4351            .unwrap();
4352        assert!(
4353            matches!(result, Err(StorageError::Timeout { .. })),
4354            "request cancellation must surface as a typed timeout; got {result:?}"
4355        );
4356        assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
4357
4358        let (_, log_after, checkpointed_after) = passive_checkpoint(&writer);
4359        assert_eq!(
4360            log_after, checkpointed_after,
4361            "cancellation must release the explicit reader's WAL snapshot"
4362        );
4363        contender
4364            .query_all(SqlStatement {
4365                sql: "SELECT 1".into(),
4366                params: vec![],
4367                label: None,
4368            })
4369            .await
4370            .expect("the sole reader permit must be reusable after rollback");
4371
4372        let stopped_at = progress.load(std::sync::atomic::Ordering::SeqCst);
4373        reader
4374            .query_all(SqlStatement {
4375                sql: "WITH RECURSIVE rows(value) AS (\
4376                      SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 10000\
4377                      ) SELECT SUM(value) FROM rows"
4378                    .into(),
4379                params: vec![],
4380                label: None,
4381            })
4382            .await
4383            .expect("same connection must remain usable after handler teardown");
4384        assert_eq!(
4385            progress.load(std::sync::atomic::Ordering::SeqCst),
4386            stopped_at,
4387            "the cancelled request's progress callback bled into the next borrower"
4388        );
4389    }
4390
4391    /// A statement that blocks inside a single SQLite VM step for longer
4392    /// than [`crate::read_cancellation::DEFAULT_SQLITE_INTERRUPT_GRACE_MS`].
4393    /// Unlike a recursive CTE (interrupt-checked every 1,000 VM
4394    /// instructions, so it stops promptly), a UDF call is one opcode: SQLite
4395    /// cannot observe the interrupt flag until the call returns. This is the
4396    /// only way to deterministically force a worker past the grace window
4397    /// rather than merely past `wait_for_progress`.
4398    fn register_khive_test_slow_udf(
4399        conn: &rusqlite::Connection,
4400        sleep_ms: u64,
4401        started: Arc<std::sync::atomic::AtomicBool>,
4402    ) {
4403        conn.create_scalar_function(
4404            "khive_test_slow_udf",
4405            0,
4406            rusqlite::functions::FunctionFlags::SQLITE_UTF8,
4407            move |_| {
4408                started.store(true, std::sync::atomic::Ordering::Release);
4409                std::thread::sleep(std::time::Duration::from_millis(sleep_ms));
4410                Ok(0i64)
4411            },
4412        )
4413        .unwrap();
4414    }
4415
4416    async fn wait_for_flag(flag: &std::sync::atomic::AtomicBool) {
4417        tokio::time::timeout(std::time::Duration::from_secs(1), async {
4418            while !flag.load(std::sync::atomic::Ordering::Acquire) {
4419                tokio::task::yield_now().await;
4420            }
4421        })
4422        .await
4423        .expect("slow UDF never started");
4424    }
4425
4426    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4427    async fn abandoned_read_past_grace_recovers_admission_after_bounded_join() {
4428        // Regression for the PR #1897 review blocker: a `spawn_blocking`
4429        // read worker that outlives `KHIVE_SQLITE_INTERRUPT_GRACE_MS` must
4430        // not be treated as reaped just because the async side detached
4431        // from it. This forces the worker past grace with a slow UDF (so
4432        // the interrupt genuinely cannot be observed mid-call) and asserts
4433        // that admission and the WAL snapshot are only reported recovered
4434        // once the real worker has actually joined.
4435        let dir = tempfile::tempdir().unwrap();
4436        let config = PoolConfig {
4437            path: Some(dir.path().join("sql_bridge_grace_exceeded.db")),
4438            max_readers: 1,
4439            checkout_timeout: std::time::Duration::from_millis(2_000),
4440            ..PoolConfig::default()
4441        };
4442        let pool = Arc::new(ConnectionPool::new(config).unwrap());
4443        let writer = open_standalone_writer(&pool).unwrap();
4444        writer
4445            .execute_batch(
4446                "CREATE TABLE grace_probe(id INTEGER PRIMARY KEY, value TEXT NOT NULL); \
4447                 INSERT INTO grace_probe(value) VALUES ('seed');",
4448            )
4449            .unwrap();
4450
4451        let mut reader = SqliteReader {
4452            handle: Some(open_cached_reader_handle(Arc::clone(&pool)).await.unwrap()),
4453            pool: Arc::clone(&pool),
4454        };
4455        let mut contender = SqliteReader {
4456            handle: Some(open_cached_reader_handle(Arc::clone(&pool)).await.unwrap()),
4457            pool: Arc::clone(&pool),
4458        };
4459        let udf_started = Arc::new(std::sync::atomic::AtomicBool::new(false));
4460        register_khive_test_slow_udf(
4461            &reader.handle.as_ref().unwrap().conn,
4462            900,
4463            Arc::clone(&udf_started),
4464        );
4465
4466        reader
4467            .query_all(SqlStatement {
4468                sql: "BEGIN DEFERRED".into(),
4469                params: vec![],
4470                label: None,
4471            })
4472            .await
4473            .expect("begin admitted transaction");
4474        reader
4475            .query_all(SqlStatement {
4476                sql: "SELECT * FROM grace_probe".into(),
4477                params: vec![],
4478                label: None,
4479            })
4480            .await
4481            .expect("materialize a real WAL snapshot");
4482        writer
4483            .execute_batch(
4484                "WITH RECURSIVE rows(value) AS (\
4485                 SELECT 1 UNION ALL SELECT value + 1 FROM rows WHERE value < 100\
4486                 ) INSERT INTO grace_probe(value) SELECT printf('row-%d', value) FROM rows;",
4487            )
4488            .unwrap();
4489        let (_, log_before, checkpointed_before) = passive_checkpoint(&writer);
4490        assert!(
4491            log_before > checkpointed_before,
4492            "the explicit reader snapshot must pin WAL frames before cancellation"
4493        );
4494
4495        let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
4496        let query = tokio::spawn(crate::scope_request_read_cancellation(
4497            cancel_rx,
4498            async move {
4499                let result = reader
4500                    .query_all(SqlStatement {
4501                        sql: "SELECT khive_test_slow_udf()".into(),
4502                        params: vec![],
4503                        label: None,
4504                    })
4505                    .await;
4506                (reader, result)
4507            },
4508        ));
4509        // Wait for the UDF to actually be running (not just scheduled) so
4510        // registration has completed and the worker is provably blocked
4511        // inside SQLite before cancelling — the same proof `wait_for_progress`
4512        // gives the recursive-CTE tests, since a progress callback (fired
4513        // between opcodes) never runs during the UDF's own blocking call.
4514        wait_for_flag(udf_started.as_ref()).await;
4515        cancel_tx.send(true).unwrap();
4516
4517        let (mut reader, result) = tokio::time::timeout(std::time::Duration::from_secs(3), query)
4518            .await
4519            .expect(
4520                "a worker that settles within the grace+hard-cap bound must not hang the caller",
4521            )
4522            .unwrap();
4523        assert!(
4524            matches!(result, Err(StorageError::Timeout { .. })),
4525            "request cancellation must still surface as a typed timeout even after grace \
4526             was exceeded; got {result:?}"
4527        );
4528
4529        // The response was only returned above because the real worker
4530        // joined — expected-fail arm: this assertion would fail (permit
4531        // still 0) under the pre-fix behavior, which detached and returned
4532        // Timeout at the grace boundary while the worker (and its permit)
4533        // were still live.
4534        assert_eq!(
4535            pool.sql_bridge_reader_slots().available_permits(),
4536            1,
4537            "the sole reader permit must be visible again once the bounded join completes"
4538        );
4539
4540        let (_, log_after, checkpointed_after) = passive_checkpoint(&writer);
4541        assert_eq!(
4542            log_after, checkpointed_after,
4543            "the abandoned explicit read transaction must release its WAL snapshot by the \
4544             time the caller observes the timeout"
4545        );
4546
4547        contender
4548            .query_all(SqlStatement {
4549                sql: "SELECT 1".into(),
4550                params: vec![],
4551                label: None,
4552            })
4553            .await
4554            .expect("a fresh reader must be admitted once the zombie worker has settled");
4555
4556        // The settled connection itself must also be reusable (not
4557        // quarantined) and must not still be carrying the slow UDF's
4558        // progress callback into a later borrower.
4559        reader
4560            .query_all(SqlStatement {
4561                sql: "SELECT 1".into(),
4562                params: vec![],
4563                label: None,
4564            })
4565            .await
4566            .expect("the interrupted connection must remain usable after settling");
4567    }
4568
4569    #[tokio::test]
4570    #[serial_test::serial(tx_registry)]
4571    async fn cached_reader_transaction_lifecycle_survives_sqlite_empty_prefixes() {
4572        let dir = tempfile::tempdir().unwrap();
4573        let config = PoolConfig {
4574            path: Some(dir.path().join("sql_bridge_reader_prefixed_tx_control.db")),
4575            write_queue_enabled: Some(true),
4576            max_readers: 1,
4577            ..PoolConfig::default()
4578        };
4579        let pool = Arc::new(ConnectionPool::new(config).unwrap());
4580        let bridge = SqlBridge::new(Arc::clone(&pool), true);
4581        let mut reader = bridge.reader().await.unwrap();
4582
4583        reader
4584            .query_all(SqlStatement {
4585                sql: " ; -- empty statement\n /* leading comment */ \u{feff} BEGIN DEFERRED".into(),
4586                params: vec![],
4587                label: None,
4588            })
4589            .await
4590            .expect("prefixed BEGIN must enter the admitted transaction state");
4591        assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 0);
4592        reader
4593            .query_all(SqlStatement {
4594                sql: " /* leading comment */ \u{feff} ; COMMIT".into(),
4595                params: vec![],
4596                label: None,
4597            })
4598            .await
4599            .expect("prefixed COMMIT must end the admitted transaction state");
4600        assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
4601
4602        let rejected = reader
4603            .query_all(SqlStatement {
4604                sql: " ; /* no active transaction */ \u{feff} COMMIT".into(),
4605                params: vec![],
4606                label: None,
4607            })
4608            .await;
4609        assert!(
4610            matches!(
4611                &rejected,
4612                Err(StorageError::InvalidInput {
4613                    operation,
4614                    message,
4615                    ..
4616                }) if operation.as_ref() == "query_all"
4617                    && message.contains("transaction control")
4618                    && message.contains("COMMIT")
4619            ),
4620            "a prefixed COMMIT without an admitted transaction must still fail closed; \
4621             got {rejected:?}"
4622        );
4623
4624        let mut queue_backed_writer = bridge.writer().await.unwrap();
4625        let rejected = queue_backed_writer
4626            .query_all(SqlStatement {
4627                sql: "-- leading comment\n \u{feff} ; /* empty */ SAVEPOINT pinned".into(),
4628                params: vec![],
4629                label: None,
4630            })
4631            .await;
4632        assert!(
4633            matches!(
4634                &rejected,
4635                Err(StorageError::InvalidInput {
4636                    operation,
4637                    message,
4638                    ..
4639                }) if operation.as_ref() == "writer.query_all"
4640                    && message.contains("transaction control")
4641                    && message.contains("SAVEPOINT")
4642            ),
4643            "a queue-backed cached reader must classify transaction control through \
4644             comments, BOMs, and empty statements; got {rejected:?}"
4645        );
4646
4647        let value = reader
4648            .query_scalar(SqlStatement {
4649                sql: "SELECT 10".into(),
4650                params: vec![],
4651                label: None,
4652            })
4653            .await
4654            .expect("prefixed transaction lifecycle must preserve the cached reader");
4655        assert!(matches!(value, Some(SqlValue::Integer(10))));
4656    }
4657
4658    #[tokio::test]
4659    async fn cached_reader_restores_autocommit_before_releasing_its_operation_permit() {
4660        let dir = tempfile::tempdir().unwrap();
4661        let config = PoolConfig {
4662            path: Some(dir.path().join("sql_bridge_reader_autocommit.db")),
4663            max_readers: 1,
4664            ..PoolConfig::default()
4665        };
4666        let pool = Arc::new(ConnectionPool::new(config).unwrap());
4667        let conn = open_standalone_reader(&pool).unwrap();
4668        conn.execute_batch("BEGIN DEFERRED; SELECT * FROM sqlite_schema")
4669            .unwrap();
4670        assert!(
4671            !conn.is_autocommit(),
4672            "the regression precondition needs a live read transaction"
4673        );
4674        let mut reader = SqliteReader {
4675            handle: Some(StandaloneHandle {
4676                conn,
4677                _retained_slot: None,
4678                read_transaction_slot: None,
4679            }),
4680            pool: Arc::clone(&pool),
4681        };
4682
4683        let rejected = reader
4684            .query_all(SqlStatement {
4685                // The stale-state cleanup must take precedence over the
4686                // ordinary transaction-control rejection. Otherwise an idle
4687                // snapshot could survive every rejected ROLLBACK attempt.
4688                sql: "ROLLBACK".into(),
4689                params: vec![],
4690                label: None,
4691            })
4692            .await;
4693        assert!(
4694            matches!(
4695                &rejected,
4696                Err(StorageError::InvalidInput {
4697                    operation,
4698                    message,
4699                    ..
4700                }) if operation.as_ref() == "query_all"
4701                    && message.contains("outside autocommit")
4702            ),
4703            "a cached reader that reaches the boundary outside autocommit must fail closed; \
4704             got {rejected:?}"
4705        );
4706        assert_eq!(
4707            pool.sql_bridge_reader_slots().available_permits(),
4708            1,
4709            "the permit may be released only after the stale transaction is gone"
4710        );
4711        assert!(
4712            reader
4713                .handle
4714                .as_ref()
4715                .expect("successful rollback should preserve the cached handle")
4716                .conn
4717                .is_autocommit(),
4718            "the restored idle connection must not retain a WAL snapshot"
4719        );
4720
4721        let value = reader
4722            .query_scalar(SqlStatement {
4723                sql: "SELECT 9".into(),
4724                params: vec![],
4725                label: None,
4726            })
4727            .await
4728            .expect("the cleaned cached reader must remain usable");
4729        assert!(matches!(value, Some(SqlValue::Integer(9))));
4730    }
4731
4732    #[tokio::test]
4733    async fn standalone_writer_read_preserves_manual_atomic_transaction() {
4734        let dir = tempfile::tempdir().unwrap();
4735        let config = PoolConfig {
4736            path: Some(dir.path().join("sql_bridge_writer_atomic_read.db")),
4737            write_queue_enabled: Some(false),
4738            max_readers: 1,
4739            ..PoolConfig::default()
4740        };
4741        let pool = Arc::new(ConnectionPool::new(config).unwrap());
4742        {
4743            let writer = pool.writer().unwrap();
4744            writer
4745                .conn()
4746                .execute_batch(
4747                    "CREATE TABLE atomic_read_test \
4748                     (id INTEGER PRIMARY KEY, value TEXT NOT NULL)",
4749                )
4750                .unwrap();
4751        }
4752        let bridge = SqlBridge::new(Arc::clone(&pool), true);
4753
4754        let observed = bridge
4755            .atomic_unit(Box::new(|writer| {
4756                Box::pin(async move {
4757                    writer
4758                        .execute(SqlStatement {
4759                            sql: "INSERT INTO atomic_read_test (id, value) VALUES (1, 'pending')"
4760                                .into(),
4761                            params: vec![],
4762                            label: None,
4763                        })
4764                        .await?;
4765                    let count = writer
4766                        .query_scalar(SqlStatement {
4767                            sql: "SELECT COUNT(*) FROM atomic_read_test".into(),
4768                            params: vec![],
4769                            label: None,
4770                        })
4771                        .await?;
4772                    Ok(Box::new(count) as Box<dyn std::any::Any + Send>)
4773                })
4774            }))
4775            .await
4776            .expect("manual atomic read must not be mistaken for an idle reader snapshot");
4777        let observed = match observed.downcast::<Option<SqlValue>>() {
4778            Ok(observed) => observed,
4779            Err(_) => panic!("unexpected atomic result type"),
4780        };
4781        assert!(matches!(*observed, Some(SqlValue::Integer(1))));
4782
4783        let mut reader = bridge.reader().await.unwrap();
4784        let committed = reader
4785            .query_scalar(SqlStatement {
4786                sql: "SELECT COUNT(*) FROM atomic_read_test".into(),
4787                params: vec![],
4788                label: None,
4789            })
4790            .await
4791            .unwrap();
4792        assert!(matches!(committed, Some(SqlValue::Integer(1))));
4793    }
4794
4795    #[tokio::test]
4796    async fn request_cancellation_preserves_file_backed_manual_atomic_read_and_commit() {
4797        let dir = tempfile::tempdir().unwrap();
4798        let pool = Arc::new(
4799            ConnectionPool::new(PoolConfig {
4800                path: Some(dir.path().join("sql_bridge_writer_tx_cancel.db")),
4801                write_queue_enabled: Some(false),
4802                ..PoolConfig::default()
4803            })
4804            .unwrap(),
4805        );
4806        pool.writer()
4807            .unwrap()
4808            .conn()
4809            .execute_batch(
4810                "CREATE TABLE writer_tx_cancel_probe(\
4811                 id INTEGER PRIMARY KEY, value TEXT NOT NULL)",
4812            )
4813            .unwrap();
4814        let bridge = SqlBridge::new(Arc::clone(&pool), true);
4815        let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
4816
4817        let observed = crate::scope_request_read_cancellation(
4818            cancel_rx,
4819            bridge.atomic_unit(Box::new(move |writer| {
4820                Box::pin(async move {
4821                    writer
4822                        .execute(SqlStatement {
4823                            sql: "INSERT INTO writer_tx_cancel_probe VALUES (1, 'before')".into(),
4824                            params: vec![],
4825                            label: None,
4826                        })
4827                        .await?;
4828                    cancel_tx.send(true).unwrap();
4829                    let count = writer
4830                        .query_scalar(SqlStatement {
4831                            sql: "SELECT COUNT(*) FROM writer_tx_cancel_probe".into(),
4832                            params: vec![],
4833                            label: None,
4834                        })
4835                        .await?;
4836                    writer
4837                        .execute(SqlStatement {
4838                            sql: "INSERT INTO writer_tx_cancel_probe VALUES (2, 'after')".into(),
4839                            params: vec![],
4840                            label: None,
4841                        })
4842                        .await?;
4843                    Ok(Box::new(count) as Box<dyn std::any::Any + Send>)
4844                })
4845            })),
4846        )
4847        .await
4848        .expect("request cancellation must not interrupt an admitted manual write transaction");
4849        let observed = match observed.downcast::<Option<SqlValue>>() {
4850            Ok(observed) => observed,
4851            Err(_) => panic!("unexpected atomic result type"),
4852        };
4853        assert!(matches!(*observed, Some(SqlValue::Integer(1))));
4854
4855        let reader = pool.reader().unwrap();
4856        let rows: i64 = reader
4857            .conn()
4858            .query_row("SELECT COUNT(*) FROM writer_tx_cancel_probe", [], |row| {
4859                row.get(0)
4860            })
4861            .unwrap();
4862        assert_eq!(rows, 2, "both writes around the SELECT must commit");
4863    }
4864
4865    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4866    async fn cancelled_standalone_writer_transaction_retains_active_reader_admission() {
4867        let dir = tempfile::tempdir().unwrap();
4868        let pool = Arc::new(
4869            ConnectionPool::new(PoolConfig {
4870                path: Some(dir.path().join("sql_bridge_writer_tx_admission.db")),
4871                write_queue_enabled: Some(false),
4872                max_readers: 1,
4873                checkout_timeout: std::time::Duration::from_millis(250),
4874                ..PoolConfig::default()
4875            })
4876            .unwrap(),
4877        );
4878        let writer_slot = pool
4879            .sql_bridge_writer_slots()
4880            .acquire_owned()
4881            .await
4882            .unwrap();
4883        let conn = open_standalone_writer(&pool).unwrap();
4884        let (entered, release, _completed) = blocking_non_interrupting_progress_gate(&conn);
4885        let mut writer = SqliteWriter {
4886            handle: Some(StandaloneHandle {
4887                conn,
4888                _retained_slot: Some(writer_slot),
4889                read_transaction_slot: None,
4890            }),
4891            writer_task: None,
4892            origin: pool.origin(),
4893            db: crate::timeout_sink::db_label(&pool),
4894            pool: Arc::clone(&pool),
4895        };
4896        khive_storage::SqlWriter::execute(
4897            &mut writer,
4898            SqlStatement {
4899                sql: "BEGIN IMMEDIATE".into(),
4900                params: vec![],
4901                label: None,
4902            },
4903        )
4904        .await
4905        .unwrap();
4906        let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
4907        cancel_tx.send(true).unwrap();
4908
4909        let query = tokio::spawn(crate::scope_request_read_cancellation(
4910            cancel_rx,
4911            async move {
4912                let result =
4913                    khive_storage::SqlReader::query_all(&mut writer, progress_gate_statement())
4914                        .await;
4915                let rollback = khive_storage::SqlWriter::execute(
4916                    &mut writer,
4917                    SqlStatement {
4918                        sql: "ROLLBACK".into(),
4919                        params: vec![],
4920                        label: None,
4921                    },
4922                )
4923                .await;
4924                (result, rollback)
4925            },
4926        ));
4927        tokio::time::timeout(std::time::Duration::from_secs(1), entered.notified())
4928            .await
4929            .expect("cancelled writer-transaction SELECT never reached SQLite");
4930        assert_eq!(
4931            pool.sql_bridge_reader_slots().available_permits(),
4932            0,
4933            "a writer-supertrait SELECT must retain ordinary active-reader admission"
4934        );
4935
4936        tokio::task::spawn_blocking(move || release.wait())
4937            .await
4938            .unwrap();
4939        let (rows, rollback) = tokio::time::timeout(std::time::Duration::from_secs(2), query)
4940            .await
4941            .expect("writer-transaction SELECT did not finish after its gate opened")
4942            .unwrap();
4943        assert_eq!(
4944            rows.expect("request cancellation interrupted the admitted writer transaction")
4945                .len(),
4946            1
4947        );
4948        rollback.expect("writer transaction did not return to autocommit");
4949        assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
4950    }
4951
4952    #[tokio::test]
4953    async fn expired_deadline_preserves_pool_backed_manual_atomic_read_and_commit() {
4954        let pool = Arc::new(ConnectionPool::new(PoolConfig::default()).unwrap());
4955        pool.writer()
4956            .unwrap()
4957            .conn()
4958            .execute_batch(
4959                "CREATE TABLE pool_writer_tx_deadline_probe(\
4960                 id INTEGER PRIMARY KEY, value TEXT NOT NULL)",
4961            )
4962            .unwrap();
4963        let bridge = SqlBridge::new(Arc::clone(&pool), false);
4964
4965        let observed = crate::scope_request_read_deadline(
4966            std::time::Duration::ZERO,
4967            bridge.atomic_unit(Box::new(|writer| {
4968                Box::pin(async move {
4969                    writer
4970                        .execute(SqlStatement {
4971                            sql: "INSERT INTO pool_writer_tx_deadline_probe VALUES (1, 'before')"
4972                                .into(),
4973                            params: vec![],
4974                            label: None,
4975                        })
4976                        .await?;
4977                    let count = writer
4978                        .query_scalar(SqlStatement {
4979                            sql: "SELECT COUNT(*) FROM pool_writer_tx_deadline_probe".into(),
4980                            params: vec![],
4981                            label: None,
4982                        })
4983                        .await?;
4984                    writer
4985                        .execute(SqlStatement {
4986                            sql: "INSERT INTO pool_writer_tx_deadline_probe VALUES (2, 'after')"
4987                                .into(),
4988                            params: vec![],
4989                            label: None,
4990                        })
4991                        .await?;
4992                    Ok(Box::new(count) as Box<dyn std::any::Any + Send>)
4993                })
4994            })),
4995        )
4996        .await
4997        .expect("an expired read deadline must not interrupt an admitted manual write transaction");
4998        let observed = match observed.downcast::<Option<SqlValue>>() {
4999            Ok(observed) => observed,
5000            Err(_) => panic!("unexpected atomic result type"),
5001        };
5002        assert!(matches!(*observed, Some(SqlValue::Integer(1))));
5003
5004        let reader = pool.reader().unwrap();
5005        let rows: i64 = reader
5006            .conn()
5007            .query_row(
5008                "SELECT COUNT(*) FROM pool_writer_tx_deadline_probe",
5009                [],
5010                |row| row.get(0),
5011            )
5012            .unwrap();
5013        assert_eq!(rows, 2, "both writes around the SELECT must commit");
5014    }
5015
5016    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5017    async fn cancelled_standalone_open_retains_slot_until_open_finishes() {
5018        let dir = tempfile::tempdir().unwrap();
5019        let config = PoolConfig {
5020            path: Some(dir.path().join("sql_bridge_cancelled_open.db")),
5021            max_readers: 1,
5022            checkout_timeout: std::time::Duration::from_millis(250),
5023            ..PoolConfig::default()
5024        };
5025        let pool = Arc::new(ConnectionPool::new(config).unwrap());
5026        let slots = pool.sql_bridge_reader_slots();
5027        let slot = Arc::clone(&slots).acquire_owned().await.unwrap();
5028        assert_eq!(slots.available_permits(), 0);
5029
5030        let (entered_tx, entered_rx) = std::sync::mpsc::channel();
5031        let (release_tx, release_rx) = std::sync::mpsc::channel();
5032        let open = tokio::spawn(open_standalone_on_blocking(
5033            Arc::clone(&pool),
5034            slot,
5035            "test_open_reader",
5036            move |pool| {
5037                entered_tx.send(()).unwrap();
5038                release_rx.recv().unwrap();
5039                open_standalone_reader(pool)
5040            },
5041        ));
5042        tokio::task::spawn_blocking(move || entered_rx.recv())
5043            .await
5044            .unwrap()
5045            .unwrap();
5046
5047        open.abort();
5048        assert!(matches!(open.await, Err(error) if error.is_cancelled()));
5049        assert_eq!(
5050            slots.available_permits(),
5051            0,
5052            "the permit must remain in the detached open closure"
5053        );
5054        let contender = tokio::time::timeout(
5055            std::time::Duration::from_millis(50),
5056            Arc::clone(&slots).acquire_owned(),
5057        )
5058        .await;
5059        assert!(contender.is_err(), "an in-flight open must retain the cap");
5060
5061        release_tx.send(()).unwrap();
5062        let recovered = tokio::time::timeout(
5063            std::time::Duration::from_secs(1),
5064            Arc::clone(&slots).acquire_owned(),
5065        )
5066        .await
5067        .expect("the detached open did not release its permit")
5068        .unwrap();
5069        assert_eq!(slots.available_permits(), 0);
5070        drop(recovered);
5071        assert_eq!(slots.available_permits(), 1);
5072    }
5073
5074    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5075    async fn abandoned_writer_read_interrupts_and_releases_writer_handle() {
5076        let dir = tempfile::tempdir().unwrap();
5077        let config = PoolConfig {
5078            path: Some(dir.path().join("sql_bridge_cancelled_writer.db")),
5079            write_queue_enabled: Some(false),
5080            checkout_timeout: std::time::Duration::from_millis(250),
5081            ..PoolConfig::default()
5082        };
5083        let pool = Arc::new(ConnectionPool::new(config).unwrap());
5084        let bridge = SqlBridge::new(Arc::clone(&pool), true);
5085
5086        let handle_slot = pool
5087            .sql_bridge_writer_slots()
5088            .acquire_owned()
5089            .await
5090            .unwrap();
5091        let conn = open_standalone_writer(&pool).unwrap();
5092        let mut writer = SqliteWriter {
5093            handle: Some(StandaloneHandle {
5094                conn,
5095                _retained_slot: Some(handle_slot),
5096                read_transaction_slot: None,
5097            }),
5098            writer_task: None,
5099            origin: pool.origin(),
5100            db: crate::timeout_sink::db_label(&pool),
5101            pool: Arc::clone(&pool),
5102        };
5103        let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
5104        let query = tokio::spawn(crate::scope_test_read_progress(
5105            Arc::clone(&progress),
5106            async move {
5107                khive_storage::SqlReader::query_all(&mut writer, deliberately_slow_read_statement())
5108                    .await
5109            },
5110        ));
5111
5112        wait_for_progress(progress.as_ref()).await;
5113        query.abort();
5114        assert!(matches!(query.await, Err(error) if error.is_cancelled()));
5115        let writer_after =
5116            tokio::time::timeout(std::time::Duration::from_millis(500), bridge.writer())
5117                .await
5118                .expect("abandoned SQLite read did not release the writer handle promptly")
5119                .expect("writer handle remained unavailable after read interruption");
5120        drop(writer_after);
5121        let stopped_at = progress.load(std::sync::atomic::Ordering::SeqCst);
5122        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
5123        assert_eq!(
5124            progress.load(std::sync::atomic::Ordering::SeqCst),
5125            stopped_at,
5126            "writer-backed SQLite read kept consuming work after cancellation"
5127        );
5128    }
5129
5130    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5131    async fn request_cancellation_never_interrupts_admitted_execute_batch() {
5132        let dir = tempfile::tempdir().unwrap();
5133        let config = PoolConfig {
5134            path: Some(dir.path().join("sql_bridge_cancelled_writer_batch.db")),
5135            write_queue_enabled: Some(false),
5136            checkout_timeout: std::time::Duration::from_millis(250),
5137            ..PoolConfig::default()
5138        };
5139        let pool = Arc::new(ConnectionPool::new(config).unwrap());
5140        let bridge = SqlBridge::new(Arc::clone(&pool), true);
5141        {
5142            let guard = pool.writer().unwrap();
5143            guard
5144                .conn()
5145                .execute_batch(
5146                    "CREATE TABLE cancellation_write_probe(\
5147                     id INTEGER PRIMARY KEY, value INTEGER NOT NULL)",
5148                )
5149                .unwrap();
5150        }
5151
5152        let handle_slot = pool
5153            .sql_bridge_writer_slots()
5154            .acquire_owned()
5155            .await
5156            .unwrap();
5157        let conn = open_standalone_writer(&pool).unwrap();
5158        let (entered, release, completed) = blocking_non_interrupting_progress_gate(&conn);
5159        let mut writer = SqliteWriter {
5160            handle: Some(StandaloneHandle {
5161                conn,
5162                _retained_slot: Some(handle_slot),
5163                read_transaction_slot: None,
5164            }),
5165            writer_task: None,
5166            origin: pool.origin(),
5167            db: crate::timeout_sink::db_label(&pool),
5168            pool: Arc::clone(&pool),
5169        };
5170        let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
5171        let query = tokio::spawn(crate::scope_request_read_cancellation(
5172            cancel_rx,
5173            async move {
5174                khive_storage::SqlWriter::execute_batch(&mut writer, vec![slow_insert_statement()])
5175                    .await
5176            },
5177        ));
5178
5179        tokio::time::timeout(std::time::Duration::from_secs(1), entered.notified())
5180            .await
5181            .expect("mutating execute_batch never reached SQLite VM work");
5182        cancel_tx.send(true).unwrap();
5183        tokio::time::sleep(std::time::Duration::from_millis(25)).await;
5184        assert!(
5185            !query.is_finished(),
5186            "request-read cancellation must not interrupt an admitted batch"
5187        );
5188
5189        let contender = bridge.writer().await;
5190        let retained_slot = matches!(
5191            &contender,
5192            Err(StorageError::AdmissionTimeout { operation, .. })
5193                if operation.as_ref() == "sql_bridge.writer_handle"
5194        );
5195        drop(contender);
5196
5197        tokio::task::spawn_blocking(move || release.wait())
5198            .await
5199            .unwrap();
5200        let affected = tokio::time::timeout(std::time::Duration::from_secs(2), query)
5201            .await
5202            .expect("admitted batch did not finish after its gate was released")
5203            .unwrap()
5204            .expect("request cancellation must preserve the batch result");
5205        assert_eq!(affected, 10_000);
5206        tokio::time::timeout(std::time::Duration::from_secs(1), completed.notified())
5207            .await
5208            .expect("completed batch did not release its connection");
5209        assert!(
5210            retained_slot,
5211            "request cancellation released the writer slot before the admitted batch stopped"
5212        );
5213        let reader = pool.reader().unwrap();
5214        let count: i64 = reader
5215            .conn()
5216            .query_row("SELECT COUNT(*) FROM cancellation_write_probe", [], |row| {
5217                row.get(0)
5218            })
5219            .unwrap();
5220        assert_eq!(count, 10_000, "the admitted batch must commit every row");
5221    }
5222
5223    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5224    async fn request_cancellation_never_interrupts_dml_returning_via_sql_reader() {
5225        let dir = tempfile::tempdir().unwrap();
5226        let config = PoolConfig {
5227            path: Some(dir.path().join("sql_bridge_dml_returning_cancel.db")),
5228            write_queue_enabled: Some(false),
5229            checkout_timeout: std::time::Duration::from_millis(250),
5230            ..PoolConfig::default()
5231        };
5232        let pool = Arc::new(ConnectionPool::new(config).unwrap());
5233        {
5234            let guard = pool.writer().unwrap();
5235            guard
5236                .conn()
5237                .execute_batch(
5238                    "CREATE TABLE returning_write_probe(\
5239                     id INTEGER PRIMARY KEY, value INTEGER NOT NULL)",
5240                )
5241                .unwrap();
5242        }
5243
5244        let handle_slot = pool
5245            .sql_bridge_writer_slots()
5246            .acquire_owned()
5247            .await
5248            .unwrap();
5249        let conn = open_standalone_writer(&pool).unwrap();
5250        let (entered, release, completed) = blocking_non_interrupting_progress_gate(&conn);
5251        let mut writer = SqliteWriter {
5252            handle: Some(StandaloneHandle {
5253                conn,
5254                _retained_slot: Some(handle_slot),
5255                read_transaction_slot: None,
5256            }),
5257            writer_task: None,
5258            origin: pool.origin(),
5259            db: crate::timeout_sink::db_label(&pool),
5260            pool: Arc::clone(&pool),
5261        };
5262        let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
5263        let query = tokio::spawn(crate::scope_request_read_cancellation(
5264            cancel_rx,
5265            async move {
5266                khive_storage::SqlReader::query_all(
5267                    &mut writer,
5268                    SqlStatement {
5269                        sql: "INSERT INTO returning_write_probe(value) \
5270                          WITH RECURSIVE rows(value) AS (\
5271                          SELECT 1 UNION ALL SELECT value + 1 FROM rows WHERE value < 10000\
5272                          ) SELECT value FROM rows RETURNING id"
5273                            .into(),
5274                        params: vec![],
5275                        label: Some("non-interruptible-returning-probe".into()),
5276                    },
5277                )
5278                .await
5279            },
5280        ));
5281
5282        tokio::time::timeout(std::time::Duration::from_secs(1), entered.notified())
5283            .await
5284            .expect("DML RETURNING never reached admitted SQLite work");
5285        cancel_tx.send(true).unwrap();
5286        tokio::time::sleep(std::time::Duration::from_millis(25)).await;
5287        assert!(
5288            !query.is_finished(),
5289            "request-read cancellation interrupted DML RETURNING"
5290        );
5291
5292        tokio::task::spawn_blocking(move || release.wait())
5293            .await
5294            .unwrap();
5295        let rows = tokio::time::timeout(std::time::Duration::from_secs(2), query)
5296            .await
5297            .expect("DML RETURNING did not finish after its gate was released")
5298            .unwrap()
5299            .expect("request cancellation must preserve DML RETURNING's result");
5300        assert_eq!(rows.len(), 10_000);
5301        tokio::time::timeout(std::time::Duration::from_secs(1), completed.notified())
5302            .await
5303            .expect("completed DML RETURNING did not release its connection");
5304
5305        let reader = pool.reader().unwrap();
5306        let count: i64 = reader
5307            .conn()
5308            .query_row("SELECT COUNT(*) FROM returning_write_probe", [], |row| {
5309                row.get(0)
5310            })
5311            .unwrap();
5312        assert_eq!(count, 10_000, "DML RETURNING must commit every row");
5313    }
5314
5315    /// Cancelling an in-flight call permanently invalidates the boxed handle:
5316    /// the call took the handle's connection into the detached blocking task,
5317    /// so every subsequent call on the SAME handle fails loudly with
5318    /// "connection already consumed" instead of silently operating on a
5319    /// connection that may still be running the cancelled statement.
5320    /// Callers that cancel or time out a bridge call must drop the handle and
5321    /// acquire a fresh one.
5322    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5323    async fn cancelled_call_invalidates_handle_reuse_fails_loud() {
5324        let dir = tempfile::tempdir().unwrap();
5325        let config = PoolConfig {
5326            path: Some(dir.path().join("sql_bridge_cancelled_reuse.db")),
5327            checkout_timeout: std::time::Duration::from_millis(250),
5328            ..PoolConfig::default()
5329        };
5330        let pool = Arc::new(ConnectionPool::new(config).unwrap());
5331
5332        let handle_slot = acquire_handle_slot(
5333            pool.sql_bridge_writer_slots(),
5334            pool.config().checkout_timeout,
5335            "sql_bridge.writer_handle",
5336            SlotTimeoutClass::Admission,
5337        )
5338        .await
5339        .unwrap();
5340        let conn = open_standalone_writer(&pool).unwrap();
5341        let (entered, release, completed) = blocking_non_interrupting_progress_gate(&conn);
5342        let writer = Arc::new(tokio::sync::Mutex::new(SqliteWriter {
5343            handle: Some(StandaloneHandle {
5344                conn,
5345                _retained_slot: Some(handle_slot),
5346                read_transaction_slot: None,
5347            }),
5348            writer_task: None,
5349            origin: pool.origin(),
5350            db: crate::timeout_sink::db_label(&pool),
5351            pool: Arc::clone(&pool),
5352        }));
5353        let writer_clone = Arc::clone(&writer);
5354        let query = tokio::spawn(async move {
5355            khive_storage::SqlWriter::execute_batch(
5356                &mut *writer_clone.lock().await,
5357                vec![progress_gate_statement()],
5358            )
5359            .await
5360        });
5361
5362        entered.notified().await;
5363        query.abort();
5364        let cancelled = matches!(query.await, Err(error) if error.is_cancelled());
5365
5366        let reuse = khive_storage::SqlWriter::execute(
5367            &mut *writer.lock().await,
5368            SqlStatement {
5369                sql: "CREATE TABLE cancelled_reuse_probe (id INTEGER PRIMARY KEY)".into(),
5370                params: vec![],
5371                label: None,
5372            },
5373        )
5374        .await;
5375        let message = match reuse {
5376            Err(StorageError::Pool { message, .. }) => message,
5377            other => panic!(
5378                "reusing a cancelled writer handle must fail loudly with \
5379                 'connection already consumed'; got {other:?}"
5380            ),
5381        };
5382        assert!(
5383            message.contains("connection already consumed"),
5384            "expected the cancelled handle's reuse error to name the pinned \
5385             failure; got {message:?}"
5386        );
5387
5388        tokio::task::spawn_blocking(move || release.wait())
5389            .await
5390            .unwrap();
5391        tokio::time::timeout(std::time::Duration::from_secs(1), completed.notified())
5392            .await
5393            .expect("cancelled writer's detached SQLite call did not finish");
5394        assert!(cancelled, "writer batch task did not report cancellation");
5395    }
5396
5397    /// Transaction-control statements (`BEGIN`/`START`/`COMMIT`/`END`/
5398    /// `ROLLBACK`/`SAVEPOINT`/`RELEASE`) in `execute_batch` input are rejected with a
5399    /// typed invalid-input error BEFORE anything executes, on the standalone
5400    /// path: a caller `COMMIT` inside the batch's own `BEGIN IMMEDIATE`
5401    /// would commit early and break the all-or-nothing contract. The
5402    /// rejection must leave the handle untouched and fully reusable, and no
5403    /// statement (not even the valid ones before the offending one) may
5404    /// have run.
5405    #[tokio::test]
5406    async fn execute_batch_rejects_transaction_control_before_executing_anything() {
5407        let dir = tempfile::tempdir().unwrap();
5408        let config = PoolConfig {
5409            path: Some(dir.path().join("sql_bridge_tx_control_reject.db")),
5410            checkout_timeout: std::time::Duration::from_millis(250),
5411            write_queue_enabled: Some(false),
5412            ..PoolConfig::default()
5413        };
5414        let pool = Arc::new(ConnectionPool::new(config).unwrap());
5415        {
5416            let guard = pool.writer().unwrap();
5417            guard
5418                .conn()
5419                .execute_batch(
5420                    "CREATE TABLE tx_reject_test (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
5421                )
5422                .unwrap();
5423        }
5424
5425        let handle_slot = acquire_handle_slot(
5426            pool.sql_bridge_writer_slots(),
5427            pool.config().checkout_timeout,
5428            "sql_bridge.writer_handle",
5429            SlotTimeoutClass::Admission,
5430        )
5431        .await
5432        .unwrap();
5433        let conn = open_standalone_writer(&pool).unwrap();
5434        let mut writer = SqliteWriter {
5435            handle: Some(StandaloneHandle {
5436                conn,
5437                _retained_slot: Some(handle_slot),
5438                read_transaction_slot: None,
5439            }),
5440            writer_task: None,
5441            origin: pool.origin(),
5442            db: crate::timeout_sink::db_label(&pool),
5443            pool: Arc::clone(&pool),
5444        };
5445
5446        for tail in ["COMMIT", "BEGIN"] {
5447            let multi = khive_storage::SqlWriter::execute_batch(
5448                &mut writer,
5449                vec![SqlStatement {
5450                    sql: format!(
5451                        "INSERT INTO tx_reject_test (id, val) VALUES (10, 'tail'); {tail}"
5452                    ),
5453                    params: vec![],
5454                    label: None,
5455                }],
5456            )
5457            .await;
5458            let message = multi
5459                .as_ref()
5460                .err()
5461                .map(ToString::to_string)
5462                .unwrap_or_default();
5463            assert!(
5464                message.contains("Multiple statements"),
5465                "a SqlStatement with trailing {tail} must be rejected before execution; got {message}"
5466            );
5467        }
5468
5469        // A valid INSERT first, then a bare COMMIT: the whole batch must be
5470        // rejected and the INSERT must NOT have run.
5471        let batch = khive_storage::SqlWriter::execute_batch(
5472            &mut writer,
5473            vec![
5474                SqlStatement {
5475                    sql: "INSERT INTO tx_reject_test (id, val) VALUES (1, 'a')".into(),
5476                    params: vec![],
5477                    label: None,
5478                },
5479                SqlStatement {
5480                    sql: "COMMIT".into(),
5481                    params: vec![],
5482                    label: None,
5483                },
5484            ],
5485        )
5486        .await;
5487        match &batch {
5488            Err(StorageError::InvalidInput {
5489                operation, message, ..
5490            }) => {
5491                assert_eq!(operation.as_ref(), "execute_batch");
5492                assert!(
5493                    message.contains("transaction control") && message.contains("COMMIT"),
5494                    "the rejection must name the offending statement head; got {message:?}"
5495                );
5496            }
5497            other => {
5498                panic!("a batch containing a bare COMMIT must be rejected up front; got {other:?}")
5499            }
5500        }
5501
5502        // Every transaction-control head is rejected, case-insensitively and
5503        // through leading whitespace and `--`/`/* */` comments.
5504        for sql in [
5505            "BEGIN IMMEDIATE",
5506            "START TRANSACTION",
5507            "commit",
5508            "End transaction",
5509            "ROLLBACK",
5510            "SAVEPOINT sp1",
5511            "RELEASE sp1",
5512            "  -- leading comment\nCOMMIT",
5513            "/* block */ rollback to savepoint sp1",
5514        ] {
5515            let rejected = khive_storage::SqlWriter::execute_batch(
5516                &mut writer,
5517                vec![SqlStatement {
5518                    sql: sql.into(),
5519                    params: vec![],
5520                    label: None,
5521                }],
5522            )
5523            .await;
5524            assert!(
5525                matches!(&rejected, Err(StorageError::InvalidInput { .. })),
5526                "transaction-control head {sql:?} must be rejected; got {rejected:?}"
5527            );
5528        }
5529
5530        // The rejection ran before the handle was taken: no statement
5531        // executed (the INSERT above did not land), and the handle is still
5532        // fully reusable.
5533        let count: i64 = {
5534            let guard = pool.reader().unwrap();
5535            guard
5536                .conn()
5537                .query_row("SELECT COUNT(*) FROM tx_reject_test", [], |r| r.get(0))
5538                .unwrap()
5539        };
5540        assert_eq!(count, 0, "a rejected batch must not have executed anything");
5541
5542        let affected = khive_storage::SqlWriter::execute(
5543            &mut writer,
5544            SqlStatement {
5545                sql: "INSERT INTO tx_reject_test (id, val) VALUES (2, 'b')".into(),
5546                params: vec![],
5547                label: None,
5548            },
5549        )
5550        .await
5551        .expect("the handle must survive a rejected batch untouched");
5552        assert_eq!(affected, 1);
5553    }
5554
5555    #[tokio::test]
5556    async fn standalone_execute_batch_rejects_prefixed_commit_before_any_write() {
5557        let dir = tempfile::tempdir().unwrap();
5558        let config = PoolConfig {
5559            path: Some(dir.path().join("sql_bridge_prefixed_commit_standalone.db")),
5560            write_queue_enabled: Some(false),
5561            ..PoolConfig::default()
5562        };
5563        let pool = Arc::new(ConnectionPool::new(config).unwrap());
5564        {
5565            let guard = pool.writer().unwrap();
5566            guard
5567                .conn()
5568                .execute_batch("CREATE TABLE prefixed_commit (id INTEGER PRIMARY KEY)")
5569                .unwrap();
5570        }
5571        let bridge = SqlBridge::new(Arc::clone(&pool), true);
5572        let mut writer = bridge.writer().await.unwrap();
5573
5574        let rejected = writer
5575            .execute_batch(vec![
5576                SqlStatement {
5577                    sql: "INSERT INTO prefixed_commit (id) VALUES (1)".into(),
5578                    params: vec![],
5579                    label: None,
5580                },
5581                SqlStatement {
5582                    sql: " ; -- empty statement\n /* leading comment */ \u{feff} ; COMMIT".into(),
5583                    params: vec![],
5584                    label: None,
5585                },
5586            ])
5587            .await;
5588        assert!(
5589            matches!(
5590                &rejected,
5591                Err(StorageError::InvalidInput {
5592                    operation,
5593                    message,
5594                    ..
5595                }) if operation.as_ref() == "execute_batch"
5596                    && message.contains("transaction control")
5597                    && message.contains("COMMIT")
5598            ),
5599            "standalone execute_batch must reject a prefixed COMMIT before the INSERT; \
5600             got {rejected:?}"
5601        );
5602
5603        let mut reader = bridge.reader().await.unwrap();
5604        let count = reader
5605            .query_scalar(SqlStatement {
5606                sql: "SELECT COUNT(*) FROM prefixed_commit".into(),
5607                params: vec![],
5608                label: None,
5609            })
5610            .await
5611            .unwrap();
5612        assert!(
5613            matches!(count, Some(SqlValue::Integer(0))),
5614            "prefixed COMMIT rejection must happen before the earlier INSERT; got {count:?}"
5615        );
5616
5617        let affected = writer
5618            .execute(SqlStatement {
5619                sql: "INSERT INTO prefixed_commit (id) VALUES (2)".into(),
5620                params: vec![],
5621                label: None,
5622            })
5623            .await
5624            .expect("prefixed COMMIT rejection must leave the standalone handle reusable");
5625        assert_eq!(affected, 1);
5626    }
5627
5628    #[tokio::test]
5629    async fn execute_batch_rejects_multi_statement_on_pool_backed_path() {
5630        let pool = Arc::new(ConnectionPool::new(PoolConfig::default()).unwrap());
5631        pool.writer()
5632            .unwrap()
5633            .conn()
5634            .execute_batch(
5635                "CREATE TABLE multi_statement_pool_test (id INTEGER PRIMARY KEY, val TEXT)",
5636            )
5637            .unwrap();
5638        let bridge = SqlBridge::new(Arc::clone(&pool), false);
5639        let mut writer = bridge.writer().await.unwrap();
5640
5641        let result = khive_storage::SqlWriter::execute_batch(
5642            &mut *writer,
5643            vec![SqlStatement {
5644                sql: "INSERT INTO multi_statement_pool_test (id, val) VALUES (1, 'x'); COMMIT"
5645                    .into(),
5646                params: vec![],
5647                label: None,
5648            }],
5649        )
5650        .await;
5651        let message = result
5652            .as_ref()
5653            .err()
5654            .map(ToString::to_string)
5655            .unwrap_or_default();
5656        assert!(
5657            message.contains("Multiple statements"),
5658            "pool-backed execute_batch must reject a trailing COMMIT; got {message}"
5659        );
5660        let count: i64 = pool
5661            .reader()
5662            .unwrap()
5663            .conn()
5664            .query_row(
5665                "SELECT COUNT(*) FROM multi_statement_pool_test",
5666                [],
5667                |row| row.get(0),
5668            )
5669            .unwrap();
5670        assert_eq!(count, 0);
5671    }
5672
5673    #[tokio::test]
5674    async fn inline_execute_batch_rejects_multi_statement_sql() {
5675        let dir = tempfile::tempdir().unwrap();
5676        let pool = Arc::new(
5677            ConnectionPool::new(PoolConfig {
5678                path: Some(dir.path().join("sql_bridge_multi_statement_inline.db")),
5679                write_queue_enabled: Some(true),
5680                write_routing_strict: true,
5681                ..PoolConfig::default()
5682            })
5683            .unwrap(),
5684        );
5685        pool.writer()
5686            .unwrap()
5687            .conn()
5688            .execute_batch(
5689                "CREATE TABLE multi_statement_inline_test (id INTEGER PRIMARY KEY, val TEXT)",
5690            )
5691            .unwrap();
5692        let bridge = SqlBridge::new(Arc::clone(&pool), true);
5693
5694        let result = bridge
5695            .atomic_unit(Box::new(|writer| {
5696                Box::pin(async move {
5697                    writer
5698                        .execute_batch(vec![SqlStatement {
5699                            sql: "INSERT INTO multi_statement_inline_test (id, val) VALUES (1, 'x'); BEGIN"
5700                                .into(),
5701                            params: vec![],
5702                            label: None,
5703                        }])
5704                        .await
5705                        .map(|_| Box::new(()) as Box<dyn Any + Send>)
5706                })
5707            }))
5708            .await;
5709        let message = result
5710            .as_ref()
5711            .err()
5712            .map(ToString::to_string)
5713            .unwrap_or_default();
5714        assert!(
5715            message.contains("Multiple statements"),
5716            "InlineWriter must reject a trailing BEGIN; got {message}"
5717        );
5718        let count: i64 = pool
5719            .reader()
5720            .unwrap()
5721            .conn()
5722            .query_row(
5723                "SELECT COUNT(*) FROM multi_statement_inline_test",
5724                [],
5725                |row| row.get(0),
5726            )
5727            .unwrap();
5728        assert_eq!(count, 0);
5729    }
5730
5731    /// Unit matrix for [`transaction_control_head`]: statement heads are
5732    /// classified case-insensitively through leading whitespace and
5733    /// comments; non-transaction-control heads (including identifiers that
5734    /// merely START with a keyword) never match.
5735    #[test]
5736    fn transaction_control_head_classification_matrix() {
5737        for (sql, expected) in [
5738            ("BEGIN", Some("BEGIN")),
5739            ("begin immediate", Some("BEGIN")),
5740            ("START TRANSACTION", Some("START")),
5741            ("start transaction", Some("START")),
5742            ("COMMIT", Some("COMMIT")),
5743            ("commit;", Some("COMMIT")),
5744            ("END", Some("END")),
5745            ("end transaction", Some("END")),
5746            ("ROLLBACK", Some("ROLLBACK")),
5747            ("rollback to savepoint sp1", Some("ROLLBACK")),
5748            ("SAVEPOINT sp1", Some("SAVEPOINT")),
5749            ("RELEASE sp1", Some("RELEASE")),
5750            ("release savepoint sp1", Some("RELEASE")),
5751            ("   \t COMMIT", Some("COMMIT")),
5752            ("\u{feff}BEGIN", Some("BEGIN")),
5753            (" ; BEGIN", Some("BEGIN")),
5754            (" ; ; -- empty\n /* comment */ COMMIT", Some("COMMIT")),
5755            ("  \u{feff} SAVEPOINT sp1", Some("SAVEPOINT")),
5756            ("/* comment */ \u{feff} ; RELEASE sp1", Some("RELEASE")),
5757            ("\u{feff} ; \u{feff} -- empty\n ROLLBACK", Some("ROLLBACK")),
5758            ("-- a comment\nCOMMIT", Some("COMMIT")),
5759            // SQLite does not nest block comments: the comment ends at the
5760            // first `*/`, leaving `*/ COMMIT`, which is not a statement head.
5761            ("/* /* nested? no */ */ COMMIT", None),
5762            ("-- one\n-- two\n  /* x */ begin", Some("BEGIN")),
5763            ("INSERT INTO t VALUES (1)", None),
5764            ("UPDATE t SET x = 1", None),
5765            ("DELETE FROM t", None),
5766            ("SELECT * FROM commit_log", None),
5767            ("CREATE TABLE rollback_audit (id INTEGER)", None),
5768            ("/* comment only */", None),
5769            (" ; /* empty statements only */ ; ", None),
5770            (" ; SELECT 1", None),
5771            ("", None),
5772        ] {
5773            assert_eq!(
5774                transaction_control_head(sql),
5775                expected,
5776                "classification mismatch for {sql:?}"
5777            );
5778        }
5779    }
5780
5781    #[test]
5782    fn cached_read_transaction_control_classification_matrix() {
5783        use CachedReadTransactionControl::{BeginDeferred, Finish, Unsupported};
5784
5785        for (sql, expected) in [
5786            ("BEGIN", Some(BeginDeferred)),
5787            ("begin transaction", Some(BeginDeferred)),
5788            (
5789                "/* p */ \u{feff} ; BEGIN /* mode */ DEFERRED",
5790                Some(BeginDeferred),
5791            ),
5792            ("BEGIN DEFERRED TRANSACTION", Some(BeginDeferred)),
5793            ("BEGIN IMMEDIATE", Some(Unsupported("BEGIN"))),
5794            ("BEGIN /* lock */ EXCLUSIVE", Some(Unsupported("BEGIN"))),
5795            ("BEGIN TRANSACTION IMMEDIATE", Some(Unsupported("BEGIN"))),
5796            ("begin transaction exclusive", Some(Unsupported("BEGIN"))),
5797            ("BEGIN TRANSACTION DEFERRED", Some(Unsupported("BEGIN"))),
5798            ("BEGIN IMMEDIATE TRANSACTION", Some(Unsupported("BEGIN"))),
5799            ("BEGIN TRANSACTION named_txn", Some(Unsupported("BEGIN"))),
5800            (
5801                "BEGIN DEFERRED TRANSACTION trailing",
5802                Some(Unsupported("BEGIN")),
5803            ),
5804            ("BEGIN DEFERRED DEFERRED", Some(Unsupported("BEGIN"))),
5805            // Non-identifier tails: the tokenizer yields no token for a
5806            // quoted, bracketed, or backticked tail, which must read as a
5807            // refused remainder, never as end-of-statement.
5808            (
5809                "BEGIN TRANSACTION \"IMMEDIATE\"",
5810                Some(Unsupported("BEGIN")),
5811            ),
5812            ("BEGIN TRANSACTION [IMMEDIATE]", Some(Unsupported("BEGIN"))),
5813            ("BEGIN TRANSACTION `IMMEDIATE`", Some(Unsupported("BEGIN"))),
5814            ("BEGIN TRANSACTION 'IMMEDIATE'", Some(Unsupported("BEGIN"))),
5815            ("BEGIN \"DEFERRED\"", Some(Unsupported("BEGIN"))),
5816            ("BEGIN; COMMIT", Some(Unsupported("BEGIN"))),
5817            // Trailing empty statements and trivia remain an accepted end.
5818            ("BEGIN;", Some(BeginDeferred)),
5819            ("BEGIN DEFERRED ; -- done", Some(BeginDeferred)),
5820            ("BEGIN TRANSACTION /* t */ ;;", Some(BeginDeferred)),
5821            ("START TRANSACTION", Some(Unsupported("START"))),
5822            ("COMMIT", Some(Finish("COMMIT"))),
5823            ("END TRANSACTION", Some(Finish("END"))),
5824            ("ROLLBACK", Some(Finish("ROLLBACK"))),
5825            ("ROLLBACK TRANSACTION", Some(Finish("ROLLBACK"))),
5826            ("ROLLBACK TO sp", Some(Unsupported("ROLLBACK"))),
5827            (
5828                "ROLLBACK /* nested */ TRANSACTION /* target */ TO sp",
5829                Some(Unsupported("ROLLBACK")),
5830            ),
5831            ("SAVEPOINT sp", Some(Unsupported("SAVEPOINT"))),
5832            ("SELECT 1", None),
5833        ] {
5834            assert_eq!(
5835                cached_read_transaction_control(sql),
5836                expected,
5837                "cached-reader transaction classification mismatch for {sql:?}"
5838            );
5839        }
5840    }
5841
5842    #[test]
5843    fn sqlite_accepts_utf8_bom_before_transaction_control() {
5844        let conn = rusqlite::Connection::open_in_memory().unwrap();
5845        conn.execute_batch("CREATE TABLE bom_transaction_test (id INTEGER)")
5846            .unwrap();
5847        conn.execute_batch("\u{feff}BEGIN IMMEDIATE").unwrap();
5848        conn.execute_batch("ROLLBACK").unwrap();
5849    }
5850
5851    /// The queue-backed `execute_batch` path rejects transaction-control
5852    /// statements too, and the rejection protects the writer task: a caller
5853    /// `COMMIT` that reached the task would close its per-request `BEGIN
5854    /// IMMEDIATE` and terminate the task permanently. After the typed
5855    /// rejection, a legitimate batch must still succeed through the SAME
5856    /// writer task (it was never touched).
5857    #[tokio::test]
5858    async fn execute_batch_rejects_transaction_control_on_queue_backed_path() {
5859        let dir = tempfile::tempdir().unwrap();
5860        let config = PoolConfig {
5861            path: Some(dir.path().join("sql_bridge_tx_reject_queue.db")),
5862            checkout_timeout: std::time::Duration::from_millis(250),
5863            write_queue_enabled: Some(true),
5864            write_routing_strict: true,
5865            ..PoolConfig::default()
5866        };
5867        let pool = Arc::new(ConnectionPool::new(config).unwrap());
5868        {
5869            let guard = pool.writer().unwrap();
5870            guard
5871                .conn()
5872                .execute_batch(
5873                    "CREATE TABLE tx_reject_queue_test (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
5874                )
5875                .unwrap();
5876        }
5877        let bridge = SqlBridge::new(Arc::clone(&pool), true);
5878        let mut writer = bridge.writer().await.unwrap();
5879
5880        let rejected = khive_storage::SqlWriter::execute_batch(
5881            &mut *writer,
5882            vec![
5883                SqlStatement {
5884                    sql: "INSERT INTO tx_reject_queue_test (id, val) VALUES (1, 'a')".into(),
5885                    params: vec![],
5886                    label: None,
5887                },
5888                SqlStatement {
5889                    sql: "COMMIT".into(),
5890                    params: vec![],
5891                    label: None,
5892                },
5893            ],
5894        )
5895        .await;
5896        assert!(
5897            matches!(&rejected, Err(StorageError::InvalidInput { .. })),
5898            "a bare COMMIT in a queue-backed batch must be rejected up front; got {rejected:?}"
5899        );
5900
5901        let prefixed = khive_storage::SqlWriter::execute_batch(
5902            &mut *writer,
5903            vec![
5904                SqlStatement {
5905                    sql: "INSERT INTO tx_reject_queue_test (id, val) VALUES (3, 'prefixed')".into(),
5906                    params: vec![],
5907                    label: None,
5908                },
5909                SqlStatement {
5910                    sql: "/* leading */ \u{feff} ; -- empty\n ; COMMIT".into(),
5911                    params: vec![],
5912                    label: None,
5913                },
5914            ],
5915        )
5916        .await;
5917        assert!(
5918            matches!(
5919                &prefixed,
5920                Err(StorageError::InvalidInput {
5921                    operation,
5922                    message,
5923                    ..
5924                }) if operation.as_ref() == "execute_batch"
5925                    && message.contains("transaction control")
5926                    && message.contains("COMMIT")
5927            ),
5928            "a prefixed COMMIT must be rejected before touching the writer task; got {prefixed:?}"
5929        );
5930
5931        let affected = khive_storage::SqlWriter::execute_batch(
5932            &mut *writer,
5933            vec![SqlStatement {
5934                sql: "INSERT INTO tx_reject_queue_test (id, val) VALUES (2, 'b')".into(),
5935                params: vec![],
5936                label: None,
5937            }],
5938        )
5939        .await
5940        .expect("the writer task must survive the rejected batch");
5941        assert_eq!(affected, 1);
5942
5943        let count: i64 = {
5944            let guard = pool.reader().unwrap();
5945            guard
5946                .conn()
5947                .query_row("SELECT COUNT(*) FROM tx_reject_queue_test", [], |r| {
5948                    r.get(0)
5949                })
5950                .unwrap()
5951        };
5952        assert_eq!(
5953            count, 1,
5954            "exactly the post-rejection batch's row may have landed"
5955        );
5956    }
5957
5958    /// A failed ROLLBACK after a statement failure poisons the handle: the
5959    /// connection may be in an unknown transaction state, so it is dropped
5960    /// instead of restored, and every subsequent call on the same handle
5961    /// fails loudly with "connection already consumed". The caller sees the
5962    /// ORIGINAL statement error with the poison context attached (the
5963    /// rollback failure is never hidden, but never replaces the original).
5964    ///
5965    /// Forcing the arm legitimately (the pre-round-2 version smuggled a
5966    /// bare `COMMIT` into the batch, which `execute_batch` now rejects up
5967    /// front): a connection authorizer denies the `ROLLBACK` transaction
5968    /// operation, so the error path's `ROLLBACK` genuinely fails while the
5969    /// batch's own `BEGIN IMMEDIATE` and the statements run normally.
5970    #[tokio::test]
5971    async fn failed_rollback_poisons_handle_reuse_fails_loud() {
5972        use rusqlite::hooks::{AuthAction, AuthContext, Authorization, TransactionOperation};
5973
5974        fn deny_rollback(ctx: AuthContext<'_>) -> Authorization {
5975            match ctx.action {
5976                AuthAction::Transaction {
5977                    operation: TransactionOperation::Rollback,
5978                } => Authorization::Deny,
5979                _ => Authorization::Allow,
5980            }
5981        }
5982
5983        let dir = tempfile::tempdir().unwrap();
5984        let config = PoolConfig {
5985            path: Some(dir.path().join("sql_bridge_rollback_poison.db")),
5986            checkout_timeout: std::time::Duration::from_millis(250),
5987            ..PoolConfig::default()
5988        };
5989        let pool = Arc::new(ConnectionPool::new(config).unwrap());
5990        {
5991            let guard = pool.writer().unwrap();
5992            guard
5993                .conn()
5994                .execute_batch(
5995                    "CREATE TABLE rollback_poison_test (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
5996                )
5997                .unwrap();
5998        }
5999
6000        let handle_slot = acquire_handle_slot(
6001            pool.sql_bridge_writer_slots(),
6002            pool.config().checkout_timeout,
6003            "sql_bridge.writer_handle",
6004            SlotTimeoutClass::Admission,
6005        )
6006        .await
6007        .unwrap();
6008        let conn = open_standalone_writer(&pool).unwrap();
6009        conn.authorizer(Some(deny_rollback)).unwrap();
6010        let mut writer = SqliteWriter {
6011            handle: Some(StandaloneHandle {
6012                conn,
6013                _retained_slot: Some(handle_slot),
6014                read_transaction_slot: None,
6015            }),
6016            writer_task: None,
6017            origin: pool.origin(),
6018            db: crate::timeout_sink::db_label(&pool),
6019            pool: Arc::clone(&pool),
6020        };
6021
6022        let batch = khive_storage::SqlWriter::execute_batch(
6023            &mut writer,
6024            vec![
6025                SqlStatement {
6026                    sql: "INSERT INTO rollback_poison_test (id, val) VALUES (1, 'a')".into(),
6027                    params: vec![],
6028                    label: None,
6029                },
6030                SqlStatement {
6031                    sql: "SELECT FROM WHERE".into(),
6032                    params: vec![],
6033                    label: None,
6034                },
6035            ],
6036        )
6037        .await;
6038        let batch_error = batch.expect_err("invalid second statement must fail the batch");
6039        let poison = match &batch_error {
6040            StorageError::Driver { source, .. } => source
6041                .downcast_ref::<PoisonedBatchError>()
6042                .expect("failed rollback must retain its typed poison wrapper"),
6043            other => panic!("failed rollback must return a driver error; got {other:?}"),
6044        };
6045        assert!(
6046            matches!(&poison.poison_reason, BatchPoisonReason::RollbackFailed(_)),
6047            "the poison cause must be compiler-checked as RollbackFailed; got {poison:?}"
6048        );
6049        let batch_message = batch_error.to_string();
6050        assert!(
6051            batch_message.contains("ROLLBACK after statement failure failed"),
6052            "the caller must see the poison context naming the failed \
6053             rollback; got {batch_message:?}"
6054        );
6055        assert!(
6056            batch_message.contains("original error"),
6057            "the original statement error must stay visible alongside the \
6058             poison context; got {batch_message:?}"
6059        );
6060
6061        let reuse = khive_storage::SqlWriter::execute(
6062            &mut writer,
6063            SqlStatement {
6064                sql: "CREATE TABLE rollback_poison_probe (id INTEGER PRIMARY KEY)".into(),
6065                params: vec![],
6066                label: None,
6067            },
6068        )
6069        .await;
6070        let message = match reuse {
6071            Err(StorageError::Pool { message, .. }) => message,
6072            other => panic!(
6073                "reusing a poisoned writer handle must fail loudly with \
6074                 'connection already consumed'; got {other:?}"
6075            ),
6076        };
6077        assert!(
6078            message.contains("connection already consumed"),
6079            "expected the poisoned handle's reuse error to name the pinned \
6080             failure; got {message:?}"
6081        );
6082    }
6083
6084    /// A NON-TRANSIENT `BEGIN IMMEDIATE` failure poisons the handle instead
6085    /// of restoring it: the connection's transaction state is suspect (here
6086    /// a caller-driven transaction is already open on the same connection,
6087    /// so SQLite answers "cannot start a transaction within a transaction"),
6088    /// and the returned error carries the poison context.
6089    #[tokio::test]
6090    async fn non_transient_begin_failure_poisons_handle() {
6091        let dir = tempfile::tempdir().unwrap();
6092        let config = PoolConfig {
6093            path: Some(dir.path().join("sql_bridge_begin_poison.db")),
6094            checkout_timeout: std::time::Duration::from_millis(250),
6095            ..PoolConfig::default()
6096        };
6097        let pool = Arc::new(ConnectionPool::new(config).unwrap());
6098
6099        let handle_slot = acquire_handle_slot(
6100            pool.sql_bridge_writer_slots(),
6101            pool.config().checkout_timeout,
6102            "sql_bridge.writer_handle",
6103            SlotTimeoutClass::Admission,
6104        )
6105        .await
6106        .unwrap();
6107        let conn = open_standalone_writer(&pool).unwrap();
6108        // A caller-driven open transaction on the same connection: the
6109        // batch's own `BEGIN IMMEDIATE` fails non-transiently.
6110        conn.execute_batch("BEGIN IMMEDIATE").unwrap();
6111        let mut writer = SqliteWriter {
6112            handle: Some(StandaloneHandle {
6113                conn,
6114                _retained_slot: Some(handle_slot),
6115                read_transaction_slot: None,
6116            }),
6117            writer_task: None,
6118            origin: pool.origin(),
6119            db: crate::timeout_sink::db_label(&pool),
6120            pool: Arc::clone(&pool),
6121        };
6122
6123        let batch = khive_storage::SqlWriter::execute_batch(
6124            &mut writer,
6125            vec![SqlStatement {
6126                sql: "SELECT 1".into(),
6127                params: vec![],
6128                label: None,
6129            }],
6130        )
6131        .await;
6132        let batch_error = batch.expect_err("BEGIN inside an open transaction must fail");
6133        let poison = match &batch_error {
6134            StorageError::Driver { source, .. } => source
6135                .downcast_ref::<PoisonedBatchError>()
6136                .expect("failed BEGIN must retain its typed poison wrapper"),
6137            other => panic!("failed BEGIN must return a driver error; got {other:?}"),
6138        };
6139        assert!(
6140            matches!(&poison.poison_reason, BatchPoisonReason::BeginFailed),
6141            "the poison cause must be compiler-checked as BeginFailed; got {poison:?}"
6142        );
6143        let batch_message = batch_error.to_string();
6144        assert!(
6145            batch_message.contains("BEGIN IMMEDIATE failed non-transiently"),
6146            "a non-transient BEGIN failure must surface the poison context; \
6147             got {batch_message:?}"
6148        );
6149        assert!(
6150            batch_message.contains("cannot start a transaction within a transaction"),
6151            "the original BEGIN error must stay visible; got {batch_message:?}"
6152        );
6153
6154        let reuse = khive_storage::SqlWriter::execute(
6155            &mut writer,
6156            SqlStatement {
6157                sql: "CREATE TABLE begin_poison_probe (id INTEGER PRIMARY KEY)".into(),
6158                params: vec![],
6159                label: None,
6160            },
6161        )
6162        .await;
6163        assert!(
6164            matches!(
6165                &reuse,
6166                Err(StorageError::Pool { message, .. })
6167                    if message.contains("connection already consumed")
6168            ),
6169            "a handle poisoned by a non-transient BEGIN failure must be \
6170             dropped, not restored; got {reuse:?}"
6171        );
6172    }
6173
6174    /// A BUSY/LOCKED `BEGIN IMMEDIATE` failure is transient contention: the
6175    /// connection itself is untouched, so the handle is restored as
6176    /// reusable, and the next call succeeds once the contending lock is
6177    /// released.
6178    #[tokio::test]
6179    async fn busy_begin_failure_restores_handle_reusable() {
6180        let dir = tempfile::tempdir().unwrap();
6181        let config = PoolConfig {
6182            path: Some(dir.path().join("sql_bridge_begin_busy.db")),
6183            checkout_timeout: std::time::Duration::from_millis(250),
6184            busy_timeout: std::time::Duration::from_millis(100),
6185            ..PoolConfig::default()
6186        };
6187        let pool = Arc::new(ConnectionPool::new(config).unwrap());
6188        {
6189            let guard = pool.writer().unwrap();
6190            guard
6191                .conn()
6192                .execute_batch("CREATE TABLE begin_busy_test (id INTEGER PRIMARY KEY)")
6193                .unwrap();
6194        }
6195
6196        // Hold SQLite's write lock from a separate connection so the batch's
6197        // `BEGIN IMMEDIATE` genuinely fails with SQLITE_BUSY after the short
6198        // busy timeout.
6199        let lock_conn = pool.open_standalone_writer().unwrap();
6200        lock_conn.execute_batch("BEGIN IMMEDIATE").unwrap();
6201
6202        let handle_slot = acquire_handle_slot(
6203            pool.sql_bridge_writer_slots(),
6204            pool.config().checkout_timeout,
6205            "sql_bridge.writer_handle",
6206            SlotTimeoutClass::Admission,
6207        )
6208        .await
6209        .unwrap();
6210        let conn = open_standalone_writer(&pool).unwrap();
6211        let mut writer = SqliteWriter {
6212            handle: Some(StandaloneHandle {
6213                conn,
6214                _retained_slot: Some(handle_slot),
6215                read_transaction_slot: None,
6216            }),
6217            writer_task: None,
6218            origin: pool.origin(),
6219            db: crate::timeout_sink::db_label(&pool),
6220            pool: Arc::clone(&pool),
6221        };
6222
6223        let batch = khive_storage::SqlWriter::execute_batch(
6224            &mut writer,
6225            vec![SqlStatement {
6226                sql: "INSERT INTO begin_busy_test (id) VALUES (1)".into(),
6227                params: vec![],
6228                label: None,
6229            }],
6230        )
6231        .await;
6232        let batch_error = batch.expect_err("BEGIN IMMEDIATE under a held write lock must fail");
6233        assert!(
6234            batch_error.to_string().contains("database is locked"),
6235            "the busy BEGIN failure must surface SQLite's busy error; got {batch_error:?}"
6236        );
6237
6238        lock_conn.execute_batch("ROLLBACK").unwrap();
6239        drop(lock_conn);
6240
6241        let affected = khive_storage::SqlWriter::execute(
6242            &mut writer,
6243            SqlStatement {
6244                sql: "INSERT INTO begin_busy_test (id) VALUES (2)".into(),
6245                params: vec![],
6246                label: None,
6247            },
6248        )
6249        .await
6250        .expect("a busy BEGIN failure must restore the handle as reusable");
6251        assert_eq!(affected, 1);
6252    }
6253
6254    /// The manual `atomic_unit` path (write queue off) shares the pool's
6255    /// one-permit writer-handle budget with `writer()`: while a boxed writer
6256    /// handle is live, `atomic_unit` times out; after the handle drops, the
6257    /// next `atomic_unit` succeeds on the same pool.
6258    #[tokio::test]
6259    async fn manual_atomic_unit_shares_writer_permit_budget_with_writer_handle() {
6260        let dir = tempfile::tempdir().unwrap();
6261        let config = PoolConfig {
6262            path: Some(dir.path().join("sql_bridge_atomic_unit_budget.db")),
6263            checkout_timeout: std::time::Duration::from_millis(50),
6264            write_queue_enabled: Some(false),
6265            ..PoolConfig::default()
6266        };
6267        let pool = Arc::new(ConnectionPool::new(config).unwrap());
6268        let bridge = SqlBridge::new(Arc::clone(&pool), true);
6269        {
6270            let guard = pool.writer().unwrap();
6271            guard
6272                .conn()
6273                .execute_batch(
6274                    "CREATE TABLE IF NOT EXISTS atomic_unit_budget_test \
6275                     (id INTEGER PRIMARY KEY, val INTEGER NOT NULL)",
6276                )
6277                .unwrap();
6278        }
6279
6280        fn insert_op(id: i64) -> AtomicUnitOp {
6281            Box::new(move |writer| {
6282                Box::pin(async move {
6283                    writer
6284                        .execute(SqlStatement {
6285                            sql: "INSERT INTO atomic_unit_budget_test (id, val) VALUES (?1, ?2)"
6286                                .into(),
6287                            params: vec![SqlValue::Integer(id), SqlValue::Integer(id)],
6288                            label: None,
6289                        })
6290                        .await
6291                        .map_err(|e| {
6292                            khive_storage::StorageError::driver(
6293                                StorageCapability::Sql,
6294                                "atomic_unit_budget_test_insert",
6295                                e,
6296                            )
6297                        })?;
6298                    Ok(Box::new(()) as Box<dyn std::any::Any + Send>)
6299                })
6300            })
6301        }
6302
6303        let writer_handle = bridge.writer().await.unwrap();
6304        let blocked = bridge.atomic_unit(insert_op(1)).await;
6305        assert!(
6306            matches!(
6307                &blocked,
6308                Err(StorageError::AdmissionTimeout { operation, .. })
6309                    if operation.as_ref() == "sql_bridge.atomic_unit_handle"
6310            ),
6311            "atomic_unit must time out on the shared writer permit while a \
6312             writer handle is live; got {blocked:?}"
6313        );
6314
6315        drop(writer_handle);
6316        let unblocked = bridge.atomic_unit(insert_op(2)).await;
6317        assert!(
6318            unblocked.is_ok(),
6319            "atomic_unit must succeed once the writer handle releases the \
6320             shared writer permit; got {unblocked:?}"
6321        );
6322
6323        let mut reader = bridge.reader().await.unwrap();
6324        let count = reader
6325            .query_scalar(SqlStatement {
6326                sql: "SELECT COUNT(*) FROM atomic_unit_budget_test".into(),
6327                params: vec![],
6328                label: None,
6329            })
6330            .await
6331            .unwrap();
6332        assert!(
6333            matches!(count, Some(SqlValue::Integer(1))),
6334            "only the post-drop atomic_unit call may have committed; got {count:?}"
6335        );
6336    }
6337
6338    /// ADR-067 Component A entry 10: with `KHIVE_WRITE_QUEUE=1`,
6339    /// `SqliteWriter::execute_batch` (reached via `SqlBridge::writer()`)
6340    /// routes the whole statement list through the WriterTask channel
6341    /// instead of opening its own `BEGIN IMMEDIATE` on the standalone
6342    /// connection, and the row is actually committed and readable back.
6343    #[tokio::test]
6344    async fn execute_batch_routes_through_writer_task_when_flag_enabled() {
6345        let dir = tempfile::tempdir().unwrap();
6346        let path = dir.path().join("write_queue_execute_batch.db");
6347        let config = PoolConfig {
6348            path: Some(path.clone()),
6349            write_queue_enabled: Some(true),
6350            ..PoolConfig::default()
6351        };
6352        let pool = Arc::new(ConnectionPool::new(config).unwrap());
6353        {
6354            let guard = pool.writer().unwrap();
6355            guard
6356                .conn()
6357                .execute_batch(
6358                    "CREATE TABLE IF NOT EXISTS write_queue_batch_test \
6359                     (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
6360                )
6361                .unwrap();
6362        }
6363
6364        let bridge = SqlBridge::new(Arc::clone(&pool), true);
6365
6366        let mut writer = bridge.writer().await.unwrap();
6367        let affected = writer
6368            .execute_batch(vec![
6369                SqlStatement {
6370                    sql: "INSERT INTO write_queue_batch_test (id, val) VALUES (?1, ?2)".into(),
6371                    params: vec![SqlValue::Integer(1), SqlValue::Text("a".into())],
6372                    label: None,
6373                },
6374                SqlStatement {
6375                    sql: "INSERT INTO write_queue_batch_test (id, val) VALUES (?1, ?2)".into(),
6376                    params: vec![SqlValue::Integer(2), SqlValue::Text("b".into())],
6377                    label: None,
6378                },
6379            ])
6380            .await
6381            .unwrap();
6382        assert_eq!(affected, 2);
6383
6384        let mut reader = bridge.reader().await.unwrap();
6385        let count = reader
6386            .query_scalar(SqlStatement {
6387                sql: "SELECT COUNT(*) FROM write_queue_batch_test".into(),
6388                params: vec![],
6389                label: None,
6390            })
6391            .await
6392            .unwrap();
6393        assert!(
6394            matches!(count, Some(SqlValue::Integer(2))),
6395            "expected 2 rows, got {count:?}"
6396        );
6397        assert_eq!(
6398            pool.writer_task_spawn_count(),
6399            1,
6400            "the flag-ON path must actually spawn and use the writer task"
6401        );
6402    }
6403
6404    /// ADR-067 Component A entry 10, atomicity: a batch whose second
6405    /// statement fails (duplicate primary key) must roll back the WHOLE
6406    /// request — including the first statement's otherwise-successful
6407    /// INSERT — because the WriterTask commits or rolls back one
6408    /// `WriteRequest` as a single unit (ADR-067 Component A). Zero rows must
6409    /// land, not one.
6410    #[tokio::test]
6411    async fn execute_batch_rolls_back_atomically_on_mid_sequence_failure() {
6412        let dir = tempfile::tempdir().unwrap();
6413        let path = dir.path().join("write_queue_execute_batch_rollback.db");
6414        let config = PoolConfig {
6415            path: Some(path.clone()),
6416            write_queue_enabled: Some(true),
6417            ..PoolConfig::default()
6418        };
6419        let pool = Arc::new(ConnectionPool::new(config).unwrap());
6420        {
6421            let guard = pool.writer().unwrap();
6422            guard
6423                .conn()
6424                .execute_batch(
6425                    "CREATE TABLE IF NOT EXISTS write_queue_rollback_test \
6426                     (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
6427                )
6428                .unwrap();
6429        }
6430
6431        let bridge = SqlBridge::new(Arc::clone(&pool), true);
6432
6433        let mut writer = bridge.writer().await.unwrap();
6434        let result = writer
6435            .execute_batch(vec![
6436                // Statement 1: succeeds on its own.
6437                SqlStatement {
6438                    sql: "INSERT INTO write_queue_rollback_test (id, val) VALUES (?1, ?2)".into(),
6439                    params: vec![SqlValue::Integer(1), SqlValue::Text("first".into())],
6440                    label: None,
6441                },
6442                // Statement 2: duplicate primary key — fails mid-sequence.
6443                SqlStatement {
6444                    sql: "INSERT INTO write_queue_rollback_test (id, val) VALUES (?1, ?2)".into(),
6445                    params: vec![SqlValue::Integer(1), SqlValue::Text("duplicate".into())],
6446                    label: None,
6447                },
6448                // Statement 3: never reached.
6449                SqlStatement {
6450                    sql: "INSERT INTO write_queue_rollback_test (id, val) VALUES (?1, ?2)".into(),
6451                    params: vec![SqlValue::Integer(2), SqlValue::Text("third".into())],
6452                    label: None,
6453                },
6454            ])
6455            .await;
6456        assert!(
6457            result.is_err(),
6458            "a batch with a mid-sequence PK conflict must return an error"
6459        );
6460
6461        let mut reader = bridge.reader().await.unwrap();
6462        let count = reader
6463            .query_scalar(SqlStatement {
6464                sql: "SELECT COUNT(*) FROM write_queue_rollback_test".into(),
6465                params: vec![],
6466                label: None,
6467            })
6468            .await
6469            .unwrap();
6470        assert!(
6471            matches!(count, Some(SqlValue::Integer(0))),
6472            "the whole request must roll back — including statement 1's \
6473             otherwise-successful INSERT — not just the failing statement; \
6474             got {count:?}"
6475        );
6476    }
6477
6478    /// ADR-067 Component A: before
6479    /// this fix, `block_on_sync` (this file) `unreachable!()`-panicked if
6480    /// an `atomic_unit` closure's future was `Pending` on its first poll.
6481    /// That panic ran inside the writer task's own `spawn_blocking` frame
6482    /// (see `atomic_unit`'s flag-on branch), and `run_writer_task` treats
6483    /// any `spawn_blocking` `JoinError` as fatal — the whole writer task
6484    /// exits, taking down every subsequent write for this pool. Proves the
6485    /// fix: an `atomic_unit` op built to suspend on first poll (via
6486    /// `std::future::pending`, never actually resolving) now returns a
6487    /// clean `Err` from `atomic_unit` — no panic — AND the writer task
6488    /// survives to serve a completely unrelated, well-behaved `atomic_unit`
6489    /// call immediately afterward.
6490    ///
6491    /// Not `#[serial]` / no env var: builds the pool directly with
6492    /// `write_queue_enabled: Some(true)` in the `PoolConfig` literal, same
6493    /// technique as this round's other new routing tests.
6494    #[tokio::test]
6495    async fn atomic_unit_pending_future_errors_without_killing_writer_task() {
6496        let dir = tempfile::tempdir().unwrap();
6497        let path = dir.path().join("atomic_unit_pending_future.db");
6498        let config = PoolConfig {
6499            path: Some(path.clone()),
6500            write_queue_enabled: Some(true),
6501            ..PoolConfig::default()
6502        };
6503        let pool = Arc::new(ConnectionPool::new(config).unwrap());
6504        {
6505            let guard = pool.writer().unwrap();
6506            guard
6507                .conn()
6508                .execute_batch(
6509                    "CREATE TABLE IF NOT EXISTS atomic_unit_pending_test \
6510                     (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
6511                )
6512                .unwrap();
6513        }
6514        assert!(
6515            pool.writer_task_handle().unwrap().is_some(),
6516            "writer task must be spawned with the flag on for a file-backed pool"
6517        );
6518
6519        let bridge = SqlBridge::new(Arc::clone(&pool), true);
6520
6521        // A closure whose future never resolves on first poll — the exact
6522        // misuse `block_on_sync` must reject instead of panicking on.
6523        let pending_op: AtomicUnitOp = Box::new(|_writer| {
6524            Box::pin(std::future::pending::<
6525                khive_storage::types::StorageResult<Box<dyn std::any::Any + Send>>,
6526            >())
6527        });
6528
6529        let pending_result = bridge.atomic_unit(pending_op).await;
6530        assert!(
6531            pending_result.is_err(),
6532            "a Pending-on-first-poll atomic_unit closure must return Err, \
6533             not panic; got {pending_result:?}"
6534        );
6535
6536        // If the panic had instead killed the writer task, every subsequent
6537        // write on this pool (including a completely unrelated, correctly
6538        // non-blocking atomic_unit call) would now fail with a channel-closed
6539        // error. Prove the task is still alive and serving requests.
6540        let ok_op: AtomicUnitOp = Box::new(|writer| {
6541            Box::pin(async move {
6542                writer
6543                    .execute(SqlStatement {
6544                        sql: "INSERT INTO atomic_unit_pending_test (id, val) VALUES (?1, ?2)"
6545                            .into(),
6546                        params: vec![SqlValue::Integer(1), SqlValue::Text("survived".into())],
6547                        label: None,
6548                    })
6549                    .await
6550                    .map_err(|e| {
6551                        khive_storage::StorageError::driver(
6552                            StorageCapability::Sql,
6553                            "atomic_unit_pending_future_test_insert",
6554                            e,
6555                        )
6556                    })?;
6557                Ok(Box::new(()) as Box<dyn std::any::Any + Send>)
6558            })
6559        });
6560        let ok_result = bridge.atomic_unit(ok_op).await;
6561        assert!(
6562            ok_result.is_ok(),
6563            "writer task must survive a Pending misuse and keep serving \
6564             subsequent well-behaved atomic_unit requests; got {ok_result:?}"
6565        );
6566
6567        let mut reader = bridge.reader().await.unwrap();
6568        let count = reader
6569            .query_scalar(SqlStatement {
6570                sql: "SELECT COUNT(*) FROM atomic_unit_pending_test".into(),
6571                params: vec![],
6572                label: None,
6573            })
6574            .await
6575            .unwrap();
6576        assert!(
6577            matches!(count, Some(SqlValue::Integer(1))),
6578            "the well-behaved atomic_unit call after the Pending misuse must \
6579             have actually committed its write; got {count:?}"
6580        );
6581    }
6582
6583    /// ADR-136 D1 gate 1/3: with `KHIVE_WRITE_ROUTING=strict` and no writer
6584    /// task available, `SqlBridge::writer()` must error instead of silently
6585    /// degrading to a standalone connection — even when the reason no handle
6586    /// exists is simply that the queue itself was never enabled. Strict
6587    /// routing without an enabled queue is a caller misconfiguration this
6588    /// gate refuses rather than silently no-ops: an operator who set
6589    /// `KHIVE_WRITE_ROUTING=strict` believing every write is single-admission
6590    /// must be told loudly if `KHIVE_WRITE_QUEUE` was never turned on, not
6591    /// left thinking strict routing is in effect when it is not.
6592    #[tokio::test]
6593    async fn writer_strict_routing_fails_closed_without_writer_task() {
6594        let dir = tempfile::tempdir().unwrap();
6595        let path = dir.path().join("strict_writer.db");
6596        let config = PoolConfig {
6597            path: Some(path),
6598            write_queue_enabled: Some(false),
6599            write_routing_strict: true,
6600            ..PoolConfig::default()
6601        };
6602        let pool = Arc::new(ConnectionPool::new(config).unwrap());
6603        let bridge = SqlBridge::new(Arc::clone(&pool), true);
6604
6605        let result = bridge.writer().await;
6606        let err = match result {
6607            Ok(_) => panic!(
6608                "KHIVE_WRITE_ROUTING=strict with no writer task must fail closed, not \
6609                 silently degrade to a standalone connection"
6610            ),
6611            Err(e) => e,
6612        };
6613        assert!(
6614            err.to_string().contains("strict"),
6615            "error must name strict routing, got: {err}"
6616        );
6617    }
6618
6619    /// ADR-136 D1 gate 3/4: with `KHIVE_WRITE_ROUTING=strict` and no writer
6620    /// task available, `SqlBridge::atomic_unit` must error instead of
6621    /// silently falling back to a manual `BEGIN IMMEDIATE` on a standalone
6622    /// connection.
6623    #[tokio::test]
6624    async fn atomic_unit_strict_routing_fails_closed_without_writer_task() {
6625        let dir = tempfile::tempdir().unwrap();
6626        let path = dir.path().join("strict_atomic_unit.db");
6627        let config = PoolConfig {
6628            path: Some(path),
6629            write_queue_enabled: Some(false),
6630            write_routing_strict: true,
6631            ..PoolConfig::default()
6632        };
6633        let pool = Arc::new(ConnectionPool::new(config).unwrap());
6634        let bridge = SqlBridge::new(Arc::clone(&pool), true);
6635
6636        let op: AtomicUnitOp = Box::new(|_writer| {
6637            Box::pin(async move { Ok(Box::new(()) as Box<dyn std::any::Any + Send>) })
6638        });
6639        let result = bridge.atomic_unit(op).await;
6640        assert!(
6641            result.is_err(),
6642            "KHIVE_WRITE_ROUTING=strict but the queue is off (no writer task handle) must \
6643             fail closed instead of falling back to a manual BEGIN IMMEDIATE; got {result:?}"
6644        );
6645        let msg = result.unwrap_err().to_string();
6646        assert!(
6647            msg.contains("strict"),
6648            "error must name strict routing, got: {msg}"
6649        );
6650    }
6651
6652    /// ADR-136 D1 gate 3 amendment: production DOES read through a
6653    /// queue-backed `writer()` handle — `khive-pack-comm`'s handlers obtain
6654    /// a writer then call `w.query_row(...)` cursor-style, and
6655    /// `khive-pack-gtd`'s bootstrap calls `w.query_all("PRAGMA
6656    /// table_info...")` on one. Exercise that exact shape under a strict,
6657    /// queue-enabled pool: write through the handle, then read the same row
6658    /// back through it before it is dropped.
6659    #[tokio::test]
6660    async fn writer_handle_supports_read_after_write_under_strict_queue() {
6661        let dir = tempfile::tempdir().unwrap();
6662        let path = dir.path().join("writer_read_after_write.db");
6663        let config = PoolConfig {
6664            path: Some(path),
6665            write_queue_enabled: Some(true),
6666            write_routing_strict: true,
6667            ..PoolConfig::default()
6668        };
6669        let pool = Arc::new(ConnectionPool::new(config).unwrap());
6670        {
6671            let guard = pool.writer().unwrap();
6672            guard
6673                .conn()
6674                .execute_batch(
6675                    "CREATE TABLE IF NOT EXISTS writer_cursor_test \
6676                     (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
6677                )
6678                .unwrap();
6679        }
6680
6681        let bridge = SqlBridge::new(Arc::clone(&pool), true);
6682
6683        let mut w = bridge.writer().await.unwrap();
6684        w.execute(SqlStatement {
6685            sql: "INSERT INTO writer_cursor_test (id, val) VALUES (?1, ?2)".into(),
6686            params: vec![SqlValue::Integer(1), SqlValue::Text("via-writer".into())],
6687            label: None,
6688        })
6689        .await
6690        .unwrap();
6691
6692        let row = w
6693            .query_row(SqlStatement {
6694                sql: "SELECT val FROM writer_cursor_test WHERE id = ?1".into(),
6695                params: vec![SqlValue::Integer(1)],
6696                label: None,
6697            })
6698            .await
6699            .unwrap()
6700            .expect("row inserted through the same writer handle must be visible to it");
6701        assert!(
6702            matches!(&row.columns[0].value, SqlValue::Text(v) if v == "via-writer"),
6703            "query_row through a queue-backed writer handle must see its own \
6704             committed write; got {:?}",
6705            row.columns[0].value
6706        );
6707    }
6708
6709    /// Queue-backed reads charge the READER permit budget for connection
6710    /// open and query operations, and a cancelled queue-backed
6711    /// read is followed by a successful lazy reopen — the documented
6712    /// contrast with standalone handles, whose consumed connection makes
6713    /// every later call fail. Arm 1 saturates the one reader permit and
6714    /// asserts the queue-backed read times out on the READER budget (not
6715    /// the writer budget). Arm 2 aborts an in-flight read and asserts the
6716    /// next read on the same handle succeeds by reopening.
6717    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6718    async fn queue_backed_read_uses_reader_budget_and_reopens_after_cancel() {
6719        let dir = tempfile::tempdir().unwrap();
6720        let path = dir.path().join("queue_backed_reader_budget.db");
6721        let config = PoolConfig {
6722            path: Some(path),
6723            write_queue_enabled: Some(true),
6724            write_routing_strict: true,
6725            max_readers: 1,
6726            checkout_timeout: std::time::Duration::from_millis(250),
6727            ..PoolConfig::default()
6728        };
6729        let pool = Arc::new(ConnectionPool::new(config).unwrap());
6730        {
6731            let guard = pool.writer().unwrap();
6732            guard
6733                .conn()
6734                .execute_batch(
6735                    "CREATE TABLE IF NOT EXISTS reopen_test \
6736                     (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
6737                )
6738                .unwrap();
6739        }
6740        let bridge = SqlBridge::new(Arc::clone(&pool), true);
6741
6742        let mut w = bridge.writer().await.unwrap();
6743        w.execute(SqlStatement {
6744            sql: "INSERT INTO reopen_test (id, val) VALUES (1, 'seed')".into(),
6745            params: vec![],
6746            label: None,
6747        })
6748        .await
6749        .unwrap();
6750
6751        // Arm 1: with the sole reader permit held, the queue-backed read
6752        // must time out on the reader budget.
6753        let held = pool
6754            .sql_bridge_reader_slots()
6755            .acquire_owned()
6756            .await
6757            .unwrap();
6758        let starved = w
6759            .query_row(SqlStatement {
6760                sql: "SELECT val FROM reopen_test WHERE id = 1".into(),
6761                params: vec![],
6762                label: None,
6763            })
6764            .await;
6765        assert!(
6766            matches!(
6767                &starved,
6768                Err(StorageError::Timeout { operation })
6769                    if operation.as_ref() == "sql_bridge.reader_open"
6770            ),
6771            "queue-backed read with reader permits saturated must time out \
6772             on the reader budget with the ADR-005 reader contract error; \
6773             got {starved:?}"
6774        );
6775        drop(held);
6776
6777        // Arm 2: a queue-backed handle in the exact post-cancelled-read
6778        // state — `handle: None` because the cancelled call took the boxed
6779        // connection out and never returned it — must serve the next read by
6780        // lazily reopening, never a hard "connection already consumed"
6781        // failure (that contract is standalone-only; the doc names the
6782        // contrast). Constructed directly so the state is deterministic
6783        // rather than racing an abort against ensure_conn.
6784        let writer_task = pool
6785            .writer_task_handle()
6786            .expect("queue-enabled file pool must offer a writer task")
6787            .expect("writer task present under write_queue_enabled");
6788        let mut post_cancel = SqliteWriter {
6789            handle: None,
6790            writer_task: Some(writer_task),
6791            origin: pool.origin(),
6792            db: crate::timeout_sink::db_label(&pool),
6793            pool: Arc::clone(&pool),
6794        };
6795        let row = post_cancel
6796            .query_row(SqlStatement {
6797                sql: "SELECT val FROM reopen_test WHERE id = 1".into(),
6798                params: vec![],
6799                label: None,
6800            })
6801            .await
6802            .expect("read on a queue-backed handle with no resident connection must reopen")
6803            .expect("seeded row must be visible");
6804        assert!(
6805            matches!(&row.columns[0].value, SqlValue::Text(v) if v == "seed"),
6806            "reopened read must return the seeded row; got {:?}",
6807            row.columns[0].value
6808        );
6809    }
6810
6811    /// ADR-136 D1 gate 3 amendment: `SqlWriter::query_row`/`query_all` carry
6812    /// no read-only restriction at the trait level — a caller could hand a
6813    /// DML-with-RETURNING statement to `query_row` expecting it to behave
6814    /// like any other query. Under a queue-backed handle, the standalone
6815    /// connection `SqliteWriter::ensure_conn` lazily opens must be
6816    /// read-only, so SQLite rejects the statement outright instead of
6817    /// quietly mutating the row on an untracked connection outside the
6818    /// `WriterTask`. Red-proof: reverting `ensure_conn` to
6819    /// `open_standalone_writer` makes this test fail (the UPDATE succeeds
6820    /// and mutates the row).
6821    #[tokio::test]
6822    async fn writer_query_row_rejects_dml_with_returning_on_queue_backed_handle() {
6823        let dir = tempfile::tempdir().unwrap();
6824        let path = dir.path().join("writer_readonly_returning.db");
6825        let config = PoolConfig {
6826            path: Some(path),
6827            write_queue_enabled: Some(true),
6828            write_routing_strict: true,
6829            ..PoolConfig::default()
6830        };
6831        let pool = Arc::new(ConnectionPool::new(config).unwrap());
6832        {
6833            let guard = pool.writer().unwrap();
6834            guard
6835                .conn()
6836                .execute_batch(
6837                    "CREATE TABLE IF NOT EXISTS writer_returning_test \
6838                     (id INTEGER PRIMARY KEY, val TEXT NOT NULL);
6839                     INSERT INTO writer_returning_test (id, val) VALUES (1, 'original');",
6840                )
6841                .unwrap();
6842        }
6843
6844        let bridge = SqlBridge::new(Arc::clone(&pool), true);
6845
6846        let mut w = bridge.writer().await.unwrap();
6847        let result = w
6848            .query_row(SqlStatement {
6849                sql: "UPDATE writer_returning_test SET val = 'mutated' \
6850                      WHERE id = ?1 RETURNING val"
6851                    .into(),
6852                params: vec![SqlValue::Integer(1)],
6853                label: None,
6854            })
6855            .await;
6856        assert!(
6857            result.is_err(),
6858            "a DML-with-RETURNING statement through query_row on a \
6859             queue-backed writer handle must be rejected, not executed on \
6860             an untracked read-write connection; got {result:?}"
6861        );
6862
6863        let mut reader = bridge.reader().await.unwrap();
6864        let val = reader
6865            .query_scalar(SqlStatement {
6866                sql: "SELECT val FROM writer_returning_test WHERE id = ?1".into(),
6867                params: vec![SqlValue::Integer(1)],
6868                label: None,
6869            })
6870            .await
6871            .unwrap();
6872        assert!(
6873            matches!(&val, Some(SqlValue::Text(v)) if v == "original"),
6874            "the rejected UPDATE...RETURNING must not have altered the row; got {val:?}"
6875        );
6876    }
6877
6878    /// ADR-136 D1 acceptance arm: a 5-op batch shaped like `[send, mark,
6879    /// mark, mark, mark]` at the storage layer — every "mark" op is a real
6880    /// `UPDATE` against its own pre-seeded row, so (like `send`) it routes
6881    /// through the writer task rather than bypassing it as a `SELECT` would
6882    /// — issued while 3 concurrent writers contend the write path, must
6883    /// complete every op — no checkout timeout — once routing is strict and
6884    /// the queue is on. An occupier holds the writer task's single drain
6885    /// slot until all 8 requests (3 contenders + send + 4 marks) are
6886    /// provably enqueued behind it (`queue_depth() >= 8`, the same
6887    /// occupier/`queue_depth()` discriminator the migrated-call-site tests
6888    /// use), so this proves genuine contention instead of a scheduler that
6889    /// happens to drain the tiny writes before the others even enqueue.
6890    /// Mirrors the measured production failure ADR-136's Context section
6891    /// documents (middle ops of a batch starving while a sibling write wins
6892    /// under the legacy fixed-deadline pool mutex). Red-proofed: reverting
6893    /// the marks back to `bridge.reader()` `SELECT`s (the pre-fix shape)
6894    /// makes the `queue_depth() >= 8` wait time out and fail, since a read
6895    /// never reaches the writer task's channel — confirming this version
6896    /// actually requires all four marks to be real writes.
6897    #[tokio::test]
6898    async fn acceptance_five_op_batch_completes_under_concurrent_write_contention() {
6899        let dir = tempfile::tempdir().unwrap();
6900        let path = dir.path().join("acceptance_batch.db");
6901        let config = PoolConfig {
6902            path: Some(path),
6903            write_queue_enabled: Some(true),
6904            write_routing_strict: true,
6905            ..PoolConfig::default()
6906        };
6907        let pool = Arc::new(ConnectionPool::new(config).unwrap());
6908        {
6909            let guard = pool.writer().unwrap();
6910            guard
6911                .conn()
6912                .execute_batch(
6913                    "CREATE TABLE IF NOT EXISTS acceptance_batch \
6914                     (id INTEGER PRIMARY KEY, val TEXT NOT NULL);
6915                     INSERT INTO acceptance_batch (id, val) VALUES \
6916                     (200, 'seed-0'), (201, 'seed-1'), (202, 'seed-2'), (203, 'seed-3');",
6917                )
6918                .unwrap();
6919        }
6920
6921        let bridge = Arc::new(SqlBridge::new(Arc::clone(&pool), true));
6922
6923        let writer_task = pool
6924            .writer_task_handle()
6925            .unwrap()
6926            .expect("writer task must be spawned for a file-backed pool with the flag on");
6927
6928        // Occupier: holds the single writer-task drain slot until released,
6929        // so every op below is provably queued behind it rather than racing
6930        // to finish before the others even enqueue (same technique as
6931        // `rename_namespace_routes_through_writer_task_when_flag_enabled` in
6932        // `stores::text_tests`).
6933        let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
6934        let (release_tx, release_rx) = tokio::sync::oneshot::channel::<()>();
6935        let occupier = {
6936            let writer_task = writer_task.clone();
6937            tokio::spawn(async move {
6938                writer_task
6939                    .send(move |_conn| {
6940                        let _ = started_tx.send(());
6941                        let _ = release_rx.blocking_recv();
6942                        Ok::<(), StorageError>(())
6943                    })
6944                    .await
6945            })
6946        };
6947        started_rx
6948            .await
6949            .expect("occupier must signal it has started running inside the writer task");
6950        assert_eq!(
6951            writer_task.queue_depth(),
6952            0,
6953            "channel must start empty once the occupier has been dequeued and is running"
6954        );
6955
6956        // 3 concurrent writers contending the write path — each a
6957        // self-contained `execute()` through `SqlBridge::writer()`, matching
6958        // the "several short acquisitions" shape ADR-136's Context section
6959        // measures (a logical write is many short holds, not one long one).
6960        let contenders: Vec<_> = (0..3)
6961            .map(|i| {
6962                let bridge = Arc::clone(&bridge);
6963                tokio::spawn(async move {
6964                    let mut writer = bridge.writer().await?;
6965                    writer
6966                        .execute(SqlStatement {
6967                            sql: "INSERT INTO acceptance_batch (id, val) VALUES (?1, ?2)".into(),
6968                            params: vec![
6969                                SqlValue::Integer(100 + i),
6970                                SqlValue::Text(format!("contender-{i}")),
6971                            ],
6972                            label: None,
6973                        })
6974                        .await
6975                })
6976            })
6977            .collect();
6978
6979        // The 5-op batch: [send, mark, mark, mark, mark].
6980        let send = {
6981            let bridge = Arc::clone(&bridge);
6982            tokio::spawn(async move {
6983                let mut writer = bridge.writer().await?;
6984                writer
6985                    .execute(SqlStatement {
6986                        sql: "INSERT INTO acceptance_batch (id, val) VALUES (?1, ?2)".into(),
6987                        params: vec![SqlValue::Integer(1), SqlValue::Text("send".into())],
6988                        label: None,
6989                    })
6990                    .await
6991            })
6992        };
6993        let marks: Vec<_> = (0..4)
6994            .map(|i| {
6995                let bridge = Arc::clone(&bridge);
6996                tokio::spawn(async move {
6997                    let mut writer = bridge.writer().await?;
6998                    writer
6999                        .execute(SqlStatement {
7000                            sql: "UPDATE acceptance_batch SET val = ?2 WHERE id = ?1".into(),
7001                            params: vec![
7002                                SqlValue::Integer(200 + i),
7003                                SqlValue::Text(format!("marked-{i}")),
7004                            ],
7005                            label: None,
7006                        })
7007                        .await
7008                })
7009            })
7010            .collect();
7011
7012        // All 8 requests must actually reach the writer task's channel
7013        // while the occupier still holds the single drain slot.
7014        let mut saw_all_enqueued = false;
7015        for _ in 0..200 {
7016            if writer_task.queue_depth() >= 8 {
7017                saw_all_enqueued = true;
7018                break;
7019            }
7020            tokio::time::sleep(std::time::Duration::from_millis(5)).await;
7021        }
7022        assert!(
7023            saw_all_enqueued,
7024            "not all 8 contending writes (3 contenders + send + 4 marks) reached \
7025             the writer task's channel while the occupier held the single drain \
7026             slot — got depth {}",
7027            writer_task.queue_depth()
7028        );
7029
7030        release_tx
7031            .send(())
7032            .expect("occupier must still be waiting on the release signal");
7033        occupier
7034            .await
7035            .expect("occupier task must not panic")
7036            .expect("occupier write must succeed");
7037
7038        for c in contenders {
7039            c.await
7040                .expect("contender task must not panic")
7041                .expect("contender write must complete without a checkout timeout");
7042        }
7043        send.await
7044            .expect("send task must not panic")
7045            .expect("send op must complete without a checkout timeout");
7046        for (i, m) in marks.into_iter().enumerate() {
7047            let affected = m
7048                .await
7049                .expect("mark task must not panic")
7050                .expect("mark op must complete without a checkout timeout — no starvation");
7051            assert_eq!(
7052                affected, 1,
7053                "mark {i} must have updated exactly its own row"
7054            );
7055        }
7056
7057        let mut reader = bridge.reader().await.unwrap();
7058        let count = reader
7059            .query_scalar(SqlStatement {
7060                sql: "SELECT COUNT(*) FROM acceptance_batch".into(),
7061                params: vec![],
7062                label: None,
7063            })
7064            .await
7065            .unwrap();
7066        assert!(
7067            matches!(count, Some(SqlValue::Integer(8))),
7068            "the 4 seeded mark rows plus 3 contenders plus the batch's own send \
7069             must all be present; got {count:?}"
7070        );
7071
7072        for i in 0..4i64 {
7073            let mut reader = bridge.reader().await.unwrap();
7074            let val = reader
7075                .query_scalar(SqlStatement {
7076                    sql: "SELECT val FROM acceptance_batch WHERE id = ?1".into(),
7077                    params: vec![SqlValue::Integer(200 + i)],
7078                    label: None,
7079                })
7080                .await
7081                .unwrap();
7082            assert!(
7083                matches!(&val, Some(SqlValue::Text(v)) if *v == format!("marked-{i}")),
7084                "mark row {i} must reflect the persisted UPDATE after release; got {val:?}"
7085            );
7086        }
7087    }
7088
7089    #[tokio::test]
7090    async fn file_backed_bridge_counts_writer_and_flag_off_atomic_unit_acquisitions() {
7091        let dir = tempfile::tempdir().unwrap();
7092        let config = PoolConfig {
7093            path: Some(dir.path().join("bridge_writer_acquisitions.db")),
7094            write_queue_enabled: Some(false),
7095            ..PoolConfig::default()
7096        };
7097        let pool = Arc::new(ConnectionPool::new(config).unwrap());
7098        let bridge = SqlBridge::new(Arc::clone(&pool), true);
7099
7100        let before = pool.writer_acquisition_snapshot();
7101
7102        drop(bridge.writer().await.unwrap());
7103        let after_writer = pool.writer_acquisition_snapshot();
7104        assert_eq!(
7105            after_writer.standalone_acquisitions,
7106            before.standalone_acquisitions + 1
7107        );
7108        assert_eq!(after_writer.acquisitions, before.acquisitions + 1);
7109        assert_eq!(after_writer.pooled_acquisitions, before.pooled_acquisitions);
7110        assert_eq!(
7111            after_writer.writer_task_acquisitions,
7112            before.writer_task_acquisitions
7113        );
7114
7115        let op: AtomicUnitOp = Box::new(|_writer| {
7116            Box::pin(async { Ok(Box::new(()) as Box<dyn std::any::Any + Send>) })
7117        });
7118        bridge.atomic_unit(op).await.unwrap();
7119
7120        let after_atomic_unit = pool.writer_acquisition_snapshot();
7121        assert_eq!(
7122            after_atomic_unit.standalone_acquisitions,
7123            before.standalone_acquisitions + 2
7124        );
7125        assert_eq!(after_atomic_unit.acquisitions, before.acquisitions + 2);
7126        assert_eq!(
7127            after_atomic_unit.pooled_acquisitions,
7128            before.pooled_acquisitions
7129        );
7130        assert_eq!(
7131            after_atomic_unit.writer_task_acquisitions,
7132            before.writer_task_acquisitions
7133        );
7134    }
7135}