Skip to main content

khive_db/
pool.rs

1//! Connection pool for SQLite: one exclusive writer, N concurrent readers.
2mod code_map;
3#[path = "pool/identity_registry.rs"]
4mod identity_registry;
5#[cfg(any(test, feature = "test-support"))]
6#[path = "pool/test_volume_lock_dir.rs"]
7mod test_volume_lock_dir;
8#[path = "pool/write_units.rs"]
9mod write_units;
10#[path = "pool/writer_acquisition.rs"]
11mod writer_acquisition;
12#[cfg(test)]
13use identity_registry::pool_identity_suffix;
14use identity_registry::PoolIdentityRegistration;
15#[cfg(any(test, feature = "test-support"))]
16use test_volume_lock_dir::test_volume_lock_dir;
17#[cfg(test)]
18use write_units::STARTUP_SPACE_PROBE;
19pub use write_units::{CheckpointGuard, CheckpointResult, WriterAcquisitionSnapshot, WriterGuard};
20pub(crate) use write_units::{
21    PooledAutocommitWriteUnit, PooledTransactionWriteUnit, StandaloneTransactionWriteUnit,
22    WriteAdmission, WriterAcquisitionCounters,
23};
24
25use crossbeam_queue::ArrayQueue;
26use parking_lot::{Condvar, Mutex};
27use rusqlite::hooks::{AuthContext, Authorization};
28use rusqlite::{Connection, OpenFlags};
29use serde::Serialize;
30use sha2::{Digest, Sha256};
31use std::cell::Cell;
32use std::collections::{BTreeMap, HashMap};
33use std::fs;
34use std::io::Read as _;
35use std::ops::{Deref, DerefMut};
36use std::path::{Path, PathBuf};
37use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
38use std::sync::{Arc, OnceLock};
39use std::thread;
40use std::time::{Duration, Instant};
41use tokio::sync::Semaphore;
42
43use crate::database_owner_identity::{DatabaseOwnerIdentity, DatabaseOwnerIdentityError};
44use crate::disk_guard::{LeaseRefusal, VolumeIdentity, VolumeLease};
45#[cfg(test)]
46use crate::disk_guard_config::DiskGuardConfigSource;
47use crate::disk_guard_config::{
48    resolve_disk_guard_config, EffectiveDiskGuardConfig, DEFAULT_DISK_GUARD_DEADLINE_MS,
49};
50use crate::error::SqliteError;
51#[cfg(windows)]
52use crate::file_identity::sqlite_opened_file_identity;
53#[cfg(any(unix, windows))]
54use crate::file_identity::{database_file_identity, DatabaseFileIdentity};
55use crate::writer_task::{execute_wrapped_transaction, WriterTaskHandle};
56use khive_storage::error::StorageError;
57use khive_storage::tx_registry::{DbIdentity, TxOrigin};
58use khive_storage::CapacityUnavailablePhase;
59use khive_storage::StorageCapability;
60
61mod claimed_file_identity;
62
63#[cfg(unix)]
64mod claimed_file_observer;
65
66#[cfg(unix)]
67pub use claimed_file_observer::initialize as initialize_claimed_file_observer;
68
69#[cfg(all(test, unix))]
70#[path = "pool/claimed_file_identity_tests.rs"]
71mod claimed_file_identity_tests;
72
73pub(crate) const CACHE_SIZE_KIB: &str = "-65536";
74const MMAP_SIZE_BYTES: &str = "1073741824";
75const DEFAULT_READER_CAP: usize = 8;
76
77const DEFAULT_JOURNAL_SIZE_LIMIT_BYTES: i64 = 67_108_864; // 64 MiB
78const DEFAULT_WRITE_QUEUE_CAPACITY: usize = 256;
79const DATABASE_ID_TABLE: &str = "_khive_database_identity";
80static NEXT_MAIN_POOL_GENERATION: AtomicU64 = AtomicU64::new(1);
81
82#[cfg(test)]
83#[derive(Clone, Copy, PartialEq, Eq)]
84enum IdentityOpenStage {
85    AfterMainOpenBeforeFirstStat,
86    AfterInitialIdentityWrite,
87    BeforeStandaloneOpen,
88    AfterStandaloneOpen,
89}
90
91#[cfg(test)]
92type IdentityOpenHook = Box<dyn Fn(&Path, IdentityOpenStage, Option<&Connection>)>;
93
94#[cfg(test)]
95thread_local! {
96    static IDENTITY_OPEN_HOOK: std::cell::RefCell<Option<IdentityOpenHook>> =
97        const { std::cell::RefCell::new(None) };
98}
99
100#[cfg(test)]
101fn run_identity_open_hook(path: &Path, stage: IdentityOpenStage, conn: Option<&Connection>) {
102    IDENTITY_OPEN_HOOK.with(|hook| {
103        if let Some(hook) = hook.borrow().as_ref() {
104            hook(path, stage, conn);
105        }
106    });
107}
108
109/// Runtime-owned SQL transactions that share the store write-routing policy.
110#[derive(Clone, Copy, Debug)]
111pub enum RuntimeWriteOperation {
112    MergeEntity,
113    MergeNote,
114    UpdateSymmetricEdge,
115}
116
117impl RuntimeWriteOperation {
118    fn operation(self) -> &'static str {
119        match self {
120            Self::MergeEntity => "merge_entity",
121            Self::MergeNote => "merge_note",
122            Self::UpdateSymmetricEdge => "update_edge",
123        }
124    }
125
126    fn fallback_site(self) -> crate::timeout_sink::Site {
127        match self {
128            Self::MergeEntity => crate::timeout_sink::Site::DirectRouteRuntimeMergeEntity,
129            Self::MergeNote => crate::timeout_sink::Site::DirectRouteRuntimeMergeNote,
130            Self::UpdateSymmetricEdge => {
131                crate::timeout_sink::Site::DirectRouteRuntimeUpdateSymmetricEdge
132            }
133        }
134    }
135}
136
137/// Bounded WAL autocheckpoint applied to writer-capable connections while no
138/// dedicated checkpoint owner has claimed the pool (4,000 pages ≈ 16 MiB at
139/// SQLite's default 4 KiB page size — SQLite's historic behaviour for this
140/// pool). Not a tuning parameter: there is no config field or environment
141/// override, and the only way to change the effective value is an actual
142/// ownership claim ([`ConnectionPool::claim_checkpoint_ownership`]), which a
143/// runtime may make only when it really runs the scheduled checkpoint task.
144pub(crate) const FALLBACK_WAL_AUTOCHECKPOINT_PAGES: u32 = 4_000;
145
146#[derive(Clone, Copy, Debug, PartialEq, Eq)]
147enum CheckpointOwnership {
148    Unclaimed,
149    Claiming,
150    Claimed,
151}
152
153struct CheckpointOwnershipState {
154    phase: CheckpointOwnership,
155    #[cfg(test)]
156    connection_waiters: usize,
157}
158
159#[cfg(test)]
160struct CheckpointConnectionConfigPause {
161    selected: std::sync::Barrier,
162    resume: std::sync::Barrier,
163}
164
165#[cfg(test)]
166impl CheckpointConnectionConfigPause {
167    fn new() -> Self {
168        Self {
169            selected: std::sync::Barrier::new(2),
170            resume: std::sync::Barrier::new(2),
171        }
172    }
173}
174
175struct CheckpointOwnershipGate {
176    state: Mutex<CheckpointOwnershipState>,
177    changed: Condvar,
178    #[cfg(test)]
179    connection_config_pause: Mutex<Option<Arc<CheckpointConnectionConfigPause>>>,
180    #[cfg(test)]
181    claim_lock_observed: Mutex<Option<std::sync::mpsc::SyncSender<bool>>>,
182}
183
184impl CheckpointOwnershipGate {
185    fn new() -> Self {
186        Self {
187            state: Mutex::new(CheckpointOwnershipState {
188                phase: CheckpointOwnership::Unclaimed,
189                #[cfg(test)]
190                connection_waiters: 0,
191            }),
192            changed: Condvar::new(),
193            #[cfg(test)]
194            connection_config_pause: Mutex::new(None),
195            #[cfg(test)]
196            claim_lock_observed: Mutex::new(None),
197        }
198    }
199
200    /// Join an in-flight claim, or become the one caller that configures it.
201    /// Returns `false` when another caller has already completed the claim.
202    fn begin_claim(&self) -> bool {
203        #[cfg(test)]
204        let claim_lock_observed = self.claim_lock_observed.lock().take();
205        #[cfg(test)]
206        let mut state = if let Some(observed) = claim_lock_observed {
207            match self.state.try_lock() {
208                Some(state) => {
209                    let _ = observed.send(false);
210                    state
211                }
212                None => {
213                    let _ = observed.send(true);
214                    self.state.lock()
215                }
216            }
217        } else {
218            self.state.lock()
219        };
220        #[cfg(not(test))]
221        let mut state = self.state.lock();
222        loop {
223            match state.phase {
224                CheckpointOwnership::Unclaimed => {
225                    state.phase = CheckpointOwnership::Claiming;
226                    self.changed.notify_all();
227                    return true;
228                }
229                CheckpointOwnership::Claiming => self.changed.wait(&mut state),
230                CheckpointOwnership::Claimed => return false,
231            }
232        }
233    }
234
235    fn finish_claim(&self, succeeded: bool) {
236        let mut state = self.state.lock();
237        debug_assert_eq!(state.phase, CheckpointOwnership::Claiming);
238        state.phase = if succeeded {
239            CheckpointOwnership::Claimed
240        } else {
241            CheckpointOwnership::Unclaimed
242        };
243        self.changed.notify_all();
244    }
245
246    fn settled_state(&self) -> parking_lot::MutexGuard<'_, CheckpointOwnershipState> {
247        let mut state = self.state.lock();
248        while state.phase == CheckpointOwnership::Claiming {
249            #[cfg(test)]
250            {
251                state.connection_waiters += 1;
252                self.changed.notify_all();
253            }
254            self.changed.wait(&mut state);
255            #[cfg(test)]
256            {
257                state.connection_waiters -= 1;
258                self.changed.notify_all();
259            }
260        }
261        state
262    }
263
264    #[cfg(test)]
265    fn wal_autocheckpoint_pages(&self) -> u32 {
266        let state = self.settled_state();
267        match state.phase {
268            CheckpointOwnership::Unclaimed => FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
269            CheckpointOwnership::Claimed => 0,
270            CheckpointOwnership::Claiming => unreachable!("claim wait must settle the state"),
271        }
272    }
273
274    /// Wait for any in-flight claim, select the resulting posture, and retain
275    /// the gate until SQLite has applied that connection-local PRAGMA. A claim
276    /// therefore linearizes entirely before or after this configuration,
277    /// never between its state sample and side effect.
278    fn configure_wal_autocheckpoint(&self, conn: &Connection) -> Result<(), SqliteError> {
279        let state = self.settled_state();
280        let pages = match state.phase {
281            CheckpointOwnership::Unclaimed => FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
282            CheckpointOwnership::Claimed => 0,
283            CheckpointOwnership::Claiming => unreachable!("claim wait must settle the state"),
284        };
285        #[cfg(test)]
286        if let Some(pause) = self.connection_config_pause.lock().take() {
287            pause.selected.wait();
288            pause.resume.wait();
289        }
290        conn.pragma_update(None, "wal_autocheckpoint", pages)?;
291        drop(state);
292        Ok(())
293    }
294}
295
296fn deny_retired_writer(_context: AuthContext<'_>) -> Authorization {
297    Authorization::Deny
298}
299
300pub(crate) const TEST_HARNESS_ENV: &str = "KHIVE_TEST_HARNESS";
301
302/// Where the effective WAL ceiling byte value was configured.
303#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize)]
304#[serde(rename_all = "snake_case")]
305pub enum WalCeilingSource {
306    BackendField,
307    Environment,
308    #[default]
309    Default,
310}
311
312/// Resolved WAL-extent policy for one SQLite backend. A zero-byte policy is
313/// explicitly disabled; it remains visible in diagnostics and config identity.
314#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
315pub struct WalCeilingPolicy {
316    pub bytes: u64,
317    pub source: WalCeilingSource,
318}
319
320impl WalCeilingPolicy {
321    /// Validate checks that do not need SQLite's page-size observation.
322    /// Read-only backends retain configured metadata but enforce no writer
323    /// policy. Invalid offset arithmetic is rejected in either mode.
324    pub fn validate_static(
325        self,
326        file_backed: bool,
327        wal_mode: bool,
328        read_only: bool,
329    ) -> Result<(), SqliteError> {
330        if self.bytes == 0 {
331            return Ok(());
332        }
333        if i64::try_from(self.bytes).is_err() {
334            return Err(SqliteError::WalCeilingOffsetOverflow { bytes: self.bytes });
335        }
336        if read_only {
337            return Ok(());
338        }
339        if !file_backed {
340            return Err(SqliteError::WalCeilingUnsupported {
341                bytes: self.bytes,
342                backend_kind: "in-memory backend",
343            });
344        }
345        if !wal_mode {
346            return Err(SqliteError::WalCeilingUnsupported {
347                bytes: self.bytes,
348                backend_kind: "non-WAL backend",
349            });
350        }
351        Ok(())
352    }
353
354    /// Bytes of writer policy that could be enforced on this backend.
355    pub fn effective_bytes(self, read_only: bool) -> u64 {
356        if read_only {
357            0
358        } else {
359            self.bytes
360        }
361    }
362}
363
364/// Configuration for the connection pool.
365#[derive(Clone, Debug)]
366pub struct PoolConfig {
367    /// Database path. None = in-memory (pool degrades to single connection).
368    pub path: Option<PathBuf>,
369    /// Registered native code-map VFS name; set only by the code-map constructor.
370    pub code_map_vfs: Option<String>,
371    /// File identity pinned by the caller before this pool opens SQLite.
372    /// A mismatch is refused before identity initialization or WAL setup.
373    #[cfg(any(unix, windows))]
374    pub expected_file_identity: Option<DatabaseFileIdentity>,
375    /// Number of reader connections (default: min(num_cpus, 8)).
376    pub max_readers: usize,
377    /// WAL mode (must be true for pooling to work; default: true).
378    pub wal_mode: bool,
379    /// Busy timeout per connection (default: 30s).
380    ///
381    /// Overridable via `KHIVE_BUSY_TIMEOUT_SECS`.
382    pub busy_timeout: Duration,
383    /// Time to wait for a reader connection before returning an error (default: 5s).
384    ///
385    /// For a writer it bounds only the pool-mutex wait. A writable file-backed
386    /// pool takes the volume lease first, bounded by the guard deadline
387    /// (`disk_guard_deadline_ms`, default 2000 ms), so the effective writer
388    /// wait under contention is the guard deadline, then `checkout_timeout`:
389    /// a 50 ms value here can still wait about 2 s, and a wait that ends at
390    /// the lease returns `CapacityUnavailable` with phase `lock` rather than a
391    /// checkout timeout. Diagnostics report both bounds and their sum.
392    ///
393    /// Overridable via `KHIVE_CHECKOUT_TIMEOUT_SECS`.
394    pub checkout_timeout: Duration,
395    /// Maximum WAL journal size in bytes before SQLite resets the WAL.
396    ///
397    /// Maps to `PRAGMA journal_size_limit`. Default: 64 MiB.
398    ///
399    /// Overridable via `KHIVE_JOURNAL_SIZE_LIMIT_BYTES`.
400    pub journal_size_limit_bytes: i64,
401    /// Open the database read-only (default: false).
402    ///
403    /// When true, the pool's writer connection is opened with
404    /// `SQLITE_OPEN_READ_ONLY` (no `SQLITE_OPEN_CREATE`, so a missing path is
405    /// rejected instead of created) and `PRAGMA query_only = ON` is set on
406    /// every connection that can execute SQL. Reader connections are already
407    /// opened read-only regardless of this flag.
408    pub read_only: bool,
409    /// ADR-194 WAL active-extent ceiling and its resolved configuration source.
410    /// Zero explicitly disables this independent policy.
411    pub wal_ceiling: WalCeilingPolicy,
412    /// Route migrated store write paths through the single-writer
413    /// `WriterTask` channel (ADR-067 Component A) instead of the legacy
414    /// per-call pool-mutex/standalone-connection path. Enabled by default
415    /// for file-backed pools when unset; explicit override always wins.
416    /// That default is a compatibility-routing posture subordinate to
417    /// ADR-135 Amendment 1 and ADR-136 D1/D2 — the strict-routing default
418    /// flip has NOT happened.
419    ///
420    /// The store layer resolves all of its routed write paths at write time;
421    /// the classification table in `writer_task.rs` remains the authoritative
422    /// inventory. This tranche does not claim the repository-wide
423    /// single-writer guarantee, and the strict default is still evidence-gated.
424    ///
425    /// `None` means the caller expressed no preference: [`ConnectionPool::new`]
426    /// resolves it once `path` is known, defaulting to `true` for file-backed
427    /// pools and `false` for in-memory ones. `Some(_)` is an explicit
428    /// preference and always wins, in both directions, over that default.
429    /// An explicit `Some(true)` on an in-memory pool is accepted DELIBERATELY
430    /// and emits a warning before degrading to the legacy path — an in-memory
431    /// pool cannot host a writer task (`writer_task::spawn`'s
432    /// standalone-connection open fails); see
433    /// `ConnectionPool::writer_task_handle` and the
434    /// `explicit_true_stays_on_for_memory_backed_pool` test.
435    ///
436    /// Overridable via `KHIVE_WRITE_QUEUE` (`"1"` or `"true"`,
437    /// case-insensitive, sets `Some(true)`; any other value sets `Some(false)`;
438    /// unset leaves it `None`).
439    pub write_queue_enabled: Option<bool>,
440    /// Bounded channel capacity for the `WriterTask` write queue.
441    ///
442    /// Overridable via `KHIVE_WRITE_QUEUE_CAPACITY`. Default: 256 pending
443    /// operations (ADR-067 Component A recommended default).
444    pub write_queue_capacity: usize,
445    /// ADR-136 D1: when `true`, every covered store write path that would
446    /// otherwise silently degrade to the legacy pool-mutex/standalone-
447    /// connection path on a missing or failed `WriterTask` handle instead
448    /// returns an error.
449    /// Exercises the store-layer routing tranche toward ADR-135 F2's
450    /// strict-routing precondition without changing behavior for callers that
451    /// never set the env var.
452    ///
453    /// Overridable via `KHIVE_WRITE_ROUTING` (value `"strict"`,
454    /// case-insensitive; anything else, or unset, leaves this `false`).
455    pub write_routing_strict: bool,
456    /// Dedicated admission deadline (ADR-131 Decision 2) bounding ONLY the
457    /// wait for capacity on the `WriterTask` write queue —
458    /// [`WriterTaskHandle::send_bounded`]/`send_top_level_bounded`'s default
459    /// timeout. Distinct from `checkout_timeout`, which bounds reader/pool
460    /// checkout instead; the two authorities used to be conflated (#1382,
461    /// #1643) before this field existed.
462    ///
463    /// Default: 2000 ms. Validated at [`ConnectionPool::new`] to fall in
464    /// `[100, 10000]` ms; a value outside that range is a configuration
465    /// error (`SqliteError::InvalidConfig`), never silently clamped into
466    /// range.
467    ///
468    /// Overridable via `KHIVE_WRITE_ADMISSION_DEADLINE_MS`.
469    pub write_admission_deadline_ms: u64,
470    /// SQLite disk-reserve and guard-deadline policy for this pool. `None`
471    /// resolves the process environment when the pool opens.
472    pub disk_guard_config: Option<EffectiveDiskGuardConfig>,
473    /// Shared per-user directory for the volume advisory lock files.
474    /// [`PoolConfig::default`] resolves it with [`crate::default_volume_lock_dir`]
475    /// and leaves `None` when no directory can be resolved.
476    pub volume_lock_dir: Option<PathBuf>,
477    /// Maximum age an explicit cached-reader read transaction
478    /// (`sql_bridge`'s `BEGIN`-then-reuse path) may reach before its next use
479    /// is refused and it is rolled back instead of extending its WAL
480    /// snapshot further (#1846). Shares `KHIVE_TX_MAX_AGE_SECS` with the
481    /// ADR-091 Plank 1 visibility sweep in `checkpoint.rs` so one knob
482    /// governs both when an operator is warned about a stale reader and when
483    /// that reader's snapshot is actually released.
484    ///
485    /// Overridable via `KHIVE_TX_MAX_AGE_SECS`. Default: 120 seconds.
486    pub read_tx_max_age: Duration,
487}
488
489/// ADR-131 Decision 2's validated range for `write_admission_deadline_ms`.
490const WRITE_ADMISSION_DEADLINE_MS_RANGE: std::ops::RangeInclusive<u64> = 100..=10_000;
491const DEFAULT_WRITE_ADMISSION_DEADLINE_MS: u64 = 2000;
492
493impl Default for PoolConfig {
494    fn default() -> Self {
495        Self {
496            path: None,
497            code_map_vfs: None,
498            #[cfg(any(unix, windows))]
499            expected_file_identity: None,
500            max_readers: std::thread::available_parallelism()
501                .map(|n| n.get())
502                .unwrap_or(1)
503                .clamp(1, DEFAULT_READER_CAP),
504            wal_mode: true,
505            busy_timeout: Duration::from_secs(
506                std::env::var("KHIVE_BUSY_TIMEOUT_SECS")
507                    .ok()
508                    .and_then(|v| v.parse::<u64>().ok())
509                    .unwrap_or(30),
510            ),
511            checkout_timeout: Duration::from_secs(
512                std::env::var("KHIVE_CHECKOUT_TIMEOUT_SECS")
513                    .ok()
514                    .and_then(|v| v.parse::<u64>().ok())
515                    .unwrap_or(5),
516            ),
517            journal_size_limit_bytes: std::env::var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES")
518                .ok()
519                .and_then(|v| v.parse::<i64>().ok())
520                .unwrap_or(DEFAULT_JOURNAL_SIZE_LIMIT_BYTES),
521            read_only: false,
522            wal_ceiling: WalCeilingPolicy::default(),
523            // `var_os`, not `var`: the documented contract is "any SET value
524            // other than 1/true means Some(false)" — a set-but-non-Unicode
525            // value must count as set (var() would return Err and silently
526            // fall through to the file-backed default of enabled).
527            write_queue_enabled: std::env::var_os("KHIVE_WRITE_QUEUE").map(|v| {
528                v.to_str()
529                    .is_some_and(|v| v == "1" || v.eq_ignore_ascii_case("true"))
530            }),
531            write_queue_capacity: std::env::var("KHIVE_WRITE_QUEUE_CAPACITY")
532                .ok()
533                .and_then(|v| v.parse::<usize>().ok())
534                .filter(|&n| n > 0)
535                .unwrap_or(DEFAULT_WRITE_QUEUE_CAPACITY),
536            write_routing_strict: std::env::var("KHIVE_WRITE_ROUTING")
537                .map(|v| v.eq_ignore_ascii_case("strict"))
538                .unwrap_or(false),
539            write_admission_deadline_ms: std::env::var("KHIVE_WRITE_ADMISSION_DEADLINE_MS")
540                .ok()
541                .and_then(|v| v.parse::<u64>().ok())
542                .unwrap_or(DEFAULT_WRITE_ADMISSION_DEADLINE_MS),
543            disk_guard_config: None,
544            #[cfg(test)]
545            volume_lock_dir: Some(test_volume_lock_dir()),
546            #[cfg(not(test))]
547            volume_lock_dir: crate::default_volume_lock_dir().ok(),
548            read_tx_max_age: crate::checkpoint::tx_age_thresholds_from_env(
549                Duration::from_secs(30),
550                Duration::from_secs(120),
551            )
552            .1,
553        }
554    }
555}
556
557#[cfg(any(test, feature = "test-support"))]
558impl PoolConfig {
559    /// A small concurrent pool for private test databases.
560    ///
561    /// Tests of reader admission or production sizing should set their required
562    /// count explicitly. Ordinary fixtures need not reserve a CPU-sized pool.
563    pub fn for_test() -> Self {
564        Self {
565            max_readers: 2,
566            volume_lock_dir: Some(test_volume_lock_dir()),
567            ..Self::default()
568        }
569    }
570}
571
572/// Prevent Cargo-launched tests and test subprocesses from opening the
573/// operator's default data tree in every build profile. Activation is solely
574/// the runtime `KHIVE_TEST_HARNESS=1` marker; production/installed binaries do
575/// not receive that workspace Cargo environment.
576///
577/// There is deliberately no environment override: any inheritable escape
578/// hatch set for one Cargo invocation leaks into the next `cargo test` in the
579/// same shell and re-opens the store the guard exists to protect. A deliberate
580/// session against the real store runs the built binary directly (for example
581/// `target/release/...` or an installed binary), which never receives the
582/// workspace Cargo environment and therefore never trips this guard.
583/// Existing path ancestors are canonicalized before comparison, resolving
584/// traversal, symlinks, and filesystem-provided case (including APFS case
585/// folding). Missing trailing components remain lexical because they have no
586/// filesystem identity yet. SQLite URI paths are rejected rather than trying
587/// to reproduce SQLite's URI normalization rules.
588fn refuse_home_data_store_in_tests(config: &PoolConfig) -> Result<(), SqliteError> {
589    if std::env::var(TEST_HARNESS_ENV).as_deref() != Ok("1") {
590        return Ok(());
591    }
592
593    let Some(path) = config.path.as_deref() else {
594        return Ok(());
595    };
596    if path
597        .as_os_str()
598        .as_encoded_bytes()
599        .get(..5)
600        .is_some_and(|prefix| prefix.eq_ignore_ascii_case(b"file:"))
601    {
602        return Err(SqliteError::InvalidData(format!(
603            "test harness refused SQLite URI database path {}; use a filesystem path outside \
604             HOME/.khive (deliberate sessions against a real store run the built binary \
605             directly, outside the Cargo test environment)",
606            path.display()
607        )));
608    }
609
610    let Some(home) = std::env::var_os("HOME") else {
611        return Ok(());
612    };
613    let canonical_path = canonicalize_deepest_existing(path)?;
614    let canonical_home_data_dir =
615        canonicalize_deepest_existing(&PathBuf::from(home).join(".khive"))?;
616    if canonical_path.starts_with(&canonical_home_data_dir) {
617        return Err(SqliteError::InvalidData(format!(
618            "test harness refused to open SQLite database under HOME/.khive: {} \
619             (deliberate sessions against a real store run the built binary directly, \
620             outside the Cargo test environment)",
621            canonical_path.display()
622        )));
623    }
624    Ok(())
625}
626
627fn canonicalize_deepest_existing(path: &Path) -> Result<PathBuf, SqliteError> {
628    let absolute = if path.is_absolute() {
629        path.to_path_buf()
630    } else {
631        std::env::current_dir().map_err(SqliteError::Io)?.join(path)
632    };
633
634    for ancestor in absolute.ancestors() {
635        match fs::canonicalize(ancestor) {
636            Ok(mut canonical) => {
637                let missing = absolute.strip_prefix(ancestor).map_err(|error| {
638                    SqliteError::InvalidData(format!(
639                        "failed to preserve missing path components for {}: {error}",
640                        absolute.display()
641                    ))
642                })?;
643                canonical.push(missing);
644                return Ok(canonical);
645            }
646            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
647            Err(error) => {
648                return Err(SqliteError::InvalidData(format!(
649                    "failed to canonicalize database path ancestor {}: {error}",
650                    ancestor.display()
651                )));
652            }
653        }
654    }
655
656    Err(SqliteError::InvalidData(format!(
657        "database path has no canonicalizable ancestor: {}",
658        absolute.display()
659    )))
660}
661
662/// Enforce ADR-131 Decision 2's `write_admission_deadline_ms` bound at
663/// configuration load: `[100, 10000]` ms, rejected rather than clamped when
664/// out of range so a misconfiguration is never silently reinterpreted as a
665/// different deadline than the operator asked for.
666fn validate_write_admission_deadline(deadline_ms: u64) -> Result<(), SqliteError> {
667    if WRITE_ADMISSION_DEADLINE_MS_RANGE.contains(&deadline_ms) {
668        return Ok(());
669    }
670    Err(SqliteError::InvalidConfig(format!(
671        "write_admission_deadline_ms must be in [{}, {}] ms, got {deadline_ms}",
672        WRITE_ADMISSION_DEADLINE_MS_RANGE.start(),
673        WRITE_ADMISSION_DEADLINE_MS_RANGE.end()
674    )))
675}
676
677/// Pool-scoped counters for ADR-166's search mechanism guards.
678#[derive(Clone, Debug, Default, PartialEq, Eq, serde::Serialize)]
679pub struct SearchMechanismSnapshot {
680    /// Coordinator calls actually issued to each registered backend, keyed by
681    /// the request's canonical kind. A backend skipped by served-kind routing
682    /// has no entry for that kind.
683    pub dispatches_by_backend_and_kind: BTreeMap<String, BTreeMap<String, u64>>,
684    /// Candidate note rows fetched after the text/vector fusion and fresh-tail
685    /// merge, including rows later filtered as deleted. Result metadata fetched
686    /// later by the KG handler is excluded.
687    pub note_candidate_hydration_rows: u64,
688}
689
690/// A read-write connection pool for SQLite.
691///
692/// Architecture:
693/// - 1 writer connection protected by a Mutex (exclusive access)
694/// - N reader connections in a lock-free queue (concurrent access)
695/// - All connections share the same database file in WAL mode
696///
697/// Writable in-memory databases, or writable file databases when WAL mode is
698/// disabled/unavailable, degrade to single-connection mode and route all
699/// operations through the writer connection. A file-backed read-only pool
700/// always retains at least one dedicated read-only connection: rollback-journal
701/// snapshots do not need WAL to support concurrent readers, and inspection must
702/// never alias a read onto the query-only writer slot.
703pub struct ConnectionPool {
704    writer: Arc<Mutex<Connection>>,
705    // Prepared before admission so terminal cleanup cannot fail while allocating a replacement.
706    retirement_connection: Mutex<Option<Connection>>,
707    #[cfg(any(test, feature = "test-support"))]
708    statement_observer: Arc<crate::statement_observer::StatementObserverHub>,
709    main_pool_generation: OnceLock<u64>,
710    /// Three-state gate for whether the ADR-091 scheduled task has claimed
711    /// routine WAL reclamation for this pool. Until claimed, every
712    /// writer-capable connection keeps a bounded SQLite autocheckpoint
713    /// ([`FALLBACK_WAL_AUTOCHECKPOINT_PAGES`]) so a writable pool without a
714    /// checkpoint task cannot grow its WAL without bound. After
715    /// [`Self::claim_checkpoint_ownership`], writer-capable connections open
716    /// with `wal_autocheckpoint = 0` and routine checkpoint I/O stays off
717    /// application commit paths.
718    checkpoint_ownership: CheckpointOwnershipGate,
719    /// Fail-closed guard for the legacy pool-mutex writer. A transaction
720    /// owner retires this connection after a body panic or when it cannot
721    /// prove that finalization restored autocommit mode; subsequent checkouts
722    /// must never reuse it.
723    pooled_writer_retired: AtomicBool,
724    /// Process-local writer acquisition counters shared with the pool's
725    /// lifetime-owned writer task. Keeping the counters at the actual
726    /// acquisition boundaries means new verbs inherit instrumentation without
727    /// per-verb classification (ADR-133 D8 / issue #1389).
728    writer_acquisition_counters: Arc<WriterAcquisitionCounters>,
729    /// Shared with the long-lived writer task so it can resample at every
730    /// dequeued request rather than only when its connection is opened.
731    write_admission: Arc<WriteAdmission>,
732    /// Effective source and values captured at open for diagnostic reporting;
733    /// the live volume identifier is deliberately resolved separately.
734    disk_guard_config: Option<EffectiveDiskGuardConfig>,
735    /// Pool-scoped reader route, saturation, and hold-lifecycle counters.
736    /// Instrumentation lives at the acquisition boundary so every typed
737    /// store and raw-SQL caller inherits it without per-verb bookkeeping
738    /// (ADR-165 Slice 2 / ADR-166 G2).
739    reader_acquisition_counters: ReaderAcquisitionCounters,
740    /// ADR-166 G4/G5 process-lifetime mechanism counters for this physical
741    /// backend. Backend IDs remain separate even when aliases share a pool.
742    search_dispatches: Mutex<BTreeMap<String, BTreeMap<String, u64>>>,
743    note_candidate_hydration_rows: AtomicU64,
744    readers: ArrayQueue<Connection>,
745    max_readers: usize,
746    config: PoolConfig,
747    /// Canonical physical target used by every connection in a file-backed
748    /// read-only pool. Classification and open must share this exact spelling:
749    /// deriving WAL sidecars from a configured symlink while SQLite follows it
750    /// to another file can hide committed frames or a live writable `-shm`.
751    /// The value is an `immutable=1` URI only for a clean, checkpointed WAL;
752    /// rollback-journal databases and frozen WAL+SHM snapshots retain the
753    /// canonical ordinary path and SQLite locking/change detection.
754    read_only_open_target: Option<PathBuf>,
755    sql_bridge_reader_slots: Arc<Semaphore>,
756    sql_bridge_writer_slots: Arc<Semaphore>,
757    /// The pool-wide ADR-067 Component A writer task, spawned lazily and at
758    /// most once per pool (per DB file) via [`Self::writer_task_handle`] —
759    /// see that method's doc comment for why this lives here rather than on
760    /// each store.
761    writer_task: OnceLock<Option<WriterTaskHandle>>,
762    /// The `tokio::spawn` JoinHandle of the writer task above, stored by
763    /// [`crate::writer_task::spawn`] so short-lived callers (batch CLI
764    /// paths) can await the task's exit — and therefore its connection's
765    /// close-time WAL checkpoint — before treating the database file state
766    /// as settled. Long-running callers never take it; dropping an untaken
767    /// JoinHandle detaches the task, which is exactly the pre-existing
768    /// behavior.
769    writer_task_join: Mutex<Option<tokio::task::JoinHandle<()>>>,
770    /// Monotonic "a writer-task JoinHandle was stored at least once" flag
771    /// backing [`Self::set_writer_task_join`]'s at-most-once guard: it holds
772    /// the invariant even after [`Self::take_writer_task_join`] empties the
773    /// slot, so a second store never re-arms it.
774    writer_task_join_stored: AtomicBool,
775    /// This pool's ADR-091 backend-scoped attribution origin, minted exactly
776    /// once at construction (see [`mint_db_identity`]): `Database(_)` for a
777    /// file-backed pool, `Memory` for an in-memory pool. Every
778    /// `tx_registry::register_scoped` call site in this crate reaches its
779    /// origin through [`Self::origin`] rather than re-deriving it.
780    origin: TxOrigin,
781    /// The canonical path `origin`'s `DbIdentity` was minted from, `None` for
782    /// an in-memory pool. `DbIdentity` is deliberately opaque (no path
783    /// accessor) — filesystem consumers that need the actual path (sidecar
784    /// derivation) use this, the same canonical value the identity was
785    /// minted from, via [`Self::canonical_path`].
786    identity_path: Option<PathBuf>,
787    /// The physical file SQLite opened, checked against the canonical path.
788    /// Later reader and standalone opens must retain this identity.
789    #[cfg(any(unix, windows))]
790    opened_file_identity: Option<DatabaseFileIdentity>,
791    /// A persistent nonce read through SQLite's opened main database, rather
792    /// than through the pathname that may have been replaced during open.
793    opened_database_id: Option<uuid::Uuid>,
794    /// Registered only after every connection opens successfully. RAII removes
795    /// the path when the last pool for it drops, including failed construction.
796    identity_registration: Option<PoolIdentityRegistration>,
797    /// Test-only instrumentation: counts how many times the writer-task
798    /// init closure actually ran. Must never exceed 1 per pool no matter how
799    /// many stores are constructed over it — that is the invariant
800    /// `OnceLock::get_or_init` exists to guarantee, and what
801    /// `pool.rs`'s and `entity_tests.rs`'s one-writer-per-pool tests assert.
802    #[cfg(test)]
803    writer_task_spawn_count: std::sync::atomic::AtomicUsize,
804}
805
806impl Drop for ConnectionPool {
807    /// Close every read-only reader before the fields below it drop in
808    /// declaration order (`writer` first, `readers` well before
809    /// `writer_task`). A read-only connection cannot take the EXCLUSIVE lock
810    /// SQLite needs to checkpoint on close, so if a reader were left to close
811    /// last, WAL mode would leave `-wal`/`-shm` behind. Draining `readers`
812    /// here, before that field-order drop runs, makes the writable `writer`
813    /// connection close after every reader instead of before it.
814    fn drop(&mut self) {
815        while let Some(conn) = self.readers.pop() {
816            drop(conn);
817        }
818    }
819}
820
821enum ReaderLease<'pool> {
822    Pooled(Connection),
823    Shared(parking_lot::MutexGuard<'pool, Connection>),
824}
825
826/// A value-extraction view of one row from a reader lease.
827///
828/// It deliberately exposes neither the prepared statement nor its connection.
829///
830/// ```compile_fail
831/// fn statement(row: &khive_db::ReaderRow<'_, '_>) {
832///     let _: &rusqlite::Statement<'_> = row.as_ref();
833/// }
834/// ```
835pub struct ReaderRow<'row, 'statement> {
836    row: &'row rusqlite::Row<'statement>,
837}
838
839impl ReaderRow<'_, '_> {
840    /// Extract a value by zero-based column index or column name.
841    pub fn get<I: rusqlite::RowIndex, T: rusqlite::types::FromSql>(
842        &self,
843        index: I,
844    ) -> rusqlite::Result<T> {
845        self.row.get(index)
846    }
847
848    /// Borrow a SQLite value without exposing statement metadata or execution.
849    pub fn get_ref<I: rusqlite::RowIndex>(
850        &self,
851        index: I,
852    ) -> rusqlite::Result<rusqlite::types::ValueRef<'_>> {
853        self.row.get_ref(index)
854    }
855}
856
857/// One public query owns this lease's connection-global progress handler.
858struct ReaderQueryInProgress<'a>(&'a Cell<bool>);
859
860impl Drop for ReaderQueryInProgress<'_> {
861    fn drop(&mut self) {
862        self.0.set(false);
863    }
864}
865
866/// One pool-wide reader admission permit acquired before any connection is
867/// selected, with the instant its checkout wait began so a single
868/// `checkout_timeout` bounds both the permit wait and the connection pick.
869/// Hand it to [`ConnectionPool::reader_with_admission`], which moves the permit
870/// into the resulting [`ReaderGuard`].
871pub(crate) struct ReaderAdmission {
872    slot: tokio::sync::OwnedSemaphorePermit,
873    started: Instant,
874}
875
876/// A reader connection checked out from the pool.
877/// Returns the connection to the pool on drop.
878pub struct ReaderGuard<'pool> {
879    lease: Option<ReaderLease<'pool>>,
880    /// One permit from the pool-wide reader budget, shared with the explicit
881    /// raw-SQL transaction exception. Returned only after the connection has
882    /// been reset/replaced and made reusable.
883    admission_slot: Option<tokio::sync::OwnedSemaphorePermit>,
884    pool: &'pool ConnectionPool,
885    reusable: Cell<bool>,
886    query_in_progress: Cell<bool>,
887    checked_out_at: Instant,
888    /// Set by [`Self::mark_dirty`] whenever this checkout ran a `SqlReader`
889    /// raw-SQL statement (`sql_bridge`'s `run_pool_reader_query`), never by a
890    /// typed store read. Gates whether `Drop` pays the TEMP-object,
891    /// attachment, and connection-setting pristine scan on return —
892    /// `reader_connection_state_is_pristine` and
893    /// `reader_connection_settings_match_baseline` never touch the hot path
894    /// of an ordinary typed checkout.
895    dirty: Cell<bool>,
896    /// Names the typed-store operation this checkout was resolved for, set
897    /// once by [`ConnectionPool::resolve_reader_checkout`]. `None` means the
898    /// checkout never passed that route (the pool-internal and raw-SQL
899    /// callers), and the diagnostics maximum reports it as unattributed
900    /// rather than guessing.
901    operation: Option<&'static str>,
902}
903
904impl<'pool> ReaderGuard<'pool> {
905    /// Access the connection from within `khive-db`. Every internal caller
906    /// is either a typed store (proven read-only by construction) or a
907    /// raw-SQL route that already calls [`Self::mark_dirty`] itself, so this
908    /// stays crate-private: it is the untracked half of the pristine-return
909    /// contract, and a caller outside the crate has no way to pay into that
910    /// contract by calling `mark_dirty` (also crate-private). A raw
911    /// `&Connection` is never handed to a caller outside `khive-db` — use
912    /// [`Self::query_row`] instead, which admits only read-shaped SQL.
913    pub(crate) fn conn(&self) -> &Connection {
914        match self
915            .lease
916            .as_ref()
917            .expect("reader guard missing connection")
918        {
919            ReaderLease::Pooled(conn) => conn,
920            ReaderLease::Shared(guard) => guard,
921        }
922    }
923
924    /// Run one read-only statement against this reader lease and map the
925    /// first resulting row, for callers outside `khive-db`.
926    ///
927    /// Unlike a raw `Connection`, this never hands out a capability that can
928    /// change connection-local or database state: `sql` is checked against
929    /// the same allow-listed read-shape classifier
930    /// (`SELECT`/`WITH ... SELECT`/`VALUES`/`EXPLAIN`/a fixed read-only
931    /// `PRAGMA` set) the pooled `SqlReader` surface admits raw SQL through,
932    /// and anything else — `BEGIN`, DML, DDL, a setting `PRAGMA`, `ATTACH`
933    /// — is refused before it ever reaches SQLite. An admitted statement
934    /// still marks the checkout dirty unconditionally, so `Drop` always pays
935    /// the pristine-state scan (or, in degraded shared-lease mode, the
936    /// settings/rollback verification) on return.
937    ///
938    /// SQL stepping cooperatively observes the current request's cancellation
939    /// and original absolute deadline. Synchronous mapper/native callback code
940    /// cannot be forcibly preempted; a post-check refuses a successful result
941    /// if the request stopped while that code ran.
942    pub fn query_row<T, P, F>(&self, sql: &str, params: P, f: F) -> Result<T, SqliteError>
943    where
944        P: rusqlite::Params,
945        F: FnOnce(&ReaderRow<'_, '_>) -> rusqlite::Result<T>,
946    {
947        crate::sql_bridge::reader_capability_admits(sql).map_err(SqliteError::InvalidData)?;
948        if !self.reusable.get() {
949            return Err(SqliteError::InvalidData(
950                "reader lease is quarantined after failed read cleanup".into(),
951            ));
952        }
953        // Params may call user-provided ToSql before stepping, and the mapper
954        // also runs user code. Neither may replace this query's progress handler
955        // by recursively querying the same lease.
956        if self.query_in_progress.replace(true) {
957            return Err(SqliteError::InvalidData(
958                "reader lease is already executing a query".into(),
959            ));
960        }
961        let _in_progress = ReaderQueryInProgress(&self.query_in_progress);
962        self.mark_dirty();
963        crate::read_cancellation::run_borrowed_reader(self, |conn, admission| {
964            conn.query_row(sql, params, |row| {
965                // Parameter conversion runs before any stepping, and a query
966                // cheap enough not to trip the progress handler never polls at
967                // all, so this is the only point that can keep the documented
968                // guarantee: a cancelled request does not enter the mapper.
969                // SQLITE_INTERRUPT is the code the scope already recognises, so
970                // the refusal converts to the same timeout error as a
971                // cancellation observed during stepping.
972                if !admission.admits() {
973                    return Err(rusqlite::Error::SqliteFailure(
974                        rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_INTERRUPT),
975                        Some("request stopped before the mapper".into()),
976                    ));
977                }
978                f(&ReaderRow { row })
979            })
980            .map_err(|error| {
981                StorageError::driver(StorageCapability::Sql, "reader_guard.query_row", error)
982            })
983        })
984        .map_err(|error| {
985            self.pool.record_reader_query_error(&error);
986            match error {
987                StorageError::Driver {
988                    capability,
989                    operation,
990                    source,
991                } => match source.downcast::<rusqlite::Error>() {
992                    Ok(error) => SqliteError::Rusqlite(*error),
993                    Err(source) => SqliteError::RequestReadStopped(StorageError::Driver {
994                        capability,
995                        operation,
996                        source,
997                    }),
998                },
999                other => SqliteError::RequestReadStopped(other),
1000            }
1001        })
1002    }
1003
1004    /// Fail closed when connection-global state could not be restored after
1005    /// a read. A pooled reader is closed and replaced on drop; a degraded
1006    /// shared-writer reader is quarantined for the lifetime of the pool.
1007    pub(crate) fn discard(&self) {
1008        self.reusable.set(false);
1009    }
1010
1011    /// Mark this checkout as having run a raw-SQL statement, so `Drop` pays
1012    /// the pristine-state scan on return instead of skipping it.
1013    pub(crate) fn mark_dirty(&self) {
1014        self.dirty.set(true);
1015    }
1016
1017    /// Name the typed-store operation this checkout serves, so a long hold
1018    /// can be attributed in diagnostics instead of arriving as a bare
1019    /// maximum with no next step (#2793).
1020    pub(crate) fn label_operation(&mut self, operation: &'static str) {
1021        self.operation = Some(operation);
1022    }
1023}
1024
1025impl<'pool> Drop for ReaderGuard<'pool> {
1026    fn drop(&mut self) {
1027        let Some(lease) = self.lease.take() else {
1028            return;
1029        };
1030
1031        match lease {
1032            ReaderLease::Pooled(conn) if self.reusable.get() => {
1033                self.pool.return_reader(conn, self.dirty.get())
1034            }
1035            ReaderLease::Pooled(conn) => {
1036                close_connection_quietly(conn);
1037                self.pool.replace_discarded_reader_slot();
1038            }
1039            ReaderLease::Shared(guard) if !self.reusable.get() => {
1040                self.pool.retire_pooled_writer(&guard);
1041            }
1042            ReaderLease::Shared(guard) => {
1043                // The shared lease IS the pool's writer connection (degraded
1044                // `max_readers == 0` mode) — there is no separate reader
1045                // connection to close and replace. A dirty return that
1046                // cannot be verifiably restored poisons the whole pool via
1047                // the same terminal-fault path a writer transaction fault
1048                // uses, since reuse and replacement are equally unavailable
1049                // here.
1050                if self.dirty.get() && !restore_shared_reader_state(&guard, &self.pool.config) {
1051                    self.pool.retire_pooled_writer(&guard);
1052                }
1053            }
1054        }
1055
1056        // A checkout remains active until reset/replacement returned the
1057        // connection to service (or the degraded writer guard was released),
1058        // not merely until the caller's query closure returned.
1059        drop(self.admission_slot.take());
1060        self.pool
1061            .reader_acquisition_counters
1062            .record_checkout_completed(self.checked_out_at.elapsed(), self.operation);
1063    }
1064}
1065
1066/// Owned, `'static` counterpart to [`ReaderGuard`]'s degraded-mode
1067/// (`max_readers == 0`) lease. `ReaderGuard` borrows `&'pool ConnectionPool`
1068/// and therefore cannot be retained across the separate `.await` points
1069/// between a caller's `SqlReader` trait-method calls — every ordinary call
1070/// through it draws a fresh checkout and returns it before the next call
1071/// begins. That is fine for an ordinary read, but wrong for the explicit
1072/// multi-call deferred read transaction (ADR-005/ADR-091): a `BEGIN
1073/// DEFERRED` checked out and returned this way releases the pool's one
1074/// physical connection to any concurrent caller — including a writer —
1075/// before its own matching `COMMIT`/`ROLLBACK` runs, and an abandoned span
1076/// (an error or cancellation between the two) leaves that connection sitting
1077/// in the pool inside an open transaction with nothing left holding it. This
1078/// type owns an `Arc`-rooted mutex guard instead of borrowing one, so it can
1079/// be held by the caller across those `.await` points and give the whole
1080/// span real, exclusive ownership of the one connection.
1081pub(crate) struct SharedReaderTransactionGuard {
1082    conn: parking_lot::ArcMutexGuard<parking_lot::RawMutex, Connection>,
1083    admission_slot: Option<tokio::sync::OwnedSemaphorePermit>,
1084    pool: Arc<ConnectionPool>,
1085    checked_out_at: Instant,
1086    /// Set once a statement executed against this connection could not
1087    /// prove its own cleanup ran (SQLite progress-handler removal failure).
1088    /// `Drop` then treats the connection exactly like a rollback failure:
1089    /// poison the pool rather than let possibly-corrupted connection state
1090    /// return to service.
1091    poison: Cell<bool>,
1092}
1093
1094impl SharedReaderTransactionGuard {
1095    pub(crate) fn conn(&self) -> &Connection {
1096        &self.conn
1097    }
1098
1099    /// Force poisoning on drop regardless of the connection's own
1100    /// autocommit/pristine state.
1101    pub(crate) fn poison(&self) {
1102        self.poison.set(true);
1103    }
1104}
1105
1106impl Drop for SharedReaderTransactionGuard {
1107    fn drop(&mut self) {
1108        // An abandoned span must never let this connection return to
1109        // service while still inside an open transaction — the next reader
1110        // or writer to draw the pool's one physical connection would
1111        // silently inherit it. Roll back first; only a verified return to
1112        // autocommit, followed by the same pristine-state scan an ordinary
1113        // dirty degraded checkout pays, allows reuse.
1114        let mut restored = !self.poison.get();
1115        if restored && !self.conn.is_autocommit() {
1116            restored = self.conn.execute_batch("ROLLBACK").is_ok() && self.conn.is_autocommit();
1117        }
1118        if restored {
1119            restored = restore_shared_reader_state(&self.conn, &self.pool.config);
1120        }
1121        if !restored {
1122            self.pool.retire_pooled_writer(&self.conn);
1123        }
1124        drop(self.admission_slot.take());
1125        self.pool
1126            .reader_acquisition_counters
1127            .record_checkout_completed(
1128                self.checked_out_at.elapsed(),
1129                Some("explicit_sql_read_transaction"),
1130            );
1131    }
1132}
1133
1134impl ConnectionPool {
1135    /// Check out the degraded-mode (`max_readers == 0`) pool's single
1136    /// physical connection as an owned [`SharedReaderTransactionGuard`]
1137    /// instead of a borrowed [`ReaderGuard`]. Mirrors the admission and
1138    /// mutex-acquisition loop in [`Self::reader_until`]'s degraded branch —
1139    /// the only difference is the guard type returned, so a caller can
1140    /// retain this one across several separate `.await` points.
1141    ///
1142    /// Exists only for the in-memory backend's explicit deferred-read-
1143    /// transaction bypass (`sql_bridge::PoolBackedReader`); ordinary reads
1144    /// use [`Self::reader_until`].
1145    pub(crate) fn checkout_shared_reader_transaction(
1146        self: &Arc<Self>,
1147        should_stop: impl Fn() -> bool,
1148    ) -> Result<Option<SharedReaderTransactionGuard>, SqliteError> {
1149        debug_assert_eq!(
1150            self.max_readers, 0,
1151            "the owned shared-reader-transaction guard exists only for the degraded, \
1152             single-connection backend"
1153        );
1154        self.ensure_pooled_writer_active()?;
1155        let started = Instant::now();
1156        let admission_slot = loop {
1157            if should_stop() {
1158                return Ok(None);
1159            }
1160            match Arc::clone(&self.sql_bridge_reader_slots).try_acquire_owned() {
1161                Ok(slot) => break slot,
1162                Err(tokio::sync::TryAcquireError::Closed) => {
1163                    return Err(SqliteError::InvalidData(
1164                        "reader admission semaphore is closed".to_string(),
1165                    ));
1166                }
1167                Err(tokio::sync::TryAcquireError::NoPermits) => {}
1168            }
1169            if started.elapsed() >= self.config.checkout_timeout {
1170                self.reader_acquisition_counters.record_checkout_timeout();
1171                return Err(pool_exhausted_error(
1172                    self.config.checkout_timeout,
1173                    self.max_readers,
1174                ));
1175            }
1176            thread::yield_now();
1177        };
1178
1179        loop {
1180            if should_stop() {
1181                return Ok(None);
1182            }
1183            let remaining = self
1184                .config
1185                .checkout_timeout
1186                .saturating_sub(started.elapsed());
1187            if remaining.is_zero() {
1188                self.reader_acquisition_counters.record_checkout_timeout();
1189                return Err(pool_exhausted_error(
1190                    self.config.checkout_timeout,
1191                    self.max_readers,
1192                ));
1193            }
1194            if let Some(conn) = self
1195                .writer
1196                .try_lock_arc_for(remaining.min(Duration::from_millis(2)))
1197            {
1198                self.ensure_pooled_writer_active()?;
1199                self.reader_acquisition_counters.record_pooled_checkout();
1200                return Ok(Some(SharedReaderTransactionGuard {
1201                    conn,
1202                    admission_slot: Some(admission_slot),
1203                    pool: Arc::clone(self),
1204                    checked_out_at: Instant::now(),
1205                    poison: Cell::new(false),
1206                }));
1207            }
1208        }
1209    }
1210}
1211
1212/// Why a caller is allowed to bypass the pooled reader queue.
1213///
1214/// This is deliberately a closed, crate-private list. Ordinary request reads
1215/// have no variant: they must use [`ConnectionPool::reader`] and surface its
1216/// bounded admission timeout without falling back to a fresh connection.
1217#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1218#[allow(dead_code)] // The closed ADR list includes infrastructure paths not instantiated today.
1219pub(crate) enum StandaloneReaderPurpose {
1220    /// ADR-005/ADR-091's explicit, multi-call deferred raw-SQL transaction.
1221    ExplicitSqlReadTransaction,
1222    /// Boot-time schema/model-registry inspection before a runtime pool can
1223    /// own the read. Kept separate from request traffic in diagnostics.
1224    BootSchemaProbe,
1225    /// An operator diagnostic that requires a physically independent
1226    /// snapshot. Kept separate from request traffic in diagnostics.
1227    DiagnosticsIndependentSnapshot,
1228}
1229
1230impl StandaloneReaderPurpose {
1231    fn is_infrastructure(self) -> bool {
1232        matches!(
1233            self,
1234            Self::BootSchemaProbe | Self::DiagnosticsIndependentSnapshot
1235        )
1236    }
1237}
1238
1239/// Process-local reader acquisition and hold lifecycle since one
1240/// [`ConnectionPool`] was constructed.
1241///
1242/// Counters are monotonic and reset only when the pool is reconstructed.
1243/// `active_pooled_checkouts` is the point-in-time number of live
1244/// [`ReaderGuard`] values. Hold duration includes connection reset or
1245/// replacement on return, because the slot is not reusable before then.
1246#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
1247pub struct ReaderAcquisitionSnapshot {
1248    /// Configured process-local admission budget shared by pooled readers and
1249    /// the explicit raw-SQL read-transaction exception.
1250    pub reader_admission_capacity: usize,
1251    /// Point-in-time permits not held by either pooled readers or explicit
1252    /// raw-SQL read transactions.
1253    pub available_reader_admission_slots: usize,
1254    /// Successful request-path reader acquisitions (pooled plus the explicit
1255    /// raw-SQL transaction exception; infrastructure opens excluded).
1256    pub acquisitions: u64,
1257    /// Successful bounded pooled-reader checkouts.
1258    pub pooled_checkouts: u64,
1259    /// Successful request-path standalone opens from the closed exception
1260    /// list (currently explicit raw-SQL deferred transactions only).
1261    pub standalone_opens: u64,
1262    /// Successful boot/diagnostic standalone opens, attributed separately so
1263    /// request traffic cannot be inferred from infrastructure activity.
1264    pub infrastructure_standalone_opens: u64,
1265    /// Reader-admission waits that exhausted `checkout_timeout` before any
1266    /// query began. Covers pooled checkout and the closed raw-SQL exception;
1267    /// cooperative request cancellation is intentionally excluded.
1268    pub checkout_timeouts: u64,
1269    /// Queries on a checked-out pooled reader that SQLite refused with
1270    /// `SQLITE_BUSY` after the connection's busy handler gave up (typed-store
1271    /// reads, pooled raw-SQL reads, and [`ReaderGuard::query_row`]). Counted
1272    /// after checkout succeeded, so it never overlaps `checkout_timeouts`;
1273    /// `SQLITE_LOCKED` and cooperative cancellation are excluded. Writer
1274    /// refusals are not counted here; the writer task's are in
1275    /// [`WriterAcquisitionSnapshot::writer_task_begin_busy`].
1276    pub busy_timeouts: u64,
1277    /// Pooled checkouts live at the instant this snapshot was taken.
1278    pub active_pooled_checkouts: u64,
1279    /// High-water mark of concurrent pooled checkouts.
1280    pub peak_active_pooled_checkouts: u64,
1281    /// Pooled guards that completed their full return/reset lifecycle.
1282    pub completed_pooled_checkouts: u64,
1283    /// Longest completed checkout hold, including return/reset, in
1284    /// microseconds. Diagnostic evidence only; never a test timing gate.
1285    pub max_completed_hold_micros: u64,
1286    /// The typed-store operation that held the checkout reported in
1287    /// `max_completed_hold_micros`. `None` when that hold came from a route
1288    /// that carries no operation name, which is itself the answer rather
1289    /// than a missing reading (#2793).
1290    pub max_completed_hold_operation: Option<&'static str>,
1291    /// A disqualified pooled-reader return (reset/pristine-check failure)
1292    /// whose replacement connection then also failed to open, permanently
1293    /// shrinking the physical pool by one slot below `max_readers`. Logged at
1294    /// `warn` when it happens; this counter makes the shrink observable in a
1295    /// snapshot too, since the pool itself never re-grows on its own.
1296    pub reader_replacement_open_failures: u64,
1297}
1298
1299/// The longest completed pooled-reader hold and the operation that held it,
1300/// kept under one lock so a snapshot cannot pair one checkout's duration with
1301/// another's name.
1302#[derive(Debug, Default, Clone, Copy)]
1303struct LongestCompletedHold {
1304    micros: u64,
1305    operation: Option<&'static str>,
1306}
1307
1308#[derive(Debug, Default)]
1309struct ReaderAcquisitionCounters {
1310    pooled_checkouts: AtomicU64,
1311    standalone_opens: AtomicU64,
1312    infrastructure_standalone_opens: AtomicU64,
1313    checkout_timeouts: AtomicU64,
1314    busy_timeouts: AtomicU64,
1315    active_pooled_checkouts: AtomicU64,
1316    peak_active_pooled_checkouts: AtomicU64,
1317    completed_pooled_checkouts: AtomicU64,
1318    longest_completed_hold: parking_lot::Mutex<LongestCompletedHold>,
1319    reader_replacement_open_failures: AtomicU64,
1320}
1321
1322impl ReaderAcquisitionCounters {
1323    fn record_pooled_checkout(&self) {
1324        self.pooled_checkouts.fetch_add(1, Ordering::Relaxed);
1325        let active = self
1326            .active_pooled_checkouts
1327            .fetch_add(1, Ordering::Relaxed)
1328            .saturating_add(1);
1329        self.peak_active_pooled_checkouts
1330            .fetch_max(active, Ordering::Relaxed);
1331    }
1332
1333    fn record_checkout_timeout(&self) {
1334        self.checkout_timeouts.fetch_add(1, Ordering::Relaxed);
1335    }
1336
1337    fn record_busy_timeout(&self) {
1338        self.busy_timeouts.fetch_add(1, Ordering::Relaxed);
1339    }
1340
1341    fn record_reader_replacement_open_failure(&self) {
1342        self.reader_replacement_open_failures
1343            .fetch_add(1, Ordering::Relaxed);
1344    }
1345
1346    fn record_standalone_open(&self, purpose: StandaloneReaderPurpose) {
1347        if purpose.is_infrastructure() {
1348            self.infrastructure_standalone_opens
1349                .fetch_add(1, Ordering::Relaxed);
1350        } else {
1351            self.standalone_opens.fetch_add(1, Ordering::Relaxed);
1352        }
1353    }
1354
1355    fn record_checkout_completed(&self, hold: Duration, operation: Option<&'static str>) {
1356        let previous = self.active_pooled_checkouts.fetch_sub(1, Ordering::Relaxed);
1357        debug_assert!(previous > 0, "reader active-checkout counter underflow");
1358        self.completed_pooled_checkouts
1359            .fetch_add(1, Ordering::Relaxed);
1360        let micros = u64::try_from(hold.as_micros()).unwrap_or(u64::MAX);
1361        // The maximum and the name of what held it are one reading: taken
1362        // apart they can report a duration from one checkout beside a label
1363        // from another, which is worse than no label at all.
1364        let mut longest = self.longest_completed_hold.lock();
1365        if micros > longest.micros {
1366            longest.micros = micros;
1367            longest.operation = operation;
1368        }
1369    }
1370
1371    fn snapshot(
1372        &self,
1373        reader_admission_capacity: usize,
1374        available_reader_admission_slots: usize,
1375    ) -> ReaderAcquisitionSnapshot {
1376        let pooled_checkouts = self.pooled_checkouts.load(Ordering::Relaxed);
1377        let standalone_opens = self.standalone_opens.load(Ordering::Relaxed);
1378        let longest_completed_hold = *self.longest_completed_hold.lock();
1379        ReaderAcquisitionSnapshot {
1380            reader_admission_capacity,
1381            available_reader_admission_slots,
1382            acquisitions: pooled_checkouts.saturating_add(standalone_opens),
1383            pooled_checkouts,
1384            standalone_opens,
1385            infrastructure_standalone_opens: self
1386                .infrastructure_standalone_opens
1387                .load(Ordering::Relaxed),
1388            checkout_timeouts: self.checkout_timeouts.load(Ordering::Relaxed),
1389            busy_timeouts: self.busy_timeouts.load(Ordering::Relaxed),
1390            active_pooled_checkouts: self.active_pooled_checkouts.load(Ordering::Relaxed),
1391            peak_active_pooled_checkouts: self.peak_active_pooled_checkouts.load(Ordering::Relaxed),
1392            completed_pooled_checkouts: self.completed_pooled_checkouts.load(Ordering::Relaxed),
1393            max_completed_hold_micros: longest_completed_hold.micros,
1394            max_completed_hold_operation: longest_completed_hold.operation,
1395            reader_replacement_open_failures: self
1396                .reader_replacement_open_failures
1397                .load(Ordering::Relaxed),
1398        }
1399    }
1400}
1401
1402impl ConnectionPool {
1403    /// Create a new connection pool.
1404    ///
1405    /// Opens 1 writer + N reader connections to the same database when pooling
1406    /// is enabled. All connections are configured consistently (busy timeout,
1407    /// foreign keys, cache, mmap, temp store). Writable in-memory databases and
1408    /// writable non-WAL files fall back to single-connection mode. Read-only
1409    /// files retain a dedicated reader regardless of journal mode.
1410    pub fn new(config: PoolConfig) -> Result<Self, SqliteError> {
1411        refuse_home_data_store_in_tests(&config)?;
1412        validate_write_admission_deadline(config.write_admission_deadline_ms)?;
1413        if let Some(policy) = config.disk_guard_config {
1414            policy.validate()?;
1415        }
1416        config.wal_ceiling.validate_static(
1417            config.path.is_some(),
1418            config.wal_mode,
1419            config.read_only,
1420        )?;
1421        code_map::validate_pool(&config)?;
1422        if config.path.is_some() && !config.read_only {
1423            let lock_dir =
1424                crate::disk_guard_config::require_volume_lock_dir(config.volume_lock_dir.clone())?;
1425            if !lock_dir.is_absolute() {
1426                return Err(SqliteError::CapacityUnavailable {
1427                    phase: CapacityUnavailablePhase::Lock,
1428                    message: "configured volume-lock directory is not absolute".to_string(),
1429                });
1430            }
1431        }
1432
1433        // Resolve "no preference" (`None`) now that `path` is known: on for
1434        // file-backed pools, off for in-memory ones. An explicit `Some(_)`
1435        // preference is left untouched and always wins.
1436        let mut config = config;
1437        let inert_memory_queue_request =
1438            config.path.is_none() && config.write_queue_enabled == Some(true);
1439        config.write_queue_enabled =
1440            Some(config.write_queue_enabled.unwrap_or(config.path.is_some()));
1441        if inert_memory_queue_request {
1442            tracing::warn!(
1443                "write queue explicitly requested for an in-memory pool; it is inert because \
1444                 in-memory pools cannot host a writer task"
1445            );
1446        }
1447
1448        // Mint the physical identity before WAL classification or SQLite open.
1449        // Every read-only connection below uses this same canonical path (or
1450        // an immutable URI derived from it), so a symlink cannot split main-file
1451        // resolution from sidecar resolution.
1452        let (origin, identity_path) = match config.path.as_ref() {
1453            Some(path) if config.code_map_vfs.is_some() => code_map::guarded_identity(path),
1454            Some(path) => {
1455                let (identity, canonical) = mint_db_identity(path)?;
1456                (TxOrigin::Database(identity), Some(canonical))
1457            }
1458            None => (TxOrigin::Memory, None),
1459        };
1460        let (write_admission, disk_guard_config) = if !config.read_only {
1461            if identity_path.as_deref().and_then(Path::parent).is_some() {
1462                let disk_guard = match config.disk_guard_config {
1463                    Some(policy) => policy,
1464                    None => resolve_disk_guard_config(None, None)?,
1465                };
1466                disk_guard.validate()?;
1467                if disk_guard.legacy_environment_present {
1468                    tracing::warn!(
1469                        "legacy SQLite reserve setting is deprecated; \
1470                         use KHIVE_SQLITE_DISK_RESERVE_BYTES"
1471                    );
1472                }
1473                if disk_guard.reserve_bytes == 0 {
1474                    tracing::warn!(
1475                        "SQLite disk reserve is explicitly zero; \
1476                         new logical writes will not be floor-refused"
1477                    );
1478                }
1479                (
1480                    Arc::new(WriteAdmission::new(
1481                        identity_path.clone(),
1482                        disk_guard.reserve_bytes,
1483                        disk_guard.guard_deadline_ms,
1484                        config.volume_lock_dir.clone(),
1485                    )?),
1486                    Some(disk_guard),
1487                )
1488            } else {
1489                (
1490                    Arc::new(WriteAdmission::new(
1491                        None,
1492                        0,
1493                        DEFAULT_DISK_GUARD_DEADLINE_MS,
1494                        None,
1495                    )?),
1496                    None,
1497                )
1498            }
1499        } else {
1500            (
1501                Arc::new(WriteAdmission::new(
1502                    None,
1503                    0,
1504                    DEFAULT_DISK_GUARD_DEADLINE_MS,
1505                    None,
1506                )?),
1507                None,
1508            )
1509        };
1510        let read_only_open_target = read_only_open_target(&config, identity_path.as_deref())?;
1511        #[cfg(any(unix, windows))]
1512        let identity_before_open = identity_path
1513            .as_deref()
1514            .map(database_file_identity_if_exists)
1515            .transpose()?
1516            .flatten();
1517        #[cfg(any(unix, windows))]
1518        claimed_file_identity::verify_before_open(&config, identity_before_open)?;
1519        let retirement_connection = Connection::open_in_memory()?;
1520        retirement_connection.authorizer(Some(deny_retired_writer))?;
1521        // Declared first so a failed bootstrap closes SQLite before releasing its lease.
1522        let mut initialization_lease = None;
1523        let mut writer = open_writer_connection(
1524            &config,
1525            read_only_open_target.as_deref(),
1526            identity_path.as_deref(),
1527        )?;
1528        validate_wal_ceiling_at_open(&writer, &config)?;
1529        // The identity bootstrap can take a write lock before the remaining
1530        // connection pragmas are configured. Honor the caller's wait bound.
1531        writer.busy_timeout(config.busy_timeout)?;
1532        // A read of main.sqlite_master forces SQLite's main file open without
1533        // changing either database. Reject an already-swapped target before
1534        // installing a nonce into a legacy or initially empty database.
1535        let initial_database_id = if identity_path.is_some() {
1536            read_database_id(&writer)?
1537        } else {
1538            None
1539        };
1540        #[cfg(test)]
1541        if let Some(path) = identity_path.as_deref() {
1542            run_identity_open_hook(
1543                path,
1544                IdentityOpenStage::AfterMainOpenBeforeFirstStat,
1545                Some(&writer),
1546            );
1547        }
1548        #[cfg(any(unix, windows))]
1549        let identity_before_write = identity_path
1550            .as_deref()
1551            .map(database_file_identity)
1552            .transpose()?;
1553        #[cfg(any(unix, windows))]
1554        if identity_before_open.is_some() && identity_before_open != identity_before_write {
1555            return Err(SqliteError::InvalidData(
1556                "database file identity changed while opening the pool".to_string(),
1557            ));
1558        }
1559        #[cfg(any(unix, windows))]
1560        if let Some(path) = identity_path.as_deref() {
1561            let opened = opened_sqlite_file_identity(&writer, path)?;
1562            if identity_before_write != Some(opened) {
1563                return Err(SqliteError::InvalidData(
1564                    "database file identity changed while opening the pool".to_string(),
1565                ));
1566            }
1567        }
1568        // A newly-created file must still belong to the volume captured from
1569        // its parent before open, and an existing file must retain that volume.
1570        write_admission.verify_current_volume()?;
1571        let opened_database_id = if identity_path.is_some()
1572            && !config.read_only
1573            && initial_database_id.is_none()
1574        {
1575            match write_admission.check() {
1576                Ok(()) => {
1577                    initialization_lease = write_admission.acquire()?;
1578                    match initialize_database_id_with_admission(&mut writer, Some(&write_admission))
1579                    {
1580                        Ok(id) => Some(id),
1581                        Err(SqliteError::CapacityFloor { .. }) => None,
1582                        Err(error) => return Err(error),
1583                    }
1584                }
1585                // Recovery must be able to open this pool and its checkpoint
1586                // connection below the reserve. Physical file pinning still
1587                // applies; a later pool open can install the nonce.
1588                Err(SqliteError::CapacityFloor { .. }) => None,
1589                Err(error) => return Err(error),
1590            }
1591        } else {
1592            initial_database_id
1593        };
1594        drop(initialization_lease.take());
1595        #[cfg(test)]
1596        if let Some(path) = identity_path.as_deref() {
1597            run_identity_open_hook(
1598                path,
1599                IdentityOpenStage::AfterInitialIdentityWrite,
1600                Some(&writer),
1601            );
1602        }
1603        #[cfg(any(unix, windows))]
1604        let opened_file_identity = identity_path
1605            .as_deref()
1606            .map(|path| opened_sqlite_file_identity(&writer, path))
1607            .transpose()?;
1608        #[cfg(any(unix, windows))]
1609        if identity_before_write != opened_file_identity {
1610            return Err(SqliteError::InvalidData(
1611                "database file identity changed while opening the pool".to_string(),
1612            ));
1613        }
1614        let wal_enabled = configure_writer_connection(&writer, &config)?;
1615        let max_readers = effective_reader_count(&config, wal_enabled);
1616
1617        let readers = ArrayQueue::new(max_readers.max(1));
1618
1619        #[cfg(any(test, feature = "test-support"))]
1620        let statement_observer = crate::statement_observer::StatementObserverHub::new()?;
1621        #[cfg(any(test, feature = "test-support"))]
1622        crate::statement_observer::install(&writer, &statement_observer)?;
1623
1624        let mut pool = Self {
1625            writer: Arc::new(Mutex::new(writer)),
1626            retirement_connection: Mutex::new(Some(retirement_connection)),
1627            #[cfg(any(test, feature = "test-support"))]
1628            statement_observer,
1629            main_pool_generation: OnceLock::new(),
1630            checkpoint_ownership: CheckpointOwnershipGate::new(),
1631            pooled_writer_retired: AtomicBool::new(false),
1632            writer_acquisition_counters: Arc::new(WriterAcquisitionCounters::default()),
1633            write_admission,
1634            disk_guard_config,
1635            reader_acquisition_counters: ReaderAcquisitionCounters::default(),
1636            search_dispatches: Mutex::new(BTreeMap::new()),
1637            note_candidate_hydration_rows: AtomicU64::new(0),
1638            readers,
1639            max_readers,
1640            config,
1641            read_only_open_target,
1642            sql_bridge_reader_slots: Arc::new(Semaphore::new(max_readers.max(1))),
1643            sql_bridge_writer_slots: Arc::new(Semaphore::new(1)),
1644            writer_task: OnceLock::new(),
1645            writer_task_join: Mutex::new(None),
1646            writer_task_join_stored: AtomicBool::new(false),
1647            origin,
1648            identity_path,
1649            #[cfg(any(unix, windows))]
1650            opened_file_identity,
1651            opened_database_id,
1652            identity_registration: None,
1653            #[cfg(test)]
1654            writer_task_spawn_count: std::sync::atomic::AtomicUsize::new(0),
1655        };
1656
1657        for _ in 0..pool.max_readers {
1658            let conn = pool.open_reader_connection()?;
1659            pool.readers
1660                .push(conn)
1661                .expect("reader queue must have capacity during pool initialization");
1662        }
1663
1664        // Best-effort, process-global diagnostics belong only to pools that
1665        // can acquire a writer. A read-only inspection pool has no writer
1666        // timeout to report and must neither mutate `<db_parent>/.khive-logs`
1667        // nor consume the global sink claim before a later writable pool.
1668        if !pool.config.read_only {
1669            crate::timeout_sink::init(
1670                pool.canonical_path().and_then(Path::parent),
1671                &crate::timeout_sink::db_label(&pool),
1672            );
1673        }
1674
1675        pool.identity_registration = pool.canonical_path().map(PoolIdentityRegistration::new);
1676        Ok(pool)
1677    }
1678
1679    /// Check out a reader connection.
1680    ///
1681    /// Tries to pop from the lock-free queue. If empty, spins briefly then
1682    /// waits with exponential backoff up to `checkout_timeout`.
1683    ///
1684    /// In degraded mode (WAL unavailable, `max_readers == 0`), this method
1685    /// checks the shared writer mutex in bounded slices and returns pool
1686    /// exhaustion after `checkout_timeout`; it never blocks indefinitely on
1687    /// the non-reentrant mutex.
1688    pub fn reader(&self) -> Result<ReaderGuard<'_>, SqliteError> {
1689        self.reader_until(|| false)?.ok_or_else(|| {
1690            SqliteError::InvalidData("uncancelled reader checkout stopped unexpectedly".into())
1691        })
1692    }
1693
1694    /// Check out a reader while cooperatively polling a request cancellation
1695    /// predicate. The predicate is evaluated before connection acquisition and
1696    /// between backoff slices, so an abandoned request does not sit through the
1697    /// full pool checkout timeout or execute a statement when a reader later
1698    /// becomes available.
1699    pub(crate) fn reader_until<C>(
1700        &self,
1701        should_stop: C,
1702    ) -> Result<Option<ReaderGuard<'_>>, SqliteError>
1703    where
1704        C: Fn() -> bool,
1705    {
1706        let started = Instant::now();
1707        let mut admission_attempt = 0u32;
1708        let admission_slot = loop {
1709            if should_stop() {
1710                return Ok(None);
1711            }
1712            match Arc::clone(&self.sql_bridge_reader_slots).try_acquire_owned() {
1713                Ok(slot) => break slot,
1714                Err(tokio::sync::TryAcquireError::Closed) => {
1715                    return Err(SqliteError::InvalidData(
1716                        "reader admission semaphore is closed".to_string(),
1717                    ));
1718                }
1719                Err(tokio::sync::TryAcquireError::NoPermits) => {}
1720            }
1721            if started.elapsed() >= self.config.checkout_timeout {
1722                self.reader_acquisition_counters.record_checkout_timeout();
1723                return Err(pool_exhausted_error(
1724                    self.config.checkout_timeout,
1725                    self.max_readers,
1726                ));
1727            }
1728            match admission_attempt {
1729                0..=7 => {
1730                    let spins = 1usize << admission_attempt;
1731                    for _ in 0..spins {
1732                        std::hint::spin_loop();
1733                    }
1734                }
1735                8..=15 => thread::yield_now(),
1736                _ => {
1737                    let remaining = self
1738                        .config
1739                        .checkout_timeout
1740                        .saturating_sub(started.elapsed());
1741                    let sleep =
1742                        Duration::from_micros(50 * (1u64 << (admission_attempt - 16).min(6)));
1743                    thread::sleep(sleep.min(remaining).min(Duration::from_millis(2)));
1744                }
1745            }
1746            admission_attempt = admission_attempt.saturating_add(1);
1747        };
1748
1749        self.reader_with_admission(
1750            ReaderAdmission {
1751                slot: admission_slot,
1752                started,
1753            },
1754            should_stop,
1755        )
1756    }
1757
1758    /// Wait for one pool-wide reader admission permit on the async side, so a
1759    /// read that has to queue costs a task and not a blocking-pool thread.
1760    ///
1761    /// The wait is bounded by `checkout_timeout` and by the current request's
1762    /// cancellation and deadline, and it yields the same `Ok(Some)` / `Ok(None)`
1763    /// / `Err` outcomes the permit loop in [`Self::reader_until`] does, resolved
1764    /// through the same refusal mapping as [`Self::resolve_reader_checkout`].
1765    /// Dropping the future abandons the wait without taking a permit.
1766    pub(crate) async fn acquire_reader_admission(
1767        &self,
1768        capability: StorageCapability,
1769        operation: &'static str,
1770    ) -> Result<ReaderAdmission, StorageError> {
1771        let context = khive_storage::capture_request_read_context();
1772        let started = Instant::now();
1773        let stopped = context.stop_reason().is_some();
1774        let outcome: Result<Option<ReaderAdmission>, SqliteError> = if stopped {
1775            Ok(None)
1776        } else {
1777            tokio::select! {
1778                biased;
1779                _ = context.wait_for_stop() => Ok(None),
1780                waited = tokio::time::timeout(
1781                    self.config.checkout_timeout,
1782                    Arc::clone(&self.sql_bridge_reader_slots).acquire_owned(),
1783                ) => match waited {
1784                    Ok(Ok(slot)) => Ok(Some(ReaderAdmission { slot, started })),
1785                    Ok(Err(_closed)) => Err(SqliteError::InvalidData(
1786                        "reader admission semaphore is closed".to_string(),
1787                    )),
1788                    Err(_elapsed) => {
1789                        self.reader_acquisition_counters.record_checkout_timeout();
1790                        Err(pool_exhausted_error(
1791                            self.config.checkout_timeout,
1792                            self.max_readers,
1793                        ))
1794                    }
1795                },
1796            }
1797        };
1798        match outcome {
1799            Ok(Some(admission)) => Ok(admission),
1800            Ok(None) => Err(self.reader_checkout_refusal(capability, operation, None)),
1801            Err(error) => Err(self.reader_checkout_refusal(capability, operation, Some(error))),
1802        }
1803    }
1804
1805    /// Select the reader connection for an admission permit already held.
1806    ///
1807    /// This is [`Self::reader_until`] after its permit loop: the same
1808    /// `should_stop` polling, the same `Ok(Some)` / `Ok(None)` / `Err`
1809    /// outcomes, and the remainder of the one `checkout_timeout` that began
1810    /// when the permit wait did.
1811    pub(crate) fn reader_with_admission<C>(
1812        &self,
1813        admission: ReaderAdmission,
1814        should_stop: C,
1815    ) -> Result<Option<ReaderGuard<'_>>, SqliteError>
1816    where
1817        C: Fn() -> bool,
1818    {
1819        let ReaderAdmission {
1820            slot: admission_slot,
1821            started,
1822        } = admission;
1823
1824        if self.max_readers == 0 {
1825            self.ensure_pooled_writer_active()?;
1826            loop {
1827                if should_stop() {
1828                    return Ok(None);
1829                }
1830                let remaining = self
1831                    .config
1832                    .checkout_timeout
1833                    .saturating_sub(started.elapsed());
1834                if remaining.is_zero() {
1835                    self.reader_acquisition_counters.record_checkout_timeout();
1836                    return Err(pool_exhausted_error(
1837                        self.config.checkout_timeout,
1838                        self.max_readers,
1839                    ));
1840                }
1841                if let Some(guard) = self
1842                    .writer
1843                    .try_lock_for(remaining.min(Duration::from_millis(2)))
1844                {
1845                    self.ensure_pooled_writer_active()?;
1846                    self.reader_acquisition_counters.record_pooled_checkout();
1847                    return Ok(Some(ReaderGuard {
1848                        lease: Some(ReaderLease::Shared(guard)),
1849                        admission_slot: Some(admission_slot),
1850                        pool: self,
1851                        reusable: Cell::new(true),
1852                        query_in_progress: Cell::new(false),
1853                        checked_out_at: Instant::now(),
1854                        dirty: Cell::new(false),
1855                        operation: None,
1856                    }));
1857                }
1858            }
1859        }
1860
1861        let mut attempt = 0u32;
1862
1863        loop {
1864            if should_stop() {
1865                return Ok(None);
1866            }
1867            if let Some(conn) = self.readers.pop() {
1868                self.reader_acquisition_counters.record_pooled_checkout();
1869                return Ok(Some(ReaderGuard {
1870                    lease: Some(ReaderLease::Pooled(conn)),
1871                    admission_slot: Some(admission_slot),
1872                    pool: self,
1873                    reusable: Cell::new(true),
1874                    query_in_progress: Cell::new(false),
1875                    checked_out_at: Instant::now(),
1876                    dirty: Cell::new(false),
1877                    operation: None,
1878                }));
1879            }
1880
1881            if started.elapsed() >= self.config.checkout_timeout {
1882                self.reader_acquisition_counters.record_checkout_timeout();
1883                return Err(pool_exhausted_error(
1884                    self.config.checkout_timeout,
1885                    self.max_readers,
1886                ));
1887            }
1888
1889            match attempt {
1890                0..=7 => {
1891                    let spins = 1usize << attempt;
1892                    for _ in 0..spins {
1893                        std::hint::spin_loop();
1894                    }
1895                }
1896                8..=15 => thread::yield_now(),
1897                _ => {
1898                    let remaining = self
1899                        .config
1900                        .checkout_timeout
1901                        .saturating_sub(started.elapsed());
1902                    let sleep = Duration::from_micros(50 * (1u64 << (attempt - 16).min(6)));
1903                    thread::sleep(sleep.min(remaining).min(Duration::from_millis(2)));
1904                }
1905            }
1906
1907            attempt = attempt.saturating_add(1);
1908        }
1909    }
1910
1911    /// Check out the writer connection.
1912    ///
1913    /// Waits up to `checkout_timeout` for the writer Mutex and returns
1914    /// `Err(SqliteError::WriterPoolCheckoutTimeout)` if the timeout is
1915    /// exceeded.
1916    #[track_caller]
1917    pub fn writer(&self) -> Result<WriterGuard<'_>, SqliteError> {
1918        self.writer_with_checkout_probe(true, false)
1919    }
1920
1921    /// Callers that own an explicit §4 logical-write boundary probe there,
1922    /// so checkout must not refuse ahead of their BEGIN/bootstrap statement.
1923    #[track_caller]
1924    pub(crate) fn writer_for_admitted_operation(&self) -> Result<WriterGuard<'_>, SqliteError> {
1925        self.writer_with_checkout_probe(false, false)
1926    }
1927
1928    /// Checkpoint ownership setup changes only SQLite's checkpoint policy.
1929    /// Its bounded checkout must remain available below the write reserve.
1930    fn writer_for_checkpoint_operation(&self) -> Result<WriterGuard<'_>, SqliteError> {
1931        self.writer_with_checkout_probe(false, true)
1932    }
1933
1934    #[track_caller]
1935    fn writer_with_checkout_probe(
1936        &self,
1937        compatibility_probe: bool,
1938        checkpoint_bypass: bool,
1939    ) -> Result<WriterGuard<'_>, SqliteError> {
1940        if !checkpoint_bypass {
1941            self.write_admission.ensure_settled()?;
1942        }
1943        self.ensure_pooled_writer_active()?;
1944        let volume_lease = if checkpoint_bypass {
1945            None
1946        } else {
1947            self.write_admission.acquire()?
1948        };
1949        let Some(guard) = self.writer.try_lock_for(self.config.checkout_timeout) else {
1950            self.writer_acquisition_counters
1951                .pooled_timeouts
1952                .fetch_add(1, Ordering::Relaxed);
1953            let message = format!(
1954                "timed out after {:?} waiting for sqlite writer connection",
1955                self.config.checkout_timeout
1956            );
1957            crate::timeout_sink::emit_timeout(
1958                &crate::timeout_sink::db_label(self),
1959                crate::timeout_sink::Site::PoolAdmission,
1960                &message,
1961                Some(
1962                    self.config
1963                        .checkout_timeout
1964                        .as_millis()
1965                        .min(u128::from(u64::MAX)) as u64,
1966                ),
1967            );
1968            return Err(SqliteError::WriterPoolCheckoutTimeout {
1969                timeout: self.config.checkout_timeout,
1970            });
1971        };
1972        self.ensure_pooled_writer_active()?;
1973        #[cfg(any(unix, windows))]
1974        if let Some(path) = self.identity_path.as_deref() {
1975            self.verify_opened_file_identity(path)?;
1976            self.verify_connection_file_identity(&guard, path)?;
1977        }
1978        // Public raw-Connection/Deref compatibility callers may issue an
1979        // autocommit write without another seam. Keep the old coarse sample
1980        // until those sites migrate; explicit admitted operations skip it.
1981        if compatibility_probe && !checkpoint_bypass {
1982            self.write_admission.check()?;
1983        }
1984        self.writer_acquisition_counters
1985            .pooled_acquisitions
1986            .fetch_add(1, Ordering::Relaxed);
1987        Ok(WriterGuard {
1988            guard,
1989            origin: self.origin(),
1990            pool: self,
1991            admission: self.write_admission.as_ref(),
1992            _volume_lease: volume_lease,
1993        })
1994    }
1995
1996    /// Non-panicking writer checkout.
1997    ///
1998    /// Returns `Err` on timeout instead of panicking. Use this in request
1999    /// handlers where a 500 is preferable to crashing the process.
2000    #[track_caller]
2001    pub fn try_writer(&self) -> Result<WriterGuard<'_>, SqliteError> {
2002        self.writer()
2003    }
2004
2005    // Retain the compatibility-probe path only for its focused regression
2006    // fixtures; production callers use the execution-boundary admission seam.
2007    #[cfg(test)]
2008    #[track_caller]
2009    fn writer_until<C>(&self, should_stop: C) -> Result<Option<WriterGuard<'_>>, SqliteError>
2010    where
2011        C: Fn() -> bool,
2012    {
2013        self.writer_until_with_checkout_probe(should_stop, true)
2014    }
2015
2016    #[track_caller]
2017    pub(crate) fn writer_until_for_admitted_operation<C>(
2018        &self,
2019        should_stop: C,
2020    ) -> Result<Option<WriterGuard<'_>>, SqliteError>
2021    where
2022        C: Fn() -> bool,
2023    {
2024        self.writer_until_with_checkout_probe(should_stop, false)
2025    }
2026
2027    #[track_caller]
2028    fn writer_until_with_checkout_probe<C>(
2029        &self,
2030        should_stop: C,
2031        compatibility_probe: bool,
2032    ) -> Result<Option<WriterGuard<'_>>, SqliteError>
2033    where
2034        C: Fn() -> bool,
2035    {
2036        self.ensure_pooled_writer_active()?;
2037        if should_stop() {
2038            return Ok(None);
2039        }
2040        let volume_lease = self.write_admission.acquire()?;
2041        let started = Instant::now();
2042        loop {
2043            if should_stop() {
2044                return Ok(None);
2045            }
2046            let remaining = self
2047                .config
2048                .checkout_timeout
2049                .saturating_sub(started.elapsed());
2050            if let Some(guard) = self
2051                .writer
2052                .try_lock_for(remaining.min(Duration::from_millis(2)))
2053            {
2054                // Cancellation may have arrived during the final wait slice.
2055                // Once this guard is returned, constructor DDL is not interrupted.
2056                if should_stop() {
2057                    return Ok(None);
2058                }
2059                self.ensure_pooled_writer_active()?;
2060                #[cfg(any(unix, windows))]
2061                if let Some(path) = self.identity_path.as_deref() {
2062                    self.verify_opened_file_identity(path)?;
2063                    self.verify_connection_file_identity(&guard, path)?;
2064                }
2065                if compatibility_probe {
2066                    self.write_admission.check()?;
2067                }
2068                self.writer_acquisition_counters
2069                    .pooled_acquisitions
2070                    .fetch_add(1, Ordering::Relaxed);
2071                return Ok(Some(WriterGuard {
2072                    guard,
2073                    origin: self.origin(),
2074                    pool: self,
2075                    admission: self.write_admission.as_ref(),
2076                    _volume_lease: volume_lease,
2077                }));
2078            }
2079            if started.elapsed() >= self.config.checkout_timeout {
2080                self.writer_acquisition_counters
2081                    .pooled_timeouts
2082                    .fetch_add(1, Ordering::Relaxed);
2083                let message = format!(
2084                    "timed out after {:?} waiting for sqlite writer connection",
2085                    self.config.checkout_timeout
2086                );
2087                crate::timeout_sink::emit_timeout(
2088                    &crate::timeout_sink::db_label(self),
2089                    crate::timeout_sink::Site::PoolAdmission,
2090                    &message,
2091                    Some(
2092                        self.config
2093                            .checkout_timeout
2094                            .as_millis()
2095                            .min(u128::from(u64::MAX)) as u64,
2096                    ),
2097                );
2098                return Err(SqliteError::WriterPoolCheckoutTimeout {
2099                    timeout: self.config.checkout_timeout,
2100                });
2101            }
2102        }
2103    }
2104
2105    /// Zero-wait checkpoint checkout for recovery maintenance.
2106    ///
2107    /// Uses `try_lock()` (no timeout, no spin) — returns `Err` immediately when
2108    /// any other caller holds the writer Mutex. The scheduled checkpoint task
2109    /// uses its own dedicated connection (ADR-091 Amendment 5); this optional
2110    /// pooled capability remains available for zero-wait recovery callers.
2111    ///
2112    /// It bypasses the disk-reserve floor because checkpoints can recover WAL
2113    /// space at or below that floor (ADR-154 §5). The returned guard exposes
2114    /// only fixed PASSIVE and TRUNCATE checkpoint operations.
2115    ///
2116    /// The former unrestricted checkout must not return:
2117    /// ```compile_fail
2118    /// use khive_db::{ConnectionPool, PoolConfig};
2119    /// let pool = ConnectionPool::new(PoolConfig::default()).unwrap();
2120    /// pool.try_writer_nowait().unwrap().execute_batch("CREATE TABLE bypass (id INTEGER)");
2121    /// ```
2122    pub fn try_checkpoint_nowait(&self) -> Result<CheckpointGuard<'_>, SqliteError> {
2123        self.ensure_pooled_writer_active()?;
2124        let guard = self.writer.try_lock().ok_or_else(|| {
2125            SqliteError::InvalidData(
2126                "writer connection busy (checkpoint skipped this tick)".to_string(),
2127            )
2128        })?;
2129        self.ensure_pooled_writer_active()?;
2130        Ok(CheckpointGuard { guard })
2131    }
2132
2133    pub(crate) fn retire_pooled_writer(&self, conn: &Connection) {
2134        self.pooled_writer_retired.store(true, Ordering::Release);
2135        if let Err(error) = conn.authorizer(Some(deny_retired_writer)) {
2136            tracing::error!(
2137                %error,
2138                "failed to install the retired pooled-writer quarantine authorizer"
2139            );
2140        }
2141    }
2142
2143    /// Exercise the retired connection's authorizer without exporting a raw
2144    /// handle from the pool. This is intentionally unavailable in production.
2145    #[cfg(test)]
2146    pub(crate) fn probe_retired_pooled_writer_for_test(&self) -> rusqlite::Result<i64> {
2147        self.writer
2148            .lock()
2149            .query_row("SELECT 1", [], |row| row.get(0))
2150    }
2151
2152    /// Leave an inherited transaction on the bare pooled connection to test
2153    /// the typed write-unit refusal. A normal WriterGuard rolls it back on
2154    /// drop, so it cannot create this poisoned state.
2155    #[cfg(test)]
2156    pub(crate) fn leave_pooled_writer_transaction_open_for_test(&self) -> rusqlite::Result<()> {
2157        self.writer.lock().execute_batch("BEGIN IMMEDIATE")
2158    }
2159
2160    fn ensure_pooled_writer_active(&self) -> Result<(), SqliteError> {
2161        if self.pooled_writer_retired.load(Ordering::Acquire) {
2162            return Err(SqliteError::InvalidData(
2163                "pooled writer connection retired after a terminal transaction fault".to_string(),
2164            ));
2165        }
2166        Ok(())
2167    }
2168
2169    /// Snapshot all instrumented writer acquisition outcomes since this pool
2170    /// was constructed.
2171    pub fn writer_acquisition_snapshot(&self) -> WriterAcquisitionSnapshot {
2172        self.writer_acquisition_counters
2173            .snapshot(self.write_admission.lease_timeouts())
2174    }
2175
2176    /// Count a coordinator call only after served-kind filtering selects this
2177    /// backend. The canonical requested kind is the granular kind when one was
2178    /// supplied, or the entity/note substrate otherwise.
2179    pub fn record_search_dispatch(&self, backend_id: &str, requested_kind: &str) {
2180        let mut dispatches = self.search_dispatches.lock();
2181        let count = dispatches
2182            .entry(backend_id.to_owned())
2183            .or_default()
2184            .entry(requested_kind.to_owned())
2185            .or_default();
2186        *count = count.saturating_add(1);
2187    }
2188
2189    /// Count one actual candidate note row returned at the post-fusion
2190    /// hydration seam, including rows later filtered as deleted. Absent rows
2191    /// and later result metadata are excluded.
2192    pub fn record_note_candidate_hydration_row(&self) {
2193        let _ = self.note_candidate_hydration_rows.fetch_update(
2194            Ordering::Relaxed,
2195            Ordering::Relaxed,
2196            |current| Some(current.saturating_add(1)),
2197        );
2198    }
2199
2200    /// Snapshot ADR-166 G4/G5 counters. Values reset only with this pool.
2201    pub fn search_mechanism_snapshot(&self) -> SearchMechanismSnapshot {
2202        SearchMechanismSnapshot {
2203            dispatches_by_backend_and_kind: self.search_dispatches.lock().clone(),
2204            note_candidate_hydration_rows: self
2205                .note_candidate_hydration_rows
2206                .load(Ordering::Relaxed),
2207        }
2208    }
2209
2210    /// Snapshot reader acquisition, saturation, and hold lifecycle outcomes
2211    /// since this pool was constructed. Counters reset only with pool
2212    /// reconstruction.
2213    pub fn reader_acquisition_snapshot(&self) -> ReaderAcquisitionSnapshot {
2214        self.reader_acquisition_counters.snapshot(
2215            self.max_readers.max(1),
2216            self.sql_bridge_reader_slots.available_permits(),
2217        )
2218    }
2219
2220    /// Record one pool-wide reader-admission wait that exhausted the configured
2221    /// checkout timeout outside [`Self::reader_until`] (currently the explicit
2222    /// raw-SQL read-transaction exception and reads on a standalone writer).
2223    pub(crate) fn record_reader_admission_timeout(&self) {
2224        self.reader_acquisition_counters.record_checkout_timeout();
2225    }
2226
2227    /// Count a query error from an already checked-out pooled reader when it is
2228    /// SQLite's `SQLITE_BUSY` surfacing after the busy handler gave up. Every
2229    /// other error, including `SQLITE_LOCKED`, is ignored. Checkout exhaustion
2230    /// is counted by [`Self::reader_until`] before any query runs, so the two
2231    /// classes cannot overlap.
2232    pub(crate) fn record_reader_query_error(&self, error: &StorageError) {
2233        if crate::read_cancellation::storage_error_sqlite_code(error)
2234            == Some(rusqlite::ErrorCode::DatabaseBusy)
2235        {
2236            self.reader_acquisition_counters.record_busy_timeout();
2237        }
2238    }
2239
2240    /// Clone the pool-scoped counter set for the lifetime-owned writer task.
2241    pub(crate) fn writer_acquisition_counters(&self) -> Arc<WriterAcquisitionCounters> {
2242        Arc::clone(&self.writer_acquisition_counters)
2243    }
2244
2245    pub(crate) fn write_admission(&self) -> Arc<WriteAdmission> {
2246        Arc::clone(&self.write_admission)
2247    }
2248
2249    pub fn effective_disk_guard_config(&self) -> Option<EffectiveDiskGuardConfig> {
2250        self.disk_guard_config
2251    }
2252
2253    #[cfg(test)]
2254    pub(crate) fn set_test_write_admission(
2255        &mut self,
2256        floor_bytes: u64,
2257        probe: impl Fn(&Path) -> std::io::Result<u64> + Send + Sync + 'static,
2258    ) {
2259        let database_path = if self.config.read_only {
2260            None
2261        } else {
2262            self.canonical_path().map(Path::to_path_buf)
2263        };
2264        let admission = Arc::new(
2265            WriteAdmission::new(
2266                database_path,
2267                floor_bytes,
2268                DEFAULT_DISK_GUARD_DEADLINE_MS,
2269                self.config.volume_lock_dir.clone(),
2270            )
2271            .expect("test admission volume identity"),
2272        );
2273        admission.set_test_space_probe(probe);
2274        self.write_admission = admission;
2275        if self.canonical_path().is_some() && !self.config.read_only {
2276            self.disk_guard_config = Some(EffectiveDiskGuardConfig {
2277                reserve_bytes: floor_bytes,
2278                guard_deadline_ms: DEFAULT_DISK_GUARD_DEADLINE_MS,
2279                reserve_source: DiskGuardConfigSource::Backend,
2280                deadline_source: DiskGuardConfigSource::Default,
2281                legacy_environment_present: false,
2282            });
2283        }
2284    }
2285
2286    /// Get the current number of available reader connections.
2287    pub fn available_readers(&self) -> usize {
2288        self.readers.len()
2289    }
2290
2291    /// Get the total number of reader connections in the pool.
2292    pub fn max_readers(&self) -> usize {
2293        self.max_readers
2294    }
2295
2296    /// Return the pool configuration.
2297    pub fn config(&self) -> &PoolConfig {
2298        &self.config
2299    }
2300
2301    /// Observe actual SQLite statement starts on this private test pool.
2302    ///
2303    /// The limit bounds retained SQL records. Failed steps count as attempts;
2304    /// preparation alone does not count. The guard observes every pool-owned
2305    /// connection, including the queued writer. Do not run unrelated background
2306    /// work on the fixture pool; see the guard documentation for limitations.
2307    /// Connection setup runs before observation begins on each connection and
2308    /// is never recorded, including for opens during an active observation.
2309    /// Reader connection opens use the existing reader acquisition counters instead.
2310    #[cfg(any(test, feature = "test-support"))]
2311    pub fn observe_test_statement_starts(
2312        &self,
2313        limit: usize,
2314    ) -> Result<crate::statement_observer::StatementStartObservation, SqliteError> {
2315        self.statement_observer.observe(limit)
2316    }
2317
2318    /// Identify this pool's counter window when it is designated as main.
2319    /// Repeated runtime handles and diagnostics reads reuse the same generation;
2320    /// constructing secondary pools does not consume main-pool generations.
2321    pub fn main_pool_generation(&self) -> u64 {
2322        *self.main_pool_generation.get_or_init(|| {
2323            NEXT_MAIN_POOL_GENERATION
2324                .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |next| {
2325                    next.checked_add(1)
2326                })
2327                .expect("main pool generation exhausted")
2328        })
2329    }
2330
2331    /// The typed admission failure for a pooled reader checkout that
2332    /// exhausted `checkout_timeout`: no reader was acquired, so the
2333    /// operation never started and a retry cannot duplicate a side effect.
2334    pub(crate) fn reader_admission_timeout(&self, operation: &'static str) -> StorageError {
2335        StorageError::AdmissionTimeout {
2336            operation: operation.into(),
2337            timeout_ms: u64::try_from(self.config.checkout_timeout.as_millis()).unwrap_or(u64::MAX),
2338            pool_identity: Some(
2339                self.identity_registration
2340                    .as_ref()
2341                    .map(PoolIdentityRegistration::label)
2342                    .unwrap_or_else(|| ":memory:".to_string()),
2343            ),
2344        }
2345    }
2346
2347    /// Resolve a [`Self::reader_until`] outcome into a checked-out guard or
2348    /// the canonical refusal. This is the single home of the checkout
2349    /// tri-state; call sites must not re-derive any arm of it:
2350    ///
2351    /// - `Ok(Some)` — a reader was checked out.
2352    /// - `Ok(None)` — `should_stop()` fired: the request was cancelled or hit
2353    ///   its deadline before checkout. NOT an admission wait, so it maps to
2354    ///   the non-retryable [`StorageError::Timeout`] — emitting the retryable
2355    ///   `AdmissionTimeout` here would invite an immediate retry of a request
2356    ///   its caller already abandoned, into a possibly saturated pool.
2357    /// - `Err` carrying the pool's own `SQLITE_BUSY` — `reader_until` executes
2358    ///   no SQL, so the only `SQLITE_BUSY` it can produce is
2359    ///   [`pool_exhausted_error`], raised when `checkout_timeout` elapses with
2360    ///   no reader available. A genuine admission wait that ended before any
2361    ///   work began maps to the retryable [`StorageError::AdmissionTimeout`].
2362    /// - any other `Err` (e.g. [`Self::ensure_pooled_writer_active`] returning
2363    ///   `InvalidData` for a retired pooled writer) — an opaque driver failure
2364    ///   under the caller's capability, non-retryable.
2365    pub(crate) fn resolve_reader_checkout<'p>(
2366        &self,
2367        capability: StorageCapability,
2368        operation: &'static str,
2369        outcome: Result<Option<ReaderGuard<'p>>, SqliteError>,
2370    ) -> Result<ReaderGuard<'p>, StorageError> {
2371        match outcome {
2372            Ok(Some(mut guard)) => {
2373                guard.label_operation(operation);
2374                Ok(guard)
2375            }
2376            Ok(None) => Err(self.reader_checkout_refusal(capability, operation, None)),
2377            Err(error) => Err(self.reader_checkout_refusal(capability, operation, Some(error))),
2378        }
2379    }
2380
2381    /// The refusal for a reader checkout that produced nothing, shared by
2382    /// [`Self::resolve_reader_checkout`] and [`Self::acquire_reader_admission`]
2383    /// so the arms documented on the former have one home. `None` is the
2384    /// stopped-request arm; `Some` carries the failed checkout's error.
2385    fn reader_checkout_refusal(
2386        &self,
2387        capability: StorageCapability,
2388        operation: &'static str,
2389        error: Option<SqliteError>,
2390    ) -> StorageError {
2391        let Some(error) = error else {
2392            return StorageError::Timeout {
2393                operation: operation.into(),
2394            };
2395        };
2396        let is_pool_exhausted = matches!(
2397            &error,
2398            SqliteError::Rusqlite(rusqlite::Error::SqliteFailure(code, _))
2399                if code.code == rusqlite::ErrorCode::DatabaseBusy
2400        );
2401        if is_pool_exhausted {
2402            self.reader_admission_timeout(operation)
2403        } else {
2404            StorageError::driver(capability, operation, error)
2405        }
2406    }
2407
2408    /// Pool-wide admission permits shared by pooled readers, the explicit
2409    /// raw-SQL read-transaction exception, and reads on standalone writers.
2410    pub(crate) fn sql_bridge_reader_slots(&self) -> Arc<Semaphore> {
2411        Arc::clone(&self.sql_bridge_reader_slots)
2412    }
2413
2414    /// Pool-wide permit for a file-backed raw-SQL writer handle.
2415    pub(crate) fn sql_bridge_writer_slots(&self) -> Arc<Semaphore> {
2416        Arc::clone(&self.sql_bridge_writer_slots)
2417    }
2418
2419    /// Current writer holds that prevent voluntary daemon retirement.
2420    ///
2421    /// The raw-SQL writer permit remains handle-scoped even in autocommit.
2422    /// This read acquires no connection and changes no admission policy.
2423    pub fn retirement_writer_holds(&self) -> usize {
2424        usize::from(self.writer.is_locked())
2425            + usize::from(self.sql_bridge_writer_slots.available_permits() == 0)
2426    }
2427
2428    /// This pool's ADR-091 backend-scoped attribution origin (ADR-091,
2429    /// backend-scoped WAL-pin attribution design note): `Database(_)` for a
2430    /// file-backed pool, `Memory` for an in-memory pool. Every
2431    /// `tx_registry::register_scoped` call site threaded in this crate
2432    /// passes this value as the span's origin.
2433    pub fn origin(&self) -> TxOrigin {
2434        self.origin.clone()
2435    }
2436
2437    /// The canonical path this pool's `origin()` identity was minted from,
2438    /// `None` for an in-memory pool. `DbIdentity` has no path accessor by
2439    /// design; sidecar derivation and other filesystem consumers use this —
2440    /// the same canonical value the identity was minted from — instead of
2441    /// re-deriving a path from the raw configured one.
2442    pub fn canonical_path(&self) -> Option<&Path> {
2443        self.identity_path.as_deref()
2444    }
2445
2446    /// Unix file identity retained for callers using the original tuple API.
2447    #[cfg(unix)]
2448    pub fn opened_file_identity(&self) -> Option<(u64, u64)> {
2449        self.opened_file_identity
2450            .map(DatabaseFileIdentity::unix_parts)
2451    }
2452
2453    /// Physical file identity pinned to SQLite's opened main file. Construction
2454    /// compares it with the path on both sides of the open before admission.
2455    #[cfg(any(unix, windows))]
2456    pub fn opened_file_identity_record(&self) -> Option<DatabaseFileIdentity> {
2457        self.opened_file_identity
2458    }
2459
2460    /// Ownership evidence captured through the database SQLite opened.
2461    /// This does not select the topology's main backend or establish a root binding.
2462    /// Legacy read-only and low-space opens without a stored UUID must reopen
2463    /// after installation; this accessor never installs or re-mints an identity.
2464    pub fn database_owner_identity(
2465        &self,
2466    ) -> Result<DatabaseOwnerIdentity, DatabaseOwnerIdentityError> {
2467        if self.identity_path.is_none() {
2468            return Err(DatabaseOwnerIdentityError::InMemory);
2469        }
2470        #[cfg(any(unix, windows))]
2471        {
2472            let durable_id = self
2473                .opened_database_id
2474                .ok_or(DatabaseOwnerIdentityError::DurableIdentityUnavailable)?;
2475            let file_identity = self
2476                .opened_file_identity
2477                .ok_or(DatabaseOwnerIdentityError::PhysicalIdentityUnavailable)?;
2478            Ok(DatabaseOwnerIdentity {
2479                durable_id,
2480                file_identity,
2481            })
2482        }
2483        #[cfg(not(any(unix, windows)))]
2484        {
2485            Err(DatabaseOwnerIdentityError::UnsupportedPlatform)
2486        }
2487    }
2488
2489    /// Require both this opened database's UUID and physical file to match.
2490    pub fn verify_database_owner(
2491        &self,
2492        expected: &DatabaseOwnerIdentity,
2493    ) -> Result<(), DatabaseOwnerIdentityError> {
2494        self.database_owner_identity()?.verify_owner(expected)
2495    }
2496
2497    /// Whether the write queue is effectively enabled for this pool: the
2498    /// resolved `write_queue_enabled` flag AND file-backed.
2499    ///
2500    /// `ConnectionPool::new` resolves the "no preference" (`None`) preference
2501    /// to a concrete `Some(..)` once `path` is known, so every reader of
2502    /// `config.write_queue_enabled` sees a resolved value; the `debug_assert`
2503    /// pins that invariant and a `None` that slipped past would read as
2504    /// disabled. Bypassing `ConnectionPool::new` to construct a pool is a
2505    /// construction-path bug. Use this instead of repeating
2506    /// `config().write_queue_enabled.unwrap_or(false) && config().path.is_some()`
2507    /// at every routing/violation site.
2508    pub fn write_queue_active(&self) -> bool {
2509        debug_assert!(
2510            self.config.write_queue_enabled.is_some(),
2511            "write_queue_enabled must be resolved to Some(..) by ConnectionPool::new \
2512             before any write_queue_active read"
2513        );
2514        self.config.write_queue_enabled.unwrap_or(false) && self.config.path.is_some()
2515    }
2516
2517    /// Whether a writer-task JoinHandle has been stored at least once.
2518    ///
2519    /// Unlike [`Self::take_writer_task_join`], this remains true after the
2520    /// one-shot handle slot is emptied, distinguishing a task that never
2521    /// spawned from a handle another caller already consumed.
2522    pub fn writer_task_join_was_stored(&self) -> bool {
2523        self.writer_task_join_stored.load(Ordering::SeqCst)
2524    }
2525
2526    /// Return the pool-wide ADR-067 Component A writer task, spawning it
2527    /// lazily on first access if `PoolConfig::write_queue_enabled` is set.
2528    /// Exactly one writer task exists per `ConnectionPool` (per DB file); see
2529    /// crates/khive-db/docs/api/pool.md#connectionpoolwriter_task_handle--single-writer-task-rationale
2530    /// for why a per-store writer task would defeat the single-writer
2531    /// guarantee.
2532    ///
2533    /// Returns `Ok(None)` if the flag is off, or if the writer task failed to
2534    /// spawn for a reason other than a missing runtime (for example, an
2535    /// in-memory pool has no standalone-connection support) — callers fall
2536    /// back to the legacy pool-mutex write path in either case. A spawn
2537    /// failure is logged once here (at first access), not once per store.
2538    ///
2539    /// Returns `Err(StorageError::WriterTaskNoRuntime)` instead of panicking
2540    /// when `write_queue_enabled` is set but this is the first access and no
2541    /// Tokio runtime is available on the calling thread (checked via
2542    /// [`tokio::runtime::Handle::try_current`]) — spawning the writer task
2543    /// requires `tokio::spawn`, which panics outside a runtime. Callers that
2544    /// already treat a missing writer task as best-effort (construction-time
2545    /// degrade to the legacy path, matching slice 1's documented policy) can
2546    /// collapse this into `None` with `.ok().flatten()`; callers that need to
2547    /// fail loud on a genuine misconfiguration (write queue requested but no
2548    /// runtime to run it on) can propagate the `Err` directly.
2549    pub fn writer_task_handle(&self) -> Result<Option<WriterTaskHandle>, StorageError> {
2550        // Same pinned invariant `write_queue_active` asserts, kept inline
2551        // here because this gate keys on the flag ALONE: an explicit
2552        // `Some(true)` on an in-memory pool must still attempt the spawn
2553        // and degrade (documented + tested in
2554        // `explicit_true_stays_on_for_memory_backed_pool`), so the
2555        // file-backed half of `write_queue_active` cannot gate this early
2556        // return.
2557        debug_assert!(
2558            self.config.write_queue_enabled.is_some(),
2559            "write_queue_enabled must be resolved to Some(..) by ConnectionPool::new \
2560             before any writer_task_handle read"
2561        );
2562        if !self.config.write_queue_enabled.unwrap_or(false) {
2563            return Ok(None);
2564        }
2565        // Fast path: already resolved (spawned, degraded, or off) by an
2566        // earlier call — no need to re-check the runtime.
2567        if let Some(existing) = self.writer_task.get() {
2568            return Ok(existing.clone());
2569        }
2570        // Not yet initialized and the flag is on: spawning requires
2571        // `tokio::spawn`, which panics outside a runtime context. Check
2572        // first and fail loud with a typed error instead.
2573        if tokio::runtime::Handle::try_current().is_err() {
2574            return Err(StorageError::WriterTaskNoRuntime);
2575        }
2576        Ok(self
2577            .writer_task
2578            .get_or_init(|| {
2579                #[cfg(test)]
2580                self.writer_task_spawn_count
2581                    .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
2582
2583                match crate::writer_task::spawn(self, self.config.write_queue_capacity) {
2584                    Ok(handle) => Some(handle),
2585                    Err(e) => {
2586                        tracing::warn!(
2587                            error = %e,
2588                            "KHIVE_WRITE_QUEUE=1 but the writer task failed to spawn; \
2589                             writes fall back to the pool-mutex path"
2590                        );
2591                        None
2592                    }
2593                }
2594            })
2595            .clone())
2596    }
2597
2598    /// Resolve the writer task for a store write at the moment the write is
2599    /// issued, rather than trusting only a handle cached by a synchronous
2600    /// store constructor. Construction can legitimately run before Tokio is
2601    /// entered, in which case `writer_task_handle()` returns
2602    /// `WriterTaskNoRuntime` without caching a terminal `None`.
2603    ///
2604    /// Strict routing makes every missing handle fail closed here. The
2605    /// caller remains responsible for recording a non-strict direct fallback
2606    /// at the exact fallback seam with [`Self::record_direct_route`].
2607    pub(crate) fn writer_task_for_write(
2608        &self,
2609        cached: Option<&WriterTaskHandle>,
2610        operation: &'static str,
2611    ) -> Result<Option<WriterTaskHandle>, StorageError> {
2612        let handle = match cached {
2613            Some(handle) => Some(handle.clone()),
2614            None => match self.writer_task_handle() {
2615                Ok(handle) => handle,
2616                Err(error) if self.config.write_routing_strict => return Err(error),
2617                Err(_) => None,
2618            },
2619        };
2620
2621        if handle.is_none() && self.config.write_routing_strict {
2622            return Err(StorageError::Pool {
2623                operation: operation.into(),
2624                message: "strict write routing requires a writer-task handle; no handle is \
2625                          available, so the direct writer fallback was refused"
2626                    .into(),
2627            });
2628        }
2629        Ok(handle)
2630    }
2631
2632    /// Record one actual compatibility fallback around the writer task. A
2633    /// file-backed pool with the queue enabled should never reach this seam
2634    /// in strict mode because [`Self::writer_task_for_write`] refuses first.
2635    pub(crate) fn record_direct_route(&self, site: crate::timeout_sink::Site) {
2636        if self.write_queue_active() {
2637            crate::timeout_sink::emit_direct_route_violation(
2638                &crate::timeout_sink::db_label(self),
2639                site,
2640            );
2641        }
2642    }
2643
2644    /// Resolve a runtime-owned transaction through the same strict/compatibility
2645    /// policy as store writes. `None` permits the caller's direct transaction and
2646    /// records its compatibility fallback when the file-backed queue is enabled.
2647    /// A strict refusal returns before any direct-writer acquisition or telemetry.
2648    pub fn writer_task_for_runtime_write(
2649        &self,
2650        operation: RuntimeWriteOperation,
2651    ) -> Result<Option<WriterTaskHandle>, StorageError> {
2652        let handle = self.writer_task_for_write(None, operation.operation())?;
2653        if handle.is_none() {
2654            self.record_direct_route(operation.fallback_site());
2655        }
2656        Ok(handle)
2657    }
2658
2659    /// Test-only: how many times the writer-task init closure actually ran.
2660    /// Must be at most 1 for the pool's whole lifetime, regardless of how
2661    /// many times [`Self::writer_task_handle`] is called or how many stores
2662    /// are constructed over this pool.
2663    #[cfg(test)]
2664    pub(crate) fn writer_task_spawn_count(&self) -> usize {
2665        self.writer_task_spawn_count
2666            .load(std::sync::atomic::Ordering::SeqCst)
2667    }
2668
2669    /// Record the writer task's `tokio::spawn` JoinHandle. Called exactly
2670    /// once, by [`crate::writer_task::spawn`], immediately after spawning —
2671    /// the same `writer_task` OnceLock init that makes spawn at-most-once
2672    /// per pool makes this write at-most-once per pool.
2673    ///
2674    /// First-wins: if a handle was ever stored (including one a caller has
2675    /// since taken — `writer_task_join_stored` remembers), the existing
2676    /// state is kept and the new handle is dropped (dropping a `JoinHandle`
2677    /// detaches its task without cancelling it). A second store violates the
2678    /// at-most-once contract and trips the debug_assert in debug builds;
2679    /// release builds keep the first handle rather than silently swapping
2680    /// the drain owner out from under whichever caller already took it.
2681    pub(crate) fn set_writer_task_join(&self, join: tokio::task::JoinHandle<()>) {
2682        // `swap(true)` returns the prior value: `true` means a handle was
2683        // stored at least once before, so this is a second store — even when
2684        // the slot itself is empty because `take_writer_task_join` already
2685        // ran (the slot alone cannot tell "never stored" from "taken").
2686        let first_store = !self.writer_task_join_stored.swap(true, Ordering::SeqCst);
2687        debug_assert!(
2688            first_store,
2689            "writer task JoinHandle stored twice (even counting a taken one); \
2690             the writer_task OnceLock is supposed to make spawn at-most-once per pool"
2691        );
2692        if first_store {
2693            *self.writer_task_join.lock() = Some(join);
2694        }
2695    }
2696
2697    /// Take the writer task's JoinHandle, if a writer task was spawned and
2698    /// the handle has not already been taken.
2699    ///
2700    /// Intended for short-lived batch callers that drop every
2701    /// [`WriterTaskHandle`] clone (closing the queue) and then need to await
2702    /// the task's exit before treating the database file as settled: the
2703    /// task's connection close fires SQLite's close-time WAL checkpoint, so
2704    /// until the task exits the file bytes can still move after the caller's
2705    /// last write returned.
2706    ///
2707    /// One-shot: `None` means either the write queue never spawned
2708    /// (disabled, or spawn degraded) or another caller already took the
2709    /// handle — in both cases there is nothing further to await here.
2710    /// Exactly one subsystem may own the drain: the single caller that
2711    /// receives `Some(_)` is the sole owner of the task-exit await (and of
2712    /// the close-time WAL checkpoint that settles the database file); every
2713    /// later caller receives `None` and must not arrange its own await.
2714    pub fn take_writer_task_join(&self) -> Option<tokio::task::JoinHandle<()>> {
2715        self.writer_task_join.lock().take()
2716    }
2717
2718    fn open_reader_connection(&self) -> Result<Connection, SqliteError> {
2719        let path = self.read_connection_path()?;
2720        #[cfg(any(unix, windows))]
2721        if let Some(identity_path) = self.identity_path.as_deref() {
2722            self.verify_opened_file_identity(identity_path)?;
2723        }
2724        let conn = open_reader_connection(path, &self.config)?;
2725        #[cfg(any(unix, windows))]
2726        if let Some(identity_path) = self.identity_path.as_deref() {
2727            self.verify_connection_file_identity(&conn, identity_path)?;
2728        }
2729        self.verify_opened_database_id(&conn)?;
2730        #[cfg(any(test, feature = "test-support"))]
2731        crate::statement_observer::install(&conn, &self.statement_observer)?;
2732        Ok(conn)
2733    }
2734
2735    fn read_connection_path(&self) -> Result<&Path, SqliteError> {
2736        self.read_only_open_target
2737            .as_deref()
2738            .or(self.identity_path.as_deref())
2739            .ok_or_else(|| {
2740                SqliteError::InvalidData(
2741                    "in-memory databases do not support standalone connections".to_string(),
2742                )
2743            })
2744    }
2745
2746    /// Open a standalone read-write connection to the same file-backed database.
2747    ///
2748    /// Stores whose trait methods take `Send + 'static` closures (executed via
2749    /// `spawn_blocking`) cannot hold the pooled `WriterGuard`'s `MutexGuard`
2750    /// across the call — it opens an independent connection instead. This
2751    /// must still honor `PoolConfig::read_only`: opening
2752    /// `SQLITE_OPEN_READ_WRITE` unconditionally here would let a read-only
2753    /// backend's graph/event/text stores bypass the flag that the pooled
2754    /// writer enforces via `query_only`. A fully configured successful open
2755    /// increments the standalone acquisition class exactly once.
2756    pub fn open_standalone_writer(&self) -> Result<Connection, SqliteError> {
2757        self.write_admission.check()?;
2758        self.open_standalone_writer_for_admitted_operation()
2759    }
2760
2761    /// Open for a caller that owns §4 admission at its transaction or
2762    /// autocommit execution boundary. Connection open alone cannot authorize
2763    /// the later write, especially when a handle is cached across requests.
2764    pub(crate) fn open_standalone_writer_for_admitted_operation(
2765        &self,
2766    ) -> Result<Connection, SqliteError> {
2767        let conn = self.open_standalone_writer_untracked()?;
2768        self.writer_acquisition_counters
2769            .standalone_acquisitions
2770            .fetch_add(1, Ordering::Relaxed);
2771        Ok(conn)
2772    }
2773
2774    /// Open an infrastructure-owned standalone writer connection without
2775    /// counting it as one write-operation acquisition.
2776    ///
2777    /// Restricted to the diagnostics PASSIVE probe, the writer task's
2778    /// one-time lifetime connection, and the checkpoint task's dedicated
2779    /// long-lived connection (opened once at startup and reused across
2780    /// ticks — see `CheckpointConnection::ensure_open`). Actual file-backed
2781    /// write paths must call [`Self::open_standalone_writer`] so their
2782    /// acquisitions are observable.
2783    pub(crate) fn open_standalone_writer_untracked(&self) -> Result<Connection, SqliteError> {
2784        let path = self.identity_path.as_deref().ok_or_else(|| {
2785            SqliteError::InvalidData(
2786                "in-memory databases do not support standalone connections".to_string(),
2787            )
2788        })?;
2789
2790        if self.config.read_only {
2791            return Err(SqliteError::InvalidData(
2792                "database is read-only: standalone write connections are not permitted".to_string(),
2793            ));
2794        }
2795
2796        // The configured spelling may be a symlink that changed since this
2797        // pool opened. Use its pinned target, and refuse replacement of that
2798        // target before SQLite can execute the diagnostics PASSIVE checkpoint.
2799        #[cfg(any(unix, windows))]
2800        self.verify_opened_file_identity(path)?;
2801
2802        #[cfg(test)]
2803        run_identity_open_hook(path, IdentityOpenStage::BeforeStandaloneOpen, None);
2804
2805        let conn = claimed_file_identity::open_connection(
2806            &self.config,
2807            path,
2808            OpenFlags::SQLITE_OPEN_READ_WRITE
2809                | OpenFlags::SQLITE_OPEN_NO_MUTEX
2810                | OpenFlags::SQLITE_OPEN_URI,
2811            self.identity_path.as_deref(),
2812        )?;
2813        #[cfg(test)]
2814        run_identity_open_hook(path, IdentityOpenStage::AfterStandaloneOpen, Some(&conn));
2815        #[cfg(any(unix, windows))]
2816        self.verify_connection_file_identity(&conn, path)?;
2817        self.verify_opened_database_id(&conn)?;
2818        #[cfg(feature = "namespace-trigram-proto")]
2819        register_namespace_trigram(&conn)?;
2820        register_writer_clock(&conn)?;
2821        // Expression indexes over these keys are maintained by every writer
2822        // (the write-queue task and per-store standalone writers included), so
2823        // each of them needs the same functions the pooled writer registers.
2824        register_rfc3339_key(&conn)?;
2825        conn.busy_timeout(self.config.busy_timeout)?;
2826        if self.config.code_map_vfs.is_none() {
2827            self.checkpoint_ownership
2828                .configure_wal_autocheckpoint(&conn)?;
2829        }
2830        conn.pragma_update(None, "foreign_keys", "ON")?;
2831        conn.pragma_update(None, "synchronous", "NORMAL")?;
2832
2833        let wal_enabled =
2834            self.config.wal_mode && current_journal_mode(&conn)?.eq_ignore_ascii_case("wal");
2835        if wal_enabled {
2836            conn.pragma_update(
2837                None,
2838                "journal_size_limit",
2839                self.config.journal_size_limit_bytes,
2840            )?;
2841        }
2842
2843        #[cfg(any(test, feature = "test-support"))]
2844        crate::statement_observer::install(&conn, &self.statement_observer)?;
2845        Ok(conn)
2846    }
2847
2848    #[cfg(any(unix, windows))]
2849    fn verify_opened_file_identity(&self, path: &Path) -> Result<(), SqliteError> {
2850        let Some(expected) = self.opened_file_identity else {
2851            return Err(SqliteError::InvalidData(
2852                "file-backed pool has no opened database file identity".to_string(),
2853            ));
2854        };
2855        let current = database_file_identity(path).ok();
2856        if current != Some(expected) {
2857            return Err(SqliteError::InvalidData(
2858                "pool database file identity changed since the first open; refusing standalone connection"
2859                    .to_string(),
2860            ));
2861        }
2862        Ok(())
2863    }
2864
2865    #[cfg(any(unix, windows))]
2866    fn verify_connection_file_identity(
2867        &self,
2868        conn: &Connection,
2869        path: &Path,
2870    ) -> Result<(), SqliteError> {
2871        let opened = opened_sqlite_file_identity(conn, path)?;
2872        if self.opened_file_identity != Some(opened) {
2873            return Err(SqliteError::InvalidData(
2874                "pool database file identity changed since the first open; refusing standalone connection"
2875                    .to_string(),
2876            ));
2877        }
2878        Ok(())
2879    }
2880
2881    fn verify_opened_database_id(&self, conn: &Connection) -> Result<(), SqliteError> {
2882        // A read-only pool, or a writable pool opened below the space reserve,
2883        // may precede installation of the nonce by another process. Without a
2884        // nonce to pin, use the existing file identity checks where available.
2885        let Some(expected) = self.opened_database_id else {
2886            return Ok(());
2887        };
2888        if self.identity_path.is_some() && read_database_id(conn)? != Some(expected) {
2889            return Err(SqliteError::InvalidData(
2890                "pool database identity changed since the first open; refusing standalone connection"
2891                    .to_string(),
2892            ));
2893        }
2894        Ok(())
2895    }
2896
2897    /// Effective `PRAGMA wal_autocheckpoint` for a writer-capable connection
2898    /// opened right now: `0` once a dedicated checkpoint owner has claimed
2899    /// the pool, the bounded fallback otherwise.
2900    #[cfg(test)]
2901    pub(crate) fn effective_wal_autocheckpoint_pages(&self) -> u32 {
2902        self.checkpoint_ownership.wal_autocheckpoint_pages()
2903    }
2904
2905    /// Claim routine WAL-checkpoint ownership for this pool.
2906    ///
2907    /// Called by the scheduled checkpoint task at startup — the one caller
2908    /// that actually replaces SQLite's per-commit autocheckpoint with
2909    /// dedicated PASSIVE checkpointing (ADR-091 Amendment 10). The claim
2910    /// makes every subsequently opened writer-capable connection set
2911    /// `PRAGMA wal_autocheckpoint = 0`, and re-applies that pragma on the
2912    /// already-open pooled writer under the writer mutex. A writer task
2913    /// spawned before the claim keeps its own long-lived connection;
2914    /// [`Self::propagate_checkpoint_claim_to_writer_task`] reaches that one.
2915    ///
2916    /// Without a claim, writer-capable connections keep the bounded
2917    /// `FALLBACK_WAL_AUTOCHECKPOINT_PAGES` threshold, so a writable pool
2918    /// in a process that never runs the checkpoint task (embedded runtimes,
2919    /// one-shot CLI executions) retains SQLite's own WAL reclamation instead
2920    /// of growing its WAL without bound.
2921    ///
2922    /// Read-only pools record the claim but have no writer-capable
2923    /// connections to reconfigure. Writable pools publish the claim only after
2924    /// the pooled writer is configured successfully; a failed attempt keeps
2925    /// the bounded fallback active and remains retryable.
2926    pub fn claim_checkpoint_ownership(&self) -> Result<(), SqliteError> {
2927        if !self.checkpoint_ownership.begin_claim() {
2928            return Ok(());
2929        }
2930        let result = (|| {
2931            if !self.config.read_only {
2932                let writer = self.writer_for_checkpoint_operation()?;
2933                writer.conn().pragma_update(None, "wal_autocheckpoint", 0)?;
2934            }
2935            Ok(())
2936        })();
2937        self.checkpoint_ownership.finish_claim(result.is_ok());
2938        result
2939    }
2940
2941    /// Flip an already-running writer task's long-lived connection to the
2942    /// claimed-owner setting.
2943    ///
2944    /// Connections opened after [`Self::claim_checkpoint_ownership`] inherit
2945    /// `wal_autocheckpoint = 0` at open; only a writer task spawned before
2946    /// the claim still holds a connection on the bounded fallback. Returns
2947    /// `Ok(())` without side effects when the pool's write queue is
2948    /// disabled.
2949    pub async fn propagate_checkpoint_claim_to_writer_task(&self) -> Result<(), StorageError> {
2950        let Some(handle) = self.writer_task_handle()? else {
2951            return Ok(());
2952        };
2953        handle
2954            .send_top_level(|conn| {
2955                conn.pragma_update(None, "wal_autocheckpoint", 0)
2956                    .map_err(|e| StorageError::Pool {
2957                        operation: "claim_checkpoint_ownership".into(),
2958                        message: e.to_string(),
2959                    })
2960            })
2961            .await
2962    }
2963
2964    /// Open a standalone read-only connection for one enumerated structural
2965    /// exception to pooled-reader routing.
2966    ///
2967    /// There is intentionally no public/generic standalone-reader fallback.
2968    /// Ordinary request reads use [`Self::reader`] and surface bounded pool
2969    /// exhaustion. Adding a purpose variant or call site is therefore a
2970    /// deliberate, scoped architecture change (ADR-165 Slice 2), not a
2971    /// routine extension.
2972    pub(crate) fn open_standalone_reader(
2973        &self,
2974        purpose: StandaloneReaderPurpose,
2975    ) -> Result<Connection, SqliteError> {
2976        let path = self.read_connection_path()?;
2977
2978        #[cfg(any(unix, windows))]
2979        if let Some(identity_path) = self.identity_path.as_deref() {
2980            self.verify_opened_file_identity(identity_path)?;
2981        }
2982
2983        let conn = claimed_file_identity::open_connection(
2984            &self.config,
2985            path,
2986            OpenFlags::SQLITE_OPEN_READ_ONLY
2987                | OpenFlags::SQLITE_OPEN_NO_MUTEX
2988                | OpenFlags::SQLITE_OPEN_URI,
2989            self.identity_path.as_deref(),
2990        )?;
2991        #[cfg(any(unix, windows))]
2992        if let Some(identity_path) = self.identity_path.as_deref() {
2993            self.verify_connection_file_identity(&conn, identity_path)?;
2994        }
2995        self.verify_opened_database_id(&conn)?;
2996        configure_reader_connection(&conn, &self.config)?;
2997        conn.pragma_update(None, "synchronous", "NORMAL")?;
2998        self.reader_acquisition_counters
2999            .record_standalone_open(purpose);
3000        #[cfg(any(test, feature = "test-support"))]
3001        crate::statement_observer::install(&conn, &self.statement_observer)?;
3002        Ok(conn)
3003    }
3004
3005    fn return_reader(&self, conn: Connection, dirty: bool) {
3006        if self.max_readers == 0 {
3007            return;
3008        }
3009
3010        if reset_reader_connection(&conn, dirty, &self.config)
3011            && reader_connection_is_healthy(&conn)
3012        {
3013            self.enqueue_reader_slot(conn);
3014            return;
3015        }
3016
3017        close_connection_quietly(conn);
3018        self.replace_discarded_reader_slot();
3019    }
3020
3021    /// Push a connection back onto the physical reader queue, discarding it
3022    /// (rather than growing the queue past its configured capacity) if the
3023    /// queue is already full.
3024    fn enqueue_reader_slot(&self, conn: Connection) {
3025        if let Err(conn) = self.readers.push(conn) {
3026            eprintln!("[sqlite-pool] reader pool queue full, discarding replacement connection");
3027            close_connection_quietly(conn);
3028        }
3029    }
3030
3031    /// Open a fresh connection to refill one physical reader slot after its
3032    /// previous occupant was closed and discarded — either disqualified on
3033    /// an ordinary return ([`Self::return_reader`]) or abandoned by a
3034    /// non-reusable checkout ([`ReaderGuard`]'s `Drop`). Both call sites
3035    /// share this so a failed replacement is recorded and logged identically
3036    /// either way, instead of one path silently shrinking the pool.
3037    fn replace_discarded_reader_slot(&self) {
3038        match self.open_reader_connection() {
3039            Ok(conn) => self.enqueue_reader_slot(conn),
3040            Err(error) => {
3041                self.reader_acquisition_counters
3042                    .record_reader_replacement_open_failure();
3043                tracing::warn!(
3044                    %error,
3045                    "sqlite-pool: reader replacement connection failed to open; the physical \
3046                     pool permanently shrinks by one slot below max_readers"
3047                );
3048            }
3049        }
3050    }
3051}
3052
3053/// Bound on the final-component symlink chain [`resolve_symlink_chain`]
3054/// follows before failing loud, mirroring the OS's own loop limit (e.g.
3055/// Linux/macOS `ELOOP`, commonly 40 hops) rather than looping forever on a
3056/// cycle.
3057const MAX_SYMLINK_DEPTH: u32 = 40;
3058
3059/// Mint the canonical [`DbIdentity`] for a configured database path.
3060///
3061/// The sole minting point (ADR-091 backend-scoped attribution design note):
3062/// `tx_registry` origin threading and `sidecar_dir_for` re-keying both
3063/// consume this function's output rather than re-deriving it. Operationally
3064/// three steps:
3065///
3066/// 1. A relative configured path is resolved against the process's current
3067///    directory BEFORE any canonicalization — a bare file name has an empty
3068///    parent, and canonicalizing an empty path fails.
3069/// 2. If the resolved path exists, canonicalize the full path: this
3070///    resolves symlinks at every level, including a symlink at the
3071///    database-file level itself (a `link.sqlite` pointing at the real file
3072///    mints the target's identity).
3073/// 3. If the resolved path does not yet exist (first open), a dangling
3074///    file-level symlink is a valid first-open state — SQLite creates the
3075///    target through the link on first write, and minting the link's own
3076///    name would diverge from a later opener using the target path
3077///    directly. The final-component symlink chain is followed to its
3078///    ultimate target first (bounded, see [`MAX_SYMLINK_DEPTH`]), then that
3079///    target's PARENT directory is canonicalized and the file name is
3080///    appended unchanged — the same pattern `FsBlobStore` uses for its
3081///    root-keyed write locks (`stores/blob.rs::write_lock_for_root`), and
3082///    for the same reason: `Path::canonicalize` requires an existing path.
3083///
3084/// A resolved target whose parent directory does not exist fails minting
3085/// exactly as the subsequent database open itself would fail.
3086///
3087/// Returns the minted [`DbIdentity`] alongside the canonical [`PathBuf`] it
3088/// was built from — `DbIdentity` has no path accessor by design, so callers
3089/// that need the filesystem path (sidecar derivation) keep this pairing
3090/// rather than re-deriving it from the raw configured path.
3091fn mint_db_identity(configured_path: &Path) -> Result<(DbIdentity, PathBuf), SqliteError> {
3092    let absolute = if configured_path.is_absolute() {
3093        configured_path.to_path_buf()
3094    } else {
3095        let cwd = std::env::current_dir().map_err(|e| {
3096            SqliteError::InvalidData(format!(
3097                "cannot mint database identity for {configured_path:?}: failed to resolve the \
3098                 process current directory: {e}"
3099            ))
3100        })?;
3101        cwd.join(configured_path)
3102    };
3103
3104    if absolute.exists() {
3105        let canonical = absolute.canonicalize().map_err(|e| {
3106            SqliteError::InvalidData(format!(
3107                "cannot mint database identity: failed to canonicalize existing path \
3108                 {absolute:?}: {e}"
3109            ))
3110        })?;
3111        return Ok((
3112            DbIdentity::new(canonical.clone().into_os_string()),
3113            canonical,
3114        ));
3115    }
3116
3117    let resolved_target = resolve_symlink_chain(&absolute)?;
3118    let parent = resolved_target.parent().ok_or_else(|| {
3119        SqliteError::InvalidData(format!(
3120            "cannot mint database identity for {resolved_target:?}: path has no parent \
3121             directory"
3122        ))
3123    })?;
3124    let file_name = resolved_target.file_name().ok_or_else(|| {
3125        SqliteError::InvalidData(format!(
3126            "cannot mint database identity for {resolved_target:?}: path has no file name"
3127        ))
3128    })?;
3129    let canonical_parent = parent.canonicalize().map_err(|e| {
3130        SqliteError::InvalidData(format!(
3131            "cannot mint database identity: parent directory {parent:?} of first-open path \
3132             {resolved_target:?} does not exist or is inaccessible: {e}"
3133        ))
3134    })?;
3135    let mut identity_path = canonical_parent;
3136    identity_path.push(file_name);
3137    Ok((
3138        DbIdentity::new(identity_path.clone().into_os_string()),
3139        identity_path,
3140    ))
3141}
3142
3143#[cfg(any(unix, windows))]
3144fn database_file_identity_if_exists(
3145    path: &Path,
3146) -> Result<Option<DatabaseFileIdentity>, SqliteError> {
3147    match database_file_identity(path) {
3148        Ok(identity) => Ok(Some(identity)),
3149        Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None),
3150        Err(error) => Err(error.into()),
3151    }
3152}
3153
3154#[cfg(any(unix, windows))]
3155pub(crate) fn opened_sqlite_file_identity(
3156    conn: &Connection,
3157    path: &Path,
3158) -> Result<DatabaseFileIdentity, SqliteError> {
3159    #[cfg(unix)]
3160    {
3161        verify_sqlite_opened_file_still_at_path(conn)?;
3162        Ok(database_file_identity(path)?)
3163    }
3164    #[cfg(windows)]
3165    {
3166        let opened = sqlite_opened_file_identity(conn)?;
3167        if database_file_identity(path).ok() != Some(opened) {
3168            return Err(SqliteError::InvalidData(
3169                "database file identity changed while SQLite held the opened file".to_string(),
3170            ));
3171        }
3172        Ok(opened)
3173    }
3174}
3175
3176/// The nonce lives in the main database and is read through the connection
3177/// SQLite actually opened. A pathname stat alone can observe a different file
3178/// when another process renames entries during `sqlite3_open_v2`.
3179fn read_database_id(conn: &Connection) -> Result<Option<uuid::Uuid>, SqliteError> {
3180    let table_exists: bool = conn.query_row(
3181        "SELECT count(*) != 0 FROM main.sqlite_master WHERE type = 'table' AND name = ?1",
3182        [DATABASE_ID_TABLE],
3183        |row| row.get(0),
3184    )?;
3185    if !table_exists {
3186        // Older read-only snapshots cannot be initialized here. Their Unix
3187        // file-control and inode checks still apply; writable opens backfill.
3188        return Ok(None);
3189    }
3190    let id: String = conn.query_row(
3191        &format!("SELECT id FROM main.{DATABASE_ID_TABLE} WHERE singleton = 1"),
3192        [],
3193        |row| row.get(0),
3194    )?;
3195    let id = uuid::Uuid::parse_str(&id).map_err(|error| {
3196        SqliteError::InvalidData(format!("invalid stored database identity: {error}"))
3197    })?;
3198    Ok(Some(id))
3199}
3200
3201#[cfg(test)]
3202fn initialize_database_id(conn: &mut Connection) -> Result<uuid::Uuid, SqliteError> {
3203    initialize_database_id_with_admission(conn, None)
3204}
3205
3206fn initialize_database_id_with_admission(
3207    conn: &mut Connection,
3208    admission: Option<&WriteAdmission>,
3209) -> Result<uuid::Uuid, SqliteError> {
3210    if let Some(id) = read_database_id(conn)? {
3211        return Ok(id);
3212    }
3213    let transaction = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
3214    if let Some(admission) = admission {
3215        if let Err(error) = admission.check() {
3216            let rollback = transaction.rollback();
3217            return Err(crate::migrations::capacity_refusal_after_rollback(
3218                conn,
3219                rollback,
3220                error,
3221                "database identity bootstrap",
3222            ));
3223        }
3224    }
3225    transaction.execute_batch(&format!(
3226        "CREATE TABLE IF NOT EXISTS main.{DATABASE_ID_TABLE} (\
3227             singleton INTEGER PRIMARY KEY CHECK (singleton = 1), \
3228             id TEXT NOT NULL\
3229         )"
3230    ))?;
3231    transaction.execute(
3232        &format!("INSERT OR IGNORE INTO main.{DATABASE_ID_TABLE} (singleton, id) VALUES (1, ?1)"),
3233        [uuid::Uuid::new_v4().to_string()],
3234    )?;
3235    let id: String = transaction.query_row(
3236        &format!("SELECT id FROM main.{DATABASE_ID_TABLE} WHERE singleton = 1"),
3237        [],
3238        |row| row.get(0),
3239    )?;
3240    let id = uuid::Uuid::parse_str(&id).map_err(|error| {
3241        SqliteError::InvalidData(format!("invalid stored database identity: {error}"))
3242    })?;
3243    transaction.commit()?;
3244    Ok(id)
3245}
3246
3247#[cfg(unix)]
3248fn verify_sqlite_opened_file_still_at_path(conn: &Connection) -> Result<(), SqliteError> {
3249    let mut moved: std::ffi::c_int = 0;
3250    // SAFETY: `conn` remains alive and exclusively borrowed for this call;
3251    // the `main` C string and writable integer out-parameter remain valid.
3252    // SQLite documents SQLITE_FCNTL_HAS_MOVED as querying the opened file,
3253    // and the bundled Unix VFS implements it using its retained inode.
3254    // https://www.sqlite.org/c3ref/c_fcntl_begin_atomic_write.html
3255    let result = unsafe {
3256        rusqlite::ffi::sqlite3_file_control(
3257            conn.handle(),
3258            c"main".as_ptr(),
3259            rusqlite::ffi::SQLITE_FCNTL_HAS_MOVED,
3260            (&mut moved as *mut std::ffi::c_int).cast(),
3261        )
3262    };
3263    if result != rusqlite::ffi::SQLITE_OK {
3264        return Err(SqliteError::InvalidData(format!(
3265            "cannot verify opened database file identity (SQLite file control {result})"
3266        )));
3267    }
3268    if moved != 0 {
3269        return Err(SqliteError::InvalidData(
3270            "database file identity changed while SQLite held the opened file".to_string(),
3271        ));
3272    }
3273    Ok(())
3274}
3275
3276/// Follow a (possibly dangling) final-component symlink chain to its
3277/// ultimate target, bounded at [`MAX_SYMLINK_DEPTH`] hops. A path that is
3278/// not itself a symlink — including one that does not exist at all —
3279/// returns unchanged on the first iteration; this is the common case, a
3280/// first-open path with no symlink involved.
3281fn resolve_symlink_chain(path: &Path) -> Result<PathBuf, SqliteError> {
3282    let mut current = path.to_path_buf();
3283    for _ in 0..MAX_SYMLINK_DEPTH {
3284        match fs::symlink_metadata(&current) {
3285            Ok(meta) if meta.file_type().is_symlink() => {
3286                let target = fs::read_link(&current).map_err(|e| {
3287                    SqliteError::InvalidData(format!(
3288                        "cannot mint database identity: failed to read symlink {current:?}: {e}"
3289                    ))
3290                })?;
3291                current = if target.is_absolute() {
3292                    target
3293                } else {
3294                    match current.parent() {
3295                        Some(parent) => parent.join(&target),
3296                        None => target,
3297                    }
3298                };
3299            }
3300            _ => return Ok(current),
3301        }
3302    }
3303    Err(SqliteError::InvalidData(format!(
3304        "cannot mint database identity for {path:?}: symlink chain exceeds \
3305         {MAX_SYMLINK_DEPTH} levels"
3306    )))
3307}
3308
3309fn effective_reader_count(config: &PoolConfig, wal_enabled: bool) -> usize {
3310    if config.path.is_some() && (config.read_only || config.code_map_vfs.is_some()) {
3311        config.max_readers.max(1)
3312    } else if config.path.is_some() && config.wal_mode && wal_enabled {
3313        config.max_readers
3314    } else {
3315        0
3316    }
3317}
3318
3319fn open_writer_connection(
3320    config: &PoolConfig,
3321    read_only_open_target: Option<&Path>,
3322    identity_path: Option<&Path>,
3323) -> Result<Connection, SqliteError> {
3324    claimed_file_identity::open_writer(config, read_only_open_target, identity_path)
3325}
3326
3327/// Validate the one-frame reset floor using this backend connection's own
3328/// page size. This runs before writer configuration changes journal mode or
3329/// performs any schema work. The WAL I/O limiter arrives in a later slice, so
3330/// a valid nonzero policy still refuses to open rather than running uncovered.
3331fn validate_wal_ceiling_at_open(conn: &Connection, config: &PoolConfig) -> Result<(), SqliteError> {
3332    let bytes = config.wal_ceiling.effective_bytes(config.read_only);
3333    if bytes == 0 {
3334        return Ok(());
3335    }
3336    let page_size: i64 = conn.pragma_query_value(None, "page_size", |row| row.get(0))?;
3337    let page_size = u64::try_from(page_size).map_err(|_| {
3338        SqliteError::InvalidData("SQLite reported a negative page size".to_string())
3339    })?;
3340    let minimum_bytes = page_size.checked_add(56).ok_or_else(|| {
3341        SqliteError::InvalidData("SQLite page size overflowed the WAL frame floor".to_string())
3342    })?;
3343    if bytes < minimum_bytes {
3344        return Err(SqliteError::WalCeilingBelowMinimum {
3345            bytes,
3346            page_size,
3347            minimum_bytes,
3348        });
3349    }
3350    Err(SqliteError::WalCapacityUnavailable {
3351        bytes,
3352        capability: "WAL I/O limiter",
3353    })
3354}
3355
3356/// Select the one case that may safely use SQLite's immutable URI contract: a
3357/// clean, checkpointed persistent-WAL snapshot with neither a shared-memory
3358/// index nor committed frames in `<db>-wal`. A normal read-only connection can
3359/// create fresh `-wal`/`-shm` files even for that clean database, while
3360/// `immutable=1` keeps the source directory untouched. We deliberately do not
3361/// apply `immutable=1` to:
3362///
3363/// - rollback-journal databases, which can read safely with normal locking and
3364///   should continue observing committed changes when an operator points a
3365///   read-only connection at a live database; or
3366/// - WAL databases with a read-only `-shm`, where ordinary read-only SQLite can
3367///   consume committed WAL frames without writing the frozen index; or
3368/// - WAL databases with a writable `-shm`, which are potentially live. Those
3369///   fail closed before SQLite is opened rather than mutating shared state or
3370///   suppressing change detection unsafely.
3371///
3372/// A non-empty WAL without `-shm` is also refused before open. Immutable SQLite
3373/// does not rebuild a missing WAL index: it ignores the WAL entirely, which can
3374/// make a committed row disappear from inspection. Ordinary read-only SQLite
3375/// would recover the frames but create `-shm`, violating the physical
3376/// read-only contract. The operator must provide the frozen read-only `-shm`
3377/// alongside that WAL (or checkpoint a writable copy first).
3378fn read_only_open_target(
3379    config: &PoolConfig,
3380    physical_path: Option<&Path>,
3381) -> Result<Option<PathBuf>, SqliteError> {
3382    if !config.read_only || config.code_map_vfs.is_some() {
3383        return Ok(None);
3384    }
3385    let Some(path) = physical_path else {
3386        return Ok(None);
3387    };
3388    read_only_wal_open_target_for_path(path).map(Some)
3389}
3390
3391fn read_only_wal_open_target_for_path(path: &Path) -> Result<PathBuf, SqliteError> {
3392    if !sqlite_header_uses_wal(path)? {
3393        return Ok(path.to_path_buf());
3394    }
3395
3396    let shm = sqlite_sidecar_path(path, "-shm");
3397    match fs::metadata(&shm) {
3398        Ok(metadata) if metadata.permissions().readonly() => {
3399            let wal = sqlite_sidecar_path(path, "-wal");
3400            match fs::metadata(&wal) {
3401                Ok(_) => Ok(path.to_path_buf()),
3402                Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
3403                    Err(SqliteError::InvalidData(format!(
3404                        "read-only WAL snapshot {} has a shared-memory sidecar {} but no WAL \
3405                         sidecar {}; refusing the inconsistent sidecar set before SQLite open",
3406                        path.display(),
3407                        shm.display(),
3408                        wal.display(),
3409                    )))
3410                }
3411                Err(error) => Err(SqliteError::Io(error)),
3412            }
3413        }
3414        Ok(_) => Err(SqliteError::InvalidData(format!(
3415            "read-only WAL snapshot {} has a writable WAL shared-memory sidecar {}; close every \
3416             live writer and remove the transient -shm file (or make a genuinely frozen snapshot) \
3417             before inspection",
3418            path.display(),
3419            shm.display(),
3420        ))),
3421        Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
3422            let wal = sqlite_sidecar_path(path, "-wal");
3423            match fs::metadata(&wal) {
3424                Ok(metadata) if metadata.len() > 0 => Err(SqliteError::InvalidData(format!(
3425                    "read-only WAL snapshot {} has a non-empty WAL sidecar {} but no read-only \
3426                     shared-memory sidecar {}; refusing before SQLite open because immutable \
3427                     mode would omit committed WAL frames and ordinary read-only mode would \
3428                     create or mutate -shm; include the frozen read-only -shm beside this \
3429                     snapshot, or checkpoint a writable copy before inspection",
3430                    path.display(),
3431                    wal.display(),
3432                    shm.display(),
3433                ))),
3434                Ok(_) => sqlite_immutable_uri(path),
3435                Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
3436                    sqlite_immutable_uri(path)
3437                }
3438                Err(error) => Err(SqliteError::Io(error)),
3439            }
3440        }
3441        Err(error) => Err(SqliteError::Io(error)),
3442    }
3443}
3444
3445pub(crate) fn open_read_only_snapshot_connection(path: &Path) -> Result<Connection, SqliteError> {
3446    let (_, physical_path) = mint_db_identity(path)?;
3447    let target = read_only_wal_open_target_for_path(&physical_path)?;
3448    let conn = Connection::open_with_flags(&target, reader_open_flags())?;
3449    #[cfg(feature = "namespace-trigram-proto")]
3450    register_namespace_trigram(&conn)?;
3451    Ok(conn)
3452}
3453
3454fn sqlite_header_uses_wal(path: &Path) -> Result<bool, SqliteError> {
3455    let mut file = fs::File::open(path)?;
3456    let mut header = [0_u8; 20];
3457    if let Err(error) = file.read_exact(&mut header) {
3458        if error.kind() == std::io::ErrorKind::UnexpectedEof {
3459            return Ok(false);
3460        }
3461        return Err(SqliteError::Io(error));
3462    }
3463    Ok(&header[..16] == b"SQLite format 3\0" && header[18] == 2 && header[19] == 2)
3464}
3465
3466fn sqlite_sidecar_path(path: &Path, suffix: &str) -> PathBuf {
3467    let mut sidecar = path.as_os_str().to_os_string();
3468    sidecar.push(suffix);
3469    PathBuf::from(sidecar)
3470}
3471
3472fn sqlite_immutable_uri(path: &Path) -> Result<PathBuf, SqliteError> {
3473    let absolute = if path.is_absolute() {
3474        path.to_path_buf()
3475    } else {
3476        std::env::current_dir()?.join(path)
3477    };
3478    let mut uri = String::from("file:");
3479
3480    #[cfg(unix)]
3481    {
3482        use std::os::unix::ffi::OsStrExt as _;
3483        push_sqlite_uri_path(&mut uri, absolute.as_os_str().as_bytes());
3484    }
3485
3486    #[cfg(not(unix))]
3487    {
3488        let path = absolute.to_str().ok_or_else(|| {
3489            SqliteError::InvalidData(format!(
3490                "read-only WAL snapshot path is not representable as a SQLite URI: {}",
3491                absolute.display()
3492            ))
3493        })?;
3494        let normalized = path.replace('\\', "/");
3495        if cfg!(windows) && !normalized.starts_with('/') {
3496            uri.push('/');
3497        }
3498        push_sqlite_uri_path(&mut uri, normalized.as_bytes());
3499    }
3500
3501    uri.push_str("?mode=ro&immutable=1");
3502    Ok(PathBuf::from(uri))
3503}
3504
3505fn push_sqlite_uri_path(uri: &mut String, bytes: &[u8]) {
3506    const HEX: &[u8; 16] = b"0123456789ABCDEF";
3507    for &byte in bytes {
3508        if byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.' | b'~' | b'/') {
3509            uri.push(byte as char);
3510        } else {
3511            uri.push('%');
3512            uri.push(HEX[(byte >> 4) as usize] as char);
3513            uri.push(HEX[(byte & 0x0f) as usize] as char);
3514        }
3515    }
3516}
3517
3518fn open_reader_connection(path: &Path, config: &PoolConfig) -> Result<Connection, SqliteError> {
3519    let conn = claimed_file_identity::open_connection(
3520        config,
3521        path,
3522        reader_open_flags(),
3523        config.path.as_deref(),
3524    )?;
3525    configure_reader_connection(&conn, config)?;
3526    Ok(conn)
3527}
3528
3529fn writer_open_flags() -> OpenFlags {
3530    OpenFlags::SQLITE_OPEN_READ_WRITE
3531        | OpenFlags::SQLITE_OPEN_CREATE
3532        | OpenFlags::SQLITE_OPEN_URI
3533        | OpenFlags::SQLITE_OPEN_NO_MUTEX
3534}
3535
3536/// Read-only writer-slot open flags: no `SQLITE_OPEN_CREATE`, so a missing
3537/// path is rejected rather than silently created.
3538fn writer_read_only_open_flags() -> OpenFlags {
3539    OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_URI | OpenFlags::SQLITE_OPEN_NO_MUTEX
3540}
3541
3542fn reader_open_flags() -> OpenFlags {
3543    OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_URI | OpenFlags::SQLITE_OPEN_NO_MUTEX
3544}
3545
3546#[cfg(feature = "namespace-trigram-proto")]
3547fn register_namespace_trigram(conn: &Connection) -> Result<(), SqliteError> {
3548    crate::namespace_trigram_proto::register(conn).map_err(SqliteError::InvalidData)
3549}
3550
3551fn register_writer_clock(conn: &Connection) -> Result<(), SqliteError> {
3552    // Evaluated by SQLite at statement execution, never deterministic: stream
3553    // observation deadlines use the same UTC microsecond source as note stamps.
3554    conn.create_scalar_function(
3555        "khive_now_micros",
3556        0,
3557        rusqlite::functions::FunctionFlags::SQLITE_UTF8,
3558        |_| Ok(chrono::Utc::now().timestamp_micros()),
3559    )?;
3560    Ok(())
3561}
3562
3563/// Order-preserving UTC key across Chrono's signed timestamp range, with
3564/// nanoseconds kept after the sign-adjusted epoch seconds.
3565pub(crate) fn rfc3339_instant_key(instant: chrono::DateTime<chrono::Utc>) -> Vec<u8> {
3566    let mut key = Vec::with_capacity(12);
3567    key.extend_from_slice(&((instant.timestamp() as u64) ^ (1_u64 << 63)).to_be_bytes());
3568    key.extend_from_slice(&instant.timestamp_subsec_nanos().to_be_bytes());
3569    key
3570}
3571
3572/// The outbox deadline grammar is stricter than the general timestamp filter.
3573/// This is shared by app-maintained stored keys, V44 backfill, and the read
3574/// residual; no schema expression calls an application-defined function.
3575pub(crate) fn strict_rfc3339_key(text: &str) -> Option<Vec<u8>> {
3576    chrono::DateTime::parse_from_rfc3339(text)
3577        .ok()
3578        .map(|instant| rfc3339_instant_key(instant.with_timezone(&chrono::Utc)))
3579}
3580
3581/// Register timestamp-key functions for read filters on pooled connections.
3582pub(crate) fn register_rfc3339_key(conn: &Connection) -> rusqlite::Result<()> {
3583    use rusqlite::functions::FunctionFlags;
3584    use rusqlite::types::ValueRef;
3585
3586    conn.create_scalar_function(
3587        "khive_rfc3339_key",
3588        1,
3589        FunctionFlags::SQLITE_UTF8
3590            | FunctionFlags::SQLITE_DETERMINISTIC
3591            | FunctionFlags::SQLITE_INNOCUOUS,
3592        |ctx| {
3593            let text = match ctx.get_raw(0) {
3594                ValueRef::Text(bytes) => std::str::from_utf8(bytes).ok(),
3595                _ => None,
3596            };
3597            let key = text
3598                .and_then(|text| text.parse::<chrono::DateTime<chrono::Utc>>().ok())
3599                .map(rfc3339_instant_key);
3600            Ok(key)
3601        },
3602    )?;
3603    // The outbox's legacy retry predicate used parse_from_rfc3339, while the
3604    // general key above accepts Chrono's relaxed DateTime FromStr grammar.
3605    // Keep the strict grammar separate so a relaxed-only future value still
3606    // fails open as malformed, instead of postponing the message forever.
3607    conn.create_scalar_function(
3608        "khive_rfc3339_strict_key",
3609        1,
3610        FunctionFlags::SQLITE_UTF8
3611            | FunctionFlags::SQLITE_DETERMINISTIC
3612            | FunctionFlags::SQLITE_INNOCUOUS,
3613        |ctx| {
3614            let text = match ctx.get_raw(0) {
3615                ValueRef::Text(bytes) => std::str::from_utf8(bytes).ok(),
3616                _ => None,
3617            };
3618            let key = text.and_then(strict_rfc3339_key);
3619            Ok(key)
3620        },
3621    )?;
3622    Ok(())
3623}
3624
3625fn configure_writer_connection(
3626    conn: &Connection,
3627    config: &PoolConfig,
3628) -> Result<bool, SqliteError> {
3629    #[cfg(feature = "namespace-trigram-proto")]
3630    register_namespace_trigram(conn)?;
3631    register_writer_clock(conn)?;
3632    register_rfc3339_key(conn)?;
3633    if config.read_only {
3634        // Read-only writer slot: skip write-intent PRAGMAs (journal_mode,
3635        // wal_autocheckpoint, journal_size_limit all require write access to
3636        // change) and lock the connection down with query_only instead.
3637        conn.pragma_update(None, "foreign_keys", "ON")?;
3638        conn.busy_timeout(config.busy_timeout)?;
3639        conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
3640        conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
3641        conn.pragma_update(None, "temp_store", "MEMORY")?;
3642        conn.pragma_update(None, "query_only", "ON")?;
3643
3644        let wal_enabled =
3645            config.wal_mode && current_journal_mode(conn)?.eq_ignore_ascii_case("wal");
3646        return Ok(wal_enabled);
3647    }
3648
3649    let wants_wal = config.path.is_some() && config.wal_mode;
3650
3651    if wants_wal {
3652        conn.pragma_update(None, "journal_mode", "WAL")?;
3653    }
3654    code_map::require_delete_journal(conn, config)?;
3655
3656    conn.pragma_update(None, "synchronous", "NORMAL")?;
3657    conn.pragma_update(None, "foreign_keys", "ON")?;
3658    conn.busy_timeout(config.busy_timeout)?;
3659    conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
3660    conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
3661    conn.pragma_update(None, "temp_store", "MEMORY")?;
3662    // The pool's startup writer always opens before any checkpoint owner can
3663    // claim the pool, so it starts on the bounded fallback;
3664    // `claim_checkpoint_ownership` re-applies the pragma on this connection
3665    // under the writer mutex when a dedicated owner attaches.
3666    if config.code_map_vfs.is_none() {
3667        conn.pragma_update(
3668            None,
3669            "wal_autocheckpoint",
3670            FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
3671        )?;
3672    }
3673
3674    let wal_enabled = wants_wal && current_journal_mode(conn)?.eq_ignore_ascii_case("wal");
3675
3676    if wal_enabled {
3677        conn.pragma_update(None, "journal_size_limit", config.journal_size_limit_bytes)?;
3678    }
3679
3680    Ok(wal_enabled)
3681}
3682
3683fn configure_reader_connection(conn: &Connection, config: &PoolConfig) -> Result<(), SqliteError> {
3684    #[cfg(feature = "namespace-trigram-proto")]
3685    register_namespace_trigram(conn)?;
3686    register_rfc3339_key(conn)?;
3687    conn.pragma_update(None, "foreign_keys", "ON")?;
3688    conn.busy_timeout(config.busy_timeout)?;
3689    conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
3690    conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
3691    conn.pragma_update(None, "temp_store", "MEMORY")?;
3692    Ok(())
3693}
3694
3695fn current_journal_mode(conn: &Connection) -> Result<String, SqliteError> {
3696    conn.pragma_query_value(None, "journal_mode", |row| row.get::<_, String>(0))
3697        .map(|mode| mode.to_ascii_lowercase())
3698        .map_err(Into::into)
3699}
3700
3701fn reset_reader_connection(conn: &Connection, dirty: bool, config: &PoolConfig) -> bool {
3702    if !conn.is_autocommit() {
3703        match conn.execute_batch("ROLLBACK") {
3704            Ok(()) => {}
3705            Err(rusqlite::Error::SqliteFailure(err, _)) => {
3706                if matches!(
3707                    err.code,
3708                    rusqlite::ErrorCode::CannotOpen
3709                        | rusqlite::ErrorCode::DatabaseCorrupt
3710                        | rusqlite::ErrorCode::NotADatabase
3711                        | rusqlite::ErrorCode::DiskFull
3712                ) {
3713                    return false;
3714                }
3715            }
3716            Err(_) => return false,
3717        }
3718        if !conn.is_autocommit() {
3719            return false;
3720        }
3721    }
3722
3723    if !dirty {
3724        return true;
3725    }
3726
3727    reader_connection_state_is_pristine(conn)
3728        && reader_connection_settings_match_baseline(conn, config, 0)
3729}
3730
3731/// A pooled reader must never carry connection-local state across logical
3732/// checkouts. Raw-SQL reads run arbitrary caller SQL against the shared
3733/// pooled connection (`sql_bridge`'s `run_pool_reader_query`), so a
3734/// `CREATE TEMP TABLE` or `ATTACH DATABASE` from one checkout would
3735/// otherwise stay visible to whichever later caller draws the same
3736/// connection back out of the pool. Dropping those objects individually is
3737/// order-sensitive (triggers and indexes depend on their tables), so their
3738/// presence is instead treated as reuse-disqualifying: the caller closes
3739/// the connection and opens a fresh replacement.
3740///
3741/// Called only when [`ReaderGuard::dirty`] is set (`reset_reader_connection`'s
3742/// `dirty` gate) — a typed store read never runs raw caller SQL and returns
3743/// without paying this scan; only a checkout that executed a `SqlReader`
3744/// raw-SQL statement (`sql_bridge`'s `run_pool_reader_query`) does.
3745fn reader_connection_state_is_pristine(conn: &Connection) -> bool {
3746    let has_temp_objects: bool = match conn.query_row(
3747        "SELECT EXISTS(SELECT 1 FROM sqlite_temp_master)",
3748        [],
3749        |row| row.get(0),
3750    ) {
3751        Ok(v) => v,
3752        Err(_) => return false,
3753    };
3754    if has_temp_objects {
3755        return false;
3756    }
3757
3758    let attached_databases: i64 = match conn.query_row(
3759        "SELECT COUNT(*) FROM pragma_database_list WHERE name NOT IN ('main', 'temp')",
3760        [],
3761        |row| row.get(0),
3762    ) {
3763        Ok(v) => v,
3764        Err(_) => return false,
3765    };
3766    attached_databases == 0
3767}
3768
3769/// The observable connection-local settings a reader capability could in
3770/// principle influence, compared against this pool's configured baseline.
3771/// Called only on a dirty return, alongside [`reader_connection_state_is_pristine`]
3772/// — the reader-capability admission gate (`sql_bridge::reader_capability_admits`)
3773/// refuses every raw-SQL form that could change these today, so this is
3774/// defense in depth against a gap in that gate, not the primary boundary.
3775///
3776/// `expected_query_only` is the caller's expected baseline for `query_only`:
3777/// pooled reader connections never set it explicitly (`configure_reader_connection`
3778/// does not touch it) regardless of `config.read_only`, so pooled-reader callers
3779/// pass `0`; the degraded shared-writer-as-reader lease (`max_readers == 0`)
3780/// inherits whatever `configure_writer_connection` set, which does depend on
3781/// `config.read_only`.
3782fn reader_connection_settings_match_baseline(
3783    conn: &Connection,
3784    config: &PoolConfig,
3785    expected_query_only: i64,
3786) -> bool {
3787    let expected_busy_timeout_ms =
3788        i64::try_from(config.busy_timeout.as_millis()).unwrap_or(i64::MAX);
3789    let expected_cache_size: i64 = CACHE_SIZE_KIB.parse().unwrap_or(-65536);
3790    let checks: [(&str, i64); 8] = [
3791        ("query_only", expected_query_only),
3792        ("writable_schema", 0),
3793        ("foreign_keys", 1),
3794        ("busy_timeout", expected_busy_timeout_ms),
3795        ("cache_size", expected_cache_size),
3796        ("temp_store", 2),
3797        ("read_uncommitted", 0),
3798        ("defer_foreign_keys", 0),
3799    ];
3800    checks.iter().all(|(pragma, expected)| {
3801        conn.pragma_query_value(None, pragma, |row| row.get::<_, i64>(0))
3802            .map(|actual| actual == *expected)
3803            .unwrap_or(false)
3804    })
3805}
3806
3807/// Best-effort recovery for a dirty shared reader-writer lease
3808/// (`max_readers == 0` degraded mode): there is no separate connection to
3809/// close and replace, so a disqualifying state is instead undone in place —
3810/// DETACH every non-main/non-temp database, drop every TEMP object in
3811/// dependency order (views and triggers before the indexes and tables they
3812/// depend on), then reapply the pool's baseline connection settings. Returns
3813/// `true` only if the connection verifiably passes the same pristine/settings
3814/// checks afterward; the caller poisons the lease on `false`.
3815fn restore_shared_reader_state(conn: &Connection, config: &PoolConfig) -> bool {
3816    if !detach_non_main_databases(conn) {
3817        return false;
3818    }
3819    if !drop_temp_objects(conn) {
3820        return false;
3821    }
3822    let expected_query_only = i64::from(config.read_only);
3823    if reset_observable_settings(conn, config, expected_query_only).is_err() {
3824        return false;
3825    }
3826    reader_connection_state_is_pristine(conn)
3827        && reader_connection_settings_match_baseline(conn, config, expected_query_only)
3828}
3829
3830fn detach_non_main_databases(conn: &Connection) -> bool {
3831    loop {
3832        let name: Option<String> = match conn.query_row(
3833            "SELECT name FROM pragma_database_list WHERE name NOT IN ('main', 'temp') LIMIT 1",
3834            [],
3835            |row| row.get(0),
3836        ) {
3837            Ok(name) => Some(name),
3838            Err(rusqlite::Error::QueryReturnedNoRows) => None,
3839            Err(_) => return false,
3840        };
3841        let Some(name) = name else {
3842            return true;
3843        };
3844        let quoted = format!("\"{}\"", name.replace('"', "\"\""));
3845        if conn
3846            .execute_batch(&format!("DETACH DATABASE {quoted}"))
3847            .is_err()
3848        {
3849            return false;
3850        }
3851    }
3852}
3853
3854fn drop_temp_objects(conn: &Connection) -> bool {
3855    // Views and triggers depend on tables/indexes but are never depended on
3856    // themselves; dropping them first means every later DROP TABLE/INDEX
3857    // never fails on a dangling dependent.
3858    for (kind, ddl_keyword) in [
3859        ("view", "VIEW"),
3860        ("trigger", "TRIGGER"),
3861        ("index", "INDEX"),
3862        ("table", "TABLE"),
3863    ] {
3864        loop {
3865            let name: Option<String> = match conn.query_row(
3866                "SELECT name FROM sqlite_temp_master WHERE type = ?1 LIMIT 1",
3867                [kind],
3868                |row| row.get(0),
3869            ) {
3870                Ok(name) => Some(name),
3871                Err(rusqlite::Error::QueryReturnedNoRows) => None,
3872                Err(_) => return false,
3873            };
3874            let Some(name) = name else {
3875                break;
3876            };
3877            let quoted = format!("\"{}\"", name.replace('"', "\"\""));
3878            if conn
3879                .execute_batch(&format!("DROP {ddl_keyword} IF EXISTS temp.{quoted}"))
3880                .is_err()
3881            {
3882                return false;
3883            }
3884        }
3885    }
3886    true
3887}
3888
3889fn reset_observable_settings(
3890    conn: &Connection,
3891    config: &PoolConfig,
3892    expected_query_only: i64,
3893) -> Result<(), rusqlite::Error> {
3894    conn.pragma_update(None, "query_only", expected_query_only)?;
3895    conn.pragma_update(None, "writable_schema", 0)?;
3896    conn.pragma_update(None, "foreign_keys", "ON")?;
3897    conn.busy_timeout(config.busy_timeout)?;
3898    conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
3899    conn.pragma_update(None, "temp_store", "MEMORY")?;
3900    conn.pragma_update(None, "read_uncommitted", 0)?;
3901    conn.pragma_update(None, "defer_foreign_keys", 0)?;
3902    Ok(())
3903}
3904
3905fn reader_connection_is_healthy(conn: &Connection) -> bool {
3906    match conn.query_row("SELECT 1", [], |row| row.get::<_, i64>(0)) {
3907        Ok(_) => true,
3908        Err(rusqlite::Error::SqliteFailure(err, _)) => !matches!(
3909            err.code,
3910            rusqlite::ErrorCode::CannotOpen
3911                | rusqlite::ErrorCode::NotADatabase
3912                | rusqlite::ErrorCode::DatabaseCorrupt
3913                | rusqlite::ErrorCode::PermissionDenied
3914                | rusqlite::ErrorCode::SystemIoFailure
3915        ),
3916        Err(_) => true,
3917    }
3918}
3919
3920fn close_connection_quietly(conn: Connection) {
3921    match conn.close() {
3922        Ok(()) => {}
3923        Err((conn, _)) => drop(conn),
3924    }
3925}
3926
3927fn pool_exhausted_error(timeout: Duration, max_readers: usize) -> SqliteError {
3928    rusqlite::Error::SqliteFailure(
3929        rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_BUSY),
3930        Some(format!(
3931            "Pool exhausted: no reader available after {timeout:?} (max_readers={max_readers})"
3932        )),
3933    )
3934    .into()
3935}
3936
3937#[cfg(test)]
3938#[path = "runtime_write_routing_tests.rs"]
3939mod runtime_write_routing_tests;
3940
3941#[cfg(test)]
3942#[path = "database_owner_identity_pool_tests.rs"]
3943mod database_owner_identity_pool_tests;
3944
3945#[cfg(test)]
3946#[path = "pool_tests.rs"]
3947mod tests;
3948
3949#[cfg(all(test, any(unix, windows)))]
3950#[path = "pool_identity_admission_tests.rs"]
3951mod identity_admission_tests;