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**: ordinary reads check out pooled readers per operation.
5//!   The only standalone-reader exception is an explicitly admitted multi-call
6//!   deferred read transaction; standalone writer handles remain capped at one.
7//!   Cross-statement write atomicity goes through `atomic_unit`, which drives a
8//!   single registered raw transaction span rather than a caller-held per-tx
9//!   connection.
10//! - **Memory**: Uses pool-backed approach (acquire pool connection per-query inside `spawn_blocking`).
11
12#[path = "sql_bridge/write_errors.rs"]
13mod write_errors;
14
15#[path = "sql_bridge/manual_atomic.rs"]
16mod manual_atomic;
17use manual_atomic::run_manual_atomic_unit;
18#[path = "sql_bridge/standalone_admission.rs"]
19mod standalone_admission;
20use standalone_admission::{
21    acquire_standalone_lease, acquire_unit_lease, admit_standalone_operation, StandaloneWriteError,
22};
23mod standalone_batch;
24#[cfg(test)]
25use standalone_batch::BatchPoisonReason;
26use standalone_batch::{
27    run_standalone_batch, run_standalone_script, run_standalone_statement,
28    run_standalone_top_level, BatchFailure, PoisonedBatchError,
29};
30
31use std::any::Any;
32use std::sync::atomic::{AtomicU64, Ordering};
33use std::sync::Arc;
34use std::time::Instant;
35
36use async_trait::async_trait;
37
38use khive_storage::error::StorageError;
39use khive_storage::types::{PageRequest, SqlColumn, SqlRow, SqlStatement, SqlValue};
40use khive_storage::{AtomicUnitOp, StorageCapability, TopLevelMaintenance};
41use tokio::sync::{OwnedSemaphorePermit, Semaphore};
42
43use crate::error::SqliteError;
44use crate::pool::{ConnectionPool, SharedReaderTransactionGuard, StandaloneReaderPurpose};
45
46// =============================================================================
47// Shared helpers
48// =============================================================================
49
50mod rows;
51
52pub(crate) use rows::bind_params;
53use rows::{
54    prepare_batch_statements, prepare_cached_sql_statement, prepare_sql_statement, row_to_sql_row,
55    AtomicEventRows, PreparedBatchStatement,
56};
57#[cfg(test)]
58use rows::{COUNTED_EVENT_INSERT_LABELS, ROW_CONVERSIONS};
59
60#[cfg(test)]
61#[path = "atomic_event_usage_tests.rs"]
62mod atomic_event_usage_tests;
63
64/// Bind and execute handles returned by [`prepare_batch_statements`].
65fn execute_prepared_batch<'conn>(
66    conn: &'conn rusqlite::Connection,
67    prepared: Vec<PreparedBatchStatement<'conn>>,
68    statements: &[SqlStatement],
69    event_rows: Option<&AtomicEventRows>,
70) -> Result<u64, rusqlite::Error> {
71    debug_assert_eq!(prepared.len(), statements.len());
72    let mut total = 0u64;
73    for (prepared, statement) in prepared.into_iter().zip(statements) {
74        let mut prepared = match prepared {
75            PreparedBatchStatement::Ready(prepared) => prepared,
76            PreparedBatchStatement::PrepareAtExecution => {
77                prepare_sql_statement(conn, &statement.sql)?
78            }
79        };
80        bind_params(&mut prepared, &statement.params)?;
81        let affected = prepared.raw_execute()? as u64;
82        if let Some(event_rows) = event_rows {
83            event_rows.observe(statement, affected);
84        }
85        total += affected;
86    }
87    Ok(total)
88}
89
90/// SQL statement heads that are transaction control. `execute_batch` owns the
91/// `BEGIN`/`COMMIT` boundary for the whole batch (the standalone path wraps
92/// the list in its own `BEGIN IMMEDIATE`, and the queue-backed path runs
93/// inside the writer task's per-request transaction), so a caller-supplied
94/// statement that itself starts, ends, or branches a transaction can commit
95/// or roll back early and break the batch's all-or-nothing contract. Cached
96/// read-only handles use the same lexical classification to drive their
97/// separately admitted single-level read-transaction state machine below.
98/// `START` is classified as the alternate transaction-opening spelling so
99/// callers get a typed boundary error; `END` is SQLite's `COMMIT` spelling.
100const TRANSACTION_CONTROL_KEYWORDS: [&str; 7] = [
101    "BEGIN",
102    "START",
103    "COMMIT",
104    "END",
105    "ROLLBACK",
106    "SAVEPOINT",
107    "RELEASE",
108];
109
110/// Skip the same leading whitespace, UTF-8 BOMs, empty statements (`;`), and
111/// line/block comments SQLite accepts before an executable statement.
112fn skip_sqlite_empty_prefix(mut rest: &[u8]) -> &[u8] {
113    loop {
114        let mut idx = 0;
115        while idx < rest.len() && rest[idx].is_ascii_whitespace() {
116            idx += 1;
117        }
118        rest = &rest[idx..];
119        if let Some(tail) = rest.strip_prefix(b"\xEF\xBB\xBF") {
120            rest = tail;
121            continue;
122        }
123        if let Some(tail) = rest.strip_prefix(b";") {
124            rest = tail;
125            continue;
126        }
127        if let Some(tail) = rest.strip_prefix(b"--") {
128            let mut idx = 0;
129            while idx < tail.len() && tail[idx] != b'\n' {
130                idx += 1;
131            }
132            rest = if idx < tail.len() {
133                &tail[idx + 1..]
134            } else {
135                &[]
136            };
137            continue;
138        }
139        if let Some(tail) = rest.strip_prefix(b"/*") {
140            let mut idx = 0;
141            while idx + 1 < tail.len() && !(tail[idx] == b'*' && tail[idx + 1] == b'/') {
142                idx += 1;
143            }
144            rest = if idx + 1 < tail.len() {
145                &tail[idx + 2..]
146            } else {
147                &[]
148            };
149            continue;
150        }
151        break;
152    }
153    rest
154}
155
156/// Return one ASCII SQL token after SQLite whitespace/comments/BOM trivia.
157/// Empty-statement separators are deliberately not trivia here: callers use
158/// this only after the executable statement head has already been consumed.
159fn next_sqlite_token(mut rest: &[u8]) -> Option<(&[u8], &[u8])> {
160    loop {
161        let mut idx = 0;
162        while idx < rest.len() && rest[idx].is_ascii_whitespace() {
163            idx += 1;
164        }
165        rest = &rest[idx..];
166        if let Some(tail) = rest.strip_prefix(b"\xEF\xBB\xBF") {
167            rest = tail;
168            continue;
169        }
170        if let Some(tail) = rest.strip_prefix(b"--") {
171            let mut idx = 0;
172            while idx < tail.len() && tail[idx] != b'\n' {
173                idx += 1;
174            }
175            rest = if idx < tail.len() {
176                &tail[idx + 1..]
177            } else {
178                &[]
179            };
180            continue;
181        }
182        if let Some(tail) = rest.strip_prefix(b"/*") {
183            let mut idx = 0;
184            while idx + 1 < tail.len() && !(tail[idx] == b'*' && tail[idx + 1] == b'/') {
185                idx += 1;
186            }
187            rest = if idx + 1 < tail.len() {
188                &tail[idx + 2..]
189            } else {
190                &[]
191            };
192            continue;
193        }
194        break;
195    }
196
197    let len = rest
198        .iter()
199        .take_while(|byte| byte.is_ascii_alphanumeric() || **byte == b'_')
200        .count();
201    (len != 0).then_some((&rest[..len], &rest[len..]))
202}
203
204/// Return a transaction-control keyword and the bytes following it, if any.
205/// Matching is case-insensitive and requires a word boundary, so an identifier
206/// that merely starts with `begin` or `commit` never matches.
207fn transaction_control_parts(sql: &str) -> Option<(&'static str, &[u8])> {
208    let rest = skip_sqlite_empty_prefix(sql.as_bytes());
209    TRANSACTION_CONTROL_KEYWORDS
210        .iter()
211        .copied()
212        .find_map(|keyword| {
213            let kw = keyword.as_bytes();
214            if rest.len() < kw.len() || !rest[..kw.len()].eq_ignore_ascii_case(kw) {
215                return None;
216            }
217            let boundary = match rest.get(kw.len()) {
218                Some(next) => !(next.is_ascii_alphanumeric() || *next == b'_'),
219                None => true,
220            };
221            boundary.then_some((keyword, &rest[kw.len()..]))
222        })
223}
224
225/// Return the transaction-control keyword heading `sql`, if any.
226fn transaction_control_head(sql: &str) -> Option<&'static str> {
227    transaction_control_parts(sql).map(|(keyword, _)| keyword)
228}
229
230#[derive(Clone, Copy, Debug, PartialEq, Eq)]
231enum CachedReadTransactionControl {
232    /// `BEGIN`, `BEGIN TRANSACTION`, or their explicitly `DEFERRED` form.
233    BeginDeferred,
234    /// A transaction-ending control statement and its diagnostic keyword.
235    Finish(&'static str),
236    /// Transaction control that cannot be represented by the single-level
237    /// admitted read-transaction state machine.
238    Unsupported(&'static str),
239}
240
241/// Classify transaction control for a cached read-only connection.
242///
243/// A cached reader may own exactly one top-level deferred read transaction.
244/// Immediate/exclusive starts could reserve write-side locks, while nested
245/// savepoint/rollback-to controls would require a second lifecycle level, so
246/// both remain rejected. Batch/write paths continue using
247/// [`transaction_control_head`] and reject every variant without exception.
248fn cached_read_transaction_control(sql: &str) -> Option<CachedReadTransactionControl> {
249    let (keyword, tail) = transaction_control_parts(sql)?;
250    match keyword {
251        "BEGIN" => {
252            // Accept exactly `BEGIN`, `BEGIN TRANSACTION`, `BEGIN DEFERRED`,
253            // or `BEGIN DEFERRED TRANSACTION`, with no trailing tokens.
254            // SQLite's grammar also admits `BEGIN TRANSACTION <name>` (the
255            // name parses as an identifier and is ignored), so a mode keyword
256            // in that trailing position — `BEGIN TRANSACTION IMMEDIATE` —
257            // still parses, and classifying it by its first token alone would
258            // launder what reads as a write-reserving start into a deferred
259            // one. Every trailing token is therefore Unsupported — and the
260            // check cannot stop at `next_sqlite_token` returning `None`,
261            // because that tokenizer returns `None` for any non-identifier
262            // byte, not only end-of-input: a quoted or bracketed tail
263            // (`BEGIN TRANSACTION "IMMEDIATE"`, `[IMMEDIATE]`) would fall
264            // out of the loop and read as the end of an accepted form. After
265            // the accepted keywords, the remainder must reduce to nothing
266            // under the same trivia/empty-statement skipping SQLite applies
267            // (whitespace, comments, `;`), or the statement is Unsupported.
268            let mut rest = tail;
269            let mut saw_deferred = false;
270            let mut saw_transaction = false;
271            while let Some((token, next)) = next_sqlite_token(rest) {
272                if !saw_deferred && !saw_transaction && token.eq_ignore_ascii_case(b"DEFERRED") {
273                    saw_deferred = true;
274                } else if !saw_transaction && token.eq_ignore_ascii_case(b"TRANSACTION") {
275                    saw_transaction = true;
276                } else {
277                    return Some(CachedReadTransactionControl::Unsupported(keyword));
278                }
279                rest = next;
280            }
281            if !skip_sqlite_empty_prefix(rest).is_empty() {
282                return Some(CachedReadTransactionControl::Unsupported(keyword));
283            }
284            Some(CachedReadTransactionControl::BeginDeferred)
285        }
286        "COMMIT" | "END" => Some(CachedReadTransactionControl::Finish(keyword)),
287        "ROLLBACK" => {
288            let first = next_sqlite_token(tail);
289            let rollback_target = match first {
290                Some((token, rest)) if token.eq_ignore_ascii_case(b"TRANSACTION") => {
291                    next_sqlite_token(rest).map(|(token, _)| token)
292                }
293                Some((token, _)) => Some(token),
294                None => None,
295            };
296            if rollback_target.is_some_and(|token| token.eq_ignore_ascii_case(b"TO")) {
297                Some(CachedReadTransactionControl::Unsupported(keyword))
298            } else {
299                Some(CachedReadTransactionControl::Finish(keyword))
300            }
301        }
302        _ => Some(CachedReadTransactionControl::Unsupported(keyword)),
303    }
304}
305
306/// Reject transaction-control statements in `statements` with a typed
307/// [`StorageError::InvalidInput`] BEFORE anything executes, preserving the
308/// batch's all-or-nothing contract (see [`TRANSACTION_CONTROL_KEYWORDS`]).
309fn reject_transaction_control_statements(
310    statements: &[SqlStatement],
311    operation: &'static str,
312) -> khive_storage::types::StorageResult<()> {
313    for (index, statement) in statements.iter().enumerate() {
314        if let Some(keyword) = transaction_control_head(&statement.sql) {
315            return Err(StorageError::InvalidInput {
316                capability: StorageCapability::Sql,
317                operation: operation.into(),
318                message: format!(
319                    "statement at index {index} is transaction control ({keyword}); \
320                     execute_batch owns the BEGIN/COMMIT boundary for the whole \
321                     batch — remove transaction-control statements from the batch"
322                ),
323            });
324        }
325    }
326    Ok(())
327}
328
329/// Settle a pooled call that left its connection inside a transaction, before
330/// the guard is released. The guard is checked out for this call only, so the
331/// transaction cannot be continued by a later call: it is rolled back here and
332/// the call fails, rather than the guard's drop rolling it back after the call
333/// has already reported success.
334///
335/// The outcome, in order: a rollback that cannot prove autocommit retires the
336/// writer and the call fails with `WriterSettlementUnknown` (side effects
337/// unknown), whatever the call itself returned; a call that failed keeps its own error once its
338/// transaction is rolled back; a call that succeeded fails with
339/// `InvalidInput`.
340fn settle_pooled_call<T>(
341    guard: &crate::pool::WriterGuard<'_>,
342    operation: &'static str,
343    result: khive_storage::types::StorageResult<T>,
344) -> khive_storage::types::StorageResult<T> {
345    if guard.is_autocommit() {
346        return result;
347    }
348    if let Err(settlement) = guard.rollback_or_retire("pooled call left its transaction open") {
349        if let Err(error) = &result {
350            tracing::warn!(
351                operation,
352                %error,
353                "pooled call failed inside a transaction it opened, and its rollback \
354                 could not prove autocommit; reporting the settlement failure"
355            );
356        }
357        return Err(settlement.into_storage_error(StorageCapability::Sql, operation));
358    }
359    result?;
360    Err(StorageError::InvalidInput {
361        capability: StorageCapability::Sql,
362        operation: operation.into(),
363        message: "the call left a transaction open; it was rolled back before the pooled \
364                  writer was released — use atomic_unit to run statements as one transaction"
365            .into(),
366    })
367}
368
369fn prepare_bound_statement<'conn>(
370    conn: &'conn rusqlite::Connection,
371    statement: &SqlStatement,
372) -> Result<rusqlite::Statement<'conn>, rusqlite::Error> {
373    let mut stmt = prepare_sql_statement(conn, &statement.sql)?;
374    bind_params(&mut stmt, &statement.params)?;
375    Ok(stmt)
376}
377
378fn execute_prepared_query(
379    mut stmt: rusqlite::Statement<'_>,
380) -> Result<Vec<SqlRow>, rusqlite::Error> {
381    let col_count = stmt.column_count();
382    let col_names: Vec<String> = (0..col_count)
383        .map(|i| stmt.column_name(i).unwrap_or("").to_string())
384        .collect();
385
386    let mut rows = Vec::new();
387    let mut raw_rows = stmt.raw_query();
388    while let Some(row) = raw_rows.next()? {
389        rows.push(row_to_sql_row(row, col_count, &col_names));
390    }
391    Ok(rows)
392}
393
394fn execute_prepared_query_row(
395    mut stmt: rusqlite::Statement<'_>,
396) -> Result<Option<SqlRow>, rusqlite::Error> {
397    let col_count = stmt.column_count();
398    let col_names: Vec<String> = (0..col_count)
399        .map(|i| stmt.column_name(i).unwrap_or("").to_string())
400        .collect();
401
402    let mut raw_rows = stmt.raw_query();
403    Ok(raw_rows
404        .next()?
405        .map(|row| row_to_sql_row(row, col_count, &col_names)))
406}
407
408fn execute_prepared_query_page(
409    mut stmt: rusqlite::Statement<'_>,
410    page: &PageRequest,
411) -> Result<Vec<SqlRow>, rusqlite::Error> {
412    // A zero-limit page still prepares and binds the statement, so invalid
413    // SQL fails identically across every limit; it skips the row cursor
414    // entirely and returns no rows.
415    if page.limit == 0 {
416        return Ok(Vec::new());
417    }
418
419    let col_count = stmt.column_count();
420    let col_names: Vec<String> = (0..col_count)
421        .map(|i| stmt.column_name(i).unwrap_or("").to_string())
422        .collect();
423
424    let mut rows = Vec::new();
425    let mut offset = page.offset;
426    let mut remaining = u64::from(page.limit);
427    let mut raw_rows = stmt.raw_query();
428    // The bound covers owned Rust rows only — this function advances past
429    // `offset`, owns at most the caller-supplied `page.limit` rows, and drops
430    // the statement cursor immediately afterward (ADR-005's bounded-
431    // materialization amendment). Callers own choosing a sane limit.
432    // Engine work is the query plan's own cost, not O(offset + limit):
433    // SQLite still produces and discards `offset` rows, and an unindexed
434    // ORDER BY can force a full sort of the result set before the first row
435    // is stepped. Callers deep-paging a large result set should prefer
436    // keyset pagination over growing offsets.
437    while remaining > 0 {
438        let Some(row) = raw_rows.next()? else {
439            break;
440        };
441        if offset > 0 {
442            offset -= 1;
443            continue;
444        }
445        rows.push(row_to_sql_row(row, col_count, &col_names));
446        remaining -= 1;
447    }
448    Ok(rows)
449}
450
451/// Execute a query on a `rusqlite::Connection` and return owned rows.
452fn execute_query(
453    conn: &rusqlite::Connection,
454    statement: &SqlStatement,
455) -> Result<Vec<SqlRow>, rusqlite::Error> {
456    execute_prepared_query(prepare_bound_statement(conn, statement)?)
457}
458
459fn execute_query_row(
460    conn: &rusqlite::Connection,
461    statement: &SqlStatement,
462) -> Result<Option<SqlRow>, rusqlite::Error> {
463    execute_prepared_query_row(prepare_bound_statement(conn, statement)?)
464}
465
466fn execute_query_page(
467    conn: &rusqlite::Connection,
468    statement: &SqlStatement,
469    page: &PageRequest,
470) -> Result<Vec<SqlRow>, rusqlite::Error> {
471    execute_prepared_query_page(prepare_bound_statement(conn, statement)?, page)
472}
473
474/// SQLite's prepared-statement classifier is authoritative for the safety
475/// boundary. Row-producing DML (`UPDATE ... RETURNING`) and transaction
476/// control may be called through `SqlReader`, but they are admitted writes and
477/// must never register an interrupt target.
478fn statement_is_cancellable_read(stmt: &rusqlite::Statement<'_>, sql: &str) -> bool {
479    stmt.readonly() && transaction_control_head(sql).is_none()
480}
481
482/// Introspection `PRAGMA`s admitted through the pooled reader capability with
483/// an optional single positional argument (`PRAGMA name` or `PRAGMA
484/// name(arg)`) — none of these has a set/assignment form in SQLite, so an
485/// argument is always a read-side filter, never a mutation.
486const READER_STRUCTURAL_PRAGMAS: [&str; 8] = [
487    "table_info",
488    "table_xinfo",
489    "table_list",
490    "index_list",
491    "index_info",
492    "index_xinfo",
493    "foreign_key_list",
494    "integrity_check",
495];
496
497/// Introspection `PRAGMA`s admitted through the pooled reader capability only
498/// in their bare, argument-less read form. Every one of these also has a
499/// `PRAGMA name = value` or `PRAGMA name(value)` assignment form in SQLite —
500/// an assigning form changes connection-local state that a pristine-state
501/// scan run only on dirty checkouts must never let back into the pool
502/// unnoticed (`reader_connection_state_is_pristine`,
503/// `reader_connection_settings_match_baseline`).
504const READER_SETTING_PRAGMAS: [&str; 10] = [
505    "database_list",
506    "collation_list",
507    "function_list",
508    "compile_options",
509    "page_count",
510    "freelist_count",
511    "user_version",
512    "schema_version",
513    "journal_mode",
514    "page_size",
515];
516
517/// Skip a fully parenthesized region starting at `rest[0] == b'('`,
518/// respecting SQL string/identifier quoting (`'...'` with `''` escapes,
519/// `"..."`/`` `...` `` with doubled-quote escapes, and `[...]` bracket
520/// identifiers) and line/block comments, so that a `)`, `(`, or `--` that
521/// merely appears inside a literal or comment is never mistaken for syntax.
522/// Returns the bytes after the matching `)`, or `None` if the region never
523/// closes (malformed or truncated SQL) — callers must fail closed on `None`.
524fn skip_balanced_parens(rest: &[u8]) -> Option<&[u8]> {
525    debug_assert_eq!(rest.first(), Some(&b'('));
526    let mut depth: u32 = 0;
527    let mut idx = 0;
528    loop {
529        match *rest.get(idx)? {
530            b'(' => {
531                depth += 1;
532                idx += 1;
533            }
534            b')' => {
535                depth -= 1;
536                idx += 1;
537                if depth == 0 {
538                    return Some(&rest[idx..]);
539                }
540            }
541            quote @ (b'\'' | b'"' | b'`') => {
542                idx += 1;
543                loop {
544                    match *rest.get(idx)? {
545                        byte if byte == quote => {
546                            idx += 1;
547                            if rest.get(idx) == Some(&quote) {
548                                idx += 1; // doubled-quote escape inside the literal
549                            } else {
550                                break;
551                            }
552                        }
553                        _ => idx += 1,
554                    }
555                }
556            }
557            b'[' => {
558                idx += 1;
559                while *rest.get(idx)? != b']' {
560                    idx += 1;
561                }
562                idx += 1;
563            }
564            b'-' if rest.get(idx + 1) == Some(&b'-') => {
565                idx += 2;
566                while idx < rest.len() && rest[idx] != b'\n' {
567                    idx += 1;
568                }
569            }
570            b'/' if rest.get(idx + 1) == Some(&b'*') => {
571                idx += 2;
572                while idx + 1 < rest.len() && !(rest[idx] == b'*' && rest[idx + 1] == b'/') {
573                    idx += 1;
574                }
575                idx = (idx + 2).min(rest.len());
576            }
577            _ => idx += 1,
578        }
579    }
580}
581
582/// Skip one SQLite identifier — unquoted (`[A-Za-z0-9_]+`, as
583/// [`next_sqlite_token`] already recognizes), double-quoted or
584/// backtick-quoted (both honoring the doubled-quote-character escape), or
585/// bracket-quoted (no escape; SQLite's bracket form ends at the first `]`) —
586/// and return the bytes after it. A quoted CTE name may contain any byte,
587/// including `(` and `)`, so the CTE-list walk must skip the identifier
588/// itself rather than tokenizing it the way an unquoted keyword is.
589fn skip_sqlite_identifier(rest: &[u8]) -> Option<&[u8]> {
590    match *rest.first()? {
591        quote @ (b'"' | b'`') => {
592            let mut idx = 1;
593            loop {
594                match *rest.get(idx)? {
595                    byte if byte == quote => {
596                        idx += 1;
597                        if rest.get(idx) == Some(&quote) {
598                            idx += 1; // doubled-quote escape inside the identifier
599                        } else {
600                            break;
601                        }
602                    }
603                    _ => idx += 1,
604                }
605            }
606            Some(&rest[idx..])
607        }
608        b'[' => {
609            let mut idx = 1;
610            while *rest.get(idx)? != b']' {
611                idx += 1;
612            }
613            Some(&rest[idx + 1..])
614        }
615        _ => next_sqlite_token(rest).map(|(_, next)| next),
616    }
617}
618
619/// Walk past the common-table-expression list following `WITH [RECURSIVE]`
620/// and return the bytes starting at the main statement's own head keyword.
621/// Each CTE body is skipped as a balanced parenthesized region
622/// ([`skip_balanced_parens`]), so nested parens, string literals, and
623/// comments inside a CTE body never confuse the walk. Returns `None` if the
624/// CTE list is not well-formed enough to walk past safely — the caller must
625/// fail closed (refuse admission) rather than guess.
626fn skip_common_table_expressions(tail: &[u8]) -> Option<&[u8]> {
627    let mut rest = skip_sqlite_empty_prefix(tail);
628    if let Some((word, next)) = next_sqlite_token(rest) {
629        if word.eq_ignore_ascii_case(b"RECURSIVE") {
630            rest = skip_sqlite_empty_prefix(next);
631        }
632    }
633    loop {
634        // CTE name — unquoted or SQLite-quoted (`"..."`, `` `...` ``, `[...]`).
635        let next = skip_sqlite_identifier(rest)?;
636        rest = skip_sqlite_empty_prefix(next);
637        // Optional column-name list.
638        if rest.first() == Some(&b'(') {
639            rest = skip_sqlite_empty_prefix(skip_balanced_parens(rest)?);
640        }
641        let (as_keyword, next) = next_sqlite_token(rest)?;
642        if !as_keyword.eq_ignore_ascii_case(b"AS") {
643            return None;
644        }
645        rest = skip_sqlite_empty_prefix(next);
646        // Optional `[NOT] MATERIALIZED` hint (SQLite 3.35+).
647        if let Some((word, next)) = next_sqlite_token(rest) {
648            if word.eq_ignore_ascii_case(b"MATERIALIZED") {
649                rest = skip_sqlite_empty_prefix(next);
650            } else if word.eq_ignore_ascii_case(b"NOT") {
651                let (materialized, next) = next_sqlite_token(skip_sqlite_empty_prefix(next))?;
652                if !materialized.eq_ignore_ascii_case(b"MATERIALIZED") {
653                    return None;
654                }
655                rest = skip_sqlite_empty_prefix(next);
656            }
657        }
658        // CTE body.
659        if rest.first() != Some(&b'(') {
660            return None;
661        }
662        rest = skip_sqlite_empty_prefix(skip_balanced_parens(rest)?);
663        if rest.first() == Some(&b',') {
664            rest = skip_sqlite_empty_prefix(&rest[1..]);
665            continue;
666        }
667        return Some(rest);
668    }
669}
670
671/// Reject any raw SQL that is not one of a small allow-listed set of
672/// read-only statement shapes before it ever reaches a pooled reader
673/// connection.
674///
675/// `stmt.readonly()` (SQLite's own classifier, used by
676/// [`statement_is_cancellable_read`] to decide interrupt eligibility) is not
677/// a sufficient admission predicate on its own: by SQLite's own definition it
678/// also returns `true` for `ATTACH`/`DETACH`, `CREATE TEMP TABLE`, and
679/// configuration `PRAGMA`s like `writable_schema`, `busy_timeout`, or
680/// `cache_size` — none of which write the main database file, but all of
681/// which leave connection-local state that persists across pooled checkouts.
682/// Classification here instead runs on the statement head, admitting only
683/// `SELECT`, `WITH ... SELECT`, `VALUES`, `EXPLAIN [QUERY PLAN] <admitted>`,
684/// and the fixed `PRAGMA` allow-lists above. A `WITH` clause is not itself a
685/// read shape — SQLite also allows `WITH ... INSERT/UPDATE/DELETE`, with or
686/// without `RETURNING` — so a `WITH` head is admitted only after walking past
687/// its CTE list ([`skip_common_table_expressions`]) and confirming the main
688/// statement underneath is itself `SELECT` or `VALUES`.
689pub(crate) fn reader_capability_admits(sql: &str) -> Result<(), String> {
690    let rest = skip_sqlite_empty_prefix(sql.as_bytes());
691    let Some((head, tail)) = next_sqlite_token(rest) else {
692        // No head token (empty/comment-only statement): let SQLite's own
693        // prepare step surface that error rather than duplicating it here.
694        return Ok(());
695    };
696    if head.eq_ignore_ascii_case(b"SELECT") || head.eq_ignore_ascii_case(b"VALUES") {
697        return Ok(());
698    }
699    if head.eq_ignore_ascii_case(b"WITH") {
700        let Some(after_ctes) = skip_common_table_expressions(tail) else {
701            return Err(
702                "WITH statement's common-table-expression list could not be parsed; refusing \
703                 to admit it through the reader capability"
704                    .into(),
705            );
706        };
707        return match next_sqlite_token(after_ctes) {
708            Some((main_head, _))
709                if main_head.eq_ignore_ascii_case(b"SELECT")
710                    || main_head.eq_ignore_ascii_case(b"VALUES") =>
711            {
712                Ok(())
713            }
714            other => Err(format!(
715                "WITH ... {:?} is not admitted through the reader capability; only a \
716                 read-only SELECT/VALUES body after the CTE list may run against a pooled \
717                 reader connection",
718                other.map_or_else(
719                    || "<none>".to_string(),
720                    |(main_head, _)| String::from_utf8_lossy(main_head).into_owned()
721                )
722            )),
723        };
724    }
725    if head.eq_ignore_ascii_case(b"EXPLAIN") {
726        let mut rest = skip_sqlite_empty_prefix(tail);
727        if let Some((query, next)) = next_sqlite_token(rest) {
728            if query.eq_ignore_ascii_case(b"QUERY") {
729                let after_query = skip_sqlite_empty_prefix(next);
730                match next_sqlite_token(after_query) {
731                    Some((plan, next2)) if plan.eq_ignore_ascii_case(b"PLAN") => {
732                        rest = skip_sqlite_empty_prefix(next2);
733                    }
734                    _ => {
735                        return Err(
736                            "EXPLAIN QUERY must be followed by PLAN through the reader capability"
737                                .into(),
738                        );
739                    }
740                }
741            }
742        }
743        return reader_capability_admits(&String::from_utf8_lossy(rest));
744    }
745    if head.eq_ignore_ascii_case(b"PRAGMA") {
746        return reader_capability_admits_pragma(tail);
747    }
748    Err(format!(
749        "statement head {:?} is not admitted through the reader capability; only \
750         SELECT/WITH/VALUES/EXPLAIN and an allow-listed set of read-only PRAGMA forms \
751         may run against a pooled reader connection",
752        String::from_utf8_lossy(head)
753    ))
754}
755
756fn reader_capability_admits_pragma(tail: &[u8]) -> Result<(), String> {
757    let rest = skip_sqlite_empty_prefix(tail);
758    let Some((mut name, mut after_name)) = next_sqlite_token(rest) else {
759        return Err("PRAGMA with no name is not admitted through the reader capability".into());
760    };
761    // `schema.pragma_name` — the schema qualifier changes which attached
762    // database the pragma targets, not which pragma runs.
763    if after_name.first() == Some(&b'.') {
764        let (qualified_name, qualified_after) =
765            next_sqlite_token(&after_name[1..]).ok_or_else(|| {
766                "PRAGMA schema-qualifier with no pragma name is not admitted through the \
767                 reader capability"
768                    .to_string()
769            })?;
770        name = qualified_name;
771        after_name = qualified_after;
772    }
773    let after = skip_sqlite_empty_prefix(after_name);
774    let is_structural = READER_STRUCTURAL_PRAGMAS
775        .iter()
776        .any(|allowed| name.eq_ignore_ascii_case(allowed.as_bytes()));
777    let is_setting = READER_SETTING_PRAGMAS
778        .iter()
779        .any(|allowed| name.eq_ignore_ascii_case(allowed.as_bytes()));
780    if !is_structural && !is_setting {
781        return Err(format!(
782            "PRAGMA {:?} is not admitted through the reader capability",
783            String::from_utf8_lossy(name)
784        ));
785    }
786    if after.first() == Some(&b'=') {
787        return Err(format!(
788            "PRAGMA {:?} may not be assigned through the reader capability",
789            String::from_utf8_lossy(name)
790        ));
791    }
792    if after.first() == Some(&b'(') && !is_structural {
793        return Err(format!(
794            "PRAGMA {:?} may not carry an argument through the reader capability",
795            String::from_utf8_lossy(name)
796        ));
797    }
798    Ok(())
799}
800
801/// Refuse a statement bound for the reader capability that is neither
802/// admitted transaction control (handled by the caller's cached
803/// read-transaction state machine) nor an admitted read shape
804/// ([`reader_capability_admits`]).
805fn admit_reader_capability_sql(
806    statement: &SqlStatement,
807    transaction_control: Option<CachedReadTransactionControl>,
808    operation: &'static str,
809) -> khive_storage::types::StorageResult<()> {
810    // Only the admitted single-level deferred-read span (`BEGIN`/`BEGIN
811    // DEFERRED` and its `COMMIT`/`END`/`ROLLBACK` counterpart) bypasses the
812    // read-shape check below; `Unsupported` forms (`BEGIN IMMEDIATE`,
813    // `SAVEPOINT`, `ROLLBACK TO ...`) fall through and are refused there like
814    // any other non-admitted statement head.
815    if matches!(
816        transaction_control,
817        Some(CachedReadTransactionControl::BeginDeferred)
818            | Some(CachedReadTransactionControl::Finish(_))
819    ) {
820        return Ok(());
821    }
822    reader_capability_admits(&statement.sql).map_err(|message| StorageError::InvalidInput {
823        capability: StorageCapability::Sql,
824        operation: operation.into(),
825        message,
826    })
827}
828
829fn execute_query_interruptibly(
830    scope: &crate::read_cancellation::InterruptibleReadScope,
831    conn: &rusqlite::Connection,
832    statement: &SqlStatement,
833    operation: &'static str,
834    rollback_interrupted_transaction: bool,
835    interruptible: bool,
836) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
837    let stmt = prepare_bound_statement(conn, statement)
838        .map_err(|error| map_rusqlite_err(error, operation))?;
839    if interruptible && statement_is_cancellable_read(&stmt, &statement.sql) {
840        scope.run_with_interrupted_cleanup(
841            conn,
842            move || {
843                execute_prepared_query(stmt).map_err(|error| map_rusqlite_err(error, operation))
844            },
845            || {
846                rollback_interrupted_read_transaction(
847                    conn,
848                    operation,
849                    rollback_interrupted_transaction,
850                )
851            },
852        )
853    } else {
854        scope.mark_write_committed()?;
855        execute_prepared_query(stmt).map_err(|error| map_rusqlite_err(error, operation))
856    }
857}
858
859fn execute_query_row_interruptibly(
860    scope: &crate::read_cancellation::InterruptibleReadScope,
861    conn: &rusqlite::Connection,
862    statement: &SqlStatement,
863    operation: &'static str,
864    rollback_interrupted_transaction: bool,
865    interruptible: bool,
866) -> khive_storage::types::StorageResult<Option<SqlRow>> {
867    let stmt = prepare_bound_statement(conn, statement)
868        .map_err(|error| map_rusqlite_err(error, operation))?;
869    if interruptible && statement_is_cancellable_read(&stmt, &statement.sql) {
870        scope.run_with_interrupted_cleanup(
871            conn,
872            move || {
873                execute_prepared_query_row(stmt).map_err(|error| map_rusqlite_err(error, operation))
874            },
875            || {
876                rollback_interrupted_read_transaction(
877                    conn,
878                    operation,
879                    rollback_interrupted_transaction,
880                )
881            },
882        )
883    } else {
884        scope.mark_write_committed()?;
885        execute_prepared_query_row(stmt).map_err(|error| map_rusqlite_err(error, operation))
886    }
887}
888
889fn execute_query_page_interruptibly(
890    scope: &crate::read_cancellation::InterruptibleReadScope,
891    conn: &rusqlite::Connection,
892    statement: &SqlStatement,
893    page: &PageRequest,
894    operation: &'static str,
895    rollback_interrupted_transaction: bool,
896    interruptible: bool,
897) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
898    let stmt = prepare_bound_statement(conn, statement)
899        .map_err(|error| map_rusqlite_err(error, operation))?;
900    if interruptible && statement_is_cancellable_read(&stmt, &statement.sql) {
901        scope.run_with_interrupted_cleanup(
902            conn,
903            move || {
904                execute_prepared_query_page(stmt, page)
905                    .map_err(|error| map_rusqlite_err(error, operation))
906            },
907            || {
908                rollback_interrupted_read_transaction(
909                    conn,
910                    operation,
911                    rollback_interrupted_transaction,
912                )
913            },
914        )
915    } else {
916        scope.mark_write_committed()?;
917        execute_prepared_query_page(stmt, page).map_err(|error| map_rusqlite_err(error, operation))
918    }
919}
920
921fn rollback_interrupted_read_transaction(
922    conn: &rusqlite::Connection,
923    operation: &'static str,
924    enabled: bool,
925) -> khive_storage::types::StorageResult<()> {
926    if !enabled || conn.is_autocommit() {
927        return Ok(());
928    }
929    conn.execute_batch("ROLLBACK")
930        .map_err(|error| map_rusqlite_err(error, operation))?;
931    if conn.is_autocommit() {
932        Ok(())
933    } else {
934        Err(StorageError::Transaction {
935            operation: operation.into(),
936            message: "interrupted read transaction rollback did not restore autocommit".into(),
937        })
938    }
939}
940
941/// Map a rusqlite error to `StorageError`.
942fn map_rusqlite_err(e: rusqlite::Error, op: &'static str) -> StorageError {
943    StorageError::driver(StorageCapability::Sql, op, e)
944}
945
946/// How an elapsed handle-slot deadline is classified. ADR-005 pins the closed
947/// raw-SQL standalone exception (`sql_bridge.reader_open`) and reads on a
948/// standalone writer to `StorageError::Timeout`; ordinary reader-pool
949/// saturation reports the typed `AdmissionTimeout` through
950/// `ConnectionPool::resolve_reader_checkout`.
951///
952/// Classification is decided by this enum and never by the operation label.
953/// The label names the caller's own read (`query_row`, `writer.query_all`, and
954/// so on) so that a recorded maximum reader hold can be attributed to the read
955/// that caused it.
956#[derive(Clone, Copy)]
957enum SlotTimeoutClass {
958    Admission,
959    ReaderContract,
960}
961
962async fn acquire_reader_handle_slot(
963    pool: &ConnectionPool,
964    operation: &'static str,
965    class: SlotTimeoutClass,
966) -> Result<OwnedSemaphorePermit, StorageError> {
967    let result = acquire_handle_slot(
968        pool.sql_bridge_reader_slots(),
969        pool.config().checkout_timeout,
970        operation,
971        class,
972    )
973    .await;
974    if matches!(
975        &result,
976        Err(StorageError::Timeout { .. } | StorageError::AdmissionTimeout { .. })
977    ) {
978        pool.record_reader_admission_timeout();
979    }
980    result
981}
982
983/// Serialize whole units on the one shared in-memory connection. Event-store
984/// transactions participate in the same budget as raw SQL atomic units.
985/// The permit must outlive the unit's final COMMIT or ROLLBACK.
986pub(crate) async fn acquire_in_memory_write_unit(
987    pool: &ConnectionPool,
988    operation: &'static str,
989) -> Result<OwnedSemaphorePermit, StorageError> {
990    acquire_handle_slot(
991        pool.sql_bridge_writer_slots(),
992        pool.config().checkout_timeout,
993        operation,
994        SlotTimeoutClass::Admission,
995    )
996    .await
997}
998
999async fn acquire_handle_slot(
1000    slots: Arc<Semaphore>,
1001    timeout: std::time::Duration,
1002    operation: &'static str,
1003    class: SlotTimeoutClass,
1004) -> Result<OwnedSemaphorePermit, StorageError> {
1005    tokio::time::timeout(timeout, slots.acquire_owned())
1006        .await
1007        .map_err(|_| match class {
1008            SlotTimeoutClass::Admission => StorageError::AdmissionTimeout {
1009                operation: operation.into(),
1010                timeout_ms: u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX),
1011                pool_identity: None,
1012            },
1013            SlotTimeoutClass::ReaderContract => StorageError::Timeout {
1014                operation: operation.into(),
1015            },
1016        })?
1017        .map_err(|error| StorageError::Pool {
1018            operation: operation.into(),
1019            message: error.to_string(),
1020        })
1021}
1022
1023// =============================================================================
1024// Standalone connection readers/writers (file-backed databases)
1025// =============================================================================
1026
1027fn open_standalone_reader(pool: &ConnectionPool) -> Result<rusqlite::Connection, StorageError> {
1028    pool.open_standalone_reader(StandaloneReaderPurpose::ExplicitSqlReadTransaction)
1029        .map_err(|error| StorageError::driver(StorageCapability::Sql, "open_reader", error))
1030}
1031
1032#[cfg(test)]
1033fn open_standalone_writer(pool: &ConnectionPool) -> Result<rusqlite::Connection, StorageError> {
1034    let conn = pool
1035        .open_standalone_writer()
1036        .map_err(|e| e.into_storage_error(StorageCapability::Sql, "open_writer"))?;
1037    configure_standalone_writer(pool, conn)
1038}
1039
1040fn open_admitted_standalone_writer(
1041    pool: &ConnectionPool,
1042) -> Result<rusqlite::Connection, StorageError> {
1043    let conn = pool
1044        .open_standalone_writer_for_admitted_operation()
1045        .map_err(|e| e.into_storage_error(StorageCapability::Sql, "open_writer"))?;
1046    configure_standalone_writer(pool, conn)
1047}
1048
1049fn configure_standalone_writer(
1050    pool: &ConnectionPool,
1051    conn: rusqlite::Connection,
1052) -> Result<rusqlite::Connection, StorageError> {
1053    let config = pool.config();
1054    conn.busy_timeout(config.busy_timeout)
1055        .map_err(|e| map_rusqlite_err(e, "open_writer"))?;
1056    conn.pragma_update(None, "cache_size", "-65536")
1057        .map_err(|e| map_rusqlite_err(e, "open_writer"))?;
1058    conn.pragma_update(None, "mmap_size", "1073741824")
1059        .map_err(|e| map_rusqlite_err(e, "open_writer"))?;
1060
1061    Ok(conn)
1062}
1063
1064/// Lift a standalone open onto the blocking thread pool while carrying the
1065/// connection-cap permit into the same closure.
1066///
1067/// If the awaiting future is cancelled, the detached blocking closure owns
1068/// both the connection result and the permit until it finishes, so the cap
1069/// cannot be released while an open is still running.
1070async fn open_standalone_on_blocking<F>(
1071    pool: Arc<ConnectionPool>,
1072    slot: OwnedSemaphorePermit,
1073    operation: &'static str,
1074    open: F,
1075) -> khive_storage::types::StorageResult<(rusqlite::Connection, OwnedSemaphorePermit)>
1076where
1077    F: FnOnce(&ConnectionPool) -> Result<rusqlite::Connection, StorageError> + Send + 'static,
1078{
1079    tokio::task::spawn_blocking(move || open(&pool).map(|conn| (conn, slot)))
1080        .await
1081        .map_err(|e| StorageError::driver(StorageCapability::Sql, operation, e))?
1082}
1083
1084/// [`open_standalone_reader`] lifted onto the blocking thread pool.
1085///
1086/// Opening a SQLite connection is filesystem I/O (open the file, read the
1087/// database header) followed by pragmas executed through SQLite. No
1088/// database lock is acquired at open itself — locks are taken on the first
1089/// statement — but filesystem latency is unbounded, and this module already
1090/// runs every other rusqlite call under `spawn_blocking`, so the open gets
1091/// the same treatment instead of blocking an async worker thread. The reader
1092/// permit is supplied to the helper and returned with the connection.
1093async fn open_standalone_reader_on_blocking(
1094    pool: Arc<ConnectionPool>,
1095    slot: OwnedSemaphorePermit,
1096) -> khive_storage::types::StorageResult<(rusqlite::Connection, OwnedSemaphorePermit)> {
1097    open_standalone_on_blocking(pool, slot, "open_reader", open_standalone_reader).await
1098}
1099
1100/// Open the writer handle without a premature capacity sample. Its later
1101/// operation owns admission; a cold handle must allow a recovery checkpoint.
1102/// See [`open_standalone_reader_on_blocking`] for the blocking rationale.
1103async fn open_standalone_writer_on_blocking(
1104    pool: Arc<ConnectionPool>,
1105    slot: OwnedSemaphorePermit,
1106) -> khive_storage::types::StorageResult<(rusqlite::Connection, OwnedSemaphorePermit)> {
1107    open_standalone_on_blocking(pool, slot, "open_writer", open_admitted_standalone_writer).await
1108}
1109
1110// =============================================================================
1111// File-backed: pooled SqliteReader with an explicit transaction exception
1112// =============================================================================
1113
1114const CACHED_READ_TRANSACTION_LABEL: &str = "sql_bridge_cached_read_transaction";
1115
1116/// Admission and observability guards for one explicit cached-reader
1117/// transaction. Both guards are installed only after SQLite accepts `BEGIN`
1118/// and are retained together until SQLite reports autocommit again or the
1119/// owning connection is closed.
1120///
1121/// Field order is deliberate: after [`StandaloneHandle::conn`] closes, the
1122/// reader permit is returned before the registry evidence disappears. There
1123/// is therefore no interval in which SQLite can still own the snapshot while
1124/// the transaction is absent from `tx_registry`.
1125struct CachedReadTransaction {
1126    _slot: OwnedSemaphorePermit,
1127    _tx_handle: khive_storage::tx_registry::TxHandle,
1128    /// When this explicit `BEGIN` was admitted. Read on every subsequent
1129    /// reuse of the owning cached-reader handle (#1846): a transaction whose
1130    /// age has crossed `read_tx_max_age` is rolled back instead of being
1131    /// extended by another call, bounding how long any one reader can pin
1132    /// the WAL snapshot regardless of how many further requests it makes.
1133    opened_at: Instant,
1134}
1135
1136struct StandaloneHandle {
1137    conn: rusqlite::Connection,
1138    /// Present only for a standalone read-write handle, whose one-permit
1139    /// connection budget remains handle-scoped. Read-only connections exist
1140    /// only for explicit deferred transactions and retain reader admission in
1141    /// `read_transaction_slot` until terminal control.
1142    _retained_slot: Option<OwnedSemaphorePermit>,
1143    /// Present only while a cached read-only connection owns one explicit
1144    /// multi-call read transaction. Field order is load-bearing: Rust drops
1145    /// `conn` before these guards, so cancellation or handle drop closes the
1146    /// SQLite transaction before returning reader admission or deregistering
1147    /// the transaction span.
1148    read_transaction_slot: Option<CachedReadTransaction>,
1149}
1150
1151impl StandaloneHandle {
1152    /// Whether this is an idle-cacheable read-only connection rather than a
1153    /// writer connection covered by the handle-scoped writer permit.
1154    fn is_cached_reader(&self) -> bool {
1155        self._retained_slot.is_none()
1156    }
1157
1158    fn has_read_transaction(&self) -> bool {
1159        self.read_transaction_slot.is_some()
1160    }
1161}
1162
1163struct SqliteReader {
1164    /// Present only for the ADR-005/ADR-091 explicitly admitted multi-call
1165    /// deferred transaction. Ordinary reads never populate this field and
1166    /// route through `ConnectionPool::reader` for each operation.
1167    handle: Option<StandaloneHandle>,
1168    pool: Arc<ConnectionPool>,
1169    /// Fail-loud compatibility state after a transaction connection could not
1170    /// be restored safely. `None` normally means "pooled ordinary route";
1171    /// this bit distinguishes that from "exceptional connection consumed".
1172    poisoned: bool,
1173}
1174
1175async fn open_explicit_read_transaction_handle(
1176    pool: Arc<ConnectionPool>,
1177) -> khive_storage::types::StorageResult<StandaloneHandle> {
1178    let open_slot = crate::await_request_read_phase(
1179        "sql_bridge.reader_open",
1180        acquire_reader_handle_slot(
1181            &pool,
1182            "sql_bridge.reader_open",
1183            SlotTimeoutClass::ReaderContract,
1184        ),
1185    )
1186    .await??;
1187    let (conn, open_slot) = crate::await_request_read_phase(
1188        "sql_bridge.reader_open",
1189        open_standalone_reader_on_blocking(pool, open_slot),
1190    )
1191    .await??;
1192    drop(open_slot);
1193    Ok(StandaloneHandle {
1194        conn,
1195        _retained_slot: None,
1196        read_transaction_slot: None,
1197    })
1198}
1199
1200impl SqliteReader {
1201    /// Select the standalone path only for a live or newly requested explicit
1202    /// deferred transaction. `false` means the caller must execute this
1203    /// operation through the bounded reader pool.
1204    async fn use_explicit_transaction_handle(
1205        &mut self,
1206        transaction_control: Option<CachedReadTransactionControl>,
1207        operation: &'static str,
1208    ) -> khive_storage::types::StorageResult<bool> {
1209        if self.poisoned {
1210            return Err(StorageError::Pool {
1211                operation: operation.into(),
1212                message: "connection already consumed".into(),
1213            });
1214        }
1215        if self.handle.is_some() {
1216            return Ok(true);
1217        }
1218        match transaction_control {
1219            None => Ok(false),
1220            Some(CachedReadTransactionControl::BeginDeferred) => {
1221                self.handle =
1222                    Some(open_explicit_read_transaction_handle(Arc::clone(&self.pool)).await?);
1223                Ok(true)
1224            }
1225            Some(CachedReadTransactionControl::Finish(keyword))
1226            | Some(CachedReadTransactionControl::Unsupported(keyword)) => {
1227                Err(StorageError::InvalidInput {
1228                    capability: StorageCapability::Sql,
1229                    operation: operation.into(),
1230                    message: format!(
1231                        "cached read-only handle has no admitted transaction for transaction \
1232                         control ({keyword})"
1233                    ),
1234                })
1235            }
1236        }
1237    }
1238
1239    /// A standalone connection exists only for one explicit transaction.
1240    /// Close it immediately after a failed BEGIN, COMMIT/ROLLBACK, age
1241    /// eviction, or cancellation cleanup restores autocommit. A later
1242    /// ordinary query must return to pooled routing instead of silently
1243    /// retaining a standalone cache.
1244    fn close_inactive_transaction_handle(&mut self) {
1245        if self
1246            .handle
1247            .as_ref()
1248            .is_some_and(|handle| handle.is_cached_reader() && !handle.has_read_transaction())
1249        {
1250            drop(self.handle.take());
1251        }
1252    }
1253}
1254
1255/// Run one file-backed read while coupling its connection state to the active
1256/// reader permit.
1257///
1258/// This path serves only the explicit deferred read-transaction exception and
1259/// reader-supertrait calls on a standalone read-write handle. Ordinary
1260/// read-only traffic uses [`run_pool_reader_query`]. A successful top-level
1261/// deferred `BEGIN` moves its operation permit into the handle; subsequent
1262/// reads reuse it until `COMMIT`/`END`/`ROLLBACK` restores autocommit. The
1263/// connection is declared before that retained permit, so dropping or
1264/// cancelling the handle closes SQLite first and releases admission second.
1265/// A standalone writer's reader calls still take active-reader admission; once
1266/// its connection is outside autocommit, the acquisition and SELECT are
1267/// completion-preserving rather than request-cancellable.
1268async fn execute_standalone_read<R, F>(
1269    handle: &mut Option<StandaloneHandle>,
1270    pool: Arc<ConnectionPool>,
1271    operation: &'static str,
1272    transaction_control: Option<CachedReadTransactionControl>,
1273    read: F,
1274) -> khive_storage::types::StorageResult<R>
1275where
1276    R: Send + 'static,
1277    F: FnOnce(
1278            &crate::read_cancellation::InterruptibleReadScope,
1279            &rusqlite::Connection,
1280            bool,
1281            bool,
1282        ) -> khive_storage::types::StorageResult<R>
1283        + Send
1284        + 'static,
1285{
1286    if handle.is_none() {
1287        return Err(StorageError::Pool {
1288            operation: operation.into(),
1289            message: "connection already consumed".into(),
1290        });
1291    }
1292    let active_read_transaction = handle
1293        .as_ref()
1294        .is_some_and(|handle| handle.is_cached_reader() && handle.has_read_transaction());
1295    let completion_preserving_writer_transaction = handle
1296        .as_ref()
1297        .is_some_and(|handle| !handle.is_cached_reader() && !handle.conn.is_autocommit());
1298    let mut operation_slot = if active_read_transaction {
1299        None
1300    } else if completion_preserving_writer_transaction {
1301        // A read inside an admitted write transaction still counts against the
1302        // reader budget, but request cancellation cannot skip that admission
1303        // and strand the transaction between statements.
1304        Some(acquire_reader_handle_slot(&pool, operation, SlotTimeoutClass::ReaderContract).await?)
1305    } else {
1306        Some(
1307            crate::await_request_read_phase(
1308                operation,
1309                acquire_reader_handle_slot(&pool, operation, SlotTimeoutClass::ReaderContract),
1310            )
1311            .await??,
1312        )
1313    };
1314    let Some(owned_handle) = handle.take() else {
1315        return Err(StorageError::Pool {
1316            operation: operation.into(),
1317            message: "connection already consumed".into(),
1318        });
1319    };
1320    let origin = pool.origin();
1321    let read_tx_max_age = pool.config().read_tx_max_age;
1322    let (owned_handle, result) = crate::read_cancellation::run_interruptible_read(
1323        StorageCapability::Sql,
1324        operation,
1325        move |scope| {
1326            let mut owned_handle = owned_handle;
1327            let cached_reader = owned_handle.is_cached_reader();
1328            let entered_with_transaction = owned_handle.has_read_transaction();
1329            let entered_autocommit = owned_handle.conn.is_autocommit();
1330            let mut restore_handle = true;
1331            let mut result = if cached_reader && entered_with_transaction && entered_autocommit {
1332                // A connection cannot pin a snapshot in autocommit. Repair the
1333                // admission state before returning the invariant failure.
1334                drop(owned_handle.read_transaction_slot.take());
1335                Err(StorageError::InvalidInput {
1336                    capability: StorageCapability::Sql,
1337                    operation: operation.into(),
1338                    message: "cached read-only handle retained transaction admission after SQLite \
1339                          had already returned to autocommit; the stale permit was released"
1340                        .into(),
1341                })
1342            } else if cached_reader && !entered_with_transaction && !entered_autocommit {
1343                Err(StorageError::InvalidInput {
1344                    capability: StorageCapability::Sql,
1345                    operation: operation.into(),
1346                    message: "cached read-only handle entered the operation outside autocommit; \
1347                          its transaction was rolled back before releasing the reader permit"
1348                        .into(),
1349                })
1350            } else if cached_reader
1351                && entered_with_transaction
1352                && owned_handle
1353                    .read_transaction_slot
1354                    .as_ref()
1355                    .is_some_and(|tx| tx.opened_at.elapsed() >= read_tx_max_age)
1356            {
1357                // #1846: this handle's admitted read transaction has pinned a
1358                // WAL snapshot for at least `read_tx_max_age` — reject the
1359                // continuation and roll it back instead of extending the pin
1360                // for another call, regardless of what the caller asked for.
1361                crate::checkpoint::note_read_tx_max_age_eviction();
1362                match owned_handle.conn.execute_batch("ROLLBACK") {
1363                    Ok(()) if owned_handle.conn.is_autocommit() => {
1364                        drop(owned_handle.read_transaction_slot.take());
1365                        Err(StorageError::ReadTransactionAgeEvicted {
1366                            operation: operation.into(),
1367                            max_age_secs: read_tx_max_age.as_secs(),
1368                        })
1369                    }
1370                    Ok(()) => {
1371                        restore_handle = false;
1372                        Err(StorageError::ReadTransactionAgeEvictionCleanupFailed {
1373                            operation: operation.into(),
1374                            max_age_secs: read_tx_max_age.as_secs(),
1375                            message: "rollback did not restore autocommit".into(),
1376                        })
1377                    }
1378                    Err(error) => {
1379                        restore_handle = false;
1380                        Err(StorageError::ReadTransactionAgeEvictionCleanupFailed {
1381                            operation: operation.into(),
1382                            max_age_secs: read_tx_max_age.as_secs(),
1383                            message: format!("rollback failed: {error}"),
1384                        })
1385                    }
1386                }
1387            } else if cached_reader && entered_with_transaction {
1388                match transaction_control {
1389                    None | Some(CachedReadTransactionControl::Finish(_)) => {
1390                        read(scope, &owned_handle.conn, true, true)
1391                    }
1392                    Some(CachedReadTransactionControl::BeginDeferred) => {
1393                        Err(StorageError::InvalidInput {
1394                            capability: StorageCapability::Sql,
1395                            operation: operation.into(),
1396                            message: "cached read-only handle already owns an admitted read \
1397                                  transaction; nested BEGIN is not supported"
1398                                .into(),
1399                        })
1400                    }
1401                    Some(CachedReadTransactionControl::Unsupported(keyword)) => {
1402                        Err(StorageError::InvalidInput {
1403                            capability: StorageCapability::Sql,
1404                            operation: operation.into(),
1405                            message: format!(
1406                                "cached read-only transaction does not support nested or \
1407                             write-locking transaction control ({keyword})"
1408                            ),
1409                        })
1410                    }
1411                }
1412            } else if cached_reader {
1413                match transaction_control {
1414                    None | Some(CachedReadTransactionControl::BeginDeferred) => {
1415                        read(scope, &owned_handle.conn, false, true)
1416                    }
1417                    Some(CachedReadTransactionControl::Finish(keyword))
1418                    | Some(CachedReadTransactionControl::Unsupported(keyword)) => {
1419                        Err(StorageError::InvalidInput {
1420                            capability: StorageCapability::Sql,
1421                            operation: operation.into(),
1422                            message: format!(
1423                                "cached read-only handle has no admitted transaction for \
1424                             transaction control ({keyword})"
1425                            ),
1426                        })
1427                    }
1428                }
1429            } else {
1430                read(scope, &owned_handle.conn, false, entered_autocommit)
1431            };
1432
1433            if scope.cleanup_failed() {
1434                // A connection-global callback that could not be removed may
1435                // fire for an unrelated future borrower. Closing this handle
1436                // is the only safe recovery; its transaction, if any, ends
1437                // before reader admission is released below.
1438                restore_handle = false;
1439            }
1440
1441            // An interrupted explicit read transaction must never be restored to
1442            // the cached handle: it may still own a WAL snapshot and SQLite's
1443            // interrupted flag applies to the transaction as a whole. Roll back
1444            // before releasing its retained admission; if rollback cannot prove
1445            // autocommit, discard the connection.
1446            if cached_reader
1447                && matches!(result, Err(StorageError::Timeout { .. }))
1448                && !owned_handle.conn.is_autocommit()
1449            {
1450                match owned_handle.conn.execute_batch("ROLLBACK") {
1451                    Ok(()) if owned_handle.conn.is_autocommit() => {
1452                        drop(owned_handle.read_transaction_slot.take());
1453                    }
1454                    Ok(()) => {
1455                        restore_handle = false;
1456                        result = Err(StorageError::Transaction {
1457                            operation: operation.into(),
1458                            message:
1459                                "interrupted read transaction rollback did not restore autocommit; \
1460                                  the connection was discarded"
1461                                    .into(),
1462                        });
1463                    }
1464                    Err(error) => {
1465                        restore_handle = false;
1466                        result = Err(StorageError::Transaction {
1467                            operation: operation.into(),
1468                            message: format!(
1469                                "failed to roll back interrupted read transaction ({error}); \
1470                             the connection was discarded"
1471                            ),
1472                        });
1473                    }
1474                }
1475            }
1476
1477            if cached_reader && entered_with_transaction {
1478                if owned_handle.conn.is_autocommit() {
1479                    // SQLite has ended the snapshot; release only after observing
1480                    // that terminal state. This also fails closed if an ordinary
1481                    // statement unexpectedly ended the transaction.
1482                    drop(owned_handle.read_transaction_slot.take());
1483                    if result.is_ok()
1484                        && !matches!(
1485                            transaction_control,
1486                            Some(CachedReadTransactionControl::Finish(_))
1487                        )
1488                    {
1489                        result = Err(StorageError::InvalidInput {
1490                            capability: StorageCapability::Sql,
1491                            operation: operation.into(),
1492                            message: "cached read-only operation unexpectedly ended its admitted \
1493                                  transaction; reader admission was released after autocommit"
1494                                .into(),
1495                        });
1496                    }
1497                } else if result.is_ok()
1498                    && matches!(
1499                        transaction_control,
1500                        Some(CachedReadTransactionControl::Finish(_))
1501                    )
1502                {
1503                    result = Err(StorageError::InvalidInput {
1504                        capability: StorageCapability::Sql,
1505                        operation: operation.into(),
1506                        message: "transaction-ending control completed but the cached reader \
1507                              remained outside autocommit; its reader permit remains retained"
1508                            .into(),
1509                    });
1510                }
1511            } else if cached_reader
1512                && entered_autocommit
1513                && matches!(
1514                    transaction_control,
1515                    Some(CachedReadTransactionControl::BeginDeferred)
1516                )
1517                && result.is_ok()
1518            {
1519                if owned_handle.conn.is_autocommit() {
1520                    result = Err(StorageError::InvalidInput {
1521                        capability: StorageCapability::Sql,
1522                        operation: operation.into(),
1523                        message: "deferred BEGIN completed without opening a read transaction"
1524                            .into(),
1525                    });
1526                } else {
1527                    match operation_slot.take() {
1528                        Some(slot) => {
1529                            let tx_handle = khive_storage::tx_registry::register_scoped(
1530                                Some(CACHED_READ_TRANSACTION_LABEL.to_string()),
1531                                origin.clone(),
1532                            );
1533                            owned_handle.read_transaction_slot = Some(CachedReadTransaction {
1534                                _slot: slot,
1535                                _tx_handle: tx_handle,
1536                                opened_at: Instant::now(),
1537                            });
1538                        }
1539                        None => {
1540                            result = Err(StorageError::Pool {
1541                                operation: operation.into(),
1542                                message: "successful cached-reader BEGIN had no operation permit; \
1543                                      its transaction was rolled back before returning"
1544                                    .into(),
1545                            });
1546                        }
1547                    }
1548                }
1549            }
1550
1551            // Any non-autocommit state without its retained admission is stale or
1552            // was opened by a statement the transaction classifier did not admit.
1553            if cached_reader
1554                && owned_handle.read_transaction_slot.is_none()
1555                && !owned_handle.conn.is_autocommit()
1556            {
1557                match owned_handle.conn.execute_batch("ROLLBACK") {
1558                    Ok(()) if owned_handle.conn.is_autocommit() => {
1559                        if result.is_ok() {
1560                            result = Err(StorageError::InvalidInput {
1561                                capability: StorageCapability::Sql,
1562                                operation: operation.into(),
1563                                message: "cached read-only operation left the connection outside \
1564                                      autocommit; its transaction was rolled back before \
1565                                      releasing the reader permit"
1566                                    .into(),
1567                            });
1568                        }
1569                    }
1570                    Ok(()) => {
1571                        restore_handle = false;
1572                        result = Err(StorageError::Transaction {
1573                            operation: operation.into(),
1574                            message: "ROLLBACK completed but the cached reader remained outside \
1575                                  autocommit; the connection was discarded before releasing \
1576                                  the reader permit"
1577                                .into(),
1578                        });
1579                    }
1580                    Err(error) => {
1581                        restore_handle = false;
1582                        result = Err(StorageError::Transaction {
1583                            operation: operation.into(),
1584                            message: format!(
1585                            "failed to roll back a cached reader outside autocommit ({error}); \
1586                             the connection was discarded before releasing the reader permit"
1587                        ),
1588                        });
1589                    }
1590                }
1591            }
1592
1593            let owned_handle = if restore_handle {
1594                Some(owned_handle)
1595            } else {
1596                // Closing the poisoned connection ends any remaining transaction.
1597                // This must precede the active-reader permit release below.
1598                drop(owned_handle);
1599                None
1600            };
1601            // For ordinary reads and rejected controls this is the operation
1602            // permit. A successful BEGIN moved it into `owned_handle`; poisoned
1603            // handles were closed above before this remaining permit is released.
1604            drop(operation_slot);
1605            Ok((owned_handle, result))
1606        },
1607    )
1608    .await?;
1609    *handle = owned_handle;
1610    result
1611}
1612
1613#[async_trait]
1614impl khive_storage::SqlReader for SqliteReader {
1615    async fn query_row(
1616        &mut self,
1617        statement: SqlStatement,
1618    ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
1619        let transaction_control = cached_read_transaction_control(&statement.sql);
1620        admit_reader_capability_sql(&statement, transaction_control, "query_row")?;
1621        if !self
1622            .use_explicit_transaction_handle(transaction_control, "query_row")
1623            .await?
1624        {
1625            return run_pool_reader_query(
1626                Arc::clone(&self.pool),
1627                "query_row",
1628                move |scope, conn| {
1629                    execute_query_row_interruptibly(
1630                        scope,
1631                        conn,
1632                        &statement,
1633                        "query_row",
1634                        false,
1635                        true,
1636                    )
1637                },
1638            )
1639            .await;
1640        }
1641        let result = execute_standalone_read(
1642            &mut self.handle,
1643            Arc::clone(&self.pool),
1644            "query_row",
1645            transaction_control,
1646            move |scope, conn, rollback, interruptible| {
1647                execute_query_row_interruptibly(
1648                    scope,
1649                    conn,
1650                    &statement,
1651                    "query_row",
1652                    rollback,
1653                    interruptible,
1654                )
1655            },
1656        )
1657        .await;
1658        if self.handle.is_none() {
1659            self.poisoned = true;
1660        }
1661        self.close_inactive_transaction_handle();
1662        result
1663    }
1664
1665    async fn query_all(
1666        &mut self,
1667        statement: SqlStatement,
1668    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1669        let transaction_control = cached_read_transaction_control(&statement.sql);
1670        admit_reader_capability_sql(&statement, transaction_control, "query_all")?;
1671        if !self
1672            .use_explicit_transaction_handle(transaction_control, "query_all")
1673            .await?
1674        {
1675            return run_pool_reader_query(
1676                Arc::clone(&self.pool),
1677                "query_all",
1678                move |scope, conn| {
1679                    execute_query_interruptibly(scope, conn, &statement, "query_all", false, true)
1680                },
1681            )
1682            .await;
1683        }
1684        let result = execute_standalone_read(
1685            &mut self.handle,
1686            Arc::clone(&self.pool),
1687            "query_all",
1688            transaction_control,
1689            move |scope, conn, rollback, interruptible| {
1690                execute_query_interruptibly(
1691                    scope,
1692                    conn,
1693                    &statement,
1694                    "query_all",
1695                    rollback,
1696                    interruptible,
1697                )
1698            },
1699        )
1700        .await;
1701        if self.handle.is_none() {
1702            self.poisoned = true;
1703        }
1704        self.close_inactive_transaction_handle();
1705        result
1706    }
1707
1708    async fn query_page(
1709        &mut self,
1710        statement: SqlStatement,
1711        page: PageRequest,
1712    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1713        let transaction_control = cached_read_transaction_control(&statement.sql);
1714        admit_reader_capability_sql(&statement, transaction_control, "query_page")?;
1715        if !self
1716            .use_explicit_transaction_handle(transaction_control, "query_page")
1717            .await?
1718        {
1719            return run_pool_reader_query(
1720                Arc::clone(&self.pool),
1721                "query_page",
1722                move |scope, conn| {
1723                    execute_query_page_interruptibly(
1724                        scope,
1725                        conn,
1726                        &statement,
1727                        &page,
1728                        "query_page",
1729                        false,
1730                        true,
1731                    )
1732                },
1733            )
1734            .await;
1735        }
1736        let result = execute_standalone_read(
1737            &mut self.handle,
1738            Arc::clone(&self.pool),
1739            "query_page",
1740            transaction_control,
1741            move |scope, conn, rollback, interruptible| {
1742                execute_query_page_interruptibly(
1743                    scope,
1744                    conn,
1745                    &statement,
1746                    &page,
1747                    "query_page",
1748                    rollback,
1749                    interruptible,
1750                )
1751            },
1752        )
1753        .await;
1754        if self.handle.is_none() {
1755            self.poisoned = true;
1756        }
1757        self.close_inactive_transaction_handle();
1758        result
1759    }
1760
1761    async fn query_scalar(
1762        &mut self,
1763        statement: SqlStatement,
1764    ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
1765        let row = self.query_row(statement).await?;
1766        Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
1767    }
1768
1769    async fn explain(
1770        &mut self,
1771        statement: SqlStatement,
1772    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1773        let explain_stmt = SqlStatement {
1774            sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
1775            params: statement.params,
1776            label: statement.label,
1777        };
1778        self.query_all(explain_stmt).await
1779    }
1780}
1781
1782// =============================================================================
1783// File-backed: SqliteWriter (standalone connection)
1784// =============================================================================
1785
1786struct SqliteWriter {
1787    // Manual atomic owners observe the final unit, excluding inner/cleanup errors.
1788    observe_direct_errors: bool,
1789    /// `None` at construction when a `WriterTaskHandle` was obtained (ADR-136
1790    /// D1 gate 1: queue-first `writer()` skips the standalone open in that
1791    /// case). Ordinary `SqlReader` supertrait calls then use pooled readers;
1792    /// this slot is populated only while such a queue-backed handle owns an
1793    /// explicit deferred read transaction. An eagerly opened read-write
1794    /// connection (the no-writer-task branch) retains its one-permit writer
1795    /// budget for the handle's lifetime and must serve its reads on that same
1796    /// connection to preserve manual-transaction visibility.
1797    handle: Option<StandaloneHandle>,
1798    /// ADR-067 Component A: when the write queue is enabled, `execute_batch`
1799    /// routes the whole caller-supplied statement list through the
1800    /// single-writer task instead of opening its own `BEGIN IMMEDIATE` on
1801    /// the standalone connection. `None` when the flag is off or no writer
1802    /// task is available
1803    /// (best-effort — degrades to the standalone-connection path below).
1804    writer_task: Option<crate::writer_task::WriterTaskHandle>,
1805    /// The origin (ADR-091 backend-scoped attribution) of the pool this
1806    /// standalone connection was opened against.
1807    origin: khive_storage::tx_registry::TxOrigin,
1808    /// This connection's pool's writer-timeout sink identity (`db_label`),
1809    /// captured at construction so the standalone-path busy/locked mapping
1810    /// below doesn't need a `&ConnectionPool` reference to report against.
1811    db: String,
1812    /// Reader-pool owner and explicit-transaction exception source.
1813    pool: Arc<ConnectionPool>,
1814    event_rows: Option<Arc<AtomicEventRows>>,
1815    /// The volume lease an enclosing manual atomic unit holds for its whole
1816    /// transaction, so the unit's own statements do not request it again.
1817    /// Declared last: the connection closes before the lease is released.
1818    held_lease: Option<crate::disk_guard::DetachedVolumeLease>,
1819}
1820
1821fn execute_top_level_maintenance(
1822    pool: &ConnectionPool,
1823    conn: &rusqlite::Connection,
1824    maintenance: TopLevelMaintenance,
1825) -> rusqlite::Result<()> {
1826    match maintenance {
1827        TopLevelMaintenance::WalCheckpointTruncate => {
1828            let result = conn.query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |row| {
1829                Ok((
1830                    row.get::<_, i64>(0)?,
1831                    row.get::<_, i64>(1)?,
1832                    row.get::<_, i64>(2)?,
1833                ))
1834            });
1835            crate::checkpoint::record_checkpoint_run_result(pool, result.as_ref().ok().copied());
1836            result.map(|_| ())
1837        }
1838        TopLevelMaintenance::Vacuum => conn.execute_batch(maintenance.as_sql()),
1839    }
1840}
1841
1842impl SqliteWriter {
1843    /// A standalone handle holds its connection across calls, so opening it
1844    /// cannot serve as admission for every later write on that connection.
1845    /// Each write admits itself on its blocking thread
1846    /// ([`admit_standalone_operation`]); this only refuses a consumed handle
1847    /// before it is taken.
1848    fn require_standalone_handle(
1849        &self,
1850        operation: &'static str,
1851    ) -> khive_storage::types::StorageResult<()> {
1852        if self.handle.is_none() {
1853            return Err(StorageError::Pool {
1854                operation: operation.into(),
1855                message: "connection already consumed".into(),
1856            });
1857        }
1858        Ok(())
1859    }
1860
1861    async fn use_queue_read_transaction_handle(
1862        &mut self,
1863        transaction_control: Option<CachedReadTransactionControl>,
1864        operation: &'static str,
1865    ) -> khive_storage::types::StorageResult<bool> {
1866        if self.handle.is_some() {
1867            return Ok(true);
1868        }
1869        match transaction_control {
1870            None => Ok(false),
1871            Some(CachedReadTransactionControl::BeginDeferred) => {
1872                self.handle =
1873                    Some(open_explicit_read_transaction_handle(Arc::clone(&self.pool)).await?);
1874                Ok(true)
1875            }
1876            Some(CachedReadTransactionControl::Finish(keyword))
1877            | Some(CachedReadTransactionControl::Unsupported(keyword)) => {
1878                Err(StorageError::InvalidInput {
1879                    capability: StorageCapability::Sql,
1880                    operation: operation.into(),
1881                    message: format!(
1882                        "cached read-only handle has no admitted transaction for transaction \
1883                         control ({keyword})"
1884                    ),
1885                })
1886            }
1887        }
1888    }
1889
1890    fn close_inactive_queue_read_transaction_handle(&mut self) {
1891        if self
1892            .handle
1893            .as_ref()
1894            .is_some_and(|handle| handle.is_cached_reader() && !handle.has_read_transaction())
1895        {
1896            drop(self.handle.take());
1897        }
1898    }
1899}
1900
1901#[async_trait]
1902impl khive_storage::SqlReader for SqliteWriter {
1903    async fn query_row(
1904        &mut self,
1905        statement: SqlStatement,
1906    ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
1907        if self.writer_task.is_some() {
1908            let transaction_control = cached_read_transaction_control(&statement.sql);
1909            if !self
1910                .use_queue_read_transaction_handle(transaction_control, "writer.query_row")
1911                .await?
1912            {
1913                admit_reader_capability_sql(&statement, transaction_control, "writer.query_row")?;
1914                return run_pool_reader_query(
1915                    Arc::clone(&self.pool),
1916                    "writer.query_row",
1917                    move |scope, conn| {
1918                        execute_query_row_interruptibly(
1919                            scope,
1920                            conn,
1921                            &statement,
1922                            "writer.query_row",
1923                            false,
1924                            true,
1925                        )
1926                    },
1927                )
1928                .await;
1929            }
1930            let result = execute_standalone_read(
1931                &mut self.handle,
1932                Arc::clone(&self.pool),
1933                "writer.query_row",
1934                transaction_control,
1935                move |scope, conn, rollback, interruptible| {
1936                    execute_query_row_interruptibly(
1937                        scope,
1938                        conn,
1939                        &statement,
1940                        "writer.query_row",
1941                        rollback,
1942                        interruptible,
1943                    )
1944                },
1945            )
1946            .await;
1947            self.close_inactive_queue_read_transaction_handle();
1948            return result;
1949        }
1950        let transaction_control = cached_read_transaction_control(&statement.sql);
1951        execute_standalone_read(
1952            &mut self.handle,
1953            Arc::clone(&self.pool),
1954            "writer.query_row",
1955            transaction_control,
1956            move |scope, conn, rollback, interruptible| {
1957                execute_query_row_interruptibly(
1958                    scope,
1959                    conn,
1960                    &statement,
1961                    "writer.query_row",
1962                    rollback,
1963                    interruptible,
1964                )
1965            },
1966        )
1967        .await
1968    }
1969
1970    async fn query_all(
1971        &mut self,
1972        statement: SqlStatement,
1973    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1974        if self.writer_task.is_some() {
1975            let transaction_control = cached_read_transaction_control(&statement.sql);
1976            if !self
1977                .use_queue_read_transaction_handle(transaction_control, "writer.query_all")
1978                .await?
1979            {
1980                admit_reader_capability_sql(&statement, transaction_control, "writer.query_all")?;
1981                return run_pool_reader_query(
1982                    Arc::clone(&self.pool),
1983                    "writer.query_all",
1984                    move |scope, conn| {
1985                        execute_query_interruptibly(
1986                            scope,
1987                            conn,
1988                            &statement,
1989                            "writer.query_all",
1990                            false,
1991                            true,
1992                        )
1993                    },
1994                )
1995                .await;
1996            }
1997            let result = execute_standalone_read(
1998                &mut self.handle,
1999                Arc::clone(&self.pool),
2000                "writer.query_all",
2001                transaction_control,
2002                move |scope, conn, rollback, interruptible| {
2003                    execute_query_interruptibly(
2004                        scope,
2005                        conn,
2006                        &statement,
2007                        "writer.query_all",
2008                        rollback,
2009                        interruptible,
2010                    )
2011                },
2012            )
2013            .await;
2014            self.close_inactive_queue_read_transaction_handle();
2015            return result;
2016        }
2017        let transaction_control = cached_read_transaction_control(&statement.sql);
2018        execute_standalone_read(
2019            &mut self.handle,
2020            Arc::clone(&self.pool),
2021            "writer.query_all",
2022            transaction_control,
2023            move |scope, conn, rollback, interruptible| {
2024                execute_query_interruptibly(
2025                    scope,
2026                    conn,
2027                    &statement,
2028                    "writer.query_all",
2029                    rollback,
2030                    interruptible,
2031                )
2032            },
2033        )
2034        .await
2035    }
2036
2037    async fn query_page(
2038        &mut self,
2039        statement: SqlStatement,
2040        page: PageRequest,
2041    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2042        if self.writer_task.is_some() {
2043            let transaction_control = cached_read_transaction_control(&statement.sql);
2044            if !self
2045                .use_queue_read_transaction_handle(transaction_control, "writer.query_page")
2046                .await?
2047            {
2048                admit_reader_capability_sql(&statement, transaction_control, "writer.query_page")?;
2049                return run_pool_reader_query(
2050                    Arc::clone(&self.pool),
2051                    "writer.query_page",
2052                    move |scope, conn| {
2053                        execute_query_page_interruptibly(
2054                            scope,
2055                            conn,
2056                            &statement,
2057                            &page,
2058                            "writer.query_page",
2059                            false,
2060                            true,
2061                        )
2062                    },
2063                )
2064                .await;
2065            }
2066            let result = execute_standalone_read(
2067                &mut self.handle,
2068                Arc::clone(&self.pool),
2069                "writer.query_page",
2070                transaction_control,
2071                move |scope, conn, rollback, interruptible| {
2072                    execute_query_page_interruptibly(
2073                        scope,
2074                        conn,
2075                        &statement,
2076                        &page,
2077                        "writer.query_page",
2078                        rollback,
2079                        interruptible,
2080                    )
2081                },
2082            )
2083            .await;
2084            self.close_inactive_queue_read_transaction_handle();
2085            return result;
2086        }
2087        let transaction_control = cached_read_transaction_control(&statement.sql);
2088        execute_standalone_read(
2089            &mut self.handle,
2090            Arc::clone(&self.pool),
2091            "writer.query_page",
2092            transaction_control,
2093            move |scope, conn, rollback, interruptible| {
2094                execute_query_page_interruptibly(
2095                    scope,
2096                    conn,
2097                    &statement,
2098                    &page,
2099                    "writer.query_page",
2100                    rollback,
2101                    interruptible,
2102                )
2103            },
2104        )
2105        .await
2106    }
2107
2108    async fn query_scalar(
2109        &mut self,
2110        statement: SqlStatement,
2111    ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
2112        let row = khive_storage::SqlReader::query_row(self, statement).await?;
2113        Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
2114    }
2115
2116    async fn explain(
2117        &mut self,
2118        statement: SqlStatement,
2119    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2120        let explain_stmt = SqlStatement {
2121            sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
2122            params: statement.params,
2123            label: statement.label,
2124        };
2125        khive_storage::SqlReader::query_all(self, explain_stmt).await
2126    }
2127}
2128
2129#[async_trait]
2130impl khive_storage::SqlWriter for SqliteWriter {
2131    async fn execute(
2132        &mut self,
2133        statement: SqlStatement,
2134    ) -> khive_storage::types::StorageResult<u64> {
2135        // ADR-067 Component A (Fork C slice 2): a single statement is
2136        // self-contained, just like `execute_batch`'s full statement list —
2137        // on the writer task, transaction-control rejection remains an
2138        // `execute_batch` contract; the standalone branch below refuses it
2139        // unless a unit holds the volume lease, because this primitive is also
2140        // used by internal atomic transaction owners.
2141        // route it through the writer task when available. `self.handle` is
2142        // left untouched so a subsequent `execute`/`execute_script` call on
2143        // this same handle still works over the standalone connection.
2144        if let Some(writer_task) = self.writer_task.clone() {
2145            let event_rows = self.event_rows.clone();
2146            return writer_task
2147                .send_bounded(move |conn| {
2148                    let mut stmt = prepare_cached_sql_statement(conn, &statement.sql)
2149                        .map_err(|e| map_rusqlite_err(e, "execute"))?;
2150                    bind_params(&mut stmt, &statement.params)
2151                        .map_err(|e| map_rusqlite_err(e, "execute"))?;
2152                    let affected = stmt
2153                        .raw_execute()
2154                        .map_err(|e| map_rusqlite_err(e, "execute"))?;
2155                    if let Some(event_rows) = event_rows.as_deref() {
2156                        event_rows.observe(&statement, affected as u64);
2157                    }
2158                    Ok(affected as u64)
2159                })
2160                .await;
2161        }
2162
2163        self.require_standalone_handle("execute")?;
2164        let unit_holds_lease = self.held_lease.is_some();
2165        // ADR-154: a standalone statement holds the volume lease only for the
2166        // call, so a transaction it opened would continue after the lease was
2167        // released. Only a unit that already holds the lease for its whole
2168        // span (`atomic_unit`'s manual transaction) may send transaction
2169        // control through `execute`.
2170        if !unit_holds_lease {
2171            if let Some(keyword) = transaction_control_head(&statement.sql) {
2172                return Err(StorageError::InvalidInput {
2173                    capability: StorageCapability::Sql,
2174                    operation: "execute".into(),
2175                    message: format!(
2176                        "statement is transaction control ({keyword}); a standalone \
2177                         statement holds the volume lease only for the call — use \
2178                         atomic_unit to run statements as one transaction"
2179                    ),
2180                });
2181            }
2182        }
2183        let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
2184            operation: "execute".into(),
2185            message: "connection already consumed".into(),
2186        })?;
2187        let event_rows = self.event_rows.clone();
2188        let pool = Arc::clone(&self.pool);
2189        let (handle, result) = tokio::task::spawn_blocking(move || {
2190            run_standalone_statement(handle, pool, unit_holds_lease, statement, event_rows)
2191        })
2192        .await
2193        .map_err(|e| StorageError::driver(StorageCapability::Sql, "execute", e))?;
2194        self.handle = Some(handle);
2195        let affected = result.map_err(|failure| match failure {
2196            StandaloneWriteError::Refused(error) => error,
2197            StandaloneWriteError::Sql(error) => self.map_direct_error(error, "execute"),
2198        })?;
2199        Ok(affected as u64)
2200    }
2201
2202    async fn execute_batch(
2203        &mut self,
2204        statements: Vec<SqlStatement>,
2205    ) -> khive_storage::types::StorageResult<u64> {
2206        // ADR-067 Component A: this call is self-contained (the full statement
2207        // list is supplied up front and the whole thing commits or rolls back
2208        // as one unit) — unlike `writer()`'s live incrementally-driven handle,
2209        // it maps cleanly onto a single `WriteRequest`. Route it through the
2210        // writer task when available; `self.handle` is left untouched so a
2211        // subsequent `execute`/`execute_script` call on this same handle still
2212        // works over the standalone connection (that dispatch is unmigrated —
2213        // see `SqlBridge::writer()`).
2214        //
2215        // Both paths reject transaction-control statements BEFORE executing
2216        // anything: the queue-backed branch runs inside the writer task's own
2217        // `BEGIN IMMEDIATE` (a caller `COMMIT` there would close the task's
2218        // transaction and terminate the writer task), and the standalone
2219        // branch below wraps the list in its own `BEGIN IMMEDIATE` (a caller
2220        // `COMMIT` would commit early and break all-or-nothing).
2221        reject_transaction_control_statements(&statements, "execute_batch")?;
2222        if let Some(writer_task) = self.writer_task.clone() {
2223            let event_rows = self.event_rows.clone();
2224            return writer_task
2225                .send_bounded(move |conn| {
2226                    let prepared = prepare_batch_statements(conn, &statements)
2227                        .map_err(|e| map_rusqlite_err(e, "execute_batch"))?;
2228                    execute_prepared_batch(conn, prepared, &statements, event_rows.as_deref())
2229                        .map_err(|e| map_rusqlite_err(e, "execute_batch"))
2230                })
2231                .await;
2232        }
2233
2234        self.require_standalone_handle("execute_batch")?;
2235        let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
2236            operation: "execute_batch".into(),
2237            message: "connection already consumed".into(),
2238        })?;
2239        let origin = self.origin.clone();
2240        let event_rows = self.event_rows.clone();
2241        let pool = Arc::clone(&self.pool);
2242        let unit_holds_lease = self.held_lease.is_some();
2243        let (handle, result) = tokio::task::spawn_blocking(move || {
2244            run_standalone_batch(
2245                handle,
2246                pool,
2247                unit_holds_lease,
2248                statements,
2249                origin,
2250                event_rows,
2251            )
2252        })
2253        .await
2254        .map_err(|e| StorageError::driver(StorageCapability::Sql, "execute_batch", e))?;
2255        self.handle = handle;
2256        result.map_err(|failure| match failure {
2257            StandaloneWriteError::Refused(error) => error,
2258            StandaloneWriteError::Sql(failure) => self.map_direct_batch_failure(failure),
2259        })
2260    }
2261
2262    async fn execute_script(&mut self, script: String) -> khive_storage::types::StorageResult<()> {
2263        // ADR-067 Component A (Fork C slice 2): the script text is
2264        // self-contained (supplied up front, runs as one unit), just like
2265        // `execute_batch` — route it through the writer task when
2266        // available. `self.handle` is left untouched so a subsequent
2267        // `execute`/`execute_script` call on this same handle still works
2268        // over the standalone connection. Callers must supply a DML-only
2269        // script (no bare `BEGIN`/`COMMIT`/`ROLLBACK`) on the flag-on path,
2270        // since it runs inside the writer task's own transaction — same
2271        // Boundary: transaction-control rejection is an `execute_batch`
2272        // contract; this raw script path is internal/migration-only. The
2273        // queue-backed branch still requires a DML-only script because it
2274        // runs inside the writer task's transaction.
2275        if let Some(writer_task) = self.writer_task.clone() {
2276            return writer_task
2277                .send_bounded(move |conn| {
2278                    conn.execute_batch(&script)
2279                        .map_err(|e| map_rusqlite_err(e, "execute_script"))
2280                })
2281                .await;
2282        }
2283
2284        self.require_standalone_handle("execute_script")?;
2285        let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
2286            operation: "execute_script".into(),
2287            message: "connection already consumed".into(),
2288        })?;
2289        let pool = Arc::clone(&self.pool);
2290        let unit_holds_lease = self.held_lease.is_some();
2291        let (handle, result) = tokio::task::spawn_blocking(move || {
2292            run_standalone_script(handle, pool, unit_holds_lease, script)
2293        })
2294        .await
2295        .map_err(|e| StorageError::driver(StorageCapability::Sql, "execute_script", e))?;
2296        self.handle = handle;
2297        result.map_err(|failure| match failure {
2298            StandaloneWriteError::Refused(error) => error,
2299            StandaloneWriteError::Sql(error) => self.map_direct_error(error, "execute_script"),
2300        })
2301    }
2302
2303    async fn execute_script_top_level(
2304        &mut self,
2305        maintenance: TopLevelMaintenance,
2306    ) -> khive_storage::types::StorageResult<()> {
2307        // Only the closed maintenance enum can supply unbound SQL here.
2308        // This is not the separate raw migration-script interface.
2309        // ADR-067 Component A: unlike
2310        // `execute_script`, this must NOT run inside the writer task's
2311        // per-request `BEGIN IMMEDIATE` — statements such as VACUUM are
2312        // rejected by SQLite inside any open transaction. Route through
2313        // `WriterTaskHandle::send_top_level`, which still serializes this
2314        // call through the single writer owner but skips the transaction
2315        // wrap entirely.
2316        if let Some(writer_task) = self.writer_task.clone() {
2317            let pool = Arc::clone(&self.pool);
2318            let execute = move |conn: &rusqlite::Connection| {
2319                execute_top_level_maintenance(&pool, conn, maintenance)
2320                    .map_err(|e| map_rusqlite_err(e, "execute_script_top_level"))
2321            };
2322            return if maintenance == TopLevelMaintenance::WalCheckpointTruncate {
2323                writer_task.send_checkpoint_bounded(execute).await
2324            } else {
2325                writer_task.send_vacuum_bounded(execute).await
2326            };
2327        }
2328
2329        // Flag off / no writer task: keep the volume lease through the
2330        // whole top-level operation. A checkpoint is a recovery bypass;
2331        // VACUUM uses a copy-sized DB/WAL metadata estimate.
2332        let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
2333            operation: "execute_script_top_level".into(),
2334            message: "connection already consumed".into(),
2335        })?;
2336        let pool = Arc::clone(&self.pool);
2337        let unit_holds_lease = self.held_lease.is_some();
2338        let (handle, result) = tokio::task::spawn_blocking(move || {
2339            run_standalone_top_level(handle, pool, unit_holds_lease, maintenance)
2340        })
2341        .await
2342        .map_err(|e| StorageError::driver(StorageCapability::Sql, "execute_script_top_level", e))?;
2343        self.handle = Some(handle);
2344        result.map_err(|failure| match failure {
2345            StandaloneWriteError::Refused(error) => error,
2346            StandaloneWriteError::Sql(error) => {
2347                self.map_direct_error(error, "execute_script_top_level")
2348            }
2349        })
2350    }
2351}
2352
2353// =============================================================================
2354// Pool-backed reader/writer (in-memory databases)
2355// =============================================================================
2356
2357async fn run_pool_reader_query<T, F>(
2358    pool: Arc<ConnectionPool>,
2359    operation: &'static str,
2360    query: F,
2361) -> khive_storage::types::StorageResult<T>
2362where
2363    T: Send + 'static,
2364    F: FnOnce(
2365            &crate::read_cancellation::InterruptibleReadScope,
2366            &rusqlite::Connection,
2367        ) -> khive_storage::types::StorageResult<T>
2368        + Send
2369        + 'static,
2370{
2371    // Await the permit before the blocking task starts: a queued read holds no thread.
2372    let admission = pool
2373        .acquire_reader_admission(StorageCapability::Sql, operation)
2374        .await?;
2375    crate::read_cancellation::run_interruptible_read(
2376        StorageCapability::Sql,
2377        operation,
2378        move |scope| {
2379            // Checkout tri-state (cancelled -> Timeout, admission expiry ->
2380            // retryable AdmissionTimeout, other -> Driver) lives in ONE place:
2381            // `ConnectionPool::resolve_reader_checkout`.
2382            let mut guard = pool.resolve_reader_checkout(
2383                StorageCapability::Sql,
2384                operation,
2385                pool.reader_with_admission(admission, || scope.should_stop()),
2386            )?;
2387            // Every caller of `run_pool_reader_query` runs a `SqlStatement`
2388            // (raw SQL) rather than a typed store's fixed query, so the
2389            // checkout pays the pristine-state scan on return regardless of
2390            // which `SqlReader` wrapper (reader or writer capability) drew it.
2391            guard.mark_dirty();
2392            let result = scope.with_pooled_reader(&mut guard, |conn| query(scope, conn));
2393            if let Err(error) = &result {
2394                pool.record_reader_query_error(error);
2395            }
2396            result
2397        },
2398    )
2399    .await
2400}
2401
2402async fn run_pool_writer_query<T, F>(
2403    pool: Arc<ConnectionPool>,
2404    operation: &'static str,
2405    query: F,
2406) -> khive_storage::types::StorageResult<T>
2407where
2408    T: Send + 'static,
2409    F: FnOnce(
2410            &crate::read_cancellation::InterruptibleReadScope,
2411            &rusqlite::Connection,
2412            bool,
2413        ) -> khive_storage::types::StorageResult<T>
2414        + Send
2415        + 'static,
2416{
2417    crate::read_cancellation::run_interruptible_read(
2418        StorageCapability::Sql,
2419        operation,
2420        move |scope| {
2421            let guard = pool.try_writer().map_err(|error: SqliteError| {
2422                error.into_storage_error(StorageCapability::Sql, operation)
2423            })?;
2424            scope.with_pooled_writer(&pool, &guard, |conn| {
2425                let interruptible = conn.is_autocommit();
2426                query(scope, conn, interruptible)
2427            })
2428        },
2429    )
2430    .await
2431}
2432
2433struct PoolBackedReader {
2434    pool: Arc<ConnectionPool>,
2435    /// Present only while an admitted explicit deferred read transaction
2436    /// (ADR-005/ADR-091) is open on the in-memory backend's single shared
2437    /// connection. Retaining it here — instead of drawing a fresh checkout
2438    /// per call, the way an ordinary read does — is what gives the span
2439    /// real connection ownership: see [`SharedReaderTransactionGuard`] and
2440    /// [`run_pool_backed_reader_query`].
2441    transaction: Option<SharedReaderTransactionGuard>,
2442}
2443
2444/// Decide what happened to a [`SharedReaderTransactionGuard`] after one
2445/// statement ran against it and either retain it in `*transaction` (still
2446/// mid-span) or let it drop naturally (span finished cleanly).
2447///
2448/// `expect_open_after` is the caller's terminal-state expectation: `true`
2449/// for an ordinary read inside an already-open span (the transaction must
2450/// still be open afterward), `false` for the statement that opens or closes
2451/// the span. A mismatch poisons the guard — it is never handed back for
2452/// reuse — and is folded into the returned error so a broken span is never
2453/// reported as a successful read.
2454fn finish_pool_backed_reader_step<T>(
2455    transaction: &mut Option<SharedReaderTransactionGuard>,
2456    guard: SharedReaderTransactionGuard,
2457    expect_open_after: bool,
2458    operation: &'static str,
2459    result: khive_storage::types::StorageResult<T>,
2460) -> khive_storage::types::StorageResult<T> {
2461    let still_open = !guard.conn().is_autocommit();
2462    if still_open == expect_open_after {
2463        if still_open {
2464            *transaction = Some(guard);
2465        }
2466        // Otherwise the span finished cleanly (or never opened); let `guard`
2467        // drop here, returning the connection to the pool.
2468        return result;
2469    }
2470    guard.poison();
2471    let message = if expect_open_after {
2472        "a read inside the pool-backed reader's admitted transaction unexpectedly ended it; \
2473         the connection was discarded"
2474    } else {
2475        "transaction-ending control completed but the pool-backed reader's connection \
2476         remained outside autocommit; the connection was discarded"
2477    };
2478    match result {
2479        Err(error) => Err(error),
2480        Ok(_) => Err(StorageError::InvalidInput {
2481            capability: StorageCapability::Sql,
2482            operation: operation.into(),
2483            message: message.into(),
2484        }),
2485    }
2486}
2487
2488/// Open the explicit deferred read-transaction span on the in-memory
2489/// backend's shared connection and run the admitted `BEGIN DEFERRED`
2490/// statement against it.
2491async fn open_pool_backed_reader_transaction<T, F>(
2492    transaction: &mut Option<SharedReaderTransactionGuard>,
2493    pool: Arc<ConnectionPool>,
2494    operation: &'static str,
2495    query: F,
2496) -> khive_storage::types::StorageResult<T>
2497where
2498    T: Send + 'static,
2499    F: FnOnce(
2500            &crate::read_cancellation::InterruptibleReadScope,
2501            &rusqlite::Connection,
2502            bool,
2503            bool,
2504        ) -> khive_storage::types::StorageResult<T>
2505        + Send
2506        + 'static,
2507{
2508    let (guard, result) = crate::read_cancellation::run_interruptible_read(
2509        StorageCapability::Sql,
2510        operation,
2511        move |scope| {
2512            let Some(guard) = pool
2513                .checkout_shared_reader_transaction(|| scope.should_stop())
2514                .map_err(|error| StorageError::driver(StorageCapability::Sql, operation, error))?
2515            else {
2516                return Err(StorageError::Timeout {
2517                    operation: operation.into(),
2518                });
2519            };
2520            let result = query(scope, guard.conn(), false, true);
2521            if scope.cleanup_failed() {
2522                guard.poison();
2523            }
2524            Ok((guard, result))
2525        },
2526    )
2527    .await?;
2528    finish_pool_backed_reader_step(transaction, guard, true, operation, result)
2529}
2530
2531/// Route one raw-SQL call through the in-memory backend's `PoolBackedReader`.
2532/// An ordinary read with no open span falls through to the regular per-call
2533/// pooled checkout ([`run_pool_reader_query`]); every other combination
2534/// (opening, continuing, or closing the explicit deferred read-transaction
2535/// span) is handled here so the span retains one connection end to end.
2536#[allow(clippy::too_many_lines)]
2537async fn run_pool_backed_reader_query<T, F>(
2538    transaction: &mut Option<SharedReaderTransactionGuard>,
2539    pool: Arc<ConnectionPool>,
2540    operation: &'static str,
2541    transaction_control: Option<CachedReadTransactionControl>,
2542    query: F,
2543) -> khive_storage::types::StorageResult<T>
2544where
2545    T: Send + 'static,
2546    F: FnOnce(
2547            &crate::read_cancellation::InterruptibleReadScope,
2548            &rusqlite::Connection,
2549            bool,
2550            bool,
2551        ) -> khive_storage::types::StorageResult<T>
2552        + Send
2553        + 'static,
2554{
2555    if transaction.is_none() {
2556        return match transaction_control {
2557            None => {
2558                run_pool_reader_query(pool, operation, move |scope, conn| {
2559                    query(scope, conn, false, true)
2560                })
2561                .await
2562            }
2563            Some(CachedReadTransactionControl::Finish(keyword))
2564            | Some(CachedReadTransactionControl::Unsupported(keyword)) => {
2565                Err(StorageError::InvalidInput {
2566                    capability: StorageCapability::Sql,
2567                    operation: operation.into(),
2568                    message: format!(
2569                        "pool-backed reader has no admitted transaction for transaction \
2570                         control ({keyword})"
2571                    ),
2572                })
2573            }
2574            Some(CachedReadTransactionControl::BeginDeferred) => {
2575                open_pool_backed_reader_transaction(transaction, pool, operation, query).await
2576            }
2577        };
2578    }
2579
2580    match transaction_control {
2581        Some(CachedReadTransactionControl::BeginDeferred) => {
2582            return Err(StorageError::InvalidInput {
2583                capability: StorageCapability::Sql,
2584                operation: operation.into(),
2585                message: "pool-backed reader already owns an admitted read transaction; \
2586                          nested BEGIN is not supported"
2587                    .into(),
2588            });
2589        }
2590        Some(CachedReadTransactionControl::Unsupported(keyword)) => {
2591            return Err(StorageError::InvalidInput {
2592                capability: StorageCapability::Sql,
2593                operation: operation.into(),
2594                message: format!(
2595                    "pool-backed reader's admitted read transaction does not support nested \
2596                     or write-locking transaction control ({keyword})"
2597                ),
2598            });
2599        }
2600        None | Some(CachedReadTransactionControl::Finish(_)) => {}
2601    }
2602
2603    let expect_open_after = transaction_control.is_none();
2604    let guard = transaction.take().expect("checked Some above");
2605    let (guard, result) = crate::read_cancellation::run_interruptible_read(
2606        StorageCapability::Sql,
2607        operation,
2608        move |scope| {
2609            let result = query(scope, guard.conn(), true, true);
2610            if scope.cleanup_failed() {
2611                guard.poison();
2612            }
2613            Ok((guard, result))
2614        },
2615    )
2616    .await?;
2617    finish_pool_backed_reader_step(transaction, guard, expect_open_after, operation, result)
2618}
2619
2620#[async_trait]
2621impl khive_storage::SqlReader for PoolBackedReader {
2622    async fn query_row(
2623        &mut self,
2624        statement: SqlStatement,
2625    ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
2626        let transaction_control = cached_read_transaction_control(&statement.sql);
2627        admit_reader_capability_sql(&statement, transaction_control, "pool_reader.query_row")?;
2628        let pool = Arc::clone(&self.pool);
2629        run_pool_backed_reader_query(
2630            &mut self.transaction,
2631            pool,
2632            "pool_reader.query_row",
2633            transaction_control,
2634            move |scope, conn, rollback, interruptible| {
2635                execute_query_row_interruptibly(
2636                    scope,
2637                    conn,
2638                    &statement,
2639                    "pool_reader.query_row",
2640                    rollback,
2641                    interruptible,
2642                )
2643            },
2644        )
2645        .await
2646    }
2647
2648    async fn query_all(
2649        &mut self,
2650        statement: SqlStatement,
2651    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2652        let transaction_control = cached_read_transaction_control(&statement.sql);
2653        admit_reader_capability_sql(&statement, transaction_control, "pool_reader.query_all")?;
2654        let pool = Arc::clone(&self.pool);
2655        run_pool_backed_reader_query(
2656            &mut self.transaction,
2657            pool,
2658            "pool_reader.query_all",
2659            transaction_control,
2660            move |scope, conn, rollback, interruptible| {
2661                execute_query_interruptibly(
2662                    scope,
2663                    conn,
2664                    &statement,
2665                    "pool_reader.query_all",
2666                    rollback,
2667                    interruptible,
2668                )
2669            },
2670        )
2671        .await
2672    }
2673
2674    async fn query_page(
2675        &mut self,
2676        statement: SqlStatement,
2677        page: PageRequest,
2678    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2679        let transaction_control = cached_read_transaction_control(&statement.sql);
2680        admit_reader_capability_sql(&statement, transaction_control, "pool_reader.query_page")?;
2681        let pool = Arc::clone(&self.pool);
2682        run_pool_backed_reader_query(
2683            &mut self.transaction,
2684            pool,
2685            "pool_reader.query_page",
2686            transaction_control,
2687            move |scope, conn, rollback, interruptible| {
2688                execute_query_page_interruptibly(
2689                    scope,
2690                    conn,
2691                    &statement,
2692                    &page,
2693                    "pool_reader.query_page",
2694                    rollback,
2695                    interruptible,
2696                )
2697            },
2698        )
2699        .await
2700    }
2701
2702    async fn query_scalar(
2703        &mut self,
2704        statement: SqlStatement,
2705    ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
2706        let row = self.query_row(statement).await?;
2707        Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
2708    }
2709
2710    async fn explain(
2711        &mut self,
2712        statement: SqlStatement,
2713    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2714        let explain_stmt = SqlStatement {
2715            sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
2716            params: statement.params,
2717            label: statement.label,
2718        };
2719        self.query_all(explain_stmt).await
2720    }
2721}
2722
2723struct PoolBackedWriter {
2724    pool: Arc<ConnectionPool>,
2725}
2726
2727#[async_trait]
2728impl khive_storage::SqlReader for PoolBackedWriter {
2729    async fn query_row(
2730        &mut self,
2731        statement: SqlStatement,
2732    ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
2733        let pool = Arc::clone(&self.pool);
2734        run_pool_writer_query(
2735            pool,
2736            "pool_writer.query_row",
2737            move |scope, conn, interruptible| {
2738                execute_query_row_interruptibly(
2739                    scope,
2740                    conn,
2741                    &statement,
2742                    "pool_writer.query_row",
2743                    false,
2744                    interruptible,
2745                )
2746            },
2747        )
2748        .await
2749    }
2750
2751    async fn query_all(
2752        &mut self,
2753        statement: SqlStatement,
2754    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2755        let pool = Arc::clone(&self.pool);
2756        run_pool_writer_query(
2757            pool,
2758            "pool_writer.query_all",
2759            move |scope, conn, interruptible| {
2760                execute_query_interruptibly(
2761                    scope,
2762                    conn,
2763                    &statement,
2764                    "pool_writer.query_all",
2765                    false,
2766                    interruptible,
2767                )
2768            },
2769        )
2770        .await
2771    }
2772
2773    async fn query_page(
2774        &mut self,
2775        statement: SqlStatement,
2776        page: PageRequest,
2777    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2778        let pool = Arc::clone(&self.pool);
2779        run_pool_writer_query(
2780            pool,
2781            "pool_writer.query_page",
2782            move |scope, conn, interruptible| {
2783                execute_query_page_interruptibly(
2784                    scope,
2785                    conn,
2786                    &statement,
2787                    &page,
2788                    "pool_writer.query_page",
2789                    false,
2790                    interruptible,
2791                )
2792            },
2793        )
2794        .await
2795    }
2796
2797    async fn query_scalar(
2798        &mut self,
2799        statement: SqlStatement,
2800    ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
2801        let row = khive_storage::SqlReader::query_row(self, statement).await?;
2802        Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
2803    }
2804
2805    async fn explain(
2806        &mut self,
2807        statement: SqlStatement,
2808    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2809        let explain_stmt = SqlStatement {
2810            sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
2811            params: statement.params,
2812            label: statement.label,
2813        };
2814        khive_storage::SqlReader::query_all(self, explain_stmt).await
2815    }
2816}
2817
2818#[async_trait]
2819impl khive_storage::SqlWriter for PoolBackedWriter {
2820    async fn execute(
2821        &mut self,
2822        statement: SqlStatement,
2823    ) -> khive_storage::types::StorageResult<u64> {
2824        // Each call checks the pooled writer out and releases it before
2825        // returning, so a transaction opened here could not outlive the call:
2826        // the guard would settle it on release while the caller believed it
2827        // open. Transaction control is refused before anything runs; a
2828        // transaction is opened through `atomic_unit`, which holds one guard
2829        // for the whole unit.
2830        if let Some(keyword) = transaction_control_head(&statement.sql) {
2831            return Err(StorageError::InvalidInput {
2832                capability: StorageCapability::Sql,
2833                operation: "pool_writer.execute".into(),
2834                message: format!(
2835                    "statement is transaction control ({keyword}); a pooled writer is \
2836                     released after every call, so it cannot hold a transaction across \
2837                     calls — use atomic_unit to run statements as one transaction"
2838                ),
2839            });
2840        }
2841        let pool = Arc::clone(&self.pool);
2842        tokio::task::spawn_blocking(move || {
2843            let guard = pool.try_writer().map_err(|e: SqliteError| {
2844                StorageError::driver(StorageCapability::Sql, "pool_writer.execute", e)
2845            })?;
2846            let result = (|| {
2847                let mut stmt = prepare_cached_sql_statement(&guard, &statement.sql)
2848                    .map_err(|e| map_rusqlite_err(e, "pool_writer.execute"))?;
2849                bind_params(&mut stmt, &statement.params)
2850                    .map_err(|e| map_rusqlite_err(e, "pool_writer.execute"))?;
2851                let rows = stmt
2852                    .raw_execute()
2853                    .map_err(|e| map_rusqlite_err(e, "pool_writer.execute"))?;
2854                Ok(rows as u64)
2855            })();
2856            settle_pooled_call(&guard, "pool_writer.execute", result)
2857                .inspect_err(|error| pool.record_direct_writer_error(error))
2858        })
2859        .await
2860        .map_err(|e| StorageError::driver(StorageCapability::Sql, "pool_writer.execute", e))?
2861    }
2862
2863    async fn execute_batch(
2864        &mut self,
2865        statements: Vec<SqlStatement>,
2866    ) -> khive_storage::types::StorageResult<u64> {
2867        // Same all-or-nothing contract as the file-backed path: this batch
2868        // wraps its list in its own `BEGIN IMMEDIATE`, so reject caller
2869        // transaction-control statements before executing anything.
2870        reject_transaction_control_statements(&statements, "pool_writer.execute_batch")?;
2871        let pool = Arc::clone(&self.pool);
2872        tokio::task::spawn_blocking(move || {
2873            let guard = pool.try_writer().map_err(|e: SqliteError| {
2874                StorageError::driver(StorageCapability::Sql, "pool_writer.execute_batch", e)
2875            })?;
2876            let result = (|| {
2877                let prepared = prepare_batch_statements(&guard, &statements)
2878                    .map_err(|e| map_rusqlite_err(e, "pool_writer.execute_batch"))?;
2879                guard
2880                    .execute_batch("BEGIN IMMEDIATE")
2881                    .map_err(|e| map_rusqlite_err(e, "pool_writer.execute_batch"))?;
2882                let _tx_handle = khive_storage::tx_registry::register_scoped(
2883                    Some("pool_writer.execute_batch".to_string()),
2884                    pool.origin(),
2885                );
2886                let result = execute_prepared_batch(&guard, prepared, &statements, None)
2887                    .map_err(|e| map_rusqlite_err(e, "pool_writer.execute_batch"));
2888                match result {
2889                    Ok(total) => {
2890                        if let Err(e) = guard.execute_batch("COMMIT") {
2891                            let _ = guard.execute_batch("ROLLBACK");
2892                            Err(map_rusqlite_err(e, "pool_writer.execute_batch"))
2893                        } else {
2894                            Ok(total)
2895                        }
2896                    }
2897                    Err(e) => {
2898                        let _ = guard.execute_batch("ROLLBACK");
2899                        Err(e)
2900                    }
2901                }
2902            })();
2903            result.inspect_err(|error| pool.record_direct_writer_error(error))
2904        })
2905        .await
2906        .map_err(|e| StorageError::driver(StorageCapability::Sql, "pool_writer.execute_batch", e))?
2907    }
2908
2909    async fn execute_script(&mut self, script: String) -> khive_storage::types::StorageResult<()> {
2910        // Boundary: raw scripts are internal/migration-only and do not inherit
2911        // `execute_batch`'s transaction-control rejection. A script may open
2912        // and close its own transaction; one it leaves open is rolled back and
2913        // refused before the pooled writer is released.
2914        let pool = Arc::clone(&self.pool);
2915        tokio::task::spawn_blocking(move || {
2916            let guard = pool.try_writer().map_err(|e: SqliteError| {
2917                StorageError::driver(StorageCapability::Sql, "pool_writer.execute_script", e)
2918            })?;
2919            let result = guard
2920                .execute_batch(&script)
2921                .map_err(|e| map_rusqlite_err(e, "pool_writer.execute_script"));
2922            settle_pooled_call(&guard, "pool_writer.execute_script", result)
2923                .inspect_err(|error| pool.record_direct_writer_error(error))
2924        })
2925        .await
2926        .map_err(|e| {
2927            StorageError::driver(StorageCapability::Sql, "pool_writer.execute_script", e)
2928        })?
2929    }
2930}
2931
2932// =============================================================================
2933// atomic_unit (ADR-067 Component A, Fork C slice 2)
2934// =============================================================================
2935
2936/// A purely-synchronous `SqlReader`/`SqlWriter` over a borrowed connection,
2937/// used to drive an [`AtomicUnitOp`] on the queued or in-memory path, where the
2938/// closure body runs inside a `spawn_blocking` (synchronous
2939/// `FnOnce(&rusqlite::Connection) -> ...`) rather than a real async context.
2940///
2941/// Every method here does plain, non-suspending rusqlite work — there is no
2942/// real `.await` point anywhere in this impl — so [`block_on_sync`] driving
2943/// the resulting future to completion with a single poll is sound, not a
2944/// hack: the future can never actually be `Pending`.
2945///
2946/// `SqlReader`/`SqlWriter` both carry a `'static` supertrait bound (they are
2947/// used as `Box<dyn ...>` elsewhere in this module), so this type cannot
2948/// hold a real `&'c Connection` borrow — it would tie `InlineWriter` to a
2949/// non-`'static` lifetime and, independently, `&Connection` is not `Send`
2950/// (`Connection` is `!Sync`), which the `#[async_trait]`-generated futures
2951/// require. A raw pointer sidesteps both: `*const Connection` is `Send` and
2952/// `'static` on its face, and the safety burden (the pointee outliving
2953/// every dereference) is upheld by construction — see `atomic_unit`, the
2954/// only call site: it builds an `InlineWriter` from `conn: &Connection`,
2955/// drives `op` to completion via `block_on_sync` synchronously, and drops
2956/// the `InlineWriter` before that borrow ends, all within one stack frame.
2957struct InlineWriter {
2958    event_rows: Option<Arc<AtomicEventRows>>,
2959    conn: *const rusqlite::Connection,
2960}
2961
2962// SAFETY: `InlineWriter` is never actually shared across a real thread
2963// boundary — it is constructed, driven to completion synchronously via
2964// `block_on_sync`, and dropped within a single call frame inside the
2965// `spawn_blocking` closure (see `atomic_unit`). The `Send`
2966// bound `async_trait` imposes on the futures below is a static
2967// over-approximation for this restricted, single-threaded usage pattern.
2968unsafe impl Send for InlineWriter {}
2969
2970impl InlineWriter {
2971    /// SAFETY: valid for the lifetime of the enclosing synchronous scope in
2972    /// `atomic_unit` (see the struct doc comment above) — the pointee is
2973    /// never dereferenced after that scope ends.
2974    fn conn(&self) -> &rusqlite::Connection {
2975        unsafe { &*self.conn }
2976    }
2977}
2978
2979#[async_trait]
2980impl khive_storage::SqlReader for InlineWriter {
2981    async fn query_row(
2982        &mut self,
2983        statement: SqlStatement,
2984    ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
2985        execute_query_row(self.conn(), &statement)
2986            .map_err(|e| map_rusqlite_err(e, "inline.query_row"))
2987    }
2988
2989    async fn query_all(
2990        &mut self,
2991        statement: SqlStatement,
2992    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2993        execute_query(self.conn(), &statement).map_err(|e| map_rusqlite_err(e, "inline.query_all"))
2994    }
2995
2996    async fn query_page(
2997        &mut self,
2998        statement: SqlStatement,
2999        page: PageRequest,
3000    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
3001        execute_query_page(self.conn(), &statement, &page)
3002            .map_err(|e| map_rusqlite_err(e, "inline.query_page"))
3003    }
3004
3005    async fn query_scalar(
3006        &mut self,
3007        statement: SqlStatement,
3008    ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
3009        let row = khive_storage::SqlReader::query_row(self, statement).await?;
3010        Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
3011    }
3012
3013    async fn explain(
3014        &mut self,
3015        statement: SqlStatement,
3016    ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
3017        let explain_stmt = SqlStatement {
3018            sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
3019            params: statement.params,
3020            label: statement.label,
3021        };
3022        khive_storage::SqlReader::query_all(self, explain_stmt).await
3023    }
3024}
3025
3026#[async_trait]
3027impl khive_storage::SqlWriter for InlineWriter {
3028    async fn execute(
3029        &mut self,
3030        statement: SqlStatement,
3031    ) -> khive_storage::types::StorageResult<u64> {
3032        // Boundary: `execute_batch` owns transaction-control rejection;
3033        // `atomic_unit` uses this one-statement primitive for its own boundary.
3034        let mut stmt = prepare_cached_sql_statement(self.conn(), &statement.sql)
3035            .map_err(|e| map_rusqlite_err(e, "inline.execute"))?;
3036        bind_params(&mut stmt, &statement.params)
3037            .map_err(|e| map_rusqlite_err(e, "inline.execute"))?;
3038        let affected = stmt
3039            .raw_execute()
3040            .map_err(|e| map_rusqlite_err(e, "inline.execute"))?;
3041        if let Some(event_rows) = self.event_rows.as_deref() {
3042            event_rows.observe(&statement, affected as u64);
3043        }
3044        Ok(affected as u64)
3045    }
3046
3047    async fn execute_batch(
3048        &mut self,
3049        statements: Vec<SqlStatement>,
3050    ) -> khive_storage::types::StorageResult<u64> {
3051        // Runs inside the writer task's per-request `BEGIN IMMEDIATE`
3052        // (atomic_unit flag-on path), so a caller `COMMIT` would close the
3053        // task's transaction — reject transaction-control statements up
3054        // front, same contract as every other `execute_batch`.
3055        reject_transaction_control_statements(&statements, "inline.execute_batch")?;
3056        let prepared = prepare_batch_statements(self.conn(), &statements)
3057            .map_err(|e| map_rusqlite_err(e, "inline.execute_batch"))?;
3058        execute_prepared_batch(
3059            self.conn(),
3060            prepared,
3061            &statements,
3062            self.event_rows.as_deref(),
3063        )
3064        .map_err(|e| map_rusqlite_err(e, "inline.execute_batch"))
3065    }
3066
3067    async fn execute_script(&mut self, script: String) -> khive_storage::types::StorageResult<()> {
3068        // Boundary: this raw script path is internal maintenance only and is
3069        // outside the `execute_batch` transaction-control contract.
3070        self.conn()
3071            .execute_batch(&script)
3072            .map_err(|e| map_rusqlite_err(e, "inline.execute_script"))
3073    }
3074}
3075
3076/// Poll `fut` exactly once with a no-op waker and return its output.
3077///
3078/// Only sound for futures that never actually suspend — every caller in
3079/// this module drives an [`InlineWriter`], whose methods are pure
3080/// synchronous rusqlite calls with no real `.await` point.
3081///
3082/// ADR-067 Component A: this used to
3083/// `unreachable!()`-panic on `Poll::Pending`, and a panicking closure
3084/// running inside the writer task's `spawn_blocking` (see
3085/// `SqlBridge::atomic_unit`'s flag-on branch) would surface as a
3086/// `JoinError` in `run_writer_task`, which is treated as fatal — the writer
3087/// task exits and every subsequent `WriterTaskHandle::send` on this pool
3088/// fails for the rest of the process. A future `atomic_unit` caller whose
3089/// closure ever gains a real suspend point (this file's own contract
3090/// already forbids it, but the invariant is enforced by convention, not the
3091/// type system) would take down the writer task for the whole daemon.
3092/// Returning `Err` instead lets `Pending` flow through the SAME error path
3093/// as any other `atomic_unit` op failure: `WriteRequest::execute_and_reply`
3094/// treats it as an ordinary `Err`, issues `ROLLBACK` on the writer task's
3095/// held transaction, replies the error to the caller, and the writer task's
3096/// `spawn_blocking` closure returns normally (not via panic) — so the task
3097/// keeps draining subsequent requests instead of dying with the whole pool.
3098fn block_on_sync<F: std::future::Future>(fut: F) -> Result<F::Output, StorageError> {
3099    use std::task::{Context, Poll, RawWaker, RawWakerVTable, Waker};
3100
3101    fn no_op(_: *const ()) {}
3102    fn clone_waker(_: *const ()) -> RawWaker {
3103        RawWaker::new(std::ptr::null(), &VTABLE)
3104    }
3105    static VTABLE: RawWakerVTable = RawWakerVTable::new(clone_waker, no_op, no_op, no_op);
3106
3107    // SAFETY: every `RawWakerVTable` function is a no-op that never
3108    // dereferences the data pointer, so a null data pointer is sound.
3109    let raw_waker = RawWaker::new(std::ptr::null(), &VTABLE);
3110    let waker = unsafe { Waker::from_raw(raw_waker) };
3111    let mut cx = Context::from_waker(&waker);
3112
3113    let mut fut = std::pin::pin!(fut);
3114    match fut.as_mut().poll(&mut cx) {
3115        Poll::Ready(v) => Ok(v),
3116        Poll::Pending => {
3117            tracing::error!(
3118                "block_on_sync: atomic_unit future suspended on its first poll — \
3119                 the closure passed to SqlAccess::atomic_unit must be non-blocking \
3120                 (synchronous InlineWriter calls only, no real .await point)"
3121            );
3122            Err(StorageError::Internal(
3123                "atomic_unit future suspended — closure must be non-blocking".to_string(),
3124            ))
3125        }
3126    }
3127}
3128
3129// =============================================================================
3130// SqlBridge: the SqlAccess implementor
3131// =============================================================================
3132
3133/// Bridges `ConnectionPool` to `khive_storage::SqlAccess`.
3134///
3135/// Dispatches based on whether the pool is file-backed or in-memory:
3136/// - File-backed: pooled ordinary reader operations and a lazy standalone
3137///   connection only for explicitly admitted multi-call read transactions,
3138///   plus standalone writer connections capped at one live handle; atomic
3139///   units drive a single registered raw transaction span.
3140/// - In-memory: pool-backed connections per query (single shared connection).
3141pub struct SqlBridge {
3142    pool: Arc<ConnectionPool>,
3143    is_file_backed: bool,
3144}
3145
3146impl SqlBridge {
3147    /// Create a new bridge wrapping the given pool.
3148    pub fn new(pool: Arc<ConnectionPool>, _is_file_backed: bool) -> Self {
3149        // The legacy hint is not an authority for choosing an unguarded
3150        // in-memory writer. A caller cannot route a file-backed pool through
3151        // PoolBackedWriter by supplying `false`.
3152        let is_file_backed = pool.canonical_path().is_some();
3153        Self {
3154            pool,
3155            is_file_backed,
3156        }
3157    }
3158}
3159
3160#[async_trait]
3161impl khive_storage::SqlAccess for SqlBridge {
3162    fn database_path(&self) -> Option<std::path::PathBuf> {
3163        self.pool.canonical_path().map(std::path::Path::to_path_buf)
3164    }
3165
3166    async fn reader(
3167        &self,
3168    ) -> khive_storage::types::StorageResult<Box<dyn khive_storage::SqlReader>> {
3169        if self.is_file_backed {
3170            Ok(Box::new(SqliteReader {
3171                handle: None,
3172                pool: Arc::clone(&self.pool),
3173                poisoned: false,
3174            }))
3175        } else {
3176            Ok(Box::new(PoolBackedReader {
3177                pool: Arc::clone(&self.pool),
3178                transaction: None,
3179            }))
3180        }
3181    }
3182
3183    async fn writer(
3184        &self,
3185    ) -> khive_storage::types::StorageResult<Box<dyn khive_storage::SqlWriter>> {
3186        if self.is_file_backed {
3187            if self.pool.config().read_only {
3188                return Err(StorageError::Pool {
3189                    operation: "writer".into(),
3190                    message: "backend is read-only".into(),
3191                });
3192            }
3193            let db = crate::timeout_sink::db_label(&self.pool);
3194            // ADR-136 D1 gate 1: queue-first. The handle lookup runs BEFORE
3195            // any standalone connection is opened, and a lookup failure is
3196            // propagated (never silently degraded) when strict routing is
3197            // on. Only the flag-off/degraded case still opens a standalone
3198            // connection.
3199            let writer_task = match self.pool.writer_task_handle() {
3200                Ok(handle) => handle,
3201                Err(e) => {
3202                    if self.pool.config().write_routing_strict {
3203                        return Err(e);
3204                    }
3205                    tracing::warn!(
3206                        error = %e,
3207                        "KHIVE_WRITE_ROUTING is not strict; writer() degrades to the \
3208                         standalone-connection path"
3209                    );
3210                    None
3211                }
3212            };
3213            if writer_task.is_none() && self.pool.config().write_routing_strict {
3214                return Err(StorageError::Pool {
3215                    operation: "writer".into(),
3216                    message: "KHIVE_WRITE_ROUTING=strict but no writer-task handle is \
3217                              available; refusing to fall back to a direct connection"
3218                        .into(),
3219                });
3220            }
3221            if writer_task.is_none() && self.pool.write_queue_active() {
3222                // The queue is enabled but this call didn't get a handle
3223                // (spawn/runtime degrade) — a direct-route violation in the
3224                // making once this writer's execute*/query* methods run.
3225                // In-memory pools are excluded: they never spawn a writer
3226                // task by documented design (explicit `Some(true)` degrades),
3227                // so a violation row there would be noise, not signal.
3228                crate::timeout_sink::emit_direct_route_violation(
3229                    &db,
3230                    crate::timeout_sink::Site::DirectRouteSqlBridgeWriter,
3231                );
3232            }
3233            // A standalone read-write connection is opened only when there is
3234            // no queue handle to route writes through — `SqliteWriter`'s
3235            // `SqlReader` methods (`query_row`/`query_all`/`query_page`)
3236            // use pooled readers in the queue-backed case and lazily open the
3237            // closed standalone exception only for an explicit deferred read
3238            // transaction. Production callers do read through a `writer()`
3239            // handle, so both routes are live.
3240            // The standalone open acquires the pool-wide one-permit writer
3241            // budget first, and the permit travels in the handle for the
3242            // handle's whole lifetime — a queue-backed handle holds no
3243            // writer permit (its writes route through the writer task), so
3244            // this budget caps exactly the standalone read-write
3245            // connections.
3246            let handle = if writer_task.is_none() {
3247                let handle_slot = acquire_handle_slot(
3248                    self.pool.sql_bridge_writer_slots(),
3249                    self.pool.config().checkout_timeout,
3250                    "sql_bridge.writer_handle",
3251                    SlotTimeoutClass::Admission,
3252                )
3253                .await?;
3254                let (conn, handle_slot) =
3255                    open_standalone_writer_on_blocking(Arc::clone(&self.pool), handle_slot).await?;
3256                Some(StandaloneHandle {
3257                    conn,
3258                    _retained_slot: Some(handle_slot),
3259                    read_transaction_slot: None,
3260                })
3261            } else {
3262                None
3263            };
3264            Ok(Box::new(SqliteWriter {
3265                observe_direct_errors: true,
3266                event_rows: None,
3267                handle,
3268                writer_task,
3269                origin: self.pool.origin(),
3270                db,
3271                pool: Arc::clone(&self.pool),
3272                held_lease: None,
3273            }))
3274        } else {
3275            Ok(Box::new(PoolBackedWriter {
3276                pool: Arc::clone(&self.pool),
3277            }))
3278        }
3279    }
3280
3281    /// Implements the trait's atomic-unit suspend-free invariant
3282    /// (`SqlAccess::atomic_unit`'s doc comment): on the queued and in-memory branches,
3283    /// `op` is driven through `block_on_sync` on an `InlineWriter` — a
3284    /// single-poll driver that returns `Err` the instant `op`'s future is
3285    /// `Pending` instead of ever actually suspending. `op` must therefore
3286    /// issue only synchronous DML; see `InlineWriter`'s and
3287    /// `block_on_sync`'s doc comments for the full mechanics and why this
3288    /// restriction is load-bearing (a suspended poll inside the writer
3289    /// task's `spawn_blocking` would otherwise block that task on external
3290    /// async work while holding the single write connection).
3291    async fn atomic_unit(
3292        &self,
3293        op: AtomicUnitOp,
3294    ) -> khive_storage::types::StorageResult<Box<dyn Any + Send>> {
3295        let event_rows = Arc::new(AtomicEventRows::default());
3296        let result = async {
3297            if self.is_file_backed {
3298                if self.pool.config().read_only {
3299                    return Err(StorageError::Pool {
3300                        operation: "atomic_unit".into(),
3301                        message: "backend is read-only".into(),
3302                    });
3303                }
3304                // Best-effort, same guard `writer()` uses: `Ok(None)` on flag-off;
3305                // `Err(WriterTaskNoRuntime)` propagates loud rather than silently
3306                // falling back to a competing connection from a sync caller. ADR-136
3307                // D1 gate 3: `Ok(None)` under strict routing is ALSO a fail-closed
3308                // error (queue was requested but unavailable), not just a degrade.
3309                let handle = self.pool.writer_task_handle()?;
3310                if handle.is_none() && self.pool.config().write_routing_strict {
3311                    return Err(StorageError::Pool {
3312                        operation: "atomic_unit".into(),
3313                        message: "KHIVE_WRITE_ROUTING=strict but no writer-task handle is \
3314                              available; refusing to fall back to a direct connection"
3315                            .into(),
3316                    });
3317                }
3318                if handle.is_none() && self.pool.write_queue_active() {
3319                    crate::timeout_sink::emit_direct_route_violation(
3320                        &crate::timeout_sink::db_label(&self.pool),
3321                        crate::timeout_sink::Site::DirectRouteAtomicUnit,
3322                    );
3323                }
3324                if let Some(writer_task) = handle {
3325                    // Flag-on: ONE queued WriteRequest. `run_writer_task` already
3326                    // has an open `BEGIN IMMEDIATE` on its dedicated connection
3327                    // before this closure runs and issues `COMMIT`/`ROLLBACK`
3328                    // after it returns — `op` must not (and, via `InlineWriter`,
3329                    // does not) issue its own transaction control.
3330                    let pending_event_rows = Arc::clone(&event_rows);
3331                    return writer_task
3332                        .send_bounded(move |conn| {
3333                            let mut inline = InlineWriter {
3334                                event_rows: Some(Arc::clone(&pending_event_rows)),
3335                                conn: conn as *const rusqlite::Connection,
3336                            };
3337                            // Flatten: `block_on_sync` now returns `Result<F::Output,
3338                            // StorageError>` (outer = "did the future actually
3339                            // resolve on first poll", inner = the op's own
3340                            // `StorageResult`) instead of panicking on `Pending`
3341                            // (ADR-067 Component A). Either
3342                            // error flows through this closure's ordinary `Err`
3343                            // return, which `WriteRequest::execute_and_reply`
3344                            // already turns into a normal ROLLBACK + error reply —
3345                            // no panic, so the writer task survives.
3346                            match block_on_sync(op(&mut inline)) {
3347                                Ok(inner) => inner,
3348                                Err(e) => Err(e),
3349                            }
3350                        })
3351                        .await;
3352                }
3353                // Flag-off (or no writer task available): manual
3354                // BEGIN IMMEDIATE/COMMIT/ROLLBACK on a standalone writer —
3355                // byte-for-byte the pre-ADR-067 shape.
3356                //
3357                // Contract: this acquire waits on the pool-wide one-permit
3358                // writer-handle budget — the same permit a live `writer()` handle
3359                // holds for its lifetime — so it times out with
3360                // `StorageError::AdmissionTimeout` after `checkout_timeout` while a writer
3361                // handle is checked out (and a `writer()` call times out while
3362                // this unit runs). Callers must not hold a boxed writer handle
3363                // across an `atomic_unit()` call on the same pool; drop the
3364                // handle first. The `writer_task` branch above never touches this
3365                // budget.
3366                let handle_slot = acquire_handle_slot(
3367                    self.pool.sql_bridge_writer_slots(),
3368                    self.pool.config().checkout_timeout,
3369                    "sql_bridge.atomic_unit_handle",
3370                    SlotTimeoutClass::Admission,
3371                )
3372                .await?;
3373                // The unit's lease spans its BEGIN, statements and COMMIT or
3374                // ROLLBACK, as one queued request does on the write-queue path.
3375                let unit_lease = acquire_unit_lease(Arc::clone(&self.pool)).await?;
3376                let (conn, handle_slot) =
3377                    open_standalone_writer_on_blocking(Arc::clone(&self.pool), handle_slot).await?;
3378                let mut writer = SqliteWriter {
3379                    observe_direct_errors: false,
3380                    event_rows: Some(Arc::clone(&event_rows)),
3381                    handle: Some(StandaloneHandle {
3382                        conn,
3383                        _retained_slot: Some(handle_slot),
3384                        read_transaction_slot: None,
3385                    }),
3386                    writer_task: None,
3387                    origin: self.pool.origin(),
3388                    db: crate::timeout_sink::db_label(&self.pool),
3389                    pool: Arc::clone(&self.pool),
3390                    held_lease: unit_lease,
3391                };
3392                run_manual_atomic_unit(&mut writer, op, self.pool.origin())
3393                    .await
3394                    .inspect_err(|error| self.pool.record_direct_writer_error(error))
3395            } else {
3396                // Every statement shares one connection. Keep its guard through
3397                // commit/rollback so other units and ordinary writes cannot join it.
3398                let pool = Arc::clone(&self.pool);
3399                let pending_event_rows = Arc::clone(&event_rows);
3400                tokio::task::spawn_blocking(move || {
3401                    let guard = pool.try_writer().map_err(|error: SqliteError| {
3402                        StorageError::driver(StorageCapability::Sql, "atomic_unit", error)
3403                    })?;
3404                    let conn = guard.conn();
3405                    if !conn.is_autocommit() {
3406                        pool.retire_pooled_writer(conn);
3407                        return Err(StorageError::writer_task_terminated(
3408                            khive_storage::WriterTaskRequestState::SideEffectsUnknown,
3409                        ));
3410                    }
3411                    if let Err(error) = conn.execute_batch("BEGIN IMMEDIATE") {
3412                        if !conn.is_autocommit() {
3413                            pool.retire_pooled_writer(conn);
3414                            return Err(StorageError::writer_task_terminated(
3415                                khive_storage::WriterTaskRequestState::SideEffectsUnknown,
3416                            ));
3417                        }
3418                        return Err(map_rusqlite_err(error, "atomic_unit.begin"))
3419                            .inspect_err(|error| pool.record_direct_writer_error(error));
3420                    }
3421                    let _tx_handle = khive_storage::tx_registry::register_scoped(
3422                        Some("atomic_unit".to_string()),
3423                        pool.origin(),
3424                    );
3425                    let (result, terminal_state) = crate::writer_task::execute_wrapped_transaction(
3426                        conn,
3427                        "atomic_unit.commit",
3428                        |conn| {
3429                            let mut inline = InlineWriter {
3430                                event_rows: Some(Arc::clone(&pending_event_rows)),
3431                                conn: conn as *const rusqlite::Connection,
3432                            };
3433                            block_on_sync(op(&mut inline)).and_then(|result| result)
3434                        },
3435                    );
3436                    if terminal_state.is_some() {
3437                        pool.retire_pooled_writer(conn);
3438                    }
3439                    result.inspect_err(|error| pool.record_direct_writer_error(error))
3440                })
3441                .await
3442                .map_err(|error| {
3443                    StorageError::driver(StorageCapability::Sql, "atomic_unit", error)
3444                })?
3445            }
3446        }
3447        .await;
3448        khive_storage::usage::account_event_write(
3449            result.as_ref().map(|_| event_rows.committed_rows()),
3450        );
3451        result
3452    }
3453}
3454
3455#[cfg(test)]
3456#[path = "sql_bridge_tests.rs"]
3457mod tests;
3458
3459#[cfg(test)]
3460#[path = "sql_bridge/direct_busy_tests.rs"]
3461mod direct_busy_tests;
3462
3463#[cfg(test)]
3464#[path = "sql_bridge/settlement_hazard_tests.rs"]
3465mod settlement_hazard_tests;