Skip to main content

khive_db/
pool.rs

1//! Connection pool for SQLite: one exclusive writer, N concurrent readers.
2use crossbeam_queue::ArrayQueue;
3use parking_lot::{Condvar, Mutex};
4use rusqlite::hooks::{AuthContext, Authorization};
5use rusqlite::{Connection, OpenFlags};
6use std::fs;
7use std::io::Read as _;
8use std::ops::{Deref, DerefMut};
9use std::path::{Path, PathBuf};
10use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
11use std::sync::{Arc, OnceLock};
12use std::thread;
13use std::time::{Duration, Instant};
14use tokio::sync::Semaphore;
15
16use crate::error::SqliteError;
17use crate::writer_task::WriterTaskHandle;
18use khive_storage::error::StorageError;
19use khive_storage::tx_registry::{DbIdentity, TxOrigin};
20use khive_storage::StorageCapability;
21
22const CACHE_SIZE_KIB: &str = "-65536";
23const MMAP_SIZE_BYTES: &str = "1073741824";
24const DEFAULT_READER_CAP: usize = 8;
25
26const DEFAULT_JOURNAL_SIZE_LIMIT_BYTES: i64 = 67_108_864; // 64 MiB
27const DEFAULT_WRITE_QUEUE_CAPACITY: usize = 256;
28
29/// Bounded WAL autocheckpoint applied to writer-capable connections while no
30/// dedicated checkpoint owner has claimed the pool (4,000 pages ≈ 16 MiB at
31/// SQLite's default 4 KiB page size — SQLite's historic behaviour for this
32/// pool). Not a tuning parameter: there is no config field or environment
33/// override, and the only way to change the effective value is an actual
34/// ownership claim ([`ConnectionPool::claim_checkpoint_ownership`]), which a
35/// runtime may make only when it really runs the scheduled checkpoint task.
36pub(crate) const FALLBACK_WAL_AUTOCHECKPOINT_PAGES: u32 = 4_000;
37
38#[derive(Clone, Copy, Debug, PartialEq, Eq)]
39enum CheckpointOwnership {
40    Unclaimed,
41    Claiming,
42    Claimed,
43}
44
45struct CheckpointOwnershipState {
46    phase: CheckpointOwnership,
47    #[cfg(test)]
48    connection_waiters: usize,
49}
50
51#[cfg(test)]
52struct CheckpointConnectionConfigPause {
53    selected: std::sync::Barrier,
54    resume: std::sync::Barrier,
55}
56
57#[cfg(test)]
58impl CheckpointConnectionConfigPause {
59    fn new() -> Self {
60        Self {
61            selected: std::sync::Barrier::new(2),
62            resume: std::sync::Barrier::new(2),
63        }
64    }
65}
66
67struct CheckpointOwnershipGate {
68    state: Mutex<CheckpointOwnershipState>,
69    changed: Condvar,
70    #[cfg(test)]
71    connection_config_pause: Mutex<Option<Arc<CheckpointConnectionConfigPause>>>,
72    #[cfg(test)]
73    claim_lock_observed: Mutex<Option<std::sync::mpsc::SyncSender<bool>>>,
74}
75
76impl CheckpointOwnershipGate {
77    fn new() -> Self {
78        Self {
79            state: Mutex::new(CheckpointOwnershipState {
80                phase: CheckpointOwnership::Unclaimed,
81                #[cfg(test)]
82                connection_waiters: 0,
83            }),
84            changed: Condvar::new(),
85            #[cfg(test)]
86            connection_config_pause: Mutex::new(None),
87            #[cfg(test)]
88            claim_lock_observed: Mutex::new(None),
89        }
90    }
91
92    /// Join an in-flight claim, or become the one caller that configures it.
93    /// Returns `false` when another caller has already completed the claim.
94    fn begin_claim(&self) -> bool {
95        #[cfg(test)]
96        let claim_lock_observed = self.claim_lock_observed.lock().take();
97        #[cfg(test)]
98        let mut state = if let Some(observed) = claim_lock_observed {
99            match self.state.try_lock() {
100                Some(state) => {
101                    let _ = observed.send(false);
102                    state
103                }
104                None => {
105                    let _ = observed.send(true);
106                    self.state.lock()
107                }
108            }
109        } else {
110            self.state.lock()
111        };
112        #[cfg(not(test))]
113        let mut state = self.state.lock();
114        loop {
115            match state.phase {
116                CheckpointOwnership::Unclaimed => {
117                    state.phase = CheckpointOwnership::Claiming;
118                    self.changed.notify_all();
119                    return true;
120                }
121                CheckpointOwnership::Claiming => self.changed.wait(&mut state),
122                CheckpointOwnership::Claimed => return false,
123            }
124        }
125    }
126
127    fn finish_claim(&self, succeeded: bool) {
128        let mut state = self.state.lock();
129        debug_assert_eq!(state.phase, CheckpointOwnership::Claiming);
130        state.phase = if succeeded {
131            CheckpointOwnership::Claimed
132        } else {
133            CheckpointOwnership::Unclaimed
134        };
135        self.changed.notify_all();
136    }
137
138    fn settled_state(&self) -> parking_lot::MutexGuard<'_, CheckpointOwnershipState> {
139        let mut state = self.state.lock();
140        while state.phase == CheckpointOwnership::Claiming {
141            #[cfg(test)]
142            {
143                state.connection_waiters += 1;
144                self.changed.notify_all();
145            }
146            self.changed.wait(&mut state);
147            #[cfg(test)]
148            {
149                state.connection_waiters -= 1;
150                self.changed.notify_all();
151            }
152        }
153        state
154    }
155
156    #[cfg(test)]
157    fn wal_autocheckpoint_pages(&self) -> u32 {
158        let state = self.settled_state();
159        match state.phase {
160            CheckpointOwnership::Unclaimed => FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
161            CheckpointOwnership::Claimed => 0,
162            CheckpointOwnership::Claiming => unreachable!("claim wait must settle the state"),
163        }
164    }
165
166    /// Wait for any in-flight claim, select the resulting posture, and retain
167    /// the gate until SQLite has applied that connection-local PRAGMA. A claim
168    /// therefore linearizes entirely before or after this configuration,
169    /// never between its state sample and side effect.
170    fn configure_wal_autocheckpoint(&self, conn: &Connection) -> Result<(), SqliteError> {
171        let state = self.settled_state();
172        let pages = match state.phase {
173            CheckpointOwnership::Unclaimed => FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
174            CheckpointOwnership::Claimed => 0,
175            CheckpointOwnership::Claiming => unreachable!("claim wait must settle the state"),
176        };
177        #[cfg(test)]
178        if let Some(pause) = self.connection_config_pause.lock().take() {
179            pause.selected.wait();
180            pause.resume.wait();
181        }
182        conn.pragma_update(None, "wal_autocheckpoint", pages)?;
183        drop(state);
184        Ok(())
185    }
186}
187
188fn deny_retired_writer(_context: AuthContext<'_>) -> Authorization {
189    Authorization::Deny
190}
191
192pub(crate) const TEST_HARNESS_ENV: &str = "KHIVE_TEST_HARNESS";
193
194/// Configuration for the connection pool.
195#[derive(Clone, Debug)]
196pub struct PoolConfig {
197    /// Database path. None = in-memory (pool degrades to single connection).
198    pub path: Option<PathBuf>,
199    /// Number of reader connections (default: min(num_cpus, 8)).
200    pub max_readers: usize,
201    /// WAL mode (must be true for pooling to work; default: true).
202    pub wal_mode: bool,
203    /// Busy timeout per connection (default: 30s).
204    ///
205    /// Overridable via `KHIVE_BUSY_TIMEOUT_SECS`.
206    pub busy_timeout: Duration,
207    /// Time to wait for a reader connection before returning an error (default: 5s).
208    ///
209    /// Overridable via `KHIVE_CHECKOUT_TIMEOUT_SECS`.
210    pub checkout_timeout: Duration,
211    /// Maximum WAL journal size in bytes before SQLite resets the WAL.
212    ///
213    /// Maps to `PRAGMA journal_size_limit`. Default: 64 MiB.
214    ///
215    /// Overridable via `KHIVE_JOURNAL_SIZE_LIMIT_BYTES`.
216    pub journal_size_limit_bytes: i64,
217    /// Open the database read-only (default: false).
218    ///
219    /// When true, the pool's writer connection is opened with
220    /// `SQLITE_OPEN_READ_ONLY` (no `SQLITE_OPEN_CREATE`, so a missing path is
221    /// rejected instead of created) and `PRAGMA query_only = ON` is set on
222    /// every connection that can execute SQL. Reader connections are already
223    /// opened read-only regardless of this flag.
224    pub read_only: bool,
225    /// Route migrated store write paths through the single-writer
226    /// `WriterTask` channel (ADR-067 Component A) instead of the legacy
227    /// per-call pool-mutex/standalone-connection path. Enabled by default
228    /// for file-backed pools when unset; explicit override always wins.
229    /// That default is a compatibility-routing posture subordinate to
230    /// ADR-135 Amendment 1 and ADR-136 D1/D2 — the strict-routing default
231    /// flip has NOT happened.
232    ///
233    /// The store layer resolves all of its routed write paths at write time;
234    /// the classification table in `writer_task.rs` remains the authoritative
235    /// inventory. This tranche does not claim the repository-wide
236    /// single-writer guarantee: direct runtime-orchestration call sites remain
237    /// #1847 follow-up work, and the strict default is still evidence-gated.
238    ///
239    /// `None` means the caller expressed no preference: [`ConnectionPool::new`]
240    /// resolves it once `path` is known, defaulting to `true` for file-backed
241    /// pools and `false` for in-memory ones. `Some(_)` is an explicit
242    /// preference and always wins, in both directions, over that default.
243    /// An explicit `Some(true)` on an in-memory pool is accepted DELIBERATELY
244    /// and emits a warning before degrading to the legacy path — an in-memory
245    /// pool cannot host a writer task (`writer_task::spawn`'s
246    /// standalone-connection open fails); see
247    /// `ConnectionPool::writer_task_handle` and the
248    /// `explicit_true_stays_on_for_memory_backed_pool` test.
249    ///
250    /// Overridable via `KHIVE_WRITE_QUEUE` (`"1"` or `"true"`,
251    /// case-insensitive, sets `Some(true)`; any other value sets `Some(false)`;
252    /// unset leaves it `None`).
253    pub write_queue_enabled: Option<bool>,
254    /// Bounded channel capacity for the `WriterTask` write queue.
255    ///
256    /// Overridable via `KHIVE_WRITE_QUEUE_CAPACITY`. Default: 256 pending
257    /// operations (ADR-067 Component A recommended default).
258    pub write_queue_capacity: usize,
259    /// ADR-136 D1: when `true`, every covered store write path that would
260    /// otherwise silently degrade to the legacy pool-mutex/standalone-
261    /// connection path on a missing or failed `WriterTask` handle instead
262    /// returns an error.
263    /// Exercises the store-layer routing tranche toward ADR-135 F2's
264    /// strict-routing precondition without changing behavior for callers that
265    /// never set the env var.
266    ///
267    /// Overridable via `KHIVE_WRITE_ROUTING` (value `"strict"`,
268    /// case-insensitive; anything else, or unset, leaves this `false`).
269    pub write_routing_strict: bool,
270    /// Dedicated admission deadline (ADR-131 Decision 2) bounding ONLY the
271    /// wait for capacity on the `WriterTask` write queue —
272    /// [`WriterTaskHandle::send_bounded`]/`send_top_level_bounded`'s default
273    /// timeout. Distinct from `checkout_timeout`, which bounds reader/pool
274    /// checkout instead; the two authorities used to be conflated (#1382,
275    /// #1643) before this field existed.
276    ///
277    /// Default: 2000 ms. Validated at [`ConnectionPool::new`] to fall in
278    /// `[100, 10000]` ms; a value outside that range is a configuration
279    /// error (`SqliteError::InvalidConfig`), never silently clamped into
280    /// range.
281    ///
282    /// Overridable via `KHIVE_WRITE_ADMISSION_DEADLINE_MS`.
283    pub write_admission_deadline_ms: u64,
284    /// Maximum age an explicit cached-reader read transaction
285    /// (`sql_bridge`'s `BEGIN`-then-reuse path) may reach before its next use
286    /// is refused and it is rolled back instead of extending its WAL
287    /// snapshot further (#1846). Shares `KHIVE_TX_MAX_AGE_SECS` with the
288    /// ADR-091 Plank 1 visibility sweep in `checkpoint.rs` so one knob
289    /// governs both when an operator is warned about a stale reader and when
290    /// that reader's snapshot is actually released.
291    ///
292    /// Overridable via `KHIVE_TX_MAX_AGE_SECS`. Default: 120 seconds.
293    pub read_tx_max_age: Duration,
294}
295
296/// ADR-131 Decision 2's validated range for `write_admission_deadline_ms`.
297const WRITE_ADMISSION_DEADLINE_MS_RANGE: std::ops::RangeInclusive<u64> = 100..=10_000;
298const DEFAULT_WRITE_ADMISSION_DEADLINE_MS: u64 = 2000;
299
300impl Default for PoolConfig {
301    fn default() -> Self {
302        Self {
303            path: None,
304            max_readers: std::thread::available_parallelism()
305                .map(|n| n.get())
306                .unwrap_or(1)
307                .clamp(1, DEFAULT_READER_CAP),
308            wal_mode: true,
309            busy_timeout: Duration::from_secs(
310                std::env::var("KHIVE_BUSY_TIMEOUT_SECS")
311                    .ok()
312                    .and_then(|v| v.parse::<u64>().ok())
313                    .unwrap_or(30),
314            ),
315            checkout_timeout: Duration::from_secs(
316                std::env::var("KHIVE_CHECKOUT_TIMEOUT_SECS")
317                    .ok()
318                    .and_then(|v| v.parse::<u64>().ok())
319                    .unwrap_or(5),
320            ),
321            journal_size_limit_bytes: std::env::var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES")
322                .ok()
323                .and_then(|v| v.parse::<i64>().ok())
324                .unwrap_or(DEFAULT_JOURNAL_SIZE_LIMIT_BYTES),
325            read_only: false,
326            // `var_os`, not `var`: the documented contract is "any SET value
327            // other than 1/true means Some(false)" — a set-but-non-Unicode
328            // value must count as set (var() would return Err and silently
329            // fall through to the file-backed default of enabled).
330            write_queue_enabled: std::env::var_os("KHIVE_WRITE_QUEUE").map(|v| {
331                v.to_str()
332                    .is_some_and(|v| v == "1" || v.eq_ignore_ascii_case("true"))
333            }),
334            write_queue_capacity: std::env::var("KHIVE_WRITE_QUEUE_CAPACITY")
335                .ok()
336                .and_then(|v| v.parse::<usize>().ok())
337                .filter(|&n| n > 0)
338                .unwrap_or(DEFAULT_WRITE_QUEUE_CAPACITY),
339            write_routing_strict: std::env::var("KHIVE_WRITE_ROUTING")
340                .map(|v| v.eq_ignore_ascii_case("strict"))
341                .unwrap_or(false),
342            write_admission_deadline_ms: std::env::var("KHIVE_WRITE_ADMISSION_DEADLINE_MS")
343                .ok()
344                .and_then(|v| v.parse::<u64>().ok())
345                .unwrap_or(DEFAULT_WRITE_ADMISSION_DEADLINE_MS),
346            read_tx_max_age: crate::checkpoint::tx_age_thresholds_from_env(
347                Duration::from_secs(30),
348                Duration::from_secs(120),
349            )
350            .1,
351        }
352    }
353}
354
355/// Prevent Cargo-launched tests and test subprocesses from opening the
356/// operator's default data tree in every build profile. Activation is solely
357/// the runtime `KHIVE_TEST_HARNESS=1` marker; production/installed binaries do
358/// not receive that workspace Cargo environment.
359///
360/// There is deliberately no environment override: any inheritable escape
361/// hatch set for one Cargo invocation leaks into the next `cargo test` in the
362/// same shell and re-opens the store the guard exists to protect. A deliberate
363/// session against the real store runs the built binary directly (for example
364/// `target/release/...` or an installed binary), which never receives the
365/// workspace Cargo environment and therefore never trips this guard.
366/// Existing path ancestors are canonicalized before comparison, resolving
367/// traversal, symlinks, and filesystem-provided case (including APFS case
368/// folding). Missing trailing components remain lexical because they have no
369/// filesystem identity yet. SQLite URI paths are rejected rather than trying
370/// to reproduce SQLite's URI normalization rules.
371fn refuse_home_data_store_in_tests(config: &PoolConfig) -> Result<(), SqliteError> {
372    if std::env::var(TEST_HARNESS_ENV).as_deref() != Ok("1") {
373        return Ok(());
374    }
375
376    let Some(path) = config.path.as_deref() else {
377        return Ok(());
378    };
379    if path
380        .as_os_str()
381        .as_encoded_bytes()
382        .get(..5)
383        .is_some_and(|prefix| prefix.eq_ignore_ascii_case(b"file:"))
384    {
385        return Err(SqliteError::InvalidData(format!(
386            "test harness refused SQLite URI database path {}; use a filesystem path outside \
387             HOME/.khive (deliberate sessions against a real store run the built binary \
388             directly, outside the Cargo test environment)",
389            path.display()
390        )));
391    }
392
393    let Some(home) = std::env::var_os("HOME") else {
394        return Ok(());
395    };
396    let canonical_path = canonicalize_deepest_existing(path)?;
397    let canonical_home_data_dir =
398        canonicalize_deepest_existing(&PathBuf::from(home).join(".khive"))?;
399    if canonical_path.starts_with(&canonical_home_data_dir) {
400        return Err(SqliteError::InvalidData(format!(
401            "test harness refused to open SQLite database under HOME/.khive: {} \
402             (deliberate sessions against a real store run the built binary directly, \
403             outside the Cargo test environment)",
404            canonical_path.display()
405        )));
406    }
407    Ok(())
408}
409
410fn canonicalize_deepest_existing(path: &Path) -> Result<PathBuf, SqliteError> {
411    let absolute = if path.is_absolute() {
412        path.to_path_buf()
413    } else {
414        std::env::current_dir().map_err(SqliteError::Io)?.join(path)
415    };
416
417    for ancestor in absolute.ancestors() {
418        match fs::canonicalize(ancestor) {
419            Ok(mut canonical) => {
420                let missing = absolute.strip_prefix(ancestor).map_err(|error| {
421                    SqliteError::InvalidData(format!(
422                        "failed to preserve missing path components for {}: {error}",
423                        absolute.display()
424                    ))
425                })?;
426                canonical.push(missing);
427                return Ok(canonical);
428            }
429            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
430            Err(error) => {
431                return Err(SqliteError::InvalidData(format!(
432                    "failed to canonicalize database path ancestor {}: {error}",
433                    ancestor.display()
434                )));
435            }
436        }
437    }
438
439    Err(SqliteError::InvalidData(format!(
440        "database path has no canonicalizable ancestor: {}",
441        absolute.display()
442    )))
443}
444
445/// Enforce ADR-131 Decision 2's `write_admission_deadline_ms` bound at
446/// configuration load: `[100, 10000]` ms, rejected rather than clamped when
447/// out of range so a misconfiguration is never silently reinterpreted as a
448/// different deadline than the operator asked for.
449fn validate_write_admission_deadline(deadline_ms: u64) -> Result<(), SqliteError> {
450    if WRITE_ADMISSION_DEADLINE_MS_RANGE.contains(&deadline_ms) {
451        return Ok(());
452    }
453    Err(SqliteError::InvalidConfig(format!(
454        "write_admission_deadline_ms must be in [{}, {}] ms, got {deadline_ms}",
455        WRITE_ADMISSION_DEADLINE_MS_RANGE.start(),
456        WRITE_ADMISSION_DEADLINE_MS_RANGE.end()
457    )))
458}
459
460/// A read-write connection pool for SQLite.
461///
462/// Architecture:
463/// - 1 writer connection protected by a Mutex (exclusive access)
464/// - N reader connections in a lock-free queue (concurrent access)
465/// - All connections share the same database file in WAL mode
466///
467/// Writable in-memory databases, or writable file databases when WAL mode is
468/// disabled/unavailable, degrade to single-connection mode and route all
469/// operations through the writer connection. A file-backed read-only pool
470/// always retains at least one dedicated read-only connection: rollback-journal
471/// snapshots do not need WAL to support concurrent readers, and inspection must
472/// never alias a read onto the query-only writer slot.
473pub struct ConnectionPool {
474    writer: Arc<Mutex<Connection>>,
475    /// Three-state gate for whether the ADR-091 scheduled task has claimed
476    /// routine WAL reclamation for this pool. Until claimed, every
477    /// writer-capable connection keeps a bounded SQLite autocheckpoint
478    /// ([`FALLBACK_WAL_AUTOCHECKPOINT_PAGES`]) so a writable pool without a
479    /// checkpoint task cannot grow its WAL without bound. After
480    /// [`Self::claim_checkpoint_ownership`], writer-capable connections open
481    /// with `wal_autocheckpoint = 0` and routine checkpoint I/O stays off
482    /// application commit paths.
483    checkpoint_ownership: CheckpointOwnershipGate,
484    /// Fail-closed guard for the legacy pool-mutex writer. A transaction
485    /// owner retires this connection after a body panic or when it cannot
486    /// prove that finalization restored autocommit mode; subsequent checkouts
487    /// must never reuse it.
488    pooled_writer_retired: AtomicBool,
489    /// Process-local writer acquisition counters shared with the pool's
490    /// lifetime-owned writer task. Keeping the counters at the actual
491    /// acquisition boundaries means new verbs inherit instrumentation without
492    /// per-verb classification (ADR-133 D8 / issue #1389).
493    writer_acquisition_counters: Arc<WriterAcquisitionCounters>,
494    readers: ArrayQueue<Connection>,
495    max_readers: usize,
496    config: PoolConfig,
497    /// Canonical physical target used by every connection in a file-backed
498    /// read-only pool. Classification and open must share this exact spelling:
499    /// deriving WAL sidecars from a configured symlink while SQLite follows it
500    /// to another file can hide committed frames or a live writable `-shm`.
501    /// The value is an `immutable=1` URI only for a clean, checkpointed WAL;
502    /// rollback-journal databases and frozen WAL+SHM snapshots retain the
503    /// canonical ordinary path and SQLite locking/change detection.
504    read_only_open_target: Option<PathBuf>,
505    sql_bridge_reader_slots: Arc<Semaphore>,
506    sql_bridge_writer_slots: Arc<Semaphore>,
507    /// The pool-wide ADR-067 Component A writer task, spawned lazily and at
508    /// most once per pool (per DB file) via [`Self::writer_task_handle`] —
509    /// see that method's doc comment for why this lives here rather than on
510    /// each store.
511    writer_task: OnceLock<Option<WriterTaskHandle>>,
512    /// The `tokio::spawn` JoinHandle of the writer task above, stored by
513    /// [`crate::writer_task::spawn`] so short-lived callers (batch CLI
514    /// paths) can await the task's exit — and therefore its connection's
515    /// close-time WAL checkpoint — before treating the database file state
516    /// as settled. Long-running callers never take it; dropping an untaken
517    /// JoinHandle detaches the task, which is exactly the pre-existing
518    /// behavior.
519    writer_task_join: Mutex<Option<tokio::task::JoinHandle<()>>>,
520    /// Monotonic "a writer-task JoinHandle was stored at least once" flag
521    /// backing [`Self::set_writer_task_join`]'s at-most-once guard: it holds
522    /// the invariant even after [`Self::take_writer_task_join`] empties the
523    /// slot, so a second store never re-arms it.
524    writer_task_join_stored: AtomicBool,
525    /// This pool's ADR-091 backend-scoped attribution origin, minted exactly
526    /// once at construction (see [`mint_db_identity`]): `Database(_)` for a
527    /// file-backed pool, `Memory` for an in-memory pool. Every
528    /// `tx_registry::register_scoped` call site in this crate reaches its
529    /// origin through [`Self::origin`] rather than re-deriving it.
530    origin: TxOrigin,
531    /// The canonical path `origin`'s `DbIdentity` was minted from, `None` for
532    /// an in-memory pool. `DbIdentity` is deliberately opaque (no path
533    /// accessor) — filesystem consumers that need the actual path (sidecar
534    /// derivation) use this, the same canonical value the identity was
535    /// minted from, via [`Self::canonical_path`].
536    identity_path: Option<PathBuf>,
537    /// Test-only instrumentation: counts how many times the writer-task
538    /// init closure actually ran. Must never exceed 1 per pool no matter how
539    /// many stores are constructed over it — that is the invariant
540    /// `OnceLock::get_or_init` exists to guarantee, and what
541    /// `pool.rs`'s and `entity_tests.rs`'s one-writer-per-pool tests assert.
542    #[cfg(test)]
543    writer_task_spawn_count: std::sync::atomic::AtomicUsize,
544}
545
546enum ReaderLease<'pool> {
547    Pooled(Connection),
548    Shared(parking_lot::MutexGuard<'pool, Connection>),
549}
550
551/// A reader connection checked out from the pool.
552/// Returns the connection to the pool on drop.
553pub struct ReaderGuard<'pool> {
554    lease: Option<ReaderLease<'pool>>,
555    pool: &'pool ConnectionPool,
556    reusable: bool,
557}
558
559impl<'pool> ReaderGuard<'pool> {
560    /// Access the connection.
561    pub fn conn(&self) -> &Connection {
562        match self
563            .lease
564            .as_ref()
565            .expect("reader guard missing connection")
566        {
567            ReaderLease::Pooled(conn) => conn,
568            ReaderLease::Shared(guard) => guard,
569        }
570    }
571
572    /// Fail closed when connection-global state could not be restored after
573    /// a read. A pooled reader is closed and replaced on drop; a degraded
574    /// shared-writer reader is quarantined for the lifetime of the pool.
575    pub(crate) fn discard(&mut self) {
576        self.reusable = false;
577    }
578}
579
580impl<'pool> Deref for ReaderGuard<'pool> {
581    type Target = Connection;
582
583    fn deref(&self) -> &Self::Target {
584        self.conn()
585    }
586}
587
588impl<'pool> Drop for ReaderGuard<'pool> {
589    fn drop(&mut self) {
590        let Some(lease) = self.lease.take() else {
591            return;
592        };
593
594        match lease {
595            ReaderLease::Pooled(conn) if self.reusable => self.pool.return_reader(conn),
596            ReaderLease::Pooled(conn) => {
597                close_connection_quietly(conn);
598                if let Ok(conn) = self.pool.open_reader_connection() {
599                    if let Err(conn) = self.pool.readers.push(conn) {
600                        close_connection_quietly(conn);
601                    }
602                }
603            }
604            ReaderLease::Shared(guard) if !self.reusable => {
605                self.pool.retire_pooled_writer(&guard);
606            }
607            ReaderLease::Shared(_guard) => {}
608        }
609    }
610}
611
612/// A writer connection checked out from the pool.
613/// The Mutex ensures only one writer at a time.
614pub struct WriterGuard<'pool> {
615    guard: parking_lot::MutexGuard<'pool, Connection>,
616    /// The origin (ADR-091 backend-scoped attribution) of the pool this
617    /// guard was checked out from, carried so `transaction` can register its
618    /// span with the correct origin without holding a `&ConnectionPool`.
619    origin: TxOrigin,
620}
621
622/// Process-local monotonic counters for every instrumented writer acquisition
623/// boundary owned by one [`ConnectionPool`].
624///
625/// The aggregate `acquisitions` is the saturating sum of its three explicit
626/// connection classes. Infrastructure-only opens (the diagnostics PASSIVE
627/// probe, the writer task's one-time lifetime connection, and the checkpoint
628/// task's dedicated long-lived connection) are excluded; zero-wait
629/// maintenance probes also remain outside these request-traffic counters.
630#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
631pub struct WriterAcquisitionSnapshot {
632    /// Successful acquisitions across pooled, standalone, and writer-task
633    /// connection classes.
634    pub acquisitions: u64,
635    /// Successful finite-wait pool-mutex writer checkouts.
636    pub pooled_acquisitions: u64,
637    /// Successful per-operation standalone writer connection opens.
638    pub standalone_acquisitions: u64,
639    /// Successful writer-task ownership acquisitions (one per dequeued
640    /// top-level request or successful `BEGIN IMMEDIATE`).
641    pub writer_task_acquisitions: u64,
642    /// Finite-wait pool writer checkouts that exhausted their deadline.
643    pub timeouts: u64,
644    /// Writer-task `BEGIN IMMEDIATE` attempts refused busy or locked. Counted
645    /// separately from `timeouts` because that counter names the pool-mutex
646    /// checkout stage; folding the two would mislabel the stage.
647    pub writer_task_begin_busy: u64,
648    /// Writer-task `BEGIN IMMEDIATE` attempts that failed for a reason other
649    /// than busy or locked, and so surface as `StorageError::Pool`.
650    pub writer_task_begin_errors: u64,
651    /// Dequeued writer-task requests that reached the writer seam (executed
652    /// or attempted to execute their operation) and terminated in error,
653    /// counted once per request regardless of the specific terminal state.
654    pub writer_task_request_failures: u64,
655    /// Subset of `writer_task_request_failures` whose terminal state was
656    /// `WriterTaskRequestState::SideEffectsUnknown` — the commit or rollback
657    /// outcome could not be established, so the request's side effects on
658    /// the database are unknown.
659    pub writer_task_side_effects_unknown: u64,
660}
661
662/// Atomics backing [`WriterAcquisitionSnapshot`]. The writer task retains an
663/// `Arc` after spawn so its per-request acquisition site can update the same
664/// pool-scoped snapshot without retaining the whole pool.
665#[derive(Debug, Default)]
666pub(crate) struct WriterAcquisitionCounters {
667    pooled_acquisitions: AtomicU64,
668    standalone_acquisitions: AtomicU64,
669    writer_task_acquisitions: AtomicU64,
670    pooled_timeouts: AtomicU64,
671    writer_task_begin_busy: AtomicU64,
672    writer_task_begin_errors: AtomicU64,
673    writer_task_request_failures: AtomicU64,
674    writer_task_side_effects_unknown: AtomicU64,
675}
676
677impl WriterAcquisitionCounters {
678    pub(crate) fn record_writer_task_acquisition(&self) {
679        self.writer_task_acquisitions
680            .fetch_add(1, Ordering::Relaxed);
681    }
682
683    /// Records one writer-task `BEGIN IMMEDIATE` refused busy or locked.
684    ///
685    /// Callers classify by matching the value `writer_task_begin_error`
686    /// returned rather than re-testing the `rusqlite::Error`, so the busy
687    /// rule has exactly one home and the counter cannot drift from the
688    /// error the caller is actually told about.
689    pub(crate) fn record_writer_task_begin_busy(&self) {
690        self.writer_task_begin_busy.fetch_add(1, Ordering::Relaxed);
691    }
692
693    /// Records one writer-task `BEGIN IMMEDIATE` that failed for any other
694    /// reason. Without this the non-busy arm reproduces, one level down, the
695    /// same silent-failure gap the busy counter closes.
696    pub(crate) fn record_writer_task_begin_error(&self) {
697        self.writer_task_begin_errors
698            .fetch_add(1, Ordering::Relaxed);
699    }
700
701    /// Records one dequeued writer-task request that reached the writer seam
702    /// and terminated in error. Called exactly once per such request,
703    /// regardless of which terminal state it produced.
704    pub(crate) fn record_writer_task_request_failure(&self) {
705        self.writer_task_request_failures
706            .fetch_add(1, Ordering::Relaxed);
707    }
708
709    /// Records the subset of [`Self::record_writer_task_request_failure`]
710    /// whose terminal state was `SideEffectsUnknown`. Callers pair this call
711    /// with a `record_writer_task_request_failure()` call for the same
712    /// request rather than in place of it.
713    pub(crate) fn record_writer_task_side_effects_unknown(&self) {
714        self.writer_task_side_effects_unknown
715            .fetch_add(1, Ordering::Relaxed);
716    }
717
718    fn snapshot(&self) -> WriterAcquisitionSnapshot {
719        let pooled_acquisitions = self.pooled_acquisitions.load(Ordering::Relaxed);
720        let standalone_acquisitions = self.standalone_acquisitions.load(Ordering::Relaxed);
721        let writer_task_acquisitions = self.writer_task_acquisitions.load(Ordering::Relaxed);
722        WriterAcquisitionSnapshot {
723            acquisitions: pooled_acquisitions
724                .saturating_add(standalone_acquisitions)
725                .saturating_add(writer_task_acquisitions),
726            pooled_acquisitions,
727            standalone_acquisitions,
728            writer_task_acquisitions,
729            timeouts: self.pooled_timeouts.load(Ordering::Relaxed),
730            writer_task_begin_busy: self.writer_task_begin_busy.load(Ordering::Relaxed),
731            writer_task_begin_errors: self.writer_task_begin_errors.load(Ordering::Relaxed),
732            writer_task_request_failures: self.writer_task_request_failures.load(Ordering::Relaxed),
733            writer_task_side_effects_unknown: self
734                .writer_task_side_effects_unknown
735                .load(Ordering::Relaxed),
736        }
737    }
738}
739
740impl<'pool> WriterGuard<'pool> {
741    /// Returns a shared reference to the underlying connection.
742    pub fn conn(&self) -> &Connection {
743        &self.guard
744    }
745
746    /// Returns a mutable reference to the underlying connection.
747    pub fn conn_mut(&mut self) -> &mut Connection {
748        &mut self.guard
749    }
750
751    /// Execute a write transaction.
752    /// Wraps the closure in BEGIN IMMEDIATE ... COMMIT.
753    pub fn transaction<F, R>(&self, f: F) -> Result<R, SqliteError>
754    where
755        F: FnOnce(&Connection) -> Result<R, SqliteError>,
756    {
757        self.guard.execute_batch("BEGIN IMMEDIATE")?;
758        let _tx_handle = khive_storage::tx_registry::register_scoped(
759            Some("writer_guard_tx".to_string()),
760            self.origin.clone(),
761        );
762
763        match f(&self.guard) {
764            Ok(result) => {
765                if let Err(err) = self.guard.execute_batch("COMMIT") {
766                    let _ = self.guard.execute_batch("ROLLBACK");
767                    return Err(err.into());
768                }
769                Ok(result)
770            }
771            Err(err) => {
772                let _ = self.guard.execute_batch("ROLLBACK");
773                Err(err)
774            }
775        }
776    }
777}
778
779impl<'pool> Deref for WriterGuard<'pool> {
780    type Target = Connection;
781
782    fn deref(&self) -> &Self::Target {
783        self.conn()
784    }
785}
786
787impl<'pool> DerefMut for WriterGuard<'pool> {
788    fn deref_mut(&mut self) -> &mut Self::Target {
789        self.conn_mut()
790    }
791}
792
793impl ConnectionPool {
794    /// Create a new connection pool.
795    ///
796    /// Opens 1 writer + N reader connections to the same database when pooling
797    /// is enabled. All connections are configured consistently (busy timeout,
798    /// foreign keys, cache, mmap, temp store). Writable in-memory databases and
799    /// writable non-WAL files fall back to single-connection mode. Read-only
800    /// files retain a dedicated reader regardless of journal mode.
801    pub fn new(config: PoolConfig) -> Result<Self, SqliteError> {
802        refuse_home_data_store_in_tests(&config)?;
803        validate_write_admission_deadline(config.write_admission_deadline_ms)?;
804
805        // Resolve "no preference" (`None`) now that `path` is known: on for
806        // file-backed pools, off for in-memory ones. An explicit `Some(_)`
807        // preference is left untouched and always wins.
808        let mut config = config;
809        let inert_memory_queue_request =
810            config.path.is_none() && config.write_queue_enabled == Some(true);
811        config.write_queue_enabled =
812            Some(config.write_queue_enabled.unwrap_or(config.path.is_some()));
813        if inert_memory_queue_request {
814            tracing::warn!(
815                "write queue explicitly requested for an in-memory pool; it is inert because \
816                 in-memory pools cannot host a writer task"
817            );
818        }
819
820        // Mint the physical identity before WAL classification or SQLite open.
821        // Every read-only connection below uses this same canonical path (or
822        // an immutable URI derived from it), so a symlink cannot split main-file
823        // resolution from sidecar resolution.
824        let (origin, identity_path) = match config.path.as_ref() {
825            Some(path) => {
826                let (identity, canonical) = mint_db_identity(path)?;
827                (TxOrigin::Database(identity), Some(canonical))
828            }
829            None => (TxOrigin::Memory, None),
830        };
831        let read_only_open_target = read_only_open_target(&config, identity_path.as_deref())?;
832        let writer = open_writer_connection(&config, read_only_open_target.as_deref())?;
833        let wal_enabled = configure_writer_connection(&writer, &config)?;
834        let max_readers = effective_reader_count(&config, wal_enabled);
835
836        let readers = ArrayQueue::new(max_readers.max(1));
837
838        let pool = Self {
839            writer: Arc::new(Mutex::new(writer)),
840            checkpoint_ownership: CheckpointOwnershipGate::new(),
841            pooled_writer_retired: AtomicBool::new(false),
842            writer_acquisition_counters: Arc::new(WriterAcquisitionCounters::default()),
843            readers,
844            max_readers,
845            config,
846            read_only_open_target,
847            sql_bridge_reader_slots: Arc::new(Semaphore::new(max_readers.max(1))),
848            sql_bridge_writer_slots: Arc::new(Semaphore::new(1)),
849            writer_task: OnceLock::new(),
850            writer_task_join: Mutex::new(None),
851            writer_task_join_stored: AtomicBool::new(false),
852            origin,
853            identity_path,
854            #[cfg(test)]
855            writer_task_spawn_count: std::sync::atomic::AtomicUsize::new(0),
856        };
857
858        for _ in 0..pool.max_readers {
859            let conn = pool.open_reader_connection()?;
860            pool.readers
861                .push(conn)
862                .expect("reader queue must have capacity during pool initialization");
863        }
864
865        // Best-effort, process-global diagnostics belong only to pools that
866        // can acquire a writer. A read-only inspection pool has no writer
867        // timeout to report and must neither mutate `<db_parent>/.khive-logs`
868        // nor consume the global sink claim before a later writable pool.
869        if !pool.config.read_only {
870            crate::timeout_sink::init(
871                pool.canonical_path().and_then(Path::parent),
872                &crate::timeout_sink::db_label(&pool),
873            );
874        }
875
876        Ok(pool)
877    }
878
879    /// Check out a reader connection.
880    ///
881    /// Tries to pop from the lock-free queue. If empty, spins briefly then
882    /// waits with exponential backoff up to `checkout_timeout`.
883    ///
884    /// In degraded mode (WAL unavailable, `max_readers == 0`), this method
885    /// checks the shared writer mutex in bounded slices and returns pool
886    /// exhaustion after `checkout_timeout`; it never blocks indefinitely on
887    /// the non-reentrant mutex.
888    pub fn reader(&self) -> Result<ReaderGuard<'_>, SqliteError> {
889        self.reader_until(|| false)?.ok_or_else(|| {
890            SqliteError::InvalidData("uncancelled reader checkout stopped unexpectedly".into())
891        })
892    }
893
894    /// Check out a reader while cooperatively polling a request cancellation
895    /// predicate. The predicate is evaluated before connection acquisition and
896    /// between backoff slices, so an abandoned request does not sit through the
897    /// full pool checkout timeout or execute a statement when a reader later
898    /// becomes available.
899    pub(crate) fn reader_until<C>(
900        &self,
901        should_stop: C,
902    ) -> Result<Option<ReaderGuard<'_>>, SqliteError>
903    where
904        C: Fn() -> bool,
905    {
906        if self.max_readers == 0 {
907            self.ensure_pooled_writer_active()?;
908            let started = Instant::now();
909            loop {
910                if should_stop() {
911                    return Ok(None);
912                }
913                let remaining = self
914                    .config
915                    .checkout_timeout
916                    .saturating_sub(started.elapsed());
917                if remaining.is_zero() {
918                    return Err(pool_exhausted_error(
919                        self.config.checkout_timeout,
920                        self.max_readers,
921                    ));
922                }
923                if let Some(guard) = self
924                    .writer
925                    .try_lock_for(remaining.min(Duration::from_millis(2)))
926                {
927                    self.ensure_pooled_writer_active()?;
928                    return Ok(Some(ReaderGuard {
929                        lease: Some(ReaderLease::Shared(guard)),
930                        pool: self,
931                        reusable: true,
932                    }));
933                }
934            }
935        }
936
937        let started = Instant::now();
938        let mut attempt = 0u32;
939
940        loop {
941            if should_stop() {
942                return Ok(None);
943            }
944            if let Some(conn) = self.readers.pop() {
945                return Ok(Some(ReaderGuard {
946                    lease: Some(ReaderLease::Pooled(conn)),
947                    pool: self,
948                    reusable: true,
949                }));
950            }
951
952            if started.elapsed() >= self.config.checkout_timeout {
953                return Err(pool_exhausted_error(
954                    self.config.checkout_timeout,
955                    self.max_readers,
956                ));
957            }
958
959            match attempt {
960                0..=7 => {
961                    let spins = 1usize << attempt;
962                    for _ in 0..spins {
963                        std::hint::spin_loop();
964                    }
965                }
966                8..=15 => thread::yield_now(),
967                _ => {
968                    let remaining = self
969                        .config
970                        .checkout_timeout
971                        .saturating_sub(started.elapsed());
972                    let sleep = Duration::from_micros(50 * (1u64 << (attempt - 16).min(6)));
973                    thread::sleep(sleep.min(remaining).min(Duration::from_millis(2)));
974                }
975            }
976
977            attempt = attempt.saturating_add(1);
978        }
979    }
980
981    /// Check out the writer connection.
982    ///
983    /// Waits up to `checkout_timeout` for the writer Mutex and returns
984    /// `Err(SqliteError::WriterPoolCheckoutTimeout)` if the timeout is
985    /// exceeded.
986    pub fn writer(&self) -> Result<WriterGuard<'_>, SqliteError> {
987        self.ensure_pooled_writer_active()?;
988        let Some(guard) = self.writer.try_lock_for(self.config.checkout_timeout) else {
989            self.writer_acquisition_counters
990                .pooled_timeouts
991                .fetch_add(1, Ordering::Relaxed);
992            let message = format!(
993                "timed out after {:?} waiting for sqlite writer connection",
994                self.config.checkout_timeout
995            );
996            crate::timeout_sink::emit_timeout(
997                &crate::timeout_sink::db_label(self),
998                crate::timeout_sink::Site::PoolAdmission,
999                &message,
1000                Some(
1001                    self.config
1002                        .checkout_timeout
1003                        .as_millis()
1004                        .min(u128::from(u64::MAX)) as u64,
1005                ),
1006            );
1007            return Err(SqliteError::WriterPoolCheckoutTimeout {
1008                timeout: self.config.checkout_timeout,
1009            });
1010        };
1011        self.ensure_pooled_writer_active()?;
1012        self.writer_acquisition_counters
1013            .pooled_acquisitions
1014            .fetch_add(1, Ordering::Relaxed);
1015        Ok(WriterGuard {
1016            guard,
1017            origin: self.origin(),
1018        })
1019    }
1020
1021    /// Non-panicking writer checkout.
1022    ///
1023    /// Returns `Err` on timeout instead of panicking. Use this in request
1024    /// handlers where a 500 is preferable to crashing the process.
1025    pub fn try_writer(&self) -> Result<WriterGuard<'_>, SqliteError> {
1026        self.writer()
1027    }
1028
1029    /// Zero-wait writer checkout for background tasks.
1030    ///
1031    /// Uses `try_lock()` (no timeout, no spin) — returns `Err` immediately when
1032    /// any other caller holds the writer Mutex. Background tasks (e.g. the WAL
1033    /// checkpoint task) MUST use this instead of `try_writer` so that a busy
1034    /// writer causes the background task to skip its current tick rather than
1035    /// stalling for up to `checkout_timeout` (default 5s) while write traffic
1036    /// is in progress.
1037    pub fn try_writer_nowait(&self) -> Result<WriterGuard<'_>, SqliteError> {
1038        self.ensure_pooled_writer_active()?;
1039        let guard = self.writer.try_lock().ok_or_else(|| {
1040            SqliteError::InvalidData(
1041                "writer connection busy (checkpoint skipped this tick)".to_string(),
1042            )
1043        })?;
1044        self.ensure_pooled_writer_active()?;
1045        Ok(WriterGuard {
1046            guard,
1047            origin: self.origin(),
1048        })
1049    }
1050
1051    pub(crate) fn retire_pooled_writer(&self, conn: &Connection) {
1052        self.pooled_writer_retired.store(true, Ordering::Release);
1053        if let Err(error) = conn.authorizer(Some(deny_retired_writer)) {
1054            tracing::error!(
1055                %error,
1056                "failed to install the retired pooled-writer quarantine authorizer"
1057            );
1058        }
1059    }
1060
1061    fn ensure_pooled_writer_active(&self) -> Result<(), SqliteError> {
1062        if self.pooled_writer_retired.load(Ordering::Acquire) {
1063            return Err(SqliteError::InvalidData(
1064                "pooled writer connection retired after a terminal transaction fault".to_string(),
1065            ));
1066        }
1067        Ok(())
1068    }
1069
1070    /// Snapshot all instrumented writer acquisition outcomes since this pool
1071    /// was constructed.
1072    pub fn writer_acquisition_snapshot(&self) -> WriterAcquisitionSnapshot {
1073        self.writer_acquisition_counters.snapshot()
1074    }
1075
1076    /// Clone the pool-scoped counter set for the lifetime-owned writer task.
1077    pub(crate) fn writer_acquisition_counters(&self) -> Arc<WriterAcquisitionCounters> {
1078        Arc::clone(&self.writer_acquisition_counters)
1079    }
1080
1081    /// Get the current number of available reader connections.
1082    pub fn available_readers(&self) -> usize {
1083        self.readers.len()
1084    }
1085
1086    /// Get the total number of reader connections in the pool.
1087    pub fn max_readers(&self) -> usize {
1088        self.max_readers
1089    }
1090
1091    /// Return the pool configuration.
1092    pub fn config(&self) -> &PoolConfig {
1093        &self.config
1094    }
1095
1096    /// The typed admission failure for a pooled reader checkout that
1097    /// exhausted `checkout_timeout`: no reader was acquired, so the
1098    /// operation never started and a retry cannot duplicate a side effect.
1099    pub(crate) fn reader_admission_timeout(&self, operation: &'static str) -> StorageError {
1100        StorageError::AdmissionTimeout {
1101            operation: operation.into(),
1102            timeout_ms: u64::try_from(self.config.checkout_timeout.as_millis()).unwrap_or(u64::MAX),
1103        }
1104    }
1105
1106    /// Resolve a [`Self::reader_until`] outcome into a checked-out guard or
1107    /// the canonical refusal. This is the single home of the checkout
1108    /// tri-state; call sites must not re-derive any arm of it:
1109    ///
1110    /// - `Ok(Some)` — a reader was checked out.
1111    /// - `Ok(None)` — `should_stop()` fired: the request was cancelled or hit
1112    ///   its deadline before checkout. NOT an admission wait, so it maps to
1113    ///   the non-retryable [`StorageError::Timeout`] — emitting the retryable
1114    ///   `AdmissionTimeout` here would invite an immediate retry of a request
1115    ///   its caller already abandoned, into a possibly saturated pool.
1116    /// - `Err` carrying the pool's own `SQLITE_BUSY` — `reader_until` executes
1117    ///   no SQL, so the only `SQLITE_BUSY` it can produce is
1118    ///   [`pool_exhausted_error`], raised when `checkout_timeout` elapses with
1119    ///   no reader available. A genuine admission wait that ended before any
1120    ///   work began maps to the retryable [`StorageError::AdmissionTimeout`].
1121    /// - any other `Err` (e.g. [`Self::ensure_pooled_writer_active`] returning
1122    ///   `InvalidData` for a retired pooled writer) — an opaque driver failure
1123    ///   under the caller's capability, non-retryable.
1124    pub(crate) fn resolve_reader_checkout<'p>(
1125        &self,
1126        capability: StorageCapability,
1127        operation: &'static str,
1128        outcome: Result<Option<ReaderGuard<'p>>, SqliteError>,
1129    ) -> Result<ReaderGuard<'p>, StorageError> {
1130        match outcome {
1131            Ok(Some(guard)) => Ok(guard),
1132            Ok(None) => Err(StorageError::Timeout {
1133                operation: operation.into(),
1134            }),
1135            Err(error) => {
1136                let is_pool_exhausted = matches!(
1137                    &error,
1138                    SqliteError::Rusqlite(rusqlite::Error::SqliteFailure(code, _))
1139                        if code.code == rusqlite::ErrorCode::DatabaseBusy
1140                );
1141                if is_pool_exhausted {
1142                    Err(self.reader_admission_timeout(operation))
1143                } else {
1144                    Err(StorageError::driver(capability, operation, error))
1145                }
1146            }
1147        }
1148    }
1149
1150    /// Pool-wide permits for file-backed raw-SQL reader opens and active reads.
1151    pub(crate) fn sql_bridge_reader_slots(&self) -> Arc<Semaphore> {
1152        Arc::clone(&self.sql_bridge_reader_slots)
1153    }
1154
1155    /// Pool-wide permit for a file-backed raw-SQL writer handle.
1156    pub(crate) fn sql_bridge_writer_slots(&self) -> Arc<Semaphore> {
1157        Arc::clone(&self.sql_bridge_writer_slots)
1158    }
1159
1160    /// This pool's ADR-091 backend-scoped attribution origin (ADR-091,
1161    /// backend-scoped WAL-pin attribution design note): `Database(_)` for a
1162    /// file-backed pool, `Memory` for an in-memory pool. Every
1163    /// `tx_registry::register_scoped` call site threaded in this crate
1164    /// passes this value as the span's origin.
1165    pub fn origin(&self) -> TxOrigin {
1166        self.origin.clone()
1167    }
1168
1169    /// The canonical path this pool's `origin()` identity was minted from,
1170    /// `None` for an in-memory pool. `DbIdentity` has no path accessor by
1171    /// design; sidecar derivation and other filesystem consumers use this —
1172    /// the same canonical value the identity was minted from — instead of
1173    /// re-deriving a path from the raw configured one.
1174    pub fn canonical_path(&self) -> Option<&Path> {
1175        self.identity_path.as_deref()
1176    }
1177
1178    /// Whether the write queue is effectively enabled for this pool: the
1179    /// resolved `write_queue_enabled` flag AND file-backed.
1180    ///
1181    /// `ConnectionPool::new` resolves the "no preference" (`None`) preference
1182    /// to a concrete `Some(..)` once `path` is known, so every reader of
1183    /// `config.write_queue_enabled` sees a resolved value; the `debug_assert`
1184    /// pins that invariant and a `None` that slipped past would read as
1185    /// disabled. Bypassing `ConnectionPool::new` to construct a pool is a
1186    /// construction-path bug. Use this instead of repeating
1187    /// `config().write_queue_enabled.unwrap_or(false) && config().path.is_some()`
1188    /// at every routing/violation site.
1189    pub fn write_queue_active(&self) -> bool {
1190        debug_assert!(
1191            self.config.write_queue_enabled.is_some(),
1192            "write_queue_enabled must be resolved to Some(..) by ConnectionPool::new \
1193             before any write_queue_active read"
1194        );
1195        self.config.write_queue_enabled.unwrap_or(false) && self.config.path.is_some()
1196    }
1197
1198    /// Whether a writer-task JoinHandle has been stored at least once.
1199    ///
1200    /// Unlike [`Self::take_writer_task_join`], this remains true after the
1201    /// one-shot handle slot is emptied, distinguishing a task that never
1202    /// spawned from a handle another caller already consumed.
1203    pub fn writer_task_join_was_stored(&self) -> bool {
1204        self.writer_task_join_stored.load(Ordering::SeqCst)
1205    }
1206
1207    /// Return the pool-wide ADR-067 Component A writer task, spawning it
1208    /// lazily on first access if `PoolConfig::write_queue_enabled` is set.
1209    /// Exactly one writer task exists per `ConnectionPool` (per DB file); see
1210    /// crates/khive-db/docs/api/pool.md#connectionpoolwriter_task_handle--single-writer-task-rationale
1211    /// for why a per-store writer task would defeat the single-writer
1212    /// guarantee.
1213    ///
1214    /// Returns `Ok(None)` if the flag is off, or if the writer task failed to
1215    /// spawn for a reason other than a missing runtime (for example, an
1216    /// in-memory pool has no standalone-connection support) — callers fall
1217    /// back to the legacy pool-mutex write path in either case. A spawn
1218    /// failure is logged once here (at first access), not once per store.
1219    ///
1220    /// Returns `Err(StorageError::WriterTaskNoRuntime)` instead of panicking
1221    /// when `write_queue_enabled` is set but this is the first access and no
1222    /// Tokio runtime is available on the calling thread (checked via
1223    /// [`tokio::runtime::Handle::try_current`]) — spawning the writer task
1224    /// requires `tokio::spawn`, which panics outside a runtime. Callers that
1225    /// already treat a missing writer task as best-effort (construction-time
1226    /// degrade to the legacy path, matching slice 1's documented policy) can
1227    /// collapse this into `None` with `.ok().flatten()`; callers that need to
1228    /// fail loud on a genuine misconfiguration (write queue requested but no
1229    /// runtime to run it on) can propagate the `Err` directly.
1230    pub fn writer_task_handle(&self) -> Result<Option<WriterTaskHandle>, StorageError> {
1231        // Same pinned invariant `write_queue_active` asserts, kept inline
1232        // here because this gate keys on the flag ALONE: an explicit
1233        // `Some(true)` on an in-memory pool must still attempt the spawn
1234        // and degrade (documented + tested in
1235        // `explicit_true_stays_on_for_memory_backed_pool`), so the
1236        // file-backed half of `write_queue_active` cannot gate this early
1237        // return.
1238        debug_assert!(
1239            self.config.write_queue_enabled.is_some(),
1240            "write_queue_enabled must be resolved to Some(..) by ConnectionPool::new \
1241             before any writer_task_handle read"
1242        );
1243        if !self.config.write_queue_enabled.unwrap_or(false) {
1244            return Ok(None);
1245        }
1246        // Fast path: already resolved (spawned, degraded, or off) by an
1247        // earlier call — no need to re-check the runtime.
1248        if let Some(existing) = self.writer_task.get() {
1249            return Ok(existing.clone());
1250        }
1251        // Not yet initialized and the flag is on: spawning requires
1252        // `tokio::spawn`, which panics outside a runtime context. Check
1253        // first and fail loud with a typed error instead.
1254        if tokio::runtime::Handle::try_current().is_err() {
1255            return Err(StorageError::WriterTaskNoRuntime);
1256        }
1257        Ok(self
1258            .writer_task
1259            .get_or_init(|| {
1260                #[cfg(test)]
1261                self.writer_task_spawn_count
1262                    .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1263
1264                match crate::writer_task::spawn(self, self.config.write_queue_capacity) {
1265                    Ok(handle) => Some(handle),
1266                    Err(e) => {
1267                        tracing::warn!(
1268                            error = %e,
1269                            "KHIVE_WRITE_QUEUE=1 but the writer task failed to spawn; \
1270                             writes fall back to the pool-mutex path"
1271                        );
1272                        None
1273                    }
1274                }
1275            })
1276            .clone())
1277    }
1278
1279    /// Resolve the writer task for a store write at the moment the write is
1280    /// issued, rather than trusting only a handle cached by a synchronous
1281    /// store constructor. Construction can legitimately run before Tokio is
1282    /// entered, in which case `writer_task_handle()` returns
1283    /// `WriterTaskNoRuntime` without caching a terminal `None`.
1284    ///
1285    /// Strict routing makes every missing handle fail closed here. The
1286    /// caller remains responsible for recording a non-strict direct fallback
1287    /// at the exact fallback seam with [`Self::record_direct_route`].
1288    pub(crate) fn writer_task_for_write(
1289        &self,
1290        cached: Option<&WriterTaskHandle>,
1291        operation: &'static str,
1292    ) -> Result<Option<WriterTaskHandle>, StorageError> {
1293        let handle = match cached {
1294            Some(handle) => Some(handle.clone()),
1295            None => match self.writer_task_handle() {
1296                Ok(handle) => handle,
1297                Err(error) if self.config.write_routing_strict => return Err(error),
1298                Err(_) => None,
1299            },
1300        };
1301
1302        if handle.is_none() && self.config.write_routing_strict {
1303            return Err(StorageError::Pool {
1304                operation: operation.into(),
1305                message: "strict write routing requires a writer-task handle; no handle is \
1306                          available, so the direct writer fallback was refused"
1307                    .into(),
1308            });
1309        }
1310        Ok(handle)
1311    }
1312
1313    /// Record one actual compatibility fallback around the writer task. A
1314    /// file-backed pool with the queue enabled should never reach this seam
1315    /// in strict mode because [`Self::writer_task_for_write`] refuses first.
1316    pub(crate) fn record_direct_route(&self, site: crate::timeout_sink::Site) {
1317        if self.write_queue_active() {
1318            crate::timeout_sink::emit_direct_route_violation(
1319                &crate::timeout_sink::db_label(self),
1320                site,
1321            );
1322        }
1323    }
1324
1325    /// Test-only: how many times the writer-task init closure actually ran.
1326    /// Must be at most 1 for the pool's whole lifetime, regardless of how
1327    /// many times [`Self::writer_task_handle`] is called or how many stores
1328    /// are constructed over this pool.
1329    #[cfg(test)]
1330    pub(crate) fn writer_task_spawn_count(&self) -> usize {
1331        self.writer_task_spawn_count
1332            .load(std::sync::atomic::Ordering::SeqCst)
1333    }
1334
1335    /// Record the writer task's `tokio::spawn` JoinHandle. Called exactly
1336    /// once, by [`crate::writer_task::spawn`], immediately after spawning —
1337    /// the same `writer_task` OnceLock init that makes spawn at-most-once
1338    /// per pool makes this write at-most-once per pool.
1339    ///
1340    /// First-wins: if a handle was ever stored (including one a caller has
1341    /// since taken — `writer_task_join_stored` remembers), the existing
1342    /// state is kept and the new handle is dropped (dropping a `JoinHandle`
1343    /// detaches its task without cancelling it). A second store violates the
1344    /// at-most-once contract and trips the debug_assert in debug builds;
1345    /// release builds keep the first handle rather than silently swapping
1346    /// the drain owner out from under whichever caller already took it.
1347    pub(crate) fn set_writer_task_join(&self, join: tokio::task::JoinHandle<()>) {
1348        // `swap(true)` returns the prior value: `true` means a handle was
1349        // stored at least once before, so this is a second store — even when
1350        // the slot itself is empty because `take_writer_task_join` already
1351        // ran (the slot alone cannot tell "never stored" from "taken").
1352        let first_store = !self.writer_task_join_stored.swap(true, Ordering::SeqCst);
1353        debug_assert!(
1354            first_store,
1355            "writer task JoinHandle stored twice (even counting a taken one); \
1356             the writer_task OnceLock is supposed to make spawn at-most-once per pool"
1357        );
1358        if first_store {
1359            *self.writer_task_join.lock() = Some(join);
1360        }
1361    }
1362
1363    /// Take the writer task's JoinHandle, if a writer task was spawned and
1364    /// the handle has not already been taken.
1365    ///
1366    /// Intended for short-lived batch callers that drop every
1367    /// [`WriterTaskHandle`] clone (closing the queue) and then need to await
1368    /// the task's exit before treating the database file as settled: the
1369    /// task's connection close fires SQLite's close-time WAL checkpoint, so
1370    /// until the task exits the file bytes can still move after the caller's
1371    /// last write returned.
1372    ///
1373    /// One-shot: `None` means either the write queue never spawned
1374    /// (disabled, or spawn degraded) or another caller already took the
1375    /// handle — in both cases there is nothing further to await here.
1376    /// Exactly one subsystem may own the drain: the single caller that
1377    /// receives `Some(_)` is the sole owner of the task-exit await (and of
1378    /// the close-time WAL checkpoint that settles the database file); every
1379    /// later caller receives `None` and must not arrange its own await.
1380    pub fn take_writer_task_join(&self) -> Option<tokio::task::JoinHandle<()>> {
1381        self.writer_task_join.lock().take()
1382    }
1383
1384    /// Compatibility method: returns the writer connection wrapped in `Arc<Mutex>`.
1385    ///
1386    /// WARNING: This exists only for backward compatibility with code that
1387    /// calls `store.conn()`. New code should use `reader()` and `writer()`.
1388    pub fn legacy_conn(&self) -> Arc<Mutex<Connection>> {
1389        Arc::clone(&self.writer)
1390    }
1391
1392    fn open_reader_connection(&self) -> Result<Connection, SqliteError> {
1393        let path = self.read_connection_path()?;
1394        open_reader_connection(path, &self.config)
1395    }
1396
1397    fn read_connection_path(&self) -> Result<&Path, SqliteError> {
1398        self.read_only_open_target
1399            .as_deref()
1400            .or(self.config.path.as_deref())
1401            .ok_or_else(|| {
1402                SqliteError::InvalidData(
1403                    "in-memory databases do not support standalone connections".to_string(),
1404                )
1405            })
1406    }
1407
1408    /// Open a standalone read-write connection to the same file-backed database.
1409    ///
1410    /// Stores whose trait methods take `Send + 'static` closures (executed via
1411    /// `spawn_blocking`) cannot hold the pooled `WriterGuard`'s `MutexGuard`
1412    /// across the call — it opens an independent connection instead. This
1413    /// must still honor `PoolConfig::read_only`: opening
1414    /// `SQLITE_OPEN_READ_WRITE` unconditionally here would let a read-only
1415    /// backend's graph/event/text stores bypass the flag that the pooled
1416    /// writer enforces via `query_only`. A fully configured successful open
1417    /// increments the standalone acquisition class exactly once.
1418    pub fn open_standalone_writer(&self) -> Result<Connection, SqliteError> {
1419        let conn = self.open_standalone_writer_untracked()?;
1420        self.writer_acquisition_counters
1421            .standalone_acquisitions
1422            .fetch_add(1, Ordering::Relaxed);
1423        Ok(conn)
1424    }
1425
1426    /// Open an infrastructure-owned standalone writer connection without
1427    /// counting it as one write-operation acquisition.
1428    ///
1429    /// Restricted to the diagnostics PASSIVE probe, the writer task's
1430    /// one-time lifetime connection, and the checkpoint task's dedicated
1431    /// long-lived connection (opened once at startup and reused across
1432    /// ticks — see `CheckpointConnection::ensure_open`). Actual file-backed
1433    /// write paths must call [`Self::open_standalone_writer`] so their
1434    /// acquisitions are observable.
1435    pub(crate) fn open_standalone_writer_untracked(&self) -> Result<Connection, SqliteError> {
1436        let path = self.config.path.as_ref().ok_or_else(|| {
1437            SqliteError::InvalidData(
1438                "in-memory databases do not support standalone connections".to_string(),
1439            )
1440        })?;
1441
1442        if self.config.read_only {
1443            return Err(SqliteError::InvalidData(
1444                "database is read-only: standalone write connections are not permitted".to_string(),
1445            ));
1446        }
1447
1448        let conn = Connection::open_with_flags(
1449            path,
1450            OpenFlags::SQLITE_OPEN_READ_WRITE
1451                | OpenFlags::SQLITE_OPEN_NO_MUTEX
1452                | OpenFlags::SQLITE_OPEN_URI,
1453        )?;
1454        conn.busy_timeout(self.config.busy_timeout)?;
1455        self.checkpoint_ownership
1456            .configure_wal_autocheckpoint(&conn)?;
1457        conn.pragma_update(None, "foreign_keys", "ON")?;
1458        conn.pragma_update(None, "synchronous", "NORMAL")?;
1459
1460        let wal_enabled =
1461            self.config.wal_mode && current_journal_mode(&conn)?.eq_ignore_ascii_case("wal");
1462        if wal_enabled {
1463            conn.pragma_update(
1464                None,
1465                "journal_size_limit",
1466                self.config.journal_size_limit_bytes,
1467            )?;
1468        }
1469
1470        Ok(conn)
1471    }
1472
1473    /// Effective `PRAGMA wal_autocheckpoint` for a writer-capable connection
1474    /// opened right now: `0` once a dedicated checkpoint owner has claimed
1475    /// the pool, the bounded fallback otherwise.
1476    #[cfg(test)]
1477    pub(crate) fn effective_wal_autocheckpoint_pages(&self) -> u32 {
1478        self.checkpoint_ownership.wal_autocheckpoint_pages()
1479    }
1480
1481    /// Claim routine WAL-checkpoint ownership for this pool.
1482    ///
1483    /// Called by the scheduled checkpoint task at startup — the one caller
1484    /// that actually replaces SQLite's per-commit autocheckpoint with
1485    /// dedicated PASSIVE checkpointing (ADR-091 Amendment 10). The claim
1486    /// makes every subsequently opened writer-capable connection set
1487    /// `PRAGMA wal_autocheckpoint = 0`, and re-applies that pragma on the
1488    /// already-open pooled writer under the writer mutex. A writer task
1489    /// spawned before the claim keeps its own long-lived connection;
1490    /// [`Self::propagate_checkpoint_claim_to_writer_task`] reaches that one.
1491    ///
1492    /// Without a claim, writer-capable connections keep the bounded
1493    /// `FALLBACK_WAL_AUTOCHECKPOINT_PAGES` threshold, so a writable pool
1494    /// in a process that never runs the checkpoint task (embedded runtimes,
1495    /// one-shot CLI executions) retains SQLite's own WAL reclamation instead
1496    /// of growing its WAL without bound.
1497    ///
1498    /// Read-only pools record the claim but have no writer-capable
1499    /// connections to reconfigure. Writable pools publish the claim only after
1500    /// the pooled writer is configured successfully; a failed attempt keeps
1501    /// the bounded fallback active and remains retryable.
1502    pub fn claim_checkpoint_ownership(&self) -> Result<(), SqliteError> {
1503        if !self.checkpoint_ownership.begin_claim() {
1504            return Ok(());
1505        }
1506        let result = (|| {
1507            if !self.config.read_only {
1508                let writer = self.writer()?;
1509                writer.conn().pragma_update(None, "wal_autocheckpoint", 0)?;
1510            }
1511            Ok(())
1512        })();
1513        self.checkpoint_ownership.finish_claim(result.is_ok());
1514        result
1515    }
1516
1517    /// Flip an already-running writer task's long-lived connection to the
1518    /// claimed-owner setting.
1519    ///
1520    /// Connections opened after [`Self::claim_checkpoint_ownership`] inherit
1521    /// `wal_autocheckpoint = 0` at open; only a writer task spawned before
1522    /// the claim still holds a connection on the bounded fallback. Returns
1523    /// `Ok(())` without side effects when the pool's write queue is
1524    /// disabled.
1525    pub async fn propagate_checkpoint_claim_to_writer_task(&self) -> Result<(), StorageError> {
1526        let Some(handle) = self.writer_task_handle()? else {
1527            return Ok(());
1528        };
1529        handle
1530            .send_top_level(|conn| {
1531                conn.pragma_update(None, "wal_autocheckpoint", 0)
1532                    .map_err(|e| StorageError::Pool {
1533                        operation: "claim_checkpoint_ownership".into(),
1534                        message: e.to_string(),
1535                    })
1536            })
1537            .await
1538    }
1539
1540    /// Open a standalone read-only connection to the same file-backed database.
1541    ///
1542    /// Companion to `open_standalone_writer` for stores that also need an
1543    /// independent reader connection outside the pooled reader queue.
1544    pub fn open_standalone_reader(&self) -> Result<Connection, SqliteError> {
1545        let path = self.read_connection_path()?;
1546
1547        let conn = Connection::open_with_flags(
1548            path,
1549            OpenFlags::SQLITE_OPEN_READ_ONLY
1550                | OpenFlags::SQLITE_OPEN_NO_MUTEX
1551                | OpenFlags::SQLITE_OPEN_URI,
1552        )?;
1553        configure_reader_connection(&conn, &self.config)?;
1554        conn.pragma_update(None, "synchronous", "NORMAL")?;
1555        Ok(conn)
1556    }
1557
1558    fn return_reader(&self, conn: Connection) {
1559        if self.max_readers == 0 {
1560            return;
1561        }
1562
1563        let conn = if reset_reader_connection(&conn) && reader_connection_is_healthy(&conn) {
1564            Some(conn)
1565        } else {
1566            close_connection_quietly(conn);
1567            self.open_reader_connection().ok()
1568        };
1569
1570        if let Some(conn) = conn {
1571            if let Err(conn) = self.readers.push(conn) {
1572                eprintln!(
1573                    "[sqlite-pool] reader pool queue full, discarding replacement connection"
1574                );
1575                close_connection_quietly(conn);
1576            }
1577        }
1578    }
1579}
1580
1581/// Bound on the final-component symlink chain [`resolve_symlink_chain`]
1582/// follows before failing loud, mirroring the OS's own loop limit (e.g.
1583/// Linux/macOS `ELOOP`, commonly 40 hops) rather than looping forever on a
1584/// cycle.
1585const MAX_SYMLINK_DEPTH: u32 = 40;
1586
1587/// Mint the canonical [`DbIdentity`] for a configured database path.
1588///
1589/// The sole minting point (ADR-091 backend-scoped attribution design note):
1590/// `tx_registry` origin threading and `sidecar_dir_for` re-keying both
1591/// consume this function's output rather than re-deriving it. Operationally
1592/// three steps:
1593///
1594/// 1. A relative configured path is resolved against the process's current
1595///    directory BEFORE any canonicalization — a bare file name has an empty
1596///    parent, and canonicalizing an empty path fails.
1597/// 2. If the resolved path exists, canonicalize the full path: this
1598///    resolves symlinks at every level, including a symlink at the
1599///    database-file level itself (a `link.sqlite` pointing at the real file
1600///    mints the target's identity).
1601/// 3. If the resolved path does not yet exist (first open), a dangling
1602///    file-level symlink is a valid first-open state — SQLite creates the
1603///    target through the link on first write, and minting the link's own
1604///    name would diverge from a later opener using the target path
1605///    directly. The final-component symlink chain is followed to its
1606///    ultimate target first (bounded, see [`MAX_SYMLINK_DEPTH`]), then that
1607///    target's PARENT directory is canonicalized and the file name is
1608///    appended unchanged — the same pattern `FsBlobStore` uses for its
1609///    root-keyed write locks (`stores/blob.rs::write_lock_for_root`), and
1610///    for the same reason: `Path::canonicalize` requires an existing path.
1611///
1612/// A resolved target whose parent directory does not exist fails minting
1613/// exactly as the subsequent database open itself would fail.
1614///
1615/// Returns the minted [`DbIdentity`] alongside the canonical [`PathBuf`] it
1616/// was built from — `DbIdentity` has no path accessor by design, so callers
1617/// that need the filesystem path (sidecar derivation) keep this pairing
1618/// rather than re-deriving it from the raw configured path.
1619fn mint_db_identity(configured_path: &Path) -> Result<(DbIdentity, PathBuf), SqliteError> {
1620    let absolute = if configured_path.is_absolute() {
1621        configured_path.to_path_buf()
1622    } else {
1623        let cwd = std::env::current_dir().map_err(|e| {
1624            SqliteError::InvalidData(format!(
1625                "cannot mint database identity for {configured_path:?}: failed to resolve the \
1626                 process current directory: {e}"
1627            ))
1628        })?;
1629        cwd.join(configured_path)
1630    };
1631
1632    if absolute.exists() {
1633        let canonical = absolute.canonicalize().map_err(|e| {
1634            SqliteError::InvalidData(format!(
1635                "cannot mint database identity: failed to canonicalize existing path \
1636                 {absolute:?}: {e}"
1637            ))
1638        })?;
1639        return Ok((
1640            DbIdentity::new(canonical.clone().into_os_string()),
1641            canonical,
1642        ));
1643    }
1644
1645    let resolved_target = resolve_symlink_chain(&absolute)?;
1646    let parent = resolved_target.parent().ok_or_else(|| {
1647        SqliteError::InvalidData(format!(
1648            "cannot mint database identity for {resolved_target:?}: path has no parent \
1649             directory"
1650        ))
1651    })?;
1652    let file_name = resolved_target.file_name().ok_or_else(|| {
1653        SqliteError::InvalidData(format!(
1654            "cannot mint database identity for {resolved_target:?}: path has no file name"
1655        ))
1656    })?;
1657    let canonical_parent = parent.canonicalize().map_err(|e| {
1658        SqliteError::InvalidData(format!(
1659            "cannot mint database identity: parent directory {parent:?} of first-open path \
1660             {resolved_target:?} does not exist or is inaccessible: {e}"
1661        ))
1662    })?;
1663    let mut identity_path = canonical_parent;
1664    identity_path.push(file_name);
1665    Ok((
1666        DbIdentity::new(identity_path.clone().into_os_string()),
1667        identity_path,
1668    ))
1669}
1670
1671/// Follow a (possibly dangling) final-component symlink chain to its
1672/// ultimate target, bounded at [`MAX_SYMLINK_DEPTH`] hops. A path that is
1673/// not itself a symlink — including one that does not exist at all —
1674/// returns unchanged on the first iteration; this is the common case, a
1675/// first-open path with no symlink involved.
1676fn resolve_symlink_chain(path: &Path) -> Result<PathBuf, SqliteError> {
1677    let mut current = path.to_path_buf();
1678    for _ in 0..MAX_SYMLINK_DEPTH {
1679        match fs::symlink_metadata(&current) {
1680            Ok(meta) if meta.file_type().is_symlink() => {
1681                let target = fs::read_link(&current).map_err(|e| {
1682                    SqliteError::InvalidData(format!(
1683                        "cannot mint database identity: failed to read symlink {current:?}: {e}"
1684                    ))
1685                })?;
1686                current = if target.is_absolute() {
1687                    target
1688                } else {
1689                    match current.parent() {
1690                        Some(parent) => parent.join(&target),
1691                        None => target,
1692                    }
1693                };
1694            }
1695            _ => return Ok(current),
1696        }
1697    }
1698    Err(SqliteError::InvalidData(format!(
1699        "cannot mint database identity for {path:?}: symlink chain exceeds \
1700         {MAX_SYMLINK_DEPTH} levels"
1701    )))
1702}
1703
1704fn effective_reader_count(config: &PoolConfig, wal_enabled: bool) -> usize {
1705    if config.path.is_some() && config.read_only {
1706        config.max_readers.max(1)
1707    } else if config.path.is_some() && config.wal_mode && wal_enabled {
1708        config.max_readers
1709    } else {
1710        0
1711    }
1712}
1713
1714fn open_writer_connection(
1715    config: &PoolConfig,
1716    read_only_open_target: Option<&Path>,
1717) -> Result<Connection, SqliteError> {
1718    match config.path.as_ref() {
1719        Some(path) => {
1720            let flags = if config.read_only {
1721                writer_read_only_open_flags()
1722            } else {
1723                writer_open_flags()
1724            };
1725            let target = if config.read_only {
1726                read_only_open_target.ok_or_else(|| {
1727                    SqliteError::InvalidData(
1728                        "file-backed read-only pool has no canonical open target".to_string(),
1729                    )
1730                })?
1731            } else {
1732                path
1733            };
1734            Connection::open_with_flags(target, flags).map_err(Into::into)
1735        }
1736        None => Connection::open_in_memory().map_err(Into::into),
1737    }
1738}
1739
1740/// Select the one case that may safely use SQLite's immutable URI contract: a
1741/// clean, checkpointed persistent-WAL snapshot with neither a shared-memory
1742/// index nor committed frames in `<db>-wal`. A normal read-only connection can
1743/// create fresh `-wal`/`-shm` files even for that clean database, while
1744/// `immutable=1` keeps the source directory untouched. We deliberately do not
1745/// apply `immutable=1` to:
1746///
1747/// - rollback-journal databases, which can read safely with normal locking and
1748///   should continue observing committed changes when an operator points a
1749///   read-only connection at a live database; or
1750/// - WAL databases with a read-only `-shm`, where ordinary read-only SQLite can
1751///   consume committed WAL frames without writing the frozen index; or
1752/// - WAL databases with a writable `-shm`, which are potentially live. Those
1753///   fail closed before SQLite is opened rather than mutating shared state or
1754///   suppressing change detection unsafely.
1755///
1756/// A non-empty WAL without `-shm` is also refused before open. Immutable SQLite
1757/// does not rebuild a missing WAL index: it ignores the WAL entirely, which can
1758/// make a committed row disappear from inspection. Ordinary read-only SQLite
1759/// would recover the frames but create `-shm`, violating the physical
1760/// read-only contract. The operator must provide the frozen read-only `-shm`
1761/// alongside that WAL (or checkpoint a writable copy first).
1762fn read_only_open_target(
1763    config: &PoolConfig,
1764    physical_path: Option<&Path>,
1765) -> Result<Option<PathBuf>, SqliteError> {
1766    if !config.read_only {
1767        return Ok(None);
1768    }
1769    let Some(path) = physical_path else {
1770        return Ok(None);
1771    };
1772    read_only_wal_open_target_for_path(path).map(Some)
1773}
1774
1775fn read_only_wal_open_target_for_path(path: &Path) -> Result<PathBuf, SqliteError> {
1776    if !sqlite_header_uses_wal(path)? {
1777        return Ok(path.to_path_buf());
1778    }
1779
1780    let shm = sqlite_sidecar_path(path, "-shm");
1781    match fs::metadata(&shm) {
1782        Ok(metadata) if metadata.permissions().readonly() => {
1783            let wal = sqlite_sidecar_path(path, "-wal");
1784            match fs::metadata(&wal) {
1785                Ok(_) => Ok(path.to_path_buf()),
1786                Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
1787                    Err(SqliteError::InvalidData(format!(
1788                        "read-only WAL snapshot {} has a shared-memory sidecar {} but no WAL \
1789                         sidecar {}; refusing the inconsistent sidecar set before SQLite open",
1790                        path.display(),
1791                        shm.display(),
1792                        wal.display(),
1793                    )))
1794                }
1795                Err(error) => Err(SqliteError::Io(error)),
1796            }
1797        }
1798        Ok(_) => Err(SqliteError::InvalidData(format!(
1799            "read-only WAL snapshot {} has a writable WAL shared-memory sidecar {}; close every \
1800             live writer and remove the transient -shm file (or make a genuinely frozen snapshot) \
1801             before inspection",
1802            path.display(),
1803            shm.display(),
1804        ))),
1805        Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
1806            let wal = sqlite_sidecar_path(path, "-wal");
1807            match fs::metadata(&wal) {
1808                Ok(metadata) if metadata.len() > 0 => Err(SqliteError::InvalidData(format!(
1809                    "read-only WAL snapshot {} has a non-empty WAL sidecar {} but no read-only \
1810                     shared-memory sidecar {}; refusing before SQLite open because immutable \
1811                     mode would omit committed WAL frames and ordinary read-only mode would \
1812                     create or mutate -shm; include the frozen read-only -shm beside this \
1813                     snapshot, or checkpoint a writable copy before inspection",
1814                    path.display(),
1815                    wal.display(),
1816                    shm.display(),
1817                ))),
1818                Ok(_) => sqlite_immutable_uri(path),
1819                Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
1820                    sqlite_immutable_uri(path)
1821                }
1822                Err(error) => Err(SqliteError::Io(error)),
1823            }
1824        }
1825        Err(error) => Err(SqliteError::Io(error)),
1826    }
1827}
1828
1829pub(crate) fn open_read_only_snapshot_connection(path: &Path) -> Result<Connection, SqliteError> {
1830    let (_, physical_path) = mint_db_identity(path)?;
1831    let target = read_only_wal_open_target_for_path(&physical_path)?;
1832    Connection::open_with_flags(&target, reader_open_flags()).map_err(Into::into)
1833}
1834
1835fn sqlite_header_uses_wal(path: &Path) -> Result<bool, SqliteError> {
1836    let mut file = fs::File::open(path)?;
1837    let mut header = [0_u8; 20];
1838    if let Err(error) = file.read_exact(&mut header) {
1839        if error.kind() == std::io::ErrorKind::UnexpectedEof {
1840            return Ok(false);
1841        }
1842        return Err(SqliteError::Io(error));
1843    }
1844    Ok(&header[..16] == b"SQLite format 3\0" && header[18] == 2 && header[19] == 2)
1845}
1846
1847fn sqlite_sidecar_path(path: &Path, suffix: &str) -> PathBuf {
1848    let mut sidecar = path.as_os_str().to_os_string();
1849    sidecar.push(suffix);
1850    PathBuf::from(sidecar)
1851}
1852
1853fn sqlite_immutable_uri(path: &Path) -> Result<PathBuf, SqliteError> {
1854    let absolute = if path.is_absolute() {
1855        path.to_path_buf()
1856    } else {
1857        std::env::current_dir()?.join(path)
1858    };
1859    let mut uri = String::from("file:");
1860
1861    #[cfg(unix)]
1862    {
1863        use std::os::unix::ffi::OsStrExt as _;
1864        push_sqlite_uri_path(&mut uri, absolute.as_os_str().as_bytes());
1865    }
1866
1867    #[cfg(not(unix))]
1868    {
1869        let path = absolute.to_str().ok_or_else(|| {
1870            SqliteError::InvalidData(format!(
1871                "read-only WAL snapshot path is not representable as a SQLite URI: {}",
1872                absolute.display()
1873            ))
1874        })?;
1875        let normalized = path.replace('\\', "/");
1876        if cfg!(windows) && !normalized.starts_with('/') {
1877            uri.push('/');
1878        }
1879        push_sqlite_uri_path(&mut uri, normalized.as_bytes());
1880    }
1881
1882    uri.push_str("?mode=ro&immutable=1");
1883    Ok(PathBuf::from(uri))
1884}
1885
1886fn push_sqlite_uri_path(uri: &mut String, bytes: &[u8]) {
1887    const HEX: &[u8; 16] = b"0123456789ABCDEF";
1888    for &byte in bytes {
1889        if byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.' | b'~' | b'/') {
1890            uri.push(byte as char);
1891        } else {
1892            uri.push('%');
1893            uri.push(HEX[(byte >> 4) as usize] as char);
1894            uri.push(HEX[(byte & 0x0f) as usize] as char);
1895        }
1896    }
1897}
1898
1899fn open_reader_connection(path: &Path, config: &PoolConfig) -> Result<Connection, SqliteError> {
1900    let conn = Connection::open_with_flags(path, reader_open_flags())?;
1901    configure_reader_connection(&conn, config)?;
1902    Ok(conn)
1903}
1904
1905fn writer_open_flags() -> OpenFlags {
1906    OpenFlags::SQLITE_OPEN_READ_WRITE
1907        | OpenFlags::SQLITE_OPEN_CREATE
1908        | OpenFlags::SQLITE_OPEN_URI
1909        | OpenFlags::SQLITE_OPEN_NO_MUTEX
1910}
1911
1912/// Read-only writer-slot open flags: no `SQLITE_OPEN_CREATE`, so a missing
1913/// path is rejected rather than silently created.
1914fn writer_read_only_open_flags() -> OpenFlags {
1915    OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_URI | OpenFlags::SQLITE_OPEN_NO_MUTEX
1916}
1917
1918fn reader_open_flags() -> OpenFlags {
1919    OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_URI | OpenFlags::SQLITE_OPEN_NO_MUTEX
1920}
1921
1922fn configure_writer_connection(
1923    conn: &Connection,
1924    config: &PoolConfig,
1925) -> Result<bool, SqliteError> {
1926    if config.read_only {
1927        // Read-only writer slot: skip write-intent PRAGMAs (journal_mode,
1928        // wal_autocheckpoint, journal_size_limit all require write access to
1929        // change) and lock the connection down with query_only instead.
1930        conn.pragma_update(None, "foreign_keys", "ON")?;
1931        conn.busy_timeout(config.busy_timeout)?;
1932        conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
1933        conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
1934        conn.pragma_update(None, "temp_store", "MEMORY")?;
1935        conn.pragma_update(None, "query_only", "ON")?;
1936
1937        let wal_enabled =
1938            config.wal_mode && current_journal_mode(conn)?.eq_ignore_ascii_case("wal");
1939        return Ok(wal_enabled);
1940    }
1941
1942    let wants_wal = config.path.is_some() && config.wal_mode;
1943
1944    if wants_wal {
1945        conn.pragma_update(None, "journal_mode", "WAL")?;
1946    }
1947
1948    conn.pragma_update(None, "synchronous", "NORMAL")?;
1949    conn.pragma_update(None, "foreign_keys", "ON")?;
1950    conn.busy_timeout(config.busy_timeout)?;
1951    conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
1952    conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
1953    conn.pragma_update(None, "temp_store", "MEMORY")?;
1954    // The pool's startup writer always opens before any checkpoint owner can
1955    // claim the pool, so it starts on the bounded fallback;
1956    // `claim_checkpoint_ownership` re-applies the pragma on this connection
1957    // under the writer mutex when a dedicated owner attaches.
1958    conn.pragma_update(
1959        None,
1960        "wal_autocheckpoint",
1961        FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
1962    )?;
1963
1964    let wal_enabled = wants_wal && current_journal_mode(conn)?.eq_ignore_ascii_case("wal");
1965
1966    if wal_enabled {
1967        conn.pragma_update(None, "journal_size_limit", config.journal_size_limit_bytes)?;
1968    }
1969
1970    Ok(wal_enabled)
1971}
1972
1973fn configure_reader_connection(conn: &Connection, config: &PoolConfig) -> Result<(), SqliteError> {
1974    conn.pragma_update(None, "foreign_keys", "ON")?;
1975    conn.busy_timeout(config.busy_timeout)?;
1976    conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
1977    conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
1978    conn.pragma_update(None, "temp_store", "MEMORY")?;
1979    Ok(())
1980}
1981
1982fn current_journal_mode(conn: &Connection) -> Result<String, SqliteError> {
1983    conn.pragma_query_value(None, "journal_mode", |row| row.get::<_, String>(0))
1984        .map(|mode| mode.to_ascii_lowercase())
1985        .map_err(Into::into)
1986}
1987
1988fn reset_reader_connection(conn: &Connection) -> bool {
1989    if conn.is_autocommit() {
1990        return true;
1991    }
1992
1993    match conn.execute_batch("ROLLBACK") {
1994        Ok(()) => conn.is_autocommit(),
1995        Err(rusqlite::Error::SqliteFailure(err, _)) => {
1996            if matches!(
1997                err.code,
1998                rusqlite::ErrorCode::CannotOpen
1999                    | rusqlite::ErrorCode::DatabaseCorrupt
2000                    | rusqlite::ErrorCode::NotADatabase
2001                    | rusqlite::ErrorCode::DiskFull
2002            ) {
2003                return false;
2004            }
2005            conn.is_autocommit()
2006        }
2007        Err(_) => false,
2008    }
2009}
2010
2011fn reader_connection_is_healthy(conn: &Connection) -> bool {
2012    match conn.query_row("SELECT 1", [], |row| row.get::<_, i64>(0)) {
2013        Ok(_) => true,
2014        Err(rusqlite::Error::SqliteFailure(err, _)) => !matches!(
2015            err.code,
2016            rusqlite::ErrorCode::CannotOpen
2017                | rusqlite::ErrorCode::NotADatabase
2018                | rusqlite::ErrorCode::DatabaseCorrupt
2019                | rusqlite::ErrorCode::PermissionDenied
2020                | rusqlite::ErrorCode::SystemIoFailure
2021        ),
2022        Err(_) => true,
2023    }
2024}
2025
2026fn close_connection_quietly(conn: Connection) {
2027    match conn.close() {
2028        Ok(()) => {}
2029        Err((conn, _)) => drop(conn),
2030    }
2031}
2032
2033fn pool_exhausted_error(timeout: Duration, max_readers: usize) -> SqliteError {
2034    rusqlite::Error::SqliteFailure(
2035        rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_BUSY),
2036        Some(format!(
2037            "Pool exhausted: no reader available after {timeout:?} (max_readers={max_readers})"
2038        )),
2039    )
2040    .into()
2041}
2042
2043#[cfg(test)]
2044mod tests {
2045    use super::*;
2046    use serial_test::serial;
2047
2048    struct WarningCapture {
2049        messages: Arc<std::sync::Mutex<Vec<String>>>,
2050    }
2051
2052    impl tracing::Subscriber for WarningCapture {
2053        fn enabled(&self, _: &tracing::Metadata<'_>) -> bool {
2054            true
2055        }
2056
2057        fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
2058            tracing::span::Id::from_u64(1)
2059        }
2060
2061        fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
2062
2063        fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
2064
2065        fn event(&self, event: &tracing::Event<'_>) {
2066            struct Visitor(Option<String>);
2067
2068            impl tracing::field::Visit for Visitor {
2069                fn record_debug(
2070                    &mut self,
2071                    field: &tracing::field::Field,
2072                    value: &dyn std::fmt::Debug,
2073                ) {
2074                    if field.name() == "message" {
2075                        self.0 = Some(format!("{value:?}"));
2076                    }
2077                }
2078            }
2079
2080            let mut visitor = Visitor(None);
2081            event.record(&mut visitor);
2082            if let Some(message) = visitor.0 {
2083                self.messages.lock().unwrap().push(message);
2084            }
2085        }
2086
2087        fn enter(&self, _: &tracing::span::Id) {}
2088
2089        fn exit(&self, _: &tracing::span::Id) {}
2090    }
2091
2092    /// Restores the process CWD on drop — including on panic — so a mid-test
2093    /// assertion failure (or an unexpected panic from the code under test)
2094    /// can never leave the process chdir'd into a `tempfile::tempdir()` that
2095    /// unwinds out from under every later test sharing this process.
2096    struct CwdGuard {
2097        original: PathBuf,
2098    }
2099
2100    impl CwdGuard {
2101        fn enter(dir: &Path) -> Self {
2102            let original = std::env::current_dir().unwrap();
2103            std::env::set_current_dir(dir).unwrap();
2104            Self { original }
2105        }
2106    }
2107
2108    impl Drop for CwdGuard {
2109        fn drop(&mut self) {
2110            let _ = std::env::set_current_dir(&self.original);
2111        }
2112    }
2113
2114    const POOL_ENV_VARS: [&str; 7] = [
2115        "KHIVE_BUSY_TIMEOUT_SECS",
2116        "KHIVE_CHECKOUT_TIMEOUT_SECS",
2117        "KHIVE_WAL_AUTOCHECKPOINT_PAGES",
2118        "KHIVE_JOURNAL_SIZE_LIMIT_BYTES",
2119        "KHIVE_WRITE_QUEUE",
2120        "KHIVE_WRITE_QUEUE_CAPACITY",
2121        "KHIVE_WRITE_ROUTING",
2122    ];
2123
2124    struct PoolEnvGuard {
2125        saved: Vec<(&'static str, Option<std::ffi::OsString>)>,
2126    }
2127
2128    impl PoolEnvGuard {
2129        fn capture() -> Self {
2130            Self {
2131                saved: POOL_ENV_VARS
2132                    .into_iter()
2133                    .map(|key| (key, std::env::var_os(key)))
2134                    .collect(),
2135            }
2136        }
2137    }
2138
2139    impl Drop for PoolEnvGuard {
2140        fn drop(&mut self) {
2141            for (key, value) in &self.saved {
2142                match value {
2143                    Some(value) => std::env::set_var(key, value),
2144                    None => std::env::remove_var(key),
2145                }
2146            }
2147        }
2148    }
2149
2150    fn clear_pool_env() -> PoolEnvGuard {
2151        let guard = PoolEnvGuard::capture();
2152        for var in POOL_ENV_VARS {
2153            std::env::remove_var(var);
2154        }
2155        guard
2156    }
2157
2158    fn wal_autocheckpoint_pages(conn: &Connection) -> u32 {
2159        conn.pragma_query_value(None, "wal_autocheckpoint", |row| row.get(0))
2160            .expect("read PRAGMA wal_autocheckpoint")
2161    }
2162
2163    fn journal_size_limit_bytes(conn: &Connection) -> i64 {
2164        conn.pragma_query_value(None, "journal_size_limit", |row| row.get(0))
2165            .expect("read PRAGMA journal_size_limit")
2166    }
2167
2168    #[test]
2169    fn read_only_rollback_journal_pool_keeps_a_dedicated_reader() {
2170        let dir = tempfile::tempdir().unwrap();
2171        let path = dir.path().join("read_only_delete_journal.db");
2172        {
2173            let conn = Connection::open(&path).unwrap();
2174            conn.execute_batch("CREATE TABLE snapshot_row(id INTEGER PRIMARY KEY);")
2175                .unwrap();
2176            let mode: String = conn
2177                .pragma_query_value(None, "journal_mode", |row| row.get(0))
2178                .unwrap();
2179            assert_eq!(mode.to_ascii_lowercase(), "delete");
2180        }
2181
2182        let pool = ConnectionPool::new(PoolConfig {
2183            path: Some(path),
2184            read_only: true,
2185            write_queue_enabled: Some(false),
2186            ..PoolConfig::default()
2187        })
2188        .unwrap();
2189
2190        assert!(
2191            pool.max_readers() > 0,
2192            "a read-only rollback-journal snapshot must use a genuine read-only reader, not \
2193             alias reader() onto the query-only writer slot"
2194        );
2195        let reader = pool.reader().expect("dedicated read-only reader checkout");
2196        let count: i64 = reader
2197            .query_row("SELECT COUNT(*) FROM snapshot_row", [], |row| row.get(0))
2198            .unwrap();
2199        assert_eq!(count, 0);
2200        drop(reader);
2201        assert_eq!(
2202            pool.writer_acquisition_snapshot(),
2203            WriterAcquisitionSnapshot::default(),
2204            "constructing and reading a rollback-journal snapshot must never acquire the writer"
2205        );
2206    }
2207
2208    fn sqlite_sidecar(path: &Path, suffix: &str) -> PathBuf {
2209        let mut sidecar = path.as_os_str().to_os_string();
2210        sidecar.push(suffix);
2211        PathBuf::from(sidecar)
2212    }
2213
2214    fn directory_entries(path: &Path) -> Vec<std::ffi::OsString> {
2215        let mut entries = std::fs::read_dir(path)
2216            .unwrap()
2217            .map(|entry| entry.unwrap().file_name())
2218            .collect::<Vec<_>>();
2219        entries.sort();
2220        entries
2221    }
2222
2223    /// A persistent-WAL snapshot can carry committed rows that exist only in
2224    /// `<db>-wal`. With no copied `-shm`, immutable SQLite silently ignores
2225    /// those frames while ordinary read-only SQLite creates a new `-shm`.
2226    /// Refuse before either open strategy can lose data or mutate the source.
2227    #[test]
2228    fn read_only_persistent_wal_without_shm_is_refused_without_mutation() {
2229        let dir = tempfile::tempdir().unwrap();
2230        let source = dir.path().join("wal-source.db");
2231        let snapshot = dir.path().join("snapshot ?#%.db");
2232        let source_wal = sqlite_sidecar(&source, "-wal");
2233        let snapshot_wal = sqlite_sidecar(&snapshot, "-wal");
2234        let snapshot_shm = sqlite_sidecar(&snapshot, "-shm");
2235
2236        let source_conn = Connection::open(&source).unwrap();
2237        let mode: String = source_conn
2238            .pragma_update_and_check(None, "journal_mode", "WAL", |row| row.get(0))
2239            .unwrap();
2240        assert_eq!(mode.to_ascii_lowercase(), "wal");
2241        source_conn
2242            .pragma_update(None, "wal_autocheckpoint", 0)
2243            .unwrap();
2244        source_conn
2245            .execute_batch(
2246                "CREATE TABLE snapshot_row(id INTEGER PRIMARY KEY, body TEXT NOT NULL);\
2247                 INSERT INTO snapshot_row(body) VALUES ('committed-only-in-wal');",
2248            )
2249            .unwrap();
2250        assert!(source_wal.exists(), "fixture must retain a WAL sidecar");
2251
2252        std::fs::copy(&source, &snapshot).unwrap();
2253        std::fs::copy(&source_wal, &snapshot_wal).unwrap();
2254        assert!(
2255            !snapshot_shm.exists(),
2256            "fixture intentionally omits the transient shared-memory index"
2257        );
2258
2259        let main_before = std::fs::read(&snapshot).unwrap();
2260        let wal_before = std::fs::read(&snapshot_wal).unwrap();
2261        let entries_before = directory_entries(dir.path());
2262
2263        let error = match ConnectionPool::new(PoolConfig {
2264            path: Some(snapshot.clone()),
2265            read_only: true,
2266            write_queue_enabled: Some(false),
2267            ..PoolConfig::default()
2268        }) {
2269            Ok(_) => panic!("a non-empty WAL without its frozen -shm must fail closed"),
2270            Err(error) => error,
2271        };
2272        assert!(
2273            error
2274                .to_string()
2275                .contains("would omit committed WAL frames"),
2276            "diagnostic must explain why neither unsafe open mode is allowed: {error}"
2277        );
2278
2279        assert_eq!(std::fs::read(&snapshot).unwrap(), main_before);
2280        assert_eq!(std::fs::read(&snapshot_wal).unwrap(), wal_before);
2281        assert_eq!(directory_entries(dir.path()), entries_before);
2282        assert!(
2283            !snapshot_shm.exists(),
2284            "read-only admission and every reader must keep the source free of -shm"
2285        );
2286
2287        drop(source_conn);
2288    }
2289
2290    /// A complete frozen WAL snapshot includes the WAL index. Once all three
2291    /// files are read-only, ordinary SQLite read-only mode consumes the
2292    /// committed WAL frames without changing the source. This is intentionally
2293    /// not `immutable=1`: immutable SQLite ignores WAL contents.
2294    #[test]
2295    fn read_only_persistent_wal_with_read_only_shm_reads_without_mutation() {
2296        let dir = tempfile::tempdir().unwrap();
2297        let source = dir.path().join("wal-source.db");
2298        let snapshot = dir.path().join("frozen-wal-snapshot.db");
2299        let source_wal = sqlite_sidecar(&source, "-wal");
2300        let source_shm = sqlite_sidecar(&source, "-shm");
2301        let snapshot_wal = sqlite_sidecar(&snapshot, "-wal");
2302        let snapshot_shm = sqlite_sidecar(&snapshot, "-shm");
2303
2304        let source_conn = Connection::open(&source).unwrap();
2305        let mode: String = source_conn
2306            .pragma_update_and_check(None, "journal_mode", "WAL", |row| row.get(0))
2307            .unwrap();
2308        assert_eq!(mode.to_ascii_lowercase(), "wal");
2309        source_conn
2310            .pragma_update(None, "wal_autocheckpoint", 0)
2311            .unwrap();
2312        source_conn
2313            .execute_batch(
2314                "CREATE TABLE snapshot_row(id INTEGER PRIMARY KEY, body TEXT NOT NULL);\
2315                 INSERT INTO snapshot_row(body) VALUES ('committed-only-in-wal');",
2316            )
2317            .unwrap();
2318        assert!(source_wal.exists() && source_shm.exists());
2319
2320        std::fs::copy(&source, &snapshot).unwrap();
2321        std::fs::copy(&source_wal, &snapshot_wal).unwrap();
2322        std::fs::copy(&source_shm, &snapshot_shm).unwrap();
2323
2324        let snapshot_paths = [&snapshot, &snapshot_wal, &snapshot_shm];
2325        let original_permissions =
2326            snapshot_paths.map(|path| std::fs::metadata(path).unwrap().permissions());
2327        for path in snapshot_paths {
2328            let mut permissions = std::fs::metadata(path).unwrap().permissions();
2329            permissions.set_readonly(true);
2330            std::fs::set_permissions(path, permissions).unwrap();
2331        }
2332
2333        let main_before = std::fs::read(&snapshot).unwrap();
2334        let wal_before = std::fs::read(&snapshot_wal).unwrap();
2335        let shm_before = std::fs::read(&snapshot_shm).unwrap();
2336        let entries_before = directory_entries(dir.path());
2337
2338        let pool = ConnectionPool::new(PoolConfig {
2339            path: Some(snapshot.clone()),
2340            read_only: true,
2341            write_queue_enabled: Some(false),
2342            ..PoolConfig::default()
2343        })
2344        .unwrap();
2345        let reader = pool.reader().unwrap();
2346        let body: String = reader
2347            .query_row("SELECT body FROM snapshot_row", [], |row| row.get(0))
2348            .unwrap();
2349        assert_eq!(body, "committed-only-in-wal");
2350        drop(reader);
2351
2352        let standalone = pool.open_standalone_reader().unwrap();
2353        let count: i64 = standalone
2354            .query_row("SELECT COUNT(*) FROM snapshot_row", [], |row| row.get(0))
2355            .unwrap();
2356        assert_eq!(count, 1);
2357        drop(standalone);
2358        drop(pool);
2359
2360        assert_eq!(std::fs::read(&snapshot).unwrap(), main_before);
2361        assert_eq!(std::fs::read(&snapshot_wal).unwrap(), wal_before);
2362        assert_eq!(std::fs::read(&snapshot_shm).unwrap(), shm_before);
2363        assert_eq!(directory_entries(dir.path()), entries_before);
2364
2365        for (path, permissions) in snapshot_paths.into_iter().zip(original_permissions) {
2366            std::fs::set_permissions(path, permissions).unwrap();
2367        }
2368        drop(source_conn);
2369    }
2370
2371    /// The configured spelling must not decide which WAL sidecars SQLite sees.
2372    /// A symlinked snapshot is classified and opened through one canonical
2373    /// physical path so committed frames beside the target remain visible and
2374    /// no sidecars are ever derived beside the alias.
2375    #[cfg(unix)]
2376    #[test]
2377    fn read_only_frozen_wal_symlink_reads_target_frames_without_mutation() {
2378        use std::os::unix::fs::symlink;
2379
2380        let dir = tempfile::tempdir().unwrap();
2381        let source = dir.path().join("wal-source.db");
2382        let snapshot = dir.path().join("frozen-target.db");
2383        let alias = dir.path().join("frozen-alias.db");
2384        let source_wal = sqlite_sidecar(&source, "-wal");
2385        let source_shm = sqlite_sidecar(&source, "-shm");
2386        let snapshot_wal = sqlite_sidecar(&snapshot, "-wal");
2387        let snapshot_shm = sqlite_sidecar(&snapshot, "-shm");
2388        let alias_wal = sqlite_sidecar(&alias, "-wal");
2389        let alias_shm = sqlite_sidecar(&alias, "-shm");
2390
2391        let source_conn = Connection::open(&source).unwrap();
2392        let mode: String = source_conn
2393            .pragma_update_and_check(None, "journal_mode", "WAL", |row| row.get(0))
2394            .unwrap();
2395        assert_eq!(mode.to_ascii_lowercase(), "wal");
2396        source_conn
2397            .pragma_update(None, "wal_autocheckpoint", 0)
2398            .unwrap();
2399        source_conn
2400            .execute_batch(
2401                "CREATE TABLE snapshot_row(id INTEGER PRIMARY KEY, body TEXT NOT NULL);\
2402                 INSERT INTO snapshot_row(body) VALUES ('visible-through-target-wal');",
2403            )
2404            .unwrap();
2405        assert!(source_wal.exists() && source_shm.exists());
2406
2407        std::fs::copy(&source, &snapshot).unwrap();
2408        std::fs::copy(&source_wal, &snapshot_wal).unwrap();
2409        std::fs::copy(&source_shm, &snapshot_shm).unwrap();
2410        symlink(&snapshot, &alias).unwrap();
2411        assert!(!alias_wal.exists() && !alias_shm.exists());
2412
2413        let snapshot_paths = [&snapshot, &snapshot_wal, &snapshot_shm];
2414        let original_permissions =
2415            snapshot_paths.map(|path| std::fs::metadata(path).unwrap().permissions());
2416        for path in snapshot_paths {
2417            let mut permissions = std::fs::metadata(path).unwrap().permissions();
2418            permissions.set_readonly(true);
2419            std::fs::set_permissions(path, permissions).unwrap();
2420        }
2421
2422        let main_before = std::fs::read(&snapshot).unwrap();
2423        let wal_before = std::fs::read(&snapshot_wal).unwrap();
2424        let shm_before = std::fs::read(&snapshot_shm).unwrap();
2425        let entries_before = directory_entries(dir.path());
2426
2427        let pool = ConnectionPool::new(PoolConfig {
2428            path: Some(alias.clone()),
2429            read_only: true,
2430            write_queue_enabled: Some(false),
2431            ..PoolConfig::default()
2432        })
2433        .unwrap();
2434        let reader = pool.reader().unwrap();
2435        let body: String = reader
2436            .query_row("SELECT body FROM snapshot_row", [], |row| row.get(0))
2437            .unwrap();
2438        assert_eq!(body, "visible-through-target-wal");
2439        drop(reader);
2440        let standalone = pool.open_standalone_reader().unwrap();
2441        let count: i64 = standalone
2442            .query_row("SELECT COUNT(*) FROM snapshot_row", [], |row| row.get(0))
2443            .unwrap();
2444        assert_eq!(count, 1);
2445        drop(standalone);
2446        drop(pool);
2447
2448        assert_eq!(std::fs::read(&snapshot).unwrap(), main_before);
2449        assert_eq!(std::fs::read(&snapshot_wal).unwrap(), wal_before);
2450        assert_eq!(std::fs::read(&snapshot_shm).unwrap(), shm_before);
2451        assert_eq!(directory_entries(dir.path()), entries_before);
2452        assert!(!alias_wal.exists() && !alias_shm.exists());
2453
2454        for (path, permissions) in snapshot_paths.into_iter().zip(original_permissions) {
2455            std::fs::set_permissions(path, permissions).unwrap();
2456        }
2457        drop(source_conn);
2458    }
2459
2460    /// A clean persistent-WAL database has no committed frames outside the
2461    /// checkpointed main file. This is the narrow case where an encoded
2462    /// `immutable=1` URI is safe and necessary to prevent SQLite from creating
2463    /// fresh sidecars. Reserved URI bytes in the filesystem path must still
2464    /// resolve to the exact database.
2465    #[test]
2466    fn read_only_clean_wal_snapshot_is_sidecar_free() {
2467        let dir = tempfile::tempdir().unwrap();
2468        let path = dir.path().join("clean snapshot ?#%.db");
2469        {
2470            let conn = Connection::open(&path).unwrap();
2471            let mode: String = conn
2472                .pragma_update_and_check(None, "journal_mode", "WAL", |row| row.get(0))
2473                .unwrap();
2474            assert_eq!(mode.to_ascii_lowercase(), "wal");
2475            conn.execute_batch(
2476                "CREATE TABLE snapshot_row(id INTEGER PRIMARY KEY, body TEXT NOT NULL);\
2477                 INSERT INTO snapshot_row(body) VALUES ('checkpointed');",
2478            )
2479            .unwrap();
2480        }
2481        let wal = sqlite_sidecar(&path, "-wal");
2482        let shm = sqlite_sidecar(&path, "-shm");
2483        assert!(!wal.exists() && !shm.exists());
2484        assert!(sqlite_header_uses_wal(&path).unwrap());
2485
2486        let original_permissions = std::fs::metadata(&path).unwrap().permissions();
2487        let mut read_only_permissions = original_permissions.clone();
2488        read_only_permissions.set_readonly(true);
2489        std::fs::set_permissions(&path, read_only_permissions).unwrap();
2490        let main_before = std::fs::read(&path).unwrap();
2491        let entries_before = directory_entries(dir.path());
2492
2493        let pool = ConnectionPool::new(PoolConfig {
2494            path: Some(path.clone()),
2495            read_only: true,
2496            write_queue_enabled: Some(false),
2497            ..PoolConfig::default()
2498        })
2499        .unwrap();
2500        let reader = pool.reader().unwrap();
2501        let body: String = reader
2502            .query_row("SELECT body FROM snapshot_row", [], |row| row.get(0))
2503            .unwrap();
2504        assert_eq!(body, "checkpointed");
2505        drop(reader);
2506        let standalone = pool.open_standalone_reader().unwrap();
2507        let count: i64 = standalone
2508            .query_row("SELECT COUNT(*) FROM snapshot_row", [], |row| row.get(0))
2509            .unwrap();
2510        assert_eq!(count, 1);
2511        drop(standalone);
2512        drop(pool);
2513
2514        assert_eq!(std::fs::read(&path).unwrap(), main_before);
2515        assert_eq!(directory_entries(dir.path()), entries_before);
2516        assert!(!wal.exists() && !shm.exists());
2517        std::fs::set_permissions(&path, original_permissions).unwrap();
2518    }
2519
2520    /// `immutable=1` is unsafe for a database that can still change and is not
2521    /// needed for rollback-journal reads. Keep ordinary SQLite locking/change
2522    /// detection there so an already-open read-only pool observes a later
2523    /// committed transaction from a live writer.
2524    #[test]
2525    fn read_only_live_rollback_journal_keeps_change_detection() {
2526        let dir = tempfile::tempdir().unwrap();
2527        let path = dir.path().join("live-delete-journal.db");
2528        let writer = Connection::open(&path).unwrap();
2529        writer
2530            .execute_batch("CREATE TABLE live_row(id INTEGER PRIMARY KEY);")
2531            .unwrap();
2532
2533        let pool = ConnectionPool::new(PoolConfig {
2534            path: Some(path),
2535            read_only: true,
2536            write_queue_enabled: Some(false),
2537            ..PoolConfig::default()
2538        })
2539        .unwrap();
2540        {
2541            let reader = pool.reader().unwrap();
2542            let count: i64 = reader
2543                .query_row("SELECT COUNT(*) FROM live_row", [], |row| row.get(0))
2544                .unwrap();
2545            assert_eq!(count, 0);
2546        }
2547
2548        writer
2549            .execute("INSERT INTO live_row DEFAULT VALUES", [])
2550            .unwrap();
2551        let reader = pool.reader().unwrap();
2552        let count: i64 = reader
2553            .query_row("SELECT COUNT(*) FROM live_row", [], |row| row.get(0))
2554            .unwrap();
2555        assert_eq!(
2556            count, 1,
2557            "rollback-journal read-only connections must retain live change detection"
2558        );
2559    }
2560
2561    /// A writable `-shm` beside a WAL database is evidence that the database is
2562    /// not a sidecar-free frozen snapshot (and may have a live writer). Refuse
2563    /// before opening SQLite rather than mutate the shared index or unsafely
2564    /// assert `immutable=1` over a live database.
2565    #[test]
2566    fn read_only_live_wal_with_writable_shm_is_refused_without_mutation() {
2567        let dir = tempfile::tempdir().unwrap();
2568        let path = dir.path().join("live-wal.db");
2569        let wal = sqlite_sidecar(&path, "-wal");
2570        let shm = sqlite_sidecar(&path, "-shm");
2571        let writer = Connection::open(&path).unwrap();
2572        writer.pragma_update(None, "journal_mode", "WAL").unwrap();
2573        writer
2574            .execute_batch(
2575                "CREATE TABLE live_row(id INTEGER PRIMARY KEY);\
2576                 INSERT INTO live_row DEFAULT VALUES;",
2577            )
2578            .unwrap();
2579        assert!(wal.exists() && shm.exists());
2580
2581        let main_before = std::fs::read(&path).unwrap();
2582        let wal_before = std::fs::read(&wal).unwrap();
2583        let shm_before = std::fs::read(&shm).unwrap();
2584        let error = match ConnectionPool::new(PoolConfig {
2585            path: Some(path.clone()),
2586            read_only: true,
2587            write_queue_enabled: Some(false),
2588            ..PoolConfig::default()
2589        }) {
2590            Ok(_) => panic!("a live WAL database with writable -shm must fail closed"),
2591            Err(error) => error,
2592        };
2593        assert!(
2594            error
2595                .to_string()
2596                .contains("writable WAL shared-memory sidecar"),
2597            "diagnostic must explain how to freeze the snapshot: {error}"
2598        );
2599        assert_eq!(std::fs::read(&path).unwrap(), main_before);
2600        assert_eq!(std::fs::read(&wal).unwrap(), wal_before);
2601        assert_eq!(std::fs::read(&shm).unwrap(), shm_before);
2602
2603        drop(writer);
2604    }
2605
2606    /// A symlink cannot hide a writable target `-shm`. Admission inspects the
2607    /// canonical target sidecar set before SQLite opens any connection and
2608    /// therefore refuses a potentially live WAL without touching either
2609    /// target or alias-adjacent paths.
2610    #[cfg(unix)]
2611    #[test]
2612    fn read_only_live_wal_symlink_rejects_target_writable_shm_without_mutation() {
2613        use std::os::unix::fs::symlink;
2614
2615        let dir = tempfile::tempdir().unwrap();
2616        let target = dir.path().join("live-target.db");
2617        let alias = dir.path().join("live-alias.db");
2618        let target_wal = sqlite_sidecar(&target, "-wal");
2619        let target_shm = sqlite_sidecar(&target, "-shm");
2620        let alias_wal = sqlite_sidecar(&alias, "-wal");
2621        let alias_shm = sqlite_sidecar(&alias, "-shm");
2622
2623        let writer = Connection::open(&target).unwrap();
2624        writer.pragma_update(None, "journal_mode", "WAL").unwrap();
2625        writer
2626            .execute_batch(
2627                "CREATE TABLE live_row(id INTEGER PRIMARY KEY);\
2628                 INSERT INTO live_row DEFAULT VALUES;",
2629            )
2630            .unwrap();
2631        assert!(target_wal.exists() && target_shm.exists());
2632        symlink(&target, &alias).unwrap();
2633        assert!(!alias_wal.exists() && !alias_shm.exists());
2634
2635        let main_before = std::fs::read(&target).unwrap();
2636        let wal_before = std::fs::read(&target_wal).unwrap();
2637        let shm_before = std::fs::read(&target_shm).unwrap();
2638        let entries_before = directory_entries(dir.path());
2639
2640        let error = match ConnectionPool::new(PoolConfig {
2641            path: Some(alias),
2642            read_only: true,
2643            write_queue_enabled: Some(false),
2644            ..PoolConfig::default()
2645        }) {
2646            Ok(_) => panic!("a symlink must not hide the target's writable -shm"),
2647            Err(error) => error,
2648        };
2649        assert!(
2650            error
2651                .to_string()
2652                .contains("writable WAL shared-memory sidecar"),
2653            "diagnostic must identify the canonical target's live sidecar: {error}"
2654        );
2655        assert_eq!(std::fs::read(&target).unwrap(), main_before);
2656        assert_eq!(std::fs::read(&target_wal).unwrap(), wal_before);
2657        assert_eq!(std::fs::read(&target_shm).unwrap(), shm_before);
2658        assert_eq!(directory_entries(dir.path()), entries_before);
2659        assert!(!alias_wal.exists() && !alias_shm.exists());
2660
2661        drop(writer);
2662    }
2663
2664    #[test]
2665    #[serial]
2666    fn pool_config_default_values_match_constants() {
2667        // Ensure defaults are not accidentally changed. The process env may
2668        // legitimately carry overrides (CI jobs set KHIVE_CHECKOUT_TIMEOUT_SECS),
2669        // so clear them first — this test asserts the constants, not the env.
2670        let _pool_env = clear_pool_env();
2671        let cfg = PoolConfig::default();
2672        assert_eq!(
2673            cfg.journal_size_limit_bytes,
2674            DEFAULT_JOURNAL_SIZE_LIMIT_BYTES
2675        );
2676        assert_eq!(cfg.busy_timeout, Duration::from_secs(30));
2677        assert_eq!(cfg.checkout_timeout, Duration::from_secs(5));
2678    }
2679
2680    #[test]
2681    #[serial]
2682    fn legacy_env_cannot_change_wal_autocheckpoint() {
2683        let _pool_env = clear_pool_env();
2684        std::env::set_var("KHIVE_WAL_AUTOCHECKPOINT_PAGES", "8000");
2685        let dir = tempfile::tempdir().unwrap();
2686        let path = dir.path().join("legacy_autocheckpoint_env.db");
2687        let pool = ConnectionPool::new(PoolConfig {
2688            path: Some(path),
2689            ..PoolConfig::default()
2690        })
2691        .expect("pool open");
2692        {
2693            let writer = pool.writer().expect("writer");
2694            assert_eq!(
2695                wal_autocheckpoint_pages(writer.conn()),
2696                FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
2697                "the removed env override must not change the unclaimed fallback"
2698            );
2699        }
2700        pool.claim_checkpoint_ownership().expect("claim ownership");
2701        let writer = pool.writer().expect("writer after claim");
2702        assert_eq!(
2703            wal_autocheckpoint_pages(writer.conn()),
2704            0,
2705            "the removed env override must not change the claimed-owner setting"
2706        );
2707        std::env::remove_var("KHIVE_WAL_AUTOCHECKPOINT_PAGES");
2708    }
2709
2710    #[test]
2711    #[serial]
2712    fn pool_config_env_override_journal_size_limit() {
2713        std::env::set_var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES", "134217728");
2714        let cfg = PoolConfig::default();
2715        std::env::remove_var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES");
2716        assert_eq!(cfg.journal_size_limit_bytes, 134_217_728);
2717    }
2718
2719    #[test]
2720    #[serial]
2721    fn pool_config_env_override_busy_timeout() {
2722        std::env::set_var("KHIVE_BUSY_TIMEOUT_SECS", "60");
2723        let cfg = PoolConfig::default();
2724        std::env::remove_var("KHIVE_BUSY_TIMEOUT_SECS");
2725        assert_eq!(cfg.busy_timeout, Duration::from_secs(60));
2726    }
2727
2728    #[test]
2729    #[serial]
2730    fn pool_config_env_override_checkout_timeout() {
2731        std::env::set_var("KHIVE_CHECKOUT_TIMEOUT_SECS", "10");
2732        let cfg = PoolConfig::default();
2733        std::env::remove_var("KHIVE_CHECKOUT_TIMEOUT_SECS");
2734        assert_eq!(cfg.checkout_timeout, Duration::from_secs(10));
2735    }
2736
2737    #[test]
2738    #[serial]
2739    fn pool_config_write_queue_defaults_unset() {
2740        let _pool_env = clear_pool_env();
2741        let cfg = PoolConfig::default();
2742        assert_eq!(cfg.write_queue_enabled, None);
2743        assert_eq!(cfg.write_queue_capacity, DEFAULT_WRITE_QUEUE_CAPACITY);
2744    }
2745
2746    #[test]
2747    #[serial]
2748    fn clear_pool_env_restores_overrides_on_drop() {
2749        let _ambient_env = PoolEnvGuard::capture();
2750        std::env::set_var("KHIVE_BUSY_TIMEOUT_SECS", "73");
2751
2752        {
2753            let _pool_env = clear_pool_env();
2754            assert_eq!(std::env::var_os("KHIVE_BUSY_TIMEOUT_SECS"), None);
2755        }
2756
2757        assert_eq!(
2758            std::env::var_os("KHIVE_BUSY_TIMEOUT_SECS"),
2759            Some(std::ffi::OsString::from("73"))
2760        );
2761    }
2762
2763    #[test]
2764    #[serial]
2765    fn pool_config_env_override_write_queue_enabled() {
2766        std::env::set_var("KHIVE_WRITE_QUEUE", "1");
2767        let cfg = PoolConfig::default();
2768        std::env::remove_var("KHIVE_WRITE_QUEUE");
2769        assert_eq!(cfg.write_queue_enabled, Some(true));
2770    }
2771
2772    #[test]
2773    #[serial]
2774    fn pool_config_env_override_write_queue_enabled_accepts_true_case_insensitive() {
2775        std::env::set_var("KHIVE_WRITE_QUEUE", "True");
2776        let cfg = PoolConfig::default();
2777        std::env::remove_var("KHIVE_WRITE_QUEUE");
2778        assert_eq!(cfg.write_queue_enabled, Some(true));
2779    }
2780
2781    #[test]
2782    #[serial]
2783    fn pool_config_env_override_write_queue_enabled_accepts_zero_as_explicit_off() {
2784        std::env::set_var("KHIVE_WRITE_QUEUE", "0");
2785        let cfg = PoolConfig::default();
2786        std::env::remove_var("KHIVE_WRITE_QUEUE");
2787        assert_eq!(cfg.write_queue_enabled, Some(false));
2788    }
2789
2790    /// A SET-but-non-Unicode `KHIVE_WRITE_QUEUE` value (invalid UTF-8 on
2791    /// unix) must count as SET — `Some(false)` ("any SET value other than
2792    /// 1/true means off"), never a fall-through to the file-backed default.
2793    /// That is why `PoolConfig::default()` reads `var_os`, not `var`.
2794    #[cfg(unix)]
2795    #[test]
2796    #[serial]
2797    fn pool_config_env_override_write_queue_non_unicode_value_is_explicit_off() {
2798        use std::os::unix::ffi::OsStrExt;
2799        let _pool_env = clear_pool_env();
2800        std::env::set_var(
2801            "KHIVE_WRITE_QUEUE",
2802            std::ffi::OsStr::from_bytes(b"\xff\xfe"),
2803        );
2804        let cfg = PoolConfig::default();
2805        assert_eq!(cfg.write_queue_enabled, Some(false));
2806    }
2807
2808    #[test]
2809    #[serial]
2810    fn pool_config_env_override_write_queue_invalid_value_is_explicit_off() {
2811        // Documented contract (`write_queue_enabled` docs): `"1"`/`"true"`
2812        // (case-insensitive) set `Some(true)`; any other value — garbage
2813        // included — sets `Some(false)`, never `None`.
2814        std::env::set_var("KHIVE_WRITE_QUEUE", "banana");
2815        let cfg = PoolConfig::default();
2816        std::env::remove_var("KHIVE_WRITE_QUEUE");
2817        assert_eq!(cfg.write_queue_enabled, Some(false));
2818    }
2819
2820    #[test]
2821    #[serial]
2822    fn pool_config_write_routing_strict_defaults_off() {
2823        let _pool_env = clear_pool_env();
2824        let cfg = PoolConfig::default();
2825        assert!(!cfg.write_routing_strict);
2826    }
2827
2828    #[test]
2829    #[serial]
2830    fn pool_config_env_override_write_routing_strict() {
2831        std::env::set_var("KHIVE_WRITE_ROUTING", "strict");
2832        let cfg = PoolConfig::default();
2833        std::env::remove_var("KHIVE_WRITE_ROUTING");
2834        assert!(cfg.write_routing_strict);
2835    }
2836
2837    #[test]
2838    #[serial]
2839    fn pool_config_env_override_write_routing_strict_case_insensitive() {
2840        std::env::set_var("KHIVE_WRITE_ROUTING", "STRICT");
2841        let cfg = PoolConfig::default();
2842        std::env::remove_var("KHIVE_WRITE_ROUTING");
2843        assert!(cfg.write_routing_strict);
2844    }
2845
2846    #[test]
2847    #[serial]
2848    fn pool_config_env_write_routing_ignores_unrecognized_value() {
2849        std::env::set_var("KHIVE_WRITE_ROUTING", "eventual");
2850        let cfg = PoolConfig::default();
2851        std::env::remove_var("KHIVE_WRITE_ROUTING");
2852        assert!(!cfg.write_routing_strict);
2853    }
2854
2855    #[test]
2856    #[serial]
2857    fn pool_config_env_override_write_queue_capacity() {
2858        std::env::set_var("KHIVE_WRITE_QUEUE_CAPACITY", "64");
2859        let cfg = PoolConfig::default();
2860        std::env::remove_var("KHIVE_WRITE_QUEUE_CAPACITY");
2861        assert_eq!(cfg.write_queue_capacity, 64);
2862    }
2863
2864    #[test]
2865    #[serial]
2866    fn pool_config_env_invalid_write_queue_capacity_falls_back_to_default() {
2867        std::env::set_var("KHIVE_WRITE_QUEUE_CAPACITY", "0");
2868        let cfg = PoolConfig::default();
2869        std::env::remove_var("KHIVE_WRITE_QUEUE_CAPACITY");
2870        assert_eq!(cfg.write_queue_capacity, DEFAULT_WRITE_QUEUE_CAPACITY);
2871    }
2872
2873    #[test]
2874    #[serial]
2875    fn pool_config_invalid_journal_size_limit_falls_back_to_default() {
2876        std::env::set_var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES", "");
2877        let cfg = PoolConfig::default();
2878        std::env::remove_var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES");
2879        assert_eq!(
2880            cfg.journal_size_limit_bytes,
2881            DEFAULT_JOURNAL_SIZE_LIMIT_BYTES
2882        );
2883    }
2884
2885    #[test]
2886    fn file_backed_pool_opens_successfully() {
2887        let dir = tempfile::tempdir().unwrap();
2888        let path = dir.path().join("test_pool.db");
2889        let cfg = PoolConfig {
2890            path: Some(path.clone()),
2891            ..PoolConfig::default()
2892        };
2893        let pool = ConnectionPool::new(cfg).expect("file-backed pool should open");
2894        assert!(path.exists());
2895        assert!(pool.max_readers() > 0);
2896    }
2897
2898    #[test]
2899    fn standalone_wal_writer_uses_configured_journal_size_limit() {
2900        let dir = tempfile::tempdir().unwrap();
2901        let path = dir.path().join("standalone_wal_journal_limit.db");
2902        let configured_limit = 12_345_678;
2903        let pool = ConnectionPool::new(PoolConfig {
2904            path: Some(path),
2905            journal_size_limit_bytes: configured_limit,
2906            write_queue_enabled: Some(false),
2907            ..PoolConfig::default()
2908        })
2909        .expect("WAL pool open");
2910
2911        let standalone = pool
2912            .open_standalone_writer_untracked()
2913            .expect("standalone WAL writer open");
2914        assert_eq!(current_journal_mode(&standalone).unwrap(), "wal");
2915        assert_eq!(journal_size_limit_bytes(&standalone), configured_limit);
2916    }
2917
2918    #[test]
2919    fn standalone_rollback_writer_keeps_sqlite_journal_size_limit() {
2920        let dir = tempfile::tempdir().unwrap();
2921        let path = dir.path().join("standalone_rollback_journal_limit.db");
2922        let sqlite_default = {
2923            let conn = Connection::open(&path).expect("seed rollback-journal database");
2924            assert_eq!(current_journal_mode(&conn).unwrap(), "delete");
2925            journal_size_limit_bytes(&conn)
2926        };
2927        let configured_limit = if sqlite_default == 12_345_678 {
2928            23_456_789
2929        } else {
2930            12_345_678
2931        };
2932        let pool = ConnectionPool::new(PoolConfig {
2933            path: Some(path),
2934            wal_mode: false,
2935            journal_size_limit_bytes: configured_limit,
2936            write_queue_enabled: Some(false),
2937            ..PoolConfig::default()
2938        })
2939        .expect("rollback-journal pool open");
2940
2941        let standalone = pool
2942            .open_standalone_writer_untracked()
2943            .expect("standalone rollback-journal writer open");
2944        assert_eq!(current_journal_mode(&standalone).unwrap(), "delete");
2945        assert_eq!(journal_size_limit_bytes(&standalone), sqlite_default);
2946    }
2947
2948    #[test]
2949    fn writer_connections_follow_checkpoint_ownership_claim() {
2950        let dir = tempfile::tempdir().unwrap();
2951        let path = dir.path().join("writer_autocheckpoint.db");
2952        let pool = ConnectionPool::new(PoolConfig {
2953            path: Some(path),
2954            write_queue_enabled: Some(false),
2955            ..PoolConfig::default()
2956        })
2957        .expect("pool open");
2958
2959        // Unclaimed: every writer-capable connection keeps the bounded
2960        // fallback, so a pool without a checkpoint task retains SQLite's own
2961        // WAL reclamation.
2962        {
2963            let writer = pool.writer().expect("pooled writer");
2964            assert_eq!(
2965                wal_autocheckpoint_pages(writer.conn()),
2966                FALLBACK_WAL_AUTOCHECKPOINT_PAGES
2967            );
2968        }
2969        let standalone = pool
2970            .open_standalone_writer()
2971            .expect("standalone writer opened before any claim");
2972        assert_eq!(
2973            wal_autocheckpoint_pages(&standalone),
2974            FALLBACK_WAL_AUTOCHECKPOINT_PAGES
2975        );
2976        drop(standalone);
2977
2978        // Claimed: the already-open pooled writer is re-configured under the
2979        // writer mutex, and every later writer-capable open disables the
2980        // autocheckpoint entirely.
2981        pool.claim_checkpoint_ownership().expect("claim ownership");
2982        {
2983            let writer = pool.writer().expect("pooled writer after claim");
2984            assert_eq!(wal_autocheckpoint_pages(writer.conn()), 0);
2985        }
2986        let claimed_standalone = pool
2987            .open_standalone_writer()
2988            .expect("standalone writer opened after the claim");
2989        assert_eq!(wal_autocheckpoint_pages(&claimed_standalone), 0);
2990        drop(claimed_standalone);
2991
2992        let later_infrastructure = pool
2993            .open_standalone_writer_untracked()
2994            .expect("later infrastructure writer");
2995        assert_eq!(wal_autocheckpoint_pages(&later_infrastructure), 0);
2996
2997        let memory_pool = ConnectionPool::new(PoolConfig {
2998            write_queue_enabled: Some(false),
2999            ..PoolConfig::default()
3000        })
3001        .expect("in-memory pool open");
3002        let memory_writer = memory_pool.writer().expect("in-memory writer");
3003        assert_eq!(
3004            wal_autocheckpoint_pages(memory_writer.conn()),
3005            FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
3006            "an unclaimed in-memory pool keeps the bounded fallback"
3007        );
3008    }
3009
3010    #[test]
3011    fn standalone_writer_waits_for_checkpoint_claim_resolution() {
3012        let dir = tempfile::tempdir().unwrap();
3013        let path = dir.path().join("checkpoint_claim_race.db");
3014        let pool = Arc::new(
3015            ConnectionPool::new(PoolConfig {
3016                path: Some(path),
3017                checkout_timeout: Duration::from_secs(5),
3018                write_queue_enabled: Some(false),
3019                ..PoolConfig::default()
3020            })
3021            .expect("pool open"),
3022        );
3023
3024        let legacy_conn = pool.legacy_conn();
3025        let held_writer = legacy_conn.lock();
3026        let claim_start = Arc::new(std::sync::Barrier::new(2));
3027        let claim_pool = Arc::clone(&pool);
3028        let claim_thread_start = Arc::clone(&claim_start);
3029        let claim_thread = thread::spawn(move || {
3030            claim_thread_start.wait();
3031            claim_pool.claim_checkpoint_ownership()
3032        });
3033        claim_start.wait();
3034
3035        {
3036            let mut state = pool.checkpoint_ownership.state.lock();
3037            while state.phase != CheckpointOwnership::Claiming {
3038                pool.checkpoint_ownership.changed.wait(&mut state);
3039            }
3040        }
3041
3042        let open_start = Arc::new(std::sync::Barrier::new(2));
3043        let open_pool = Arc::clone(&pool);
3044        let open_thread_start = Arc::clone(&open_start);
3045        let open_thread = thread::spawn(move || {
3046            open_thread_start.wait();
3047            let conn = open_pool
3048                .open_standalone_writer()
3049                .expect("standalone writer after claim resolution");
3050            wal_autocheckpoint_pages(&conn)
3051        });
3052        open_start.wait();
3053
3054        {
3055            let mut state = pool.checkpoint_ownership.state.lock();
3056            while state.connection_waiters == 0 {
3057                pool.checkpoint_ownership.changed.wait(&mut state);
3058            }
3059            assert_eq!(state.phase, CheckpointOwnership::Claiming);
3060        }
3061
3062        drop(held_writer);
3063        claim_thread
3064            .join()
3065            .expect("claim thread joins")
3066            .expect("claim succeeds");
3067        assert_eq!(
3068            open_thread.join().expect("standalone-open thread joins"),
3069            0,
3070            "a writer open concurrent with a successful claim must inherit claimed ownership"
3071        );
3072    }
3073
3074    #[test]
3075    fn standalone_fallback_application_linearizes_before_claim_publication() {
3076        let dir = tempfile::tempdir().unwrap();
3077        let path = dir.path().join("checkpoint_open_before_claim.db");
3078        let pool = Arc::new(
3079            ConnectionPool::new(PoolConfig {
3080                path: Some(path),
3081                checkout_timeout: Duration::from_secs(5),
3082                write_queue_enabled: Some(false),
3083                ..PoolConfig::default()
3084            })
3085            .expect("pool open"),
3086        );
3087        let pause = Arc::new(CheckpointConnectionConfigPause::new());
3088        *pool.checkpoint_ownership.connection_config_pause.lock() = Some(Arc::clone(&pause));
3089
3090        let open_pool = Arc::clone(&pool);
3091        let open_thread = thread::spawn(move || {
3092            let conn = open_pool
3093                .open_standalone_writer()
3094                .expect("standalone writer opens");
3095            wal_autocheckpoint_pages(&conn)
3096        });
3097        pause.selected.wait();
3098        assert!(
3099            pool.checkpoint_ownership.state.try_lock().is_none(),
3100            "standalone selection must retain the ownership gate until its PRAGMA is applied"
3101        );
3102
3103        let (claim_observed_tx, claim_observed_rx) = std::sync::mpsc::sync_channel(0);
3104        *pool.checkpoint_ownership.claim_lock_observed.lock() = Some(claim_observed_tx);
3105        let claim_pool = Arc::clone(&pool);
3106        let claim_thread = thread::spawn(move || claim_pool.claim_checkpoint_ownership());
3107        assert!(
3108            claim_observed_rx
3109                .recv()
3110                .expect("claim reports whether it observed gate contention"),
3111            "the claim must attempt the gate between fallback selection and PRAGMA application"
3112        );
3113        pause.resume.wait();
3114
3115        assert_eq!(
3116            open_thread.join().expect("standalone-open thread joins"),
3117            FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
3118            "an open linearized before the claim keeps the fallback"
3119        );
3120        claim_thread
3121            .join()
3122            .expect("claim thread joins")
3123            .expect("claim succeeds after standalone configuration");
3124        assert_eq!(pool.effective_wal_autocheckpoint_pages(), 0);
3125    }
3126
3127    #[test]
3128    fn failed_checkpoint_ownership_claim_keeps_fallback_and_can_be_retried() {
3129        let dir = tempfile::tempdir().unwrap();
3130        let path = dir.path().join("checkpoint_claim_retry.db");
3131        let pool = ConnectionPool::new(PoolConfig {
3132            path: Some(path),
3133            checkout_timeout: Duration::from_millis(1),
3134            write_queue_enabled: Some(false),
3135            ..PoolConfig::default()
3136        })
3137        .expect("pool open");
3138
3139        let legacy_conn = pool.legacy_conn();
3140        let held_writer = legacy_conn.lock();
3141        let error = pool
3142            .claim_checkpoint_ownership()
3143            .expect_err("the held pooled writer must make the claim time out");
3144        assert!(matches!(
3145            error,
3146            SqliteError::WriterPoolCheckoutTimeout { .. }
3147        ));
3148        assert_eq!(
3149            pool.effective_wal_autocheckpoint_pages(),
3150            FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
3151            "a failed claim must leave later writer connections fallback-safe"
3152        );
3153
3154        let fallback_writer = pool
3155            .open_standalone_writer()
3156            .expect("standalone writer after failed claim");
3157        assert_eq!(
3158            wal_autocheckpoint_pages(&fallback_writer),
3159            FALLBACK_WAL_AUTOCHECKPOINT_PAGES
3160        );
3161        drop(fallback_writer);
3162
3163        drop(held_writer);
3164        pool.claim_checkpoint_ownership()
3165            .expect("the ownership claim remains retryable");
3166        assert_eq!(pool.effective_wal_autocheckpoint_pages(), 0);
3167        let writer = pool.writer().expect("pooled writer after successful retry");
3168        assert_eq!(wal_autocheckpoint_pages(writer.conn()), 0);
3169    }
3170
3171    #[test]
3172    fn threshold_crossing_commits_do_not_run_an_implicit_checkpoint_once_claimed() {
3173        const FORMER_AUTOCHECKPOINT_THRESHOLD_PAGES: i64 = FALLBACK_WAL_AUTOCHECKPOINT_PAGES as i64;
3174
3175        let dir = tempfile::tempdir().unwrap();
3176        let path = dir.path().join("no_implicit_checkpoint.db");
3177        let pool = ConnectionPool::new(PoolConfig {
3178            path: Some(path),
3179            write_queue_enabled: Some(false),
3180            ..PoolConfig::default()
3181        })
3182        .expect("pool open");
3183        pool.claim_checkpoint_ownership()
3184            .expect("claim ownership for the dedicated-owner posture");
3185        let writer = pool.writer().expect("pooled writer");
3186        writer
3187            .execute_batch("CREATE TABLE blobs (value BLOB NOT NULL)")
3188            .expect("create fixture table");
3189
3190        let page_size: i64 = writer
3191            .pragma_query_value(None, "page_size", |row| row.get(0))
3192            .expect("read page size");
3193        let payload_bytes = page_size * 32;
3194        for _ in 0..160 {
3195            writer
3196                .execute(
3197                    "INSERT INTO blobs (value) VALUES (zeroblob(?1))",
3198                    [payload_bytes],
3199                )
3200                .expect("autocommit fixture row");
3201        }
3202
3203        let log_frames: i64 = writer
3204            .query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| row.get(1))
3205            .expect("observe WAL frame count");
3206        assert!(
3207            log_frames > FORMER_AUTOCHECKPOINT_THRESHOLD_PAGES,
3208            "the commit sequence must retain more than the former automatic threshold; \
3209             observed {log_frames} frames"
3210        );
3211    }
3212
3213    /// The other half of the ownership model: a writable pool that no
3214    /// checkpoint task ever claims must retain SQLite's own bounded WAL
3215    /// reclamation. The same commit sequence that retains >4,000 frames under
3216    /// a claimed owner must NOT accumulate them here — an implicit
3217    /// autocheckpoint fires on the threshold-crossing commit and drains the
3218    /// WAL, which is the regression guard against unbounded WAL growth (and
3219    /// eventual disk exhaustion) on embedded / one-shot writable pools.
3220    #[test]
3221    fn unclaimed_pool_retains_bounded_autocheckpoint_reclamation() {
3222        let dir = tempfile::tempdir().unwrap();
3223        let path = dir.path().join("bounded_fallback_reclamation.db");
3224        let pool = ConnectionPool::new(PoolConfig {
3225            path: Some(path),
3226            write_queue_enabled: Some(false),
3227            ..PoolConfig::default()
3228        })
3229        .expect("pool open");
3230        let writer = pool.writer().expect("pooled writer");
3231        writer
3232            .execute_batch("CREATE TABLE blobs (value BLOB NOT NULL)")
3233            .expect("create fixture table");
3234
3235        let page_size: i64 = writer
3236            .pragma_query_value(None, "page_size", |row| row.get(0))
3237            .expect("read page size");
3238        let payload_bytes = page_size * 32;
3239        for _ in 0..160 {
3240            writer
3241                .execute(
3242                    "INSERT INTO blobs (value) VALUES (zeroblob(?1))",
3243                    [payload_bytes],
3244                )
3245                .expect("autocommit fixture row");
3246        }
3247
3248        // No PASSIVE pass here — read the frame count via wal_checkpoint's
3249        // log column only after the fixture, exactly as the claimed-owner
3250        // test does. With the bounded fallback live, the autocheckpoint that
3251        // fired on a threshold-crossing commit already drained the WAL, so
3252        // far fewer than the threshold's frames remain.
3253        let log_frames: i64 = writer
3254            .query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| row.get(1))
3255            .expect("observe WAL frame count");
3256        assert!(
3257            log_frames < FALLBACK_WAL_AUTOCHECKPOINT_PAGES as i64,
3258            "an unclaimed pool must reclaim WAL frames via the bounded autocheckpoint; \
3259             observed {log_frames} retained frames"
3260        );
3261    }
3262
3263    #[tokio::test]
3264    #[serial]
3265    async fn unset_write_queue_resolves_on_for_file_backed_pool() {
3266        let _pool_env = clear_pool_env();
3267        let dir = tempfile::tempdir().unwrap();
3268        let path = dir.path().join("unset_file_backed.db");
3269        let pool = ConnectionPool::new(PoolConfig {
3270            path: Some(path),
3271            write_queue_enabled: None,
3272            ..PoolConfig::default()
3273        })
3274        .expect("file-backed pool should open");
3275        assert_eq!(pool.config().write_queue_enabled, Some(true));
3276        // Behavioral half: the resolved value actually routes — a writer
3277        // task spawns for this pool, not merely a config field flipping.
3278        assert!(
3279            pool.writer_task_handle()
3280                .expect("spawn inside a runtime context must not error")
3281                .is_some(),
3282            "resolved-on file-backed pool must actually spawn the writer task"
3283        );
3284    }
3285
3286    #[tokio::test]
3287    #[serial]
3288    async fn unset_write_queue_resolves_off_for_memory_backed_pool() {
3289        let _pool_env = clear_pool_env();
3290        let pool = ConnectionPool::new(PoolConfig {
3291            path: None,
3292            write_queue_enabled: None,
3293            ..PoolConfig::default()
3294        })
3295        .expect("in-memory pool should open");
3296        assert_eq!(pool.config().write_queue_enabled, Some(false));
3297        // Behavioral half: resolved-off means no writer task, even inside a
3298        // runtime context where one could spawn.
3299        assert!(
3300            pool.writer_task_handle()
3301                .expect("disabled queue must resolve without error")
3302                .is_none(),
3303            "resolved-off in-memory pool must not spawn a writer task"
3304        );
3305    }
3306
3307    #[test]
3308    #[serial]
3309    fn explicit_false_stays_off_for_file_backed_pool() {
3310        let _pool_env = clear_pool_env();
3311        let dir = tempfile::tempdir().unwrap();
3312        let path = dir.path().join("explicit_false_file_backed.db");
3313        let pool = ConnectionPool::new(PoolConfig {
3314            path: Some(path),
3315            write_queue_enabled: Some(false),
3316            ..PoolConfig::default()
3317        })
3318        .expect("file-backed pool should open");
3319        assert_eq!(pool.config().write_queue_enabled, Some(false));
3320    }
3321
3322    #[tokio::test]
3323    #[serial]
3324    async fn explicit_true_stays_on_for_memory_backed_pool() {
3325        let _pool_env = clear_pool_env();
3326        let pool = ConnectionPool::new(PoolConfig {
3327            path: None,
3328            write_queue_enabled: Some(true),
3329            ..PoolConfig::default()
3330        })
3331        .expect("in-memory pool should open");
3332        assert_eq!(pool.config().write_queue_enabled, Some(true));
3333        // Pinned behavioral contract: the explicit-on preference survives in
3334        // the stored config, but an in-memory pool cannot host a writer
3335        // task — `writer_task::spawn` fails its standalone-connection open
3336        // and degrades to no writer task, so callers fall back to the
3337        // legacy pool-mutex write path and there is no JoinHandle to drain.
3338        assert!(
3339            pool.writer_task_handle()
3340                .expect("spawn degrade must resolve without error")
3341                .is_none(),
3342            "explicit-on in-memory pool must degrade to no writer task"
3343        );
3344        assert_eq!(
3345            pool.writer_task_spawn_count(),
3346            1,
3347            "the spawn attempt must happen exactly once and degrade, not retry"
3348        );
3349        assert!(
3350            pool.take_writer_task_join().is_none(),
3351            "a degraded spawn stores no JoinHandle to drain"
3352        );
3353    }
3354
3355    #[test]
3356    #[serial]
3357    fn explicit_true_on_memory_pool_warns_but_false_and_none_do_not() {
3358        let _pool_env = clear_pool_env();
3359        let messages = Arc::new(std::sync::Mutex::new(Vec::new()));
3360        let subscriber = WarningCapture {
3361            messages: Arc::clone(&messages),
3362        };
3363
3364        tracing::subscriber::with_default(subscriber, || {
3365            let _explicit_true = ConnectionPool::new(PoolConfig {
3366                path: None,
3367                write_queue_enabled: Some(true),
3368                ..PoolConfig::default()
3369            })
3370            .expect("in-memory pool should open");
3371            let _explicit_false = ConnectionPool::new(PoolConfig {
3372                path: None,
3373                write_queue_enabled: Some(false),
3374                ..PoolConfig::default()
3375            })
3376            .expect("in-memory pool should open");
3377            let _unset = ConnectionPool::new(PoolConfig {
3378                path: None,
3379                write_queue_enabled: None,
3380                ..PoolConfig::default()
3381            })
3382            .expect("in-memory pool should open");
3383        });
3384
3385        let messages = messages.lock().unwrap();
3386        assert_eq!(
3387            messages
3388                .iter()
3389                .filter(|message| message.contains("write queue explicitly requested"))
3390                .count(),
3391            1,
3392            "only an explicit in-memory queue request should warn: {messages:?}"
3393        );
3394        let warning = messages
3395            .iter()
3396            .find(|message| message.contains("write queue explicitly requested"))
3397            .expect("explicit in-memory queue warning should be captured");
3398        assert!(
3399            warning.contains("in-memory pools cannot host a writer task"),
3400            "warning must explain why the request is inert: {messages:?}"
3401        );
3402    }
3403
3404    #[test]
3405    fn standalone_writer_open_counts_its_connection_class_once() {
3406        let dir = tempfile::tempdir().unwrap();
3407        let path = dir.path().join("standalone_writer_counter.db");
3408        let pool = ConnectionPool::new(PoolConfig {
3409            path: Some(path),
3410            ..PoolConfig::default()
3411        })
3412        .expect("file-backed pool");
3413
3414        let _standalone = pool
3415            .open_standalone_writer()
3416            .expect("standalone writer opens");
3417
3418        assert_eq!(
3419            pool.writer_acquisition_snapshot(),
3420            WriterAcquisitionSnapshot {
3421                acquisitions: 1,
3422                pooled_acquisitions: 0,
3423                standalone_acquisitions: 1,
3424                writer_task_acquisitions: 0,
3425                timeouts: 0,
3426                writer_task_begin_busy: 0,
3427                writer_task_begin_errors: 0,
3428                writer_task_request_failures: 0,
3429                writer_task_side_effects_unknown: 0,
3430            },
3431            "the public standalone boundary must contribute to the aggregate exactly once"
3432        );
3433    }
3434
3435    #[test]
3436    fn in_memory_pool_degrades_to_single_connection() {
3437        let cfg = PoolConfig {
3438            path: None,
3439            ..PoolConfig::default()
3440        };
3441        let pool = ConnectionPool::new(cfg).expect("in-memory pool should open");
3442        assert_eq!(pool.max_readers(), 0);
3443    }
3444
3445    #[test]
3446    fn writer_checkout_and_release_works() {
3447        let cfg = PoolConfig {
3448            path: None,
3449            ..PoolConfig::default()
3450        };
3451        let pool = ConnectionPool::new(cfg).unwrap();
3452        {
3453            let _writer = pool.writer().expect("writer checkout should succeed");
3454        }
3455        // After drop, writer should be re-acquirable.
3456        let _writer2 = pool
3457            .writer()
3458            .expect("second writer checkout should succeed");
3459    }
3460
3461    #[test]
3462    fn writer_checkout_snapshot_counts_successes_and_timeouts_at_the_pool_boundary() {
3463        let cfg = PoolConfig {
3464            path: None,
3465            checkout_timeout: Duration::from_millis(1),
3466            ..PoolConfig::default()
3467        };
3468        let pool = ConnectionPool::new(cfg).unwrap();
3469
3470        assert_eq!(
3471            pool.writer_acquisition_snapshot(),
3472            WriterAcquisitionSnapshot::default()
3473        );
3474
3475        let held = pool.writer().expect("first checkout succeeds");
3476        let error = match pool.writer() {
3477            Ok(_) => panic!("the held pool mutex must force a finite-wait timeout"),
3478            Err(error) => error,
3479        };
3480        assert!(
3481            matches!(
3482                &error,
3483                SqliteError::WriterPoolCheckoutTimeout { timeout }
3484                    if *timeout == Duration::from_millis(1)
3485            ),
3486            "timeout must have a stable, structurally matchable stage: {error}"
3487        );
3488        assert_eq!(
3489            pool.writer_acquisition_snapshot(),
3490            WriterAcquisitionSnapshot {
3491                acquisitions: 1,
3492                pooled_acquisitions: 1,
3493                standalone_acquisitions: 0,
3494                writer_task_acquisitions: 0,
3495                timeouts: 1,
3496                // A pool-mutex checkout timeout must NOT bleed into the
3497                // writer-task BEGIN counters: separate stages, separate
3498                // counters. This is the mislabeling guard in assertion form.
3499                writer_task_begin_busy: 0,
3500                writer_task_begin_errors: 0,
3501                writer_task_request_failures: 0,
3502                writer_task_side_effects_unknown: 0,
3503            }
3504        );
3505
3506        drop(held);
3507        let _reacquired = pool.writer().expect("checkout succeeds after release");
3508        assert_eq!(
3509            pool.writer_acquisition_snapshot(),
3510            WriterAcquisitionSnapshot {
3511                acquisitions: 2,
3512                pooled_acquisitions: 2,
3513                standalone_acquisitions: 0,
3514                writer_task_acquisitions: 0,
3515                timeouts: 1,
3516                writer_task_begin_busy: 0,
3517                writer_task_begin_errors: 0,
3518                writer_task_request_failures: 0,
3519                writer_task_side_effects_unknown: 0,
3520            }
3521        );
3522    }
3523
3524    #[test]
3525    fn zero_wait_maintenance_skip_is_not_reported_as_a_checkout_timeout() {
3526        let pool = ConnectionPool::new(PoolConfig::default()).unwrap();
3527        let held = pool.writer().expect("finite-wait checkout succeeds");
3528        let before = pool.writer_acquisition_snapshot();
3529
3530        assert!(
3531            pool.try_writer_nowait().is_err(),
3532            "zero-wait maintenance checkout must skip while held"
3533        );
3534
3535        assert_eq!(
3536            pool.writer_acquisition_snapshot(),
3537            before,
3538            "a checkpoint-style zero-wait skip is not a finite-wait checkout timeout"
3539        );
3540        drop(held);
3541    }
3542
3543    /// ADR-091 Plank 0: `WriterGuard::transaction` registers/deregisters a
3544    /// tx_registry entry around the closure. See
3545    /// crates/khive-db/docs/api/pool.md#writer_guard_transaction_registers_during_closure_only
3546    #[test]
3547    #[serial(tx_registry)]
3548    fn writer_guard_transaction_registers_during_closure_only() {
3549        let cfg = PoolConfig {
3550            path: None,
3551            ..PoolConfig::default()
3552        };
3553        let pool = ConnectionPool::new(cfg).unwrap();
3554        let guard = pool.writer().unwrap();
3555
3556        let mut seen_during_closure = false;
3557        let result: Result<(), SqliteError> = guard.transaction(|_conn| {
3558            seen_during_closure = khive_storage::tx_registry::snapshot()
3559                .iter()
3560                .any(|(_, label)| label.as_deref() == Some("writer_guard_tx"));
3561            Ok(())
3562        });
3563        result.expect("transaction should commit");
3564
3565        assert!(
3566            seen_during_closure,
3567            "expected a writer_guard_tx entry visible inside the closure"
3568        );
3569        assert!(
3570            !khive_storage::tx_registry::snapshot()
3571                .iter()
3572                .any(|(_, label)| label.as_deref() == Some("writer_guard_tx")),
3573            "expected the entry to be gone after the transaction completes"
3574        );
3575    }
3576
3577    /// ADR-067 Component A: `writer_task_handle` must fail loud (typed
3578    /// error, not panic) with no Tokio runtime available. See
3579    /// crates/khive-db/docs/api/pool.md#writer_task_handle_fails_loud_without_tokio_runtime
3580    #[test]
3581    fn writer_task_handle_fails_loud_without_tokio_runtime() {
3582        let dir = tempfile::tempdir().unwrap();
3583        let path = dir.path().join("writer_task_no_runtime.db");
3584        let cfg = PoolConfig {
3585            path: Some(path),
3586            write_queue_enabled: Some(true),
3587            ..PoolConfig::default()
3588        };
3589        let pool = ConnectionPool::new(cfg).expect("file-backed pool should open");
3590
3591        let result = pool.writer_task_handle();
3592
3593        assert!(
3594            matches!(result, Err(StorageError::WriterTaskNoRuntime)),
3595            "expected Err(StorageError::WriterTaskNoRuntime) outside a Tokio \
3596             runtime, got {result:?}"
3597        );
3598        assert_eq!(
3599            pool.writer_task_spawn_count(),
3600            0,
3601            "the guard must reject before ever attempting tokio::spawn"
3602        );
3603    }
3604
3605    /// #1847: strict store routing must preserve the typed missing-runtime
3606    /// failure instead of collapsing it into a direct-writer fallback.
3607    #[test]
3608    fn strict_writer_task_for_write_preserves_missing_runtime_error() {
3609        let dir = tempfile::tempdir().unwrap();
3610        let pool = ConnectionPool::new(PoolConfig {
3611            path: Some(dir.path().join("strict_writer_task_no_runtime.db")),
3612            write_queue_enabled: Some(true),
3613            write_routing_strict: true,
3614            ..PoolConfig::default()
3615        })
3616        .expect("file-backed pool should open");
3617
3618        let result = pool.writer_task_for_write(None, "strict_test_write");
3619
3620        assert!(
3621            matches!(result, Err(StorageError::WriterTaskNoRuntime)),
3622            "strict routing must preserve WriterTaskNoRuntime, got {result:?}"
3623        );
3624        assert_eq!(pool.writer_task_spawn_count(), 0);
3625    }
3626
3627    /// Join-handle lifecycle: a spawn-configured pool stores exactly one
3628    /// JoinHandle — the first `take_writer_task_join` after spawn returns
3629    /// it, and every later take returns `None` (the one-shot contract that
3630    /// lets exactly one subsystem own the drain).
3631    #[tokio::test]
3632    async fn take_writer_task_join_returns_some_once_then_none() {
3633        let dir = tempfile::tempdir().unwrap();
3634        let path = dir.path().join("join_lifecycle.db");
3635        let pool = ConnectionPool::new(PoolConfig {
3636            path: Some(path),
3637            write_queue_enabled: Some(true),
3638            ..PoolConfig::default()
3639        })
3640        .expect("file-backed pool should open");
3641
3642        // Spawning is lazy: nothing to take before the first
3643        // `writer_task_handle()` call actually spawns the task.
3644        assert!(
3645            pool.take_writer_task_join().is_none(),
3646            "before spawn there is no JoinHandle to take"
3647        );
3648        assert!(!pool.writer_task_join_was_stored());
3649        pool.writer_task_handle()
3650            .expect("runtime is present")
3651            .expect("write queue enabled must spawn a writer task");
3652        assert!(pool.writer_task_join_was_stored());
3653
3654        let join = pool
3655            .take_writer_task_join()
3656            .expect("the first take must return the spawned task's JoinHandle");
3657        assert!(
3658            pool.take_writer_task_join().is_none(),
3659            "the second take must return None — the handle is one-shot"
3660        );
3661
3662        // Await the taken handle before the test exits instead of dropping
3663        // it detached. The writer task only exits once every
3664        // `WriterTaskHandle` clone (the mpsc senders) is gone, and the pool's
3665        // own `writer_task` OnceLock holds one, so the pool must drop first —
3666        // the same drop-then-await order the batch-ingest drain relies on.
3667        drop(pool);
3668        tokio::time::timeout(Duration::from_secs(5), join)
3669            .await
3670            .expect("the writer task must exit once every handle clone is dropped")
3671            .expect("the writer task must not panic");
3672    }
3673
3674    /// Debug half of the first-wins contract: a second
3675    /// `set_writer_task_join` call is a construction bug, and debug builds
3676    /// trip the method's debug_assert loudly instead of carrying on.
3677    #[cfg(debug_assertions)]
3678    #[tokio::test]
3679    #[should_panic(expected = "writer task JoinHandle stored twice")]
3680    async fn set_writer_task_join_second_store_trips_debug_assert() {
3681        let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool should open");
3682        pool.set_writer_task_join(tokio::spawn(async {}));
3683        pool.set_writer_task_join(tokio::spawn(async {}));
3684    }
3685
3686    /// The at-most-once guard holds across the TAKEN state too: once the
3687    /// handle has been taken, the slot is empty, but a second store is still
3688    /// a construction bug and must trip the same debug_assert (the
3689    /// `writer_task_join_stored` flag remembers the first store).
3690    #[cfg(debug_assertions)]
3691    #[tokio::test]
3692    #[should_panic(expected = "writer task JoinHandle stored twice")]
3693    async fn set_writer_task_join_second_store_after_take_trips_debug_assert() {
3694        let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool should open");
3695        pool.set_writer_task_join(tokio::spawn(async {}));
3696        assert!(pool.take_writer_task_join().is_some());
3697        pool.set_writer_task_join(tokio::spawn(async {}));
3698    }
3699
3700    /// Release half of the first-wins contract: with the debug_assert
3701    /// compiled out, a second `set_writer_task_join` call keeps the
3702    /// EXISTING handle and drops the new one. The stored handle is
3703    /// therefore the first task's, so awaiting the taken handle completes
3704    /// the FIRST task's observable effect.
3705    #[cfg(not(debug_assertions))]
3706    #[tokio::test]
3707    async fn set_writer_task_join_first_wins_keeps_existing_handle() {
3708        let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool should open");
3709
3710        // First task: completes promptly and signals completion — the
3711        // observable effect the bounded await below asserts on.
3712        let (first_done_tx, first_done_rx) = tokio::sync::oneshot::channel::<()>();
3713        let first = tokio::spawn(async move {
3714            let _ = first_done_tx.send(());
3715        });
3716        // Second task: parks on a receiver nobody sends to, so it never
3717        // completes on its own. If first-wins failed and this task's handle
3718        // were the stored one, the bounded await below would time out.
3719        let (_never_sent, never_rx) = tokio::sync::oneshot::channel::<()>();
3720        let second = tokio::spawn(async move {
3721            let _ = never_rx.await;
3722        });
3723
3724        pool.set_writer_task_join(first);
3725        pool.set_writer_task_join(second);
3726
3727        let taken = pool
3728            .take_writer_task_join()
3729            .expect("the first handle must still be stored");
3730        tokio::time::timeout(Duration::from_secs(5), taken)
3731            .await
3732            .expect("stored handle must be the first task's; the second never completes")
3733            .expect("the first task must not panic");
3734        assert!(
3735            first_done_rx.await.is_ok(),
3736            "completing the taken handle must mean the FIRST task ran to completion"
3737        );
3738    }
3739
3740    /// ADR-091 backend-scoped attribution: the real path, a directory
3741    /// symlink, a file-level symlink, a relative spelling, and a bare file
3742    /// name (resolved against the current directory) must all mint an
3743    /// identical `DbIdentity` and canonical path for the same database.
3744    #[test]
3745    #[serial(pool_cwd)]
3746    fn mint_db_identity_alias_convergence() {
3747        let dir = tempfile::tempdir().unwrap();
3748        let real_dir = dir.path().join("real");
3749        fs::create_dir(&real_dir).unwrap();
3750        let db_path = real_dir.join("khive.db");
3751        fs::write(&db_path, b"").unwrap();
3752
3753        #[cfg(unix)]
3754        let dir_symlink = dir.path().join("dir_link");
3755        #[cfg(unix)]
3756        let file_symlink = dir.path().join("file_link.db");
3757        #[cfg(unix)]
3758        {
3759            std::os::unix::fs::symlink(&real_dir, &dir_symlink).unwrap();
3760            std::os::unix::fs::symlink(&db_path, &file_symlink).unwrap();
3761        }
3762
3763        let (via_real, canonical_real) = mint_db_identity(&db_path).unwrap();
3764
3765        // Relative spelling: resolved against the process CWD (step 1).
3766        let relative_result = {
3767            let _cwd = CwdGuard::enter(&real_dir);
3768            mint_db_identity(&PathBuf::from("khive.db"))
3769        };
3770        let (via_relative, canonical_relative) = relative_result.unwrap();
3771        assert_eq!(canonical_real, canonical_relative);
3772        assert_eq!(via_real, via_relative);
3773
3774        #[cfg(unix)]
3775        {
3776            let (via_dir_symlink, canonical_dir_symlink) =
3777                mint_db_identity(&dir_symlink.join("khive.db")).unwrap();
3778            assert_eq!(canonical_real, canonical_dir_symlink);
3779            assert_eq!(via_real, via_dir_symlink);
3780
3781            let (via_file_symlink, canonical_file_symlink) =
3782                mint_db_identity(&file_symlink).unwrap();
3783            assert_eq!(canonical_real, canonical_file_symlink);
3784            assert_eq!(via_real, via_file_symlink);
3785        }
3786
3787        // Bare file name: resolved against the current directory (step 1).
3788        let bare_name_result = {
3789            let _cwd = CwdGuard::enter(&real_dir);
3790            mint_db_identity(&PathBuf::from("khive.db"))
3791        };
3792        let (via_bare_name, canonical_bare_name) = bare_name_result.unwrap();
3793        assert_eq!(canonical_real, canonical_bare_name);
3794        assert_eq!(via_real, via_bare_name);
3795    }
3796
3797    /// ADR-091 backend-scoped attribution: `DbIdentity`/canonical-path
3798    /// equality across alias spellings (proven above by
3799    /// `mint_db_identity_alias_convergence`) does not by itself prove the
3800    /// walpin sidecar re-key — `sidecar_dir_for` is a separate, purely
3801    /// lexical derivation (`walpin::sidecar_dir_for`) that must be fed the
3802    /// *minted* canonical path, never the raw configured one. This test
3803    /// opens a real `ConnectionPool` (not the private `mint_db_identity` free
3804    /// function) through each alias spelling and asserts
3805    /// `sidecar_dir_for(pool.canonical_path())` converges to one directory —
3806    /// exercising the actual `ConnectionPool::new` → `canonical_path()` wiring
3807    /// every sidecar consumer (`checkpoint.rs`) reads from.
3808    #[test]
3809    #[serial(pool_cwd)]
3810    fn sidecar_dir_for_alias_convergence() {
3811        let dir = tempfile::tempdir().unwrap();
3812        let real_dir = dir.path().join("real");
3813        fs::create_dir(&real_dir).unwrap();
3814        let db_path = real_dir.join("khive.db");
3815        fs::write(&db_path, b"").unwrap();
3816
3817        #[cfg(unix)]
3818        let dir_symlink = dir.path().join("dir_link");
3819        #[cfg(unix)]
3820        let file_symlink = dir.path().join("file_link.db");
3821        #[cfg(unix)]
3822        {
3823            std::os::unix::fs::symlink(&real_dir, &dir_symlink).unwrap();
3824            std::os::unix::fs::symlink(&db_path, &file_symlink).unwrap();
3825        }
3826
3827        let pool_for = |path: &Path| -> Arc<ConnectionPool> {
3828            let cfg = PoolConfig {
3829                path: Some(path.to_path_buf()),
3830                ..PoolConfig::default()
3831            };
3832            Arc::new(ConnectionPool::new(cfg).expect("file-backed pool should open"))
3833        };
3834        let sidecar_of = |pool: &ConnectionPool| -> PathBuf {
3835            crate::walpin::sidecar_dir_for(pool.canonical_path().expect("file-backed pool"))
3836        };
3837
3838        let via_real = pool_for(&db_path);
3839        let sidecar_real = sidecar_of(&via_real);
3840
3841        let via_relative = {
3842            let _cwd = CwdGuard::enter(&real_dir);
3843            pool_for(Path::new("khive.db"))
3844        };
3845        assert_eq!(
3846            sidecar_real,
3847            sidecar_of(&via_relative),
3848            "a relative spelling of the same database must derive the same sidecar directory"
3849        );
3850
3851        #[cfg(unix)]
3852        {
3853            let via_dir_symlink = pool_for(&dir_symlink.join("khive.db"));
3854            assert_eq!(
3855                sidecar_real,
3856                sidecar_of(&via_dir_symlink),
3857                "opening through a directory symlink must derive the same sidecar directory"
3858            );
3859
3860            let via_file_symlink = pool_for(&file_symlink);
3861            assert_eq!(
3862                sidecar_real,
3863                sidecar_of(&via_file_symlink),
3864                "opening through a file-level symlink must derive the same sidecar directory"
3865            );
3866        }
3867
3868        let via_bare_name = {
3869            let _cwd = CwdGuard::enter(&real_dir);
3870            pool_for(Path::new("khive.db"))
3871        };
3872        assert_eq!(
3873            sidecar_real,
3874            sidecar_of(&via_bare_name),
3875            "a bare file name resolved against the current directory must derive the same \
3876             sidecar directory"
3877        );
3878    }
3879
3880    /// ADR-091 backend-scoped attribution: opening via a file-level symlink
3881    /// whose target does not exist yet (a valid first-open state), then
3882    /// after the target is created, opening via the target path directly,
3883    /// must mint identical `DbIdentity` values — the first-open path
3884    /// resolves the final component before canonicalizing the parent.
3885    #[cfg(unix)]
3886    #[test]
3887    fn mint_db_identity_dangling_symlink_first_open_convergence() {
3888        let dir = tempfile::tempdir().unwrap();
3889        let target = dir.path().join("target.db");
3890        let link = dir.path().join("link.db");
3891        std::os::unix::fs::symlink(&target, &link).unwrap();
3892        assert!(!target.exists(), "target must not exist yet (dangling)");
3893
3894        let (via_dangling_link, canonical_via_link) = mint_db_identity(&link).unwrap();
3895
3896        // Now create the target (as SQLite would on first write) and mint
3897        // again directly against the target path.
3898        fs::write(&target, b"").unwrap();
3899        let (via_target, canonical_via_target) = mint_db_identity(&target).unwrap();
3900
3901        assert_eq!(canonical_via_link, canonical_via_target);
3902        assert_eq!(via_dangling_link, via_target);
3903    }
3904
3905    /// A resolved target whose parent directory does not exist must fail
3906    /// minting exactly as the subsequent database open itself would fail.
3907    #[test]
3908    fn mint_db_identity_missing_parent_fails() {
3909        let dir = tempfile::tempdir().unwrap();
3910        let missing = dir.path().join("nonexistent_subdir").join("khive.db");
3911        let result = mint_db_identity(&missing);
3912        assert!(
3913            result.is_err(),
3914            "minting must fail when the parent directory does not exist"
3915        );
3916    }
3917
3918    /// Non-UTF-8 database paths (Unix) must round-trip through
3919    /// `DbIdentity`/canonicalization without loss.
3920    #[cfg(unix)]
3921    #[test]
3922    fn mint_db_identity_non_utf8_path_round_trips() {
3923        use std::ffi::OsStr;
3924        use std::os::unix::ffi::OsStrExt;
3925
3926        let dir = tempfile::tempdir().unwrap();
3927        // 0xFF is not valid UTF-8 as a standalone byte.
3928        let raw_name = OsStr::from_bytes(b"khive-\xffdb.sqlite");
3929        let db_path = dir.path().join(raw_name);
3930        // Some Unix filesystems (notably macOS's APFS) reject non-UTF-8
3931        // names outright at the syscall level — that is a filesystem
3932        // limitation, not a `mint_db_identity` bug, so skip rather than
3933        // fail where the underlying `write` itself cannot succeed.
3934        if let Err(e) = fs::write(&db_path, b"") {
3935            eprintln!(
3936                "skipping mint_db_identity_non_utf8_path_round_trips: filesystem rejected a \
3937                 non-UTF-8 file name ({e}); this platform's filesystem does not support the \
3938                 case under test"
3939            );
3940            return;
3941        }
3942
3943        let (identity, canonical) = mint_db_identity(&db_path).unwrap();
3944        assert_eq!(canonical.file_name().unwrap(), raw_name);
3945
3946        let (identity_again, canonical_again) = mint_db_identity(&db_path).unwrap();
3947        assert_eq!(identity, identity_again);
3948        assert_eq!(canonical, canonical_again);
3949    }
3950
3951    /// The checkout tri-state, arm by arm, at its single home. Each refusal
3952    /// arm asserts the classification it must NOT collapse into, because the
3953    /// historical defect was exactly a pairwise swap: cancellation surfaced as
3954    /// the retryable `AdmissionTimeout` while genuine pool exhaustion surfaced
3955    /// as a non-retryable `Driver` failure.
3956    #[test]
3957    fn resolve_reader_checkout_maps_each_arm_distinctly() {
3958        let pool = ConnectionPool::new(PoolConfig {
3959            path: None,
3960            ..PoolConfig::default()
3961        })
3962        .unwrap();
3963
3964        let guard = pool
3965            .resolve_reader_checkout(
3966                StorageCapability::Sql,
3967                "arm_checked_out",
3968                pool.reader_until(|| false),
3969            )
3970            .expect("an uncontended checkout must pass the guard through");
3971        drop(guard);
3972
3973        let Err(cancelled) =
3974            pool.resolve_reader_checkout(StorageCapability::Sql, "arm_cancelled", Ok(None))
3975        else {
3976            panic!("a cancelled checkout must be refused");
3977        };
3978        assert!(
3979            matches!(cancelled, StorageError::Timeout { .. }),
3980            "cancellation/deadline before checkout must be the non-retryable \
3981             Timeout, got {cancelled:?}"
3982        );
3983
3984        let Err(exhausted) = pool.resolve_reader_checkout(
3985            StorageCapability::Sql,
3986            "arm_exhausted",
3987            Err(pool_exhausted_error(Duration::from_millis(5), 1)),
3988        ) else {
3989            panic!("an exhausted checkout must be refused");
3990        };
3991        assert!(
3992            matches!(exhausted, StorageError::AdmissionTimeout { .. }),
3993            "the pool's own SQLITE_BUSY (checkout_timeout exhausted) must be \
3994             the retryable AdmissionTimeout, got {exhausted:?}"
3995        );
3996
3997        let Err(opaque) = pool.resolve_reader_checkout(
3998            StorageCapability::Entities,
3999            "arm_driver",
4000            Err(SqliteError::InvalidData("retired pooled writer".into())),
4001        ) else {
4002            panic!("an opaque checkout error must be refused");
4003        };
4004        assert!(
4005            matches!(
4006                &opaque,
4007                StorageError::Driver { capability, .. }
4008                    if *capability == StorageCapability::Entities
4009            ),
4010            "any other checkout error must stay a non-retryable Driver failure \
4011             under the caller's capability, got {opaque:?}"
4012        );
4013    }
4014}