Skip to main content

khive_db/
diagnostics.rs

1//! Non-mutating-by-intent database, WAL, and checkpoint diagnostics (read-only
2//! operator surface).
3//!
4//! Answers "is a reader pinning the checkpoint / why is the WAL at 64MiB"
5//! without raw SQL against a production store.
6//!
7//! NOT read-only. [`checkpoint_probe`](crate::diagnostics::checkpoint_probe) issues a real
8//! `PRAGMA wal_checkpoint(PASSIVE)`, and a PASSIVE checkpoint that succeeds
9//! BACKFILLS WAL frames into the main database — ordinary database-page
10//! writes, on the happy path. That I/O is the point: the busy/log_frames/
11//! checkpointed_frames triple reports a one-row backfill gap; a sustained
12//! pinned-frame report additionally requires a matching checkpoint run.
13//! SQLite exposes no read-only API for those counters. The
14//! guarantee is that nothing here changes logical state or destroys
15//! evidence, not that nothing touches the disk.
16//!
17//! What IS guaranteed, and what the narrowings below buy:
18//!
19//! * never creates a missing database file,
20//! * never escalates to TRUNCATE,
21//! * never perturbs the counters it reports,
22//! * never increments write-traffic acquisition counters,
23//! * never deletes a walpin sidecar entry.
24//!
25//! Deliberate narrowings make those claims true rather than aspirational:
26//!
27//! 1. [`checkpoint_probe`](crate::diagnostics::checkpoint_probe) runs a
28//!    single `PRAGMA wal_checkpoint(PASSIVE)` —
29//!    which never blocks readers or writers — and does NOT route through
30//!    [`crate::checkpoint::checkpoint_once`]: that path mutates
31//!    `TruncateState`, may escalate to TRUNCATE, and double-counts the
32//!    ADR-091 process-global counters. A verb that reports state must not
33//!    perturb the state it reports.
34//! 2. The probe's connection comes from
35//!    `ConnectionPool::open_standalone_writer_untracked`, opened without
36//!    `SQLITE_OPEN_CREATE`. A missing database yields `checkpoint_probe:
37//!    null` plus a `checkpoint_probe_error`, never a freshly created file.
38//! 3. WAL-pin attribution combines the read-only OS holder census with
39//!    `walpin::inspect_live`, a separate bounded sidecar enumerator whose
40//!    purpose flag prohibits every unlink. It applies the same descriptor-
41//!    bound trust and liveness checks as checkpoint attribution, but reports
42//!    stale cleanup candidates rather than consuming them. A complete census
43//!    plus a complete, conclusive sidecar walk can therefore report complete;
44//!    truncation, unknown entries, or holders absent from the sidecar degrade
45//!    explicitly.
46//! 4. Graph-edge integrity uses three scalar SELECTs on the same guarded
47//!    standalone connection. It exposes the exact pre-V14 duplicate-ID group
48//!    count plus raw live-edge/list-ledger counts and never repairs or deletes
49//!    data.
50//! 5. FTS5 segment diagnostics decode each index's documented one-row
51//!    structure record (`%_data.id = 10`). They never run `COUNT(*)` over or
52//!    scan a `%_idx` table, so observing segment health is bounded even when
53//!    the corpus itself is large.
54//!
55//! Checkpoint counters are process-global, while reader and writer acquisition
56//! counters belong to the supplied pool and can reset when it is reconstructed.
57//! Every payload carries [`BuildIdentity`](crate::diagnostics::BuildIdentity) and
58//! [`ProcessIdentity`](crate::diagnostics::ProcessIdentity) for the
59//! process producing the reading. The PID, OS start time, and main-pool generation
60//! identify the main pool's counter window; a secondary pool's reader/writer
61//! counters have their own reconstruction window. Checkpoint counters remain global.
62
63#[path = "diagnostics/disk_guard.rs"]
64mod disk_guard;
65#[path = "diagnostics/writer_contention.rs"]
66mod writer_contention;
67
68use std::path::{Path, PathBuf};
69use std::sync::atomic::{AtomicBool, Ordering};
70use std::sync::Arc;
71use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
72
73use khive_storage::error::StorageError;
74use khive_storage::types::StorageResult;
75use khive_storage::StorageCapability;
76use rusqlite::Connection;
77use serde::Serialize;
78
79use crate::checkpoint;
80use crate::pool::{ConnectionPool, WalCeilingSource};
81
82pub use disk_guard::DiskGuardDiagnostics;
83
84/// Raw `PRAGMA wal_checkpoint(PASSIVE)` return row.
85///
86/// SQLite returns three columns: `busy` (1 when the checkpoint could not
87/// obtain the checkpoint lock or read the WAL-index header), `log` (frames
88/// currently in the WAL), and `checkpointed` (frames moved into the database
89/// by THIS call). The one-row backfill gap is `log - checkpointed`; that
90/// difference alone does not establish a reader pin.
91#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
92pub struct CheckpointProbe {
93    pub busy: i64,
94    pub log_frames: i64,
95    pub checkpointed_frames: i64,
96}
97
98impl CheckpointProbe {
99    /// Frames still awaiting backfill in this row. Negative frame counts
100    /// (an in-memory or non-WAL database reports `-1`) return 0. Even when
101    /// busy, populated frame counts describe only a gap, never a reader pin.
102    pub fn backfill_gap_frames(&self) -> i64 {
103        if self.log_frames < 0 || self.checkpointed_frames < 0 {
104            return 0;
105        }
106        self.log_frames
107            .saturating_sub(self.checkpointed_frames)
108            .max(0)
109    }
110}
111
112/// Issue one PASSIVE checkpoint on `conn` and return the raw triple.
113///
114/// PASSIVE never blocks writers or readers, so this is safe against a live
115/// daemon. Touches none of the ADR-091 counters — unlike the periodic
116/// checkpoint task's own observation path, this does not mirror into the
117/// process-global gauges.
118///
119/// This WRITES. A PASSIVE checkpoint that makes progress backfills WAL
120/// frames into the main database file; that is normal checkpoint I/O, and on
121/// the happy path it is what the reported `checkpointed_frames` counts. The
122/// call is non-mutating in the sense that matters for a diagnostic — no
123/// logical state changes, no escalation, no evidence destroyed — but it is
124/// not write-free, and callers must not describe it as such.
125pub fn checkpoint_probe(conn: &Connection) -> rusqlite::Result<CheckpointProbe> {
126    conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
127        Ok(CheckpointProbe {
128            busy: row.get(0)?,
129            log_frames: row.get(1)?,
130            checkpointed_frames: row.get(2)?,
131        })
132    })
133}
134
135/// ADR-091 checkpoint counters, read as one snapshot.
136///
137/// The two `Option` fields carry the `u64::MAX` "never observed" sentinel as
138/// `None`, so a caller serializing this never sees `18446744073709551615`
139/// where it means "no observation yet".
140#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
141pub struct CheckpointCounters {
142    pub last_observed_wal_pages: Option<u64>,
143    pub truncate_attempts: u64,
144    pub truncate_consecutive_failures: u64,
145    pub checkpoint_skipped_ticks: u64,
146    pub checkpoint_consecutive_skips: u64,
147    pub checkpoint_last_skip_wal_pages: Option<u64>,
148    pub checkpoint_pressure_elevated_ticks: u64,
149    pub checkpoint_pressure_episodes_started: u64,
150    pub checkpoint_pressure_episodes_recovered: u64,
151    pub checkpoint_lifecycle_append_attempts: u64,
152    pub checkpoint_lifecycle_append_failures: u64,
153    pub checkpoint_lifecycle_enqueue_drops: u64,
154    /// Cached-reader read transactions rolled back on reuse for exceeding
155    /// `read_tx_max_age` (#1846) — the count of WAL snapshots actually
156    /// released by that bound, not merely logged as stale.
157    pub read_tx_max_age_evictions: u64,
158}
159
160/// Snapshot the process-global ADR-091 counters.
161pub fn checkpoint_counters() -> CheckpointCounters {
162    CheckpointCounters {
163        last_observed_wal_pages: checkpoint::last_observed_wal_pages(),
164        truncate_attempts: checkpoint::truncate_attempts(),
165        truncate_consecutive_failures: checkpoint::truncate_consecutive_failures(),
166        checkpoint_skipped_ticks: checkpoint::checkpoint_skipped_ticks(),
167        checkpoint_consecutive_skips: checkpoint::checkpoint_consecutive_skips(),
168        checkpoint_last_skip_wal_pages: checkpoint::checkpoint_last_skip_wal_pages(),
169        checkpoint_pressure_elevated_ticks: checkpoint::checkpoint_pressure_elevated_ticks(),
170        checkpoint_pressure_episodes_started: checkpoint::checkpoint_pressure_episodes_started(),
171        checkpoint_pressure_episodes_recovered: checkpoint::checkpoint_pressure_episodes_recovered(
172        ),
173        checkpoint_lifecycle_append_attempts: checkpoint::checkpoint_lifecycle_append_attempts(),
174        checkpoint_lifecycle_append_failures: checkpoint::checkpoint_lifecycle_append_failures(),
175        checkpoint_lifecycle_enqueue_drops: checkpoint::checkpoint_lifecycle_enqueue_drops(),
176        read_tx_max_age_evictions: checkpoint::read_tx_max_age_evictions(),
177    }
178}
179
180/// Which build produced this reading.
181///
182/// The counters above are process-global: a stale daemon reports a stale
183/// build's state. `build_hash` is `None` unless the binary was stamped with
184/// one — this crate does not introduce a build-metadata mechanism, it only
185/// reports whatever the caller already has.
186#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
187pub struct BuildIdentity {
188    pub version: String,
189    pub build_hash: Option<String>,
190}
191
192impl BuildIdentity {
193    /// Identity of the crate that compiled this call site.
194    pub fn from_env(version: &str, build_hash: Option<&str>) -> Self {
195        Self {
196            version: version.to_string(),
197            build_hash: build_hash.map(str::to_string),
198        }
199    }
200}
201
202/// The serving OS process and the main pool's reader/writer counter generation.
203#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
204pub struct ProcessIdentity {
205    pub pid: u32,
206    /// OS-reported process creation time in whole Unix epoch seconds (UTC).
207    /// This is neither pool creation time nor the time of the first request.
208    pub started_at: Option<i64>,
209    /// Present when the OS lookup is unsupported or unavailable. An unknown
210    /// start time is never replaced by zero or the current wall clock.
211    pub started_at_unavailable_reason: Option<String>,
212    pub pool_generation: u64,
213}
214
215impl ProcessIdentity {
216    pub fn current(pool: &ConnectionPool) -> Self {
217        let pid = std::process::id();
218        Self::from_start_time(
219            pid,
220            crate::walpin::process_start_time_secs(pid),
221            pool.main_pool_generation(),
222        )
223    }
224
225    fn from_start_time(pid: u32, started_at: Option<i64>, pool_generation: u64) -> Self {
226        Self {
227            pid,
228            started_at,
229            started_at_unavailable_reason: started_at.is_none().then(|| {
230                format!(
231                    "OS process start time is unsupported or unavailable on {}",
232                    std::env::consts::OS
233                )
234            }),
235            pool_generation,
236        }
237    }
238}
239
240/// Absolute-free WAL file state for one database path.
241#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
242pub struct WalFileState {
243    pub wal_path: String,
244    /// `None` when the `-wal` sidecar does not exist or could not be stat'd —
245    /// the reason is in `unavailable_reason`.
246    pub wal_size_bytes: Option<u64>,
247    pub unavailable_reason: Option<String>,
248}
249
250/// Stat `<db_path>-wal` and report its size in bytes.
251pub fn wal_file_state(db_path: &Path) -> WalFileState {
252    let wal_path = wal_sidecar_path(db_path);
253    match std::fs::metadata(&wal_path) {
254        Ok(md) => WalFileState {
255            wal_path: wal_path.display().to_string(),
256            wal_size_bytes: Some(md.len()),
257            unavailable_reason: None,
258        },
259        Err(e) => WalFileState {
260            wal_path: wal_path.display().to_string(),
261            wal_size_bytes: None,
262            unavailable_reason: Some(e.to_string()),
263        },
264    }
265}
266
267/// `<db_path>-wal`, built by suffixing the file name (not by replacing an
268/// extension — `khive.db` must map to `khive.db-wal`).
269fn wal_sidecar_path(db_path: &Path) -> PathBuf {
270    let mut s = db_path.as_os_str().to_os_string();
271    s.push("-wal");
272    PathBuf::from(s)
273}
274
275/// WAL-pin attribution: who currently holds the database open, and how
276/// complete that answer is.
277///
278/// `status`, `status_reasons`, and the tagged `census` field are the
279/// authoritative wire contract. The older sibling booleans and PID arrays
280/// remain available to Rust callers but are not serialized. Sidecar evidence
281/// is collected through the handle-checked, bounded, non-mutating diagnostics
282/// enumerator and reconciled with the OS holder census.
283#[derive(Debug, Clone, PartialEq, Serialize)]
284pub struct WalPinAttribution {
285    /// Authoritative quality of the complete attribution answer.
286    pub status: WalPinAttributionStatus,
287    /// Machine-adjacent reasons why the answer is degraded or unavailable.
288    pub status_reasons: Vec<String>,
289    /// Authoritative tagged result of the OS holder census.
290    pub census: WalPinCensus,
291    /// PID spelling used by the process census for the reporting process.
292    pub reporting_pid: u32,
293    /// Whether the reporting process appears among confirmed file holders.
294    pub reporting_process_is_holder: Option<bool>,
295    /// Why the census could not determine whether the reporter is a holder.
296    pub reporting_process_is_holder_unavailable_reason: Option<String>,
297    /// Raw start-time values for every confirmed holder PID.
298    pub census_process_start_times: Vec<WalPinCensusProcessStart>,
299    /// Resolution of the operating system's process start-time value.
300    pub start_time_resolution_secs: Option<u64>,
301    /// Why the platform does not provide process start-time values.
302    pub start_time_resolution_unavailable_reason: Option<String>,
303    /// `false` plus an `unavailable_reason` whenever the OS census failed,
304    /// or when only a partial (census-only) answer is available.
305    pub available: bool,
306    pub unavailable_reason: Option<String>,
307    /// OS-derived census of every PID holding the DB file open.
308    #[serde(skip_serializing)]
309    pub census_holder_pids: Vec<u32>,
310    #[serde(skip_serializing)]
311    pub census_uninspectable_pids: Vec<u32>,
312    #[serde(skip_serializing)]
313    pub census_truncated: bool,
314    #[serde(skip_serializing)]
315    pub census_is_complete: bool,
316    /// Live, identity-matched heartbeats found in the sidecar.
317    pub reporting: Vec<WalPinHolder>,
318    /// Live, identity-matched beacons with no over-threshold heartbeat.
319    pub registered_silent_pids: Vec<u32>,
320    /// Sidecar entries whose identity or freshness could not be established.
321    pub unknown_pids: Vec<u32>,
322    /// OS-confirmed holders absent from every sidecar classification.
323    pub census_pids_without_attribution: Vec<u32>,
324    /// Whether both evidence sources completed and every holder was present
325    /// in a conclusive sidecar classification.
326    pub fully_attributed: bool,
327    /// Machine-readable sidecar classifications, ordered by PID and status.
328    pub sidecar_entries: Vec<serde_json::Value>,
329    /// Present only when this request actually enumerated the sidecar.
330    #[serde(skip_serializing_if = "Option::is_none")]
331    pub sidecar_listing_truncated: Option<bool>,
332    /// Number of stale producer temps a housekeeping pass would reap. The
333    /// diagnostic pass itself never performs that cleanup.
334    #[serde(skip_serializing_if = "Option::is_none")]
335    pub sidecar_entries_cleanup_would_reap: Option<usize>,
336}
337
338/// OS-reported start time for one holder confirmed by the process census.
339#[derive(Debug, Clone, PartialEq, Serialize)]
340pub struct WalPinCensusProcessStart {
341    /// PID confirmed by the OS holder census.
342    pub pid: u32,
343    /// Process start time as epoch seconds, if the OS reported it.
344    pub process_start_time_secs: Option<i64>,
345    /// Why the process start time is unavailable.
346    pub process_start_time_unavailable_reason: Option<String>,
347}
348
349fn census_process_start_times(holder_pids: &[u32]) -> Vec<WalPinCensusProcessStart> {
350    holder_pids
351        .iter()
352        .map(|pid| {
353            let process_start_time_secs = crate::walpin::process_start_time_secs(*pid);
354            let process_start_time_unavailable_reason =
355                process_start_time_secs.is_none().then(|| {
356                    if crate::walpin::start_time_resolution_secs().is_none() {
357                        "process start time is unavailable on this platform".to_string()
358                    } else {
359                        "the operating system did not report a start time for this process"
360                            .to_string()
361                    }
362                });
363            WalPinCensusProcessStart {
364                pid: *pid,
365                process_start_time_secs,
366                process_start_time_unavailable_reason,
367            }
368        })
369        .collect()
370}
371
372/// Overall quality of the WAL-pin attribution answer.
373///
374/// A result is `complete` only when both the OS census and read-only sidecar
375/// enumeration completed and every holder has sidecar evidence.
376#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
377#[serde(rename_all = "snake_case")]
378pub enum WalPinAttributionStatus {
379    /// Holder census and sidecar reconciliation both completed.
380    Complete,
381    /// Some useful evidence is present, but full attribution is impossible.
382    Degraded,
383    /// No holder-census evidence could be collected.
384    Unavailable,
385}
386
387/// Tagged OS holder-census result.
388///
389/// An incomplete scan retains the partial holder evidence, but its wire shape
390/// cannot be mistaken for a complete census without ignoring the explicit
391/// `status` tag. This is the fail-loud direction required by ADR-091: a
392/// truncated walk or an uninspectable PID is inconclusive, never evidence
393/// that no additional holder exists.
394#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
395#[serde(tag = "status", rename_all = "snake_case")]
396pub enum WalPinCensus {
397    /// Every visible process was inspected and the walk was not truncated.
398    Complete {
399        /// PIDs confirmed to hold the database file open.
400        holder_pids: Vec<u32>,
401    },
402    /// Partial evidence from an inconclusive process walk.
403    Incomplete {
404        /// PIDs confirmed to hold the database file open.
405        holder_pids: Vec<u32>,
406        /// PIDs for which inspection failed outright.
407        uninspectable_pids: Vec<u32>,
408        /// Whether process enumeration had positive evidence of truncation.
409        truncated: bool,
410        /// Why additional holders cannot be ruled out.
411        reason: String,
412    },
413    /// The platform or census operation supplied no holder evidence.
414    Unavailable {
415        /// Why the census could not run.
416        reason: String,
417    },
418}
419
420/// Name the reason a holder census stopped short.
421///
422/// A budget stop and an enumeration failure both set `truncated`, and an
423/// operator reading one sentence has to be able to tell them apart: the first
424/// is a configurable trade this process made on purpose and is expected on a
425/// busy box, the second is evidence that process enumeration itself
426/// misbehaved. Describing both as "truncated" would make the report cry wolf
427/// on every bounded call, which is how the sentence that matters stops being
428/// read.
429fn census_truncation_cause(budget_exhausted: bool) -> String {
430    if budget_exhausted {
431        "the OS process walk stopped at its wall-clock budget (see          collection_cost.wal_pin_census_budget_ms)"
432            .to_string()
433    } else {
434        "the OS process walk was truncated".to_string()
435    }
436}
437
438/// One PID's live heartbeat as reported to an operator.
439#[derive(Debug, Clone, PartialEq, Serialize)]
440pub struct WalPinHolder {
441    pub pid: u32,
442    pub process_role: String,
443    pub current_oldest_tx_age_secs: f64,
444    pub oldest_tx_label: Option<String>,
445    pub attribution_is_evidence_backed: bool,
446}
447
448impl WalPinAttribution {
449    fn unavailable(reason: impl Into<String>) -> Self {
450        let reason = reason.into();
451        let start_time_resolution_secs = crate::walpin::start_time_resolution_secs();
452        Self {
453            status: WalPinAttributionStatus::Unavailable,
454            status_reasons: vec![reason.clone()],
455            census: WalPinCensus::Unavailable {
456                reason: reason.clone(),
457            },
458            reporting_pid: crate::walpin::reporting_pid(),
459            reporting_process_is_holder: None,
460            reporting_process_is_holder_unavailable_reason: Some(
461                "the OS holder census supplied no holder evidence".to_string(),
462            ),
463            census_process_start_times: Vec::new(),
464            start_time_resolution_secs,
465            start_time_resolution_unavailable_reason: start_time_resolution_secs
466                .is_none()
467                .then(|| "process start time is unavailable on this platform".to_string()),
468            available: false,
469            unavailable_reason: Some(reason),
470            census_holder_pids: Vec::new(),
471            census_uninspectable_pids: Vec::new(),
472            census_truncated: false,
473            census_is_complete: false,
474            reporting: Vec::new(),
475            registered_silent_pids: Vec::new(),
476            unknown_pids: Vec::new(),
477            census_pids_without_attribution: Vec::new(),
478            fully_attributed: false,
479            sidecar_entries: Vec::new(),
480            sidecar_listing_truncated: None,
481            sidecar_entries_cleanup_would_reap: None,
482        }
483    }
484}
485
486#[cfg(all(unix, test))]
487fn wal_pin_attribution_from_census(census: crate::walpin::CensusResult) -> WalPinAttribution {
488    wal_pin_attribution_without_sidecar(
489        census,
490        "read-only sidecar enumeration did not run for this attribution snapshot".to_string(),
491    )
492}
493
494#[cfg(unix)]
495fn wal_pin_attribution_without_sidecar(
496    census: crate::walpin::CensusResult,
497    sidecar_reason: String,
498) -> WalPinAttribution {
499    let census_is_complete = census.is_complete();
500    let mut census_holder_pids: Vec<u32> = census.holders.iter().copied().collect();
501    census_holder_pids.sort_unstable();
502    let reporting_pid = crate::walpin::reporting_pid();
503    let reporting_process_is_holder = census_holder_pids
504        .contains(&reporting_pid)
505        .then_some(true)
506        .or_else(|| census_is_complete.then_some(false));
507    let reporting_process_is_holder_unavailable_reason =
508        reporting_process_is_holder.is_none().then(|| {
509            "the reporting process was absent from the incomplete OS holder census".to_string()
510        });
511    let census_process_start_times = census_process_start_times(&census_holder_pids);
512    let start_time_resolution_secs = crate::walpin::start_time_resolution_secs();
513    let mut census_uninspectable_pids = census.uninspectable_pids;
514    census_uninspectable_pids.sort_unstable();
515    census_uninspectable_pids.dedup();
516    let census_truncated = census.truncated;
517    let census_budget_exhausted = census.budget_exhausted;
518
519    let mut status_reasons = vec![sidecar_reason];
520    let census = if census_is_complete {
521        WalPinCensus::Complete {
522            holder_pids: census_holder_pids.clone(),
523        }
524    } else {
525        let mut causes = Vec::new();
526        if census_truncated {
527            causes.push(census_truncation_cause(census_budget_exhausted));
528        }
529        if !census_uninspectable_pids.is_empty() {
530            causes.push(format!(
531                "{} PID(s) could not be inspected",
532                census_uninspectable_pids.len()
533            ));
534        }
535        let reason = format!(
536            "OS holder census is incomplete: {}; additional database holders cannot be ruled out",
537            causes.join("; ")
538        );
539        status_reasons.push(reason.clone());
540        WalPinCensus::Incomplete {
541            holder_pids: census_holder_pids.clone(),
542            uninspectable_pids: census_uninspectable_pids.clone(),
543            truncated: census_truncated,
544            reason,
545        }
546    };
547
548    WalPinAttribution {
549        status: WalPinAttributionStatus::Degraded,
550        unavailable_reason: Some(status_reasons.join("; ")),
551        status_reasons,
552        census,
553        reporting_pid,
554        reporting_process_is_holder,
555        reporting_process_is_holder_unavailable_reason,
556        census_process_start_times,
557        start_time_resolution_secs,
558        start_time_resolution_unavailable_reason: start_time_resolution_secs
559            .is_none()
560            .then(|| "process start time is unavailable on this platform".to_string()),
561        available: false,
562        census_holder_pids,
563        census_uninspectable_pids,
564        census_truncated,
565        census_is_complete,
566        reporting: Vec::new(),
567        registered_silent_pids: Vec::new(),
568        unknown_pids: Vec::new(),
569        census_pids_without_attribution: Vec::new(),
570        fully_attributed: false,
571        sidecar_entries: Vec::new(),
572        sidecar_listing_truncated: None,
573        sidecar_entries_cleanup_would_reap: None,
574    }
575}
576
577#[cfg(unix)]
578fn wal_pin_attribution_from_evidence(
579    census: crate::walpin::CensusResult,
580    sidecar: crate::walpin::WalpinReport,
581) -> WalPinAttribution {
582    use std::collections::BTreeSet;
583
584    let census_is_complete = census.is_complete();
585    let mut census_holder_pids: Vec<u32> = census.holders.iter().copied().collect();
586    census_holder_pids.sort_unstable();
587    let reporting_pid = crate::walpin::reporting_pid();
588    let reporting_process_is_holder = census_holder_pids
589        .contains(&reporting_pid)
590        .then_some(true)
591        .or_else(|| census_is_complete.then_some(false));
592    let reporting_process_is_holder_unavailable_reason =
593        reporting_process_is_holder.is_none().then(|| {
594            "the reporting process was absent from the incomplete OS holder census".to_string()
595        });
596    let census_process_start_times = census_process_start_times(&census_holder_pids);
597    let start_time_resolution_secs = crate::walpin::start_time_resolution_secs();
598    let mut census_uninspectable_pids = census.uninspectable_pids;
599    census_uninspectable_pids.sort_unstable();
600    census_uninspectable_pids.dedup();
601    let census_truncated = census.truncated;
602    let census_budget_exhausted = census.budget_exhausted;
603    let census_carrier = if census_is_complete {
604        WalPinCensus::Complete {
605            holder_pids: census_holder_pids.clone(),
606        }
607    } else {
608        let mut causes = Vec::new();
609        if census_truncated {
610            causes.push(census_truncation_cause(census_budget_exhausted));
611        }
612        if !census_uninspectable_pids.is_empty() {
613            causes.push(format!(
614                "{} PID(s) could not be inspected",
615                census_uninspectable_pids.len()
616            ));
617        }
618        let reason = format!(
619            "OS holder census is incomplete: {}; additional database holders cannot be ruled out",
620            causes.join("; ")
621        );
622        WalPinCensus::Incomplete {
623            holder_pids: census_holder_pids.clone(),
624            uninspectable_pids: census_uninspectable_pids.clone(),
625            truncated: census_truncated,
626            reason,
627        }
628    };
629
630    let now_epoch_secs = SystemTime::now()
631        .duration_since(UNIX_EPOCH)
632        .map(|duration| duration.as_secs() as i64)
633        .unwrap_or(0);
634    let sidecar_listing_truncated = sidecar.sidecar_listing_truncated;
635    let sidecar_entries_cleanup_would_reap = sidecar.cleanup_would_reap;
636    let mut reporting = Vec::new();
637    let mut registered_silent_pids = Vec::new();
638    let mut unknown_pids = Vec::new();
639    let mut sidecar_entries = Vec::new();
640    let mut sidecar_known_pids = BTreeSet::new();
641
642    for entry in sidecar.entries {
643        match entry {
644            crate::walpin::WalpinPidHealth::Reporting(heartbeat) => {
645                let current_oldest_tx_age_secs =
646                    heartbeat.current_oldest_tx_age_secs(now_epoch_secs);
647                let attribution_is_evidence_backed = heartbeat.attribution_is_evidence_backed();
648                sidecar_known_pids.insert(heartbeat.pid);
649                reporting.push(WalPinHolder {
650                    pid: heartbeat.pid,
651                    process_role: heartbeat.process_role.clone(),
652                    current_oldest_tx_age_secs,
653                    oldest_tx_label: heartbeat.oldest_tx_label.clone(),
654                    attribution_is_evidence_backed,
655                });
656                sidecar_entries.push((
657                    heartbeat.pid,
658                    0u8,
659                    serde_json::json!({
660                        "pid": heartbeat.pid,
661                        "status": "reporting",
662                        "process_role": heartbeat.process_role,
663                        "current_oldest_tx_age_secs": current_oldest_tx_age_secs,
664                        "oldest_tx_label": heartbeat.oldest_tx_label,
665                        "attribution_is_evidence_backed": attribution_is_evidence_backed,
666                    }),
667                ));
668            }
669            crate::walpin::WalpinPidHealth::RegisteredSilent { pid } => {
670                sidecar_known_pids.insert(pid);
671                registered_silent_pids.push(pid);
672                sidecar_entries.push((
673                    pid,
674                    1u8,
675                    serde_json::json!({"pid": pid, "status": "registered_silent"}),
676                ));
677            }
678            crate::walpin::WalpinPidHealth::Unknown { pid, reason } => {
679                sidecar_known_pids.insert(pid);
680                unknown_pids.push(pid);
681                sidecar_entries.push((
682                    pid,
683                    2u8,
684                    serde_json::json!({"pid": pid, "status": "unknown", "reason": reason}),
685                ));
686            }
687        }
688    }
689
690    reporting.sort_by_key(|holder| holder.pid);
691    reporting.dedup_by_key(|holder| holder.pid);
692    registered_silent_pids.sort_unstable();
693    registered_silent_pids.dedup();
694    unknown_pids.sort_unstable();
695    unknown_pids.dedup();
696    sidecar_entries.sort_by_key(|(pid, status_rank, _)| (*pid, *status_rank));
697    let sidecar_entries = sidecar_entries
698        .into_iter()
699        .map(|(_, _, entry)| entry)
700        .collect();
701    let census_pids_without_attribution: Vec<u32> = census_holder_pids
702        .iter()
703        .copied()
704        .filter(|pid| !sidecar_known_pids.contains(pid))
705        .collect();
706
707    let mut status_reasons = Vec::new();
708    if let WalPinCensus::Incomplete { reason, .. } = &census_carrier {
709        status_reasons.push(reason.clone());
710    }
711    if sidecar_listing_truncated {
712        status_reasons.push(
713            "read-only sidecar enumeration reached its entry cap; additional entries may exist"
714                .to_string(),
715        );
716    }
717    if !unknown_pids.is_empty() {
718        status_reasons.push(format!(
719            "{} sidecar PID(s) could not be classified conclusively",
720            unknown_pids.len()
721        ));
722    }
723    if !census_pids_without_attribution.is_empty() {
724        status_reasons.push(format!(
725            "{} OS-confirmed holder(s) have no sidecar attribution",
726            census_pids_without_attribution.len()
727        ));
728    }
729
730    let fully_attributed = census_is_complete
731        && !sidecar_listing_truncated
732        && unknown_pids.is_empty()
733        && census_pids_without_attribution.is_empty();
734    let status = if fully_attributed {
735        WalPinAttributionStatus::Complete
736    } else {
737        WalPinAttributionStatus::Degraded
738    };
739    let unavailable_reason = (!fully_attributed).then(|| status_reasons.join("; "));
740
741    WalPinAttribution {
742        status,
743        status_reasons,
744        census: census_carrier,
745        reporting_pid,
746        reporting_process_is_holder,
747        reporting_process_is_holder_unavailable_reason,
748        census_process_start_times,
749        start_time_resolution_secs,
750        start_time_resolution_unavailable_reason: start_time_resolution_secs
751            .is_none()
752            .then(|| "process start time is unavailable on this platform".to_string()),
753        available: fully_attributed,
754        unavailable_reason,
755        census_holder_pids,
756        census_uninspectable_pids,
757        census_truncated,
758        census_is_complete,
759        reporting,
760        registered_silent_pids,
761        unknown_pids,
762        census_pids_without_attribution,
763        fully_attributed,
764        sidecar_entries,
765        sidecar_listing_truncated: Some(sidecar_listing_truncated),
766        sidecar_entries_cleanup_would_reap: Some(sidecar_entries_cleanup_would_reap),
767    }
768}
769
770/// Build the WAL-pin attribution for `db_path`.
771///
772/// Unix-only: the OS census requires it. Everywhere else this degrades to
773/// `available: false` with a reason rather than failing the whole report.
774/// An operator who has explicitly disabled the sidecar (`KHIVE_WALPIN_SIDECAR`)
775/// also disables this request's sidecar collection, per ADR-091 Amendment 6 —
776/// the census still runs, but there is no sidecar evidence to reconcile it
777/// against.
778#[cfg(unix)]
779pub fn wal_pin_attribution(db_path: &Path, sweep_interval: Duration) -> WalPinAttribution {
780    use crate::walpin;
781
782    let census = match walpin::census_holders(db_path) {
783        Ok(c) => c,
784        Err(e) => return WalPinAttribution::unavailable(format!("census_holders failed: {e}")),
785    };
786    if !walpin::sidecar_enabled(true) {
787        return wal_pin_attribution_without_sidecar(census, SIDECAR_DISABLED_REASON.to_string());
788    }
789    match walpin::inspect_live(&walpin::sidecar_dir_for(db_path), sweep_interval) {
790        Ok(sidecar) => wal_pin_attribution_from_evidence(census, sidecar),
791        Err(error) => wal_pin_attribution_without_sidecar(
792            census,
793            format!("read-only sidecar enumeration failed: {error}"),
794        ),
795    }
796}
797
798/// Shared with the async collection path in `inspect_file_state_interruptibly`.
799#[cfg(unix)]
800const SIDECAR_DISABLED_REASON: &str =
801    "walpin sidecar is explicitly disabled (KHIVE_WALPIN_SIDECAR); attribution has no sidecar \
802     evidence to reconcile against the OS holder census";
803
804#[cfg(not(unix))]
805pub fn wal_pin_attribution(_db_path: &Path, _sweep_interval: Duration) -> WalPinAttribution {
806    WalPinAttribution::unavailable("WAL-pin attribution requires a Unix platform")
807}
808
809/// One typed snapshot of reader route, saturation, and hold-lifecycle signals.
810///
811/// Every field is pool-scoped. Monotonic counters reset only when the
812/// [`ConnectionPool`] is reconstructed; the active value is point-in-time.
813/// Infrastructure standalone opens are kept separate from request traffic so
814/// a boot/schema probe cannot make a hot path look like it churned readers.
815#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
816pub struct ReaderContentionDiagnostics {
817    /// Reader connections requested in this pool's configuration. The
818    /// effective admission capacity below can be smaller in degraded mode.
819    pub configured_reader_cap: usize,
820    /// Configured wait for a reader slot, in milliseconds.
821    pub configured_checkout_timeout_ms: u64,
822    /// Configured SQLite busy-handler wait, in milliseconds. This is separate
823    /// from the reader-slot checkout deadline.
824    pub configured_busy_timeout_ms: u64,
825    /// Configured total reader admission budget shared by pooled readers and
826    /// explicit raw-SQL read transactions.
827    pub reader_admission_capacity: usize,
828    /// Point-in-time permits not held when this snapshot was captured.
829    pub available_reader_admission_slots: usize,
830    /// Successful request-path acquisitions (pooled plus enumerated
831    /// standalone request exceptions; infrastructure excluded).
832    pub reader_acquisitions: u64,
833    /// Successful bounded pooled-reader checkouts.
834    pub pooled_reader_checkouts: u64,
835    /// Successful request-path standalone reader opens. Under ADR-165 Slice
836    /// 2 this is limited to the explicit raw-SQL read-transaction exception.
837    pub standalone_reader_opens: u64,
838    /// Successful standalone opens owned by an enumerated boot/diagnostic
839    /// infrastructure exception.
840    pub infrastructure_standalone_reader_opens: u64,
841    /// Pool-wide reader-admission waits that exhausted `checkout_timeout`
842    /// before work began. Cooperative request cancellation is excluded.
843    pub reader_checkout_timeouts: u64,
844    /// Queries on a checked-out pooled reader that SQLite refused with
845    /// `SQLITE_BUSY` after `configured_busy_timeout_ms` elapsed. Counted after
846    /// checkout succeeded, so it is disjoint from `reader_checkout_timeouts`.
847    /// Writer refusals are not included; see `direct_writer_busy_refusals` and
848    /// `writer_task_begin_busy`.
849    pub reader_busy_timeouts: u64,
850    /// Pooled reader guards live when the snapshot was captured.
851    pub active_pooled_reader_checkouts: u64,
852    /// Highest observed concurrent pooled-reader guard count.
853    pub peak_active_pooled_reader_checkouts: u64,
854    /// Pooled guards that completed return/reset.
855    pub completed_pooled_reader_checkouts: u64,
856    /// Longest completed hold, including return/reset, in microseconds.
857    pub max_completed_reader_hold_micros: u64,
858    /// The typed-store operation that held the checkout this maximum came
859    /// from. `None` means that hold came through a route carrying no
860    /// operation name — the pool's own internal checkout — which is the
861    /// answer to "which read was it", not a gap in the reading (#2793).
862    pub max_completed_reader_hold_operation: Option<&'static str>,
863    /// A disqualified pooled-reader return whose replacement connection then
864    /// also failed to open, permanently shrinking the physical pool by one
865    /// slot below `max_readers`. Non-zero here means the pool has fewer
866    /// physical reader connections than configured.
867    pub reader_replacement_open_failures: u64,
868}
869
870impl ReaderContentionDiagnostics {
871    fn snapshot(pool: &ConnectionPool) -> Self {
872        let reader = pool.reader_acquisition_snapshot();
873        Self {
874            configured_reader_cap: pool.config().max_readers,
875            configured_checkout_timeout_ms: u64::try_from(
876                pool.config().checkout_timeout.as_millis(),
877            )
878            .unwrap_or(u64::MAX),
879            configured_busy_timeout_ms: u64::try_from(pool.config().busy_timeout.as_millis())
880                .unwrap_or(u64::MAX),
881            reader_admission_capacity: reader.reader_admission_capacity,
882            available_reader_admission_slots: reader.available_reader_admission_slots,
883            reader_acquisitions: reader.acquisitions,
884            pooled_reader_checkouts: reader.pooled_checkouts,
885            standalone_reader_opens: reader.standalone_opens,
886            infrastructure_standalone_reader_opens: reader.infrastructure_standalone_opens,
887            reader_checkout_timeouts: reader.checkout_timeouts,
888            reader_busy_timeouts: reader.busy_timeouts,
889            active_pooled_reader_checkouts: reader.active_pooled_checkouts,
890            peak_active_pooled_reader_checkouts: reader.peak_active_pooled_checkouts,
891            completed_pooled_reader_checkouts: reader.completed_pooled_checkouts,
892            max_completed_reader_hold_micros: reader.max_completed_hold_micros,
893            max_completed_reader_hold_operation: reader.max_completed_hold_operation,
894            reader_replacement_open_failures: reader.reader_replacement_open_failures,
895        }
896    }
897}
898
899/// One typed snapshot of writer-contention signals.
900///
901/// `writer_acquisitions` is the aggregate of the three explicit connection
902/// classes below. `writer_acquisition_timeouts` remains specific to the
903/// finite-wait pool-mutex stage; standalone SQLite failures and writer-task
904/// `BEGIN` failures have different ADR-135 F6 stages and are not mislabeled as
905/// pool checkout timeouts. Those stages now carry their OWN failure counters
906/// (`writer_task_begin_busy`, `writer_task_begin_busy_absorbed`,
907/// `writer_task_begin_errors`) rather than being absent: refusing to mislabel
908/// a failure is not a reason to omit it, and an omitted failure counter fails
909/// toward looking healthy, which is the reading an operator believes.
910/// `audit_append_failures` is supplied by the runtime
911/// because the audit store lives above `khive-db`; direct `khive-db` callers
912/// receive `None` plus an explicit reason instead of a fabricated zero.
913#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
914pub struct WriterContentionDiagnostics {
915    /// Successful acquisitions across pooled, standalone, and writer-task
916    /// connection classes.
917    pub writer_acquisitions: u64,
918    /// Successful finite-wait main-pool mutex checkouts.
919    pub pooled_writer_acquisitions: u64,
920    /// Successful per-operation file-backed standalone writer opens.
921    pub standalone_writer_acquisitions: u64,
922    /// Successful writer-task ownership acquisitions.
923    pub writer_task_acquisitions: u64,
924    /// Main-pool writer checkouts that exhausted their finite deadline.
925    pub writer_acquisition_timeouts: u64,
926    /// Writer acquisitions refused because another writer held the volume
927    /// lease past the guard deadline (`CapacityUnavailable`, phase `lock`).
928    /// Every writer class takes the lease before its connection, so ordinary
929    /// same-volume writer contention is counted here, not in
930    /// `writer_acquisition_timeouts`.
931    pub writer_lease_timeouts: u64,
932    /// The guard deadline a writer waits for the volume lease, in ms; `None`
933    /// when the pool takes no lease (in-memory or read-only).
934    pub configured_guard_deadline_ms: Option<u64>,
935    /// `checkout_timeout`, which bounds only the pool-mutex wait that comes
936    /// after the lease.
937    pub configured_checkout_timeout_ms: u64,
938    /// The effective pooled-writer wait bound under contention: the guard
939    /// deadline for the lease, then `checkout_timeout` for the pool mutex. A
940    /// pool configured with a 50 ms `checkout_timeout` and the default 2000 ms
941    /// guard deadline can wait about 2050 ms before refusing.
942    pub effective_writer_wait_bound_ms: u64,
943    /// Final instrumented direct execution refusals retaining primary SQLITE_BUSY,
944    /// once per operation; excludes SQLITE_LOCKED and task/reader/open/admission errors.
945    pub direct_writer_busy_refusals: u64,
946    /// Every writer-task `BEGIN IMMEDIATE` attempt refused busy or locked,
947    /// whether or not a subsequent bounded retry absorbed it. A refusal not
948    /// absorbed by a retry also surfaces to the caller as the retryable
949    /// `writer_task_begin_busy` stage.
950    pub writer_task_begin_busy: u64,
951    /// Subset of `writer_task_begin_busy` absorbed by the bounded retry
952    /// before the request closure ran, so the caller never observed that
953    /// particular refusal. `writer_task_begin_busy - writer_task_begin_busy_absorbed`
954    /// is the count of refusals a caller actually observed.
955    pub writer_task_begin_busy_absorbed: u64,
956    /// Writer-task `BEGIN IMMEDIATE` attempts that failed for a reason other
957    /// than busy or locked.
958    pub writer_task_begin_errors: u64,
959    /// Dequeued writer-task requests that reached the writer seam and
960    /// returned error, counted once per request. Sourced directly from the
961    /// pool's own acquisition-site counters, so — unlike the runtime-supplied
962    /// fields below — it is populated identically for every caller.
963    pub writer_task_request_failures: u64,
964    /// Subset of `writer_task_request_failures` whose terminal state was
965    /// `WriterTaskRequestState::SideEffectsUnknown`.
966    pub writer_task_side_effects_unknown: u64,
967    /// Process-wide audit appends whose errors were logged and swallowed —
968    /// pure-observability rows only (config-lock rows, best-effort recall
969    /// telemetry). An obligation-bearing row's commit failure (a gate
970    /// denial's own audit row, a dispatch outcome, an unknown-verb row, or a
971    /// `git.digest` receipt) is never counted here: those either fail the
972    /// dispatch that produced them directly (visible to the caller as an
973    /// error, not as this counter moving) or, for a denial whose dispatch
974    /// already fails independent of the row, are tracked by the runtime's
975    /// own separate obligation-failure counter instead. Summing this field
976    /// with `audit_batch_flush_failures` therefore does not double-count an
977    /// obligation-bearing generation failure against this one.
978    pub audit_append_failures: Option<u64>,
979    /// Why `audit_append_failures` is unavailable to this caller.
980    pub audit_append_failures_unavailable_reason: Option<String>,
981    /// Process-wide obligation-bearing audit commit failures, including a
982    /// denial's audit failure even though the denial already refuses the call.
983    /// This counts failed producer submissions, not failed batch generations;
984    /// do not add it to `audit_batch_flush_failures` as a disjoint total.
985    pub audit_obligation_append_failures: Option<u64>,
986    /// Why the runtime-owned obligation counter is unavailable to this caller.
987    pub audit_obligation_append_failures_unavailable_reason: Option<String>,
988    /// Accepted audit-batch generations that reached a terminal non-commit
989    /// outcome after retry, including driver death; excludes preflight and
990    /// admission rejection. `None` for a direct `khive-db` caller, or when no
991    /// runtime audit-batch control has been wired into diagnostics.
992    pub audit_batch_flush_failures: Option<u64>,
993    /// Why `audit_batch_flush_failures` is unavailable to this caller.
994    pub audit_batch_flush_failures_unavailable_reason: Option<String>,
995    /// Pure-observability audit rows released without a commit. `None` under
996    /// the same conditions as `audit_batch_flush_failures`.
997    pub audit_degraded_rows: Option<u64>,
998    /// Why `audit_degraded_rows` is unavailable to this caller.
999    pub audit_degraded_rows_unavailable_reason: Option<String>,
1000    /// Monotonic process-lifetime flag set once any accepted row has been
1001    /// released without a commit: a pure-observability row released degraded,
1002    /// or a generation that failed to flush, whose rows are absent from the
1003    /// audit trail whether or not their callers were told. `None` under the
1004    /// same conditions as `audit_batch_flush_failures`.
1005    pub audit_degraded: Option<bool>,
1006    /// Why `audit_degraded` is unavailable to this caller.
1007    pub audit_degraded_unavailable_reason: Option<String>,
1008    /// Per-dispatch audit rows for an explicitly allowlisted, domain-write-free
1009    /// read verb (`VerbRegistry::ADMISSION_DEGRADE_SAFE_VERBS`) that were
1010    /// **refused before they could be enqueued** on the audit lane
1011    /// (`AuditTerminalReason::QueueAdmissionExhausted`) while the dispatch
1012    /// still reported its own successful result (ADR-103 Amendment 3, ADR-133
1013    /// Amendment 1). This is a confirmed, terminal accounting loss: the row
1014    /// never shared a generation with anyone and will never commit, so it
1015    /// undercounts `brain.event_counts`'s cost totals for exactly the rows
1016    /// counted here. Disjoint from `audit_degraded_rows` (a different reason:
1017    /// persistent commit failure of a pure-observability row, not admission
1018    /// pressure) and from `audit_admission_unresolved_obligations` (a row that
1019    /// was enqueued and may still commit). `None` under the same conditions as
1020    /// `audit_batch_flush_failures`.
1021    ///
1022    /// CUMULATIVE since process start. Nothing decrements it, so a value that
1023    /// holds steady under traffic means no refusal occurred in that window
1024    /// rather than a stalled subsystem (#2791); read
1025    /// `audit_admission_refused_obligations_last_at_ms` beside it to tell the
1026    /// two apart.
1027    pub audit_admission_refused_obligations: Option<u64>,
1028    /// Wall-clock milliseconds at which `audit_admission_refused_obligations`
1029    /// last moved in the serving process, or `None` if it has never moved.
1030    /// `None` alongside a count of zero is the ordinary quiet case; `None`
1031    /// alongside a non-zero count cannot occur and would indicate the two are
1032    /// being produced from different processes.
1033    pub audit_admission_refused_obligations_last_at_ms: Option<u64>,
1034    /// Why `audit_admission_refused_obligations` is unavailable to this caller.
1035    pub audit_admission_refused_obligations_unavailable_reason: Option<String>,
1036    /// Per-dispatch audit rows for an explicitly allowlisted, domain-write-free
1037    /// read verb (`VerbRegistry::ADMISSION_DEGRADE_SAFE_VERBS`) that were
1038    /// **already enqueued but had not resolved when the caller's admission
1039    /// wait deadline elapsed** (`AuditTerminalReason::AdmissionDeadlineExpired`)
1040    /// while the dispatch still reported its own successful result (ADR-103
1041    /// Amendment 3, ADR-133 Amendment 1). Unlike
1042    /// `audit_admission_refused_obligations`, a row counted here is not a
1043    /// confirmed loss — it may still be committed by the generation driver
1044    /// independently of the caller's timeout — so this field is an upper
1045    /// bound on the eventual undercount, not the undercount itself. `None`
1046    /// under the same conditions as `audit_batch_flush_failures`.
1047    ///
1048    /// CUMULATIVE since process start, despite the set-shaped name. There is no
1049    /// live set of unresolved obligations and nothing resolves this counter:
1050    /// each increment records one past deadline expiry, and the row it counted
1051    /// most likely committed afterwards. A steady value under traffic means no
1052    /// deadline expired in that window, which is the healthy reading (#2791).
1053    /// Read `audit_admission_unresolved_obligations_last_at_ms` beside it: a
1054    /// non-zero count whose mark is old is history, the same count with a
1055    /// recent mark is an active condition.
1056    pub audit_admission_unresolved_obligations: Option<u64>,
1057    /// Wall-clock milliseconds at which
1058    /// `audit_admission_unresolved_obligations` last moved in the serving
1059    /// process, or `None` if it has never moved.
1060    pub audit_admission_unresolved_obligations_last_at_ms: Option<u64>,
1061    /// Why `audit_admission_unresolved_obligations` is unavailable to this
1062    /// caller.
1063    pub audit_admission_unresolved_obligations_unavailable_reason: Option<String>,
1064}
1065
1066/// Process-wide audit-batch health counters, supplied by the runtime layer
1067/// that owns the audit-batch control. `khive-db` never produces these itself
1068/// — a direct `khive-db` caller always sees the corresponding
1069/// `WriterContentionDiagnostics` fields as `None` plus an explicit reason.
1070#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1071pub struct RuntimeAuditBatchMetrics {
1072    /// Accepted generations reaching terminal non-commit after retry,
1073    /// including driver death; excludes preflight/admission rejection.
1074    pub flush_failures: u64,
1075    /// Pure-observability rows released without commit.
1076    pub degraded_rows: u64,
1077    /// Monotonic process-lifetime degradation flag.
1078    pub degraded: bool,
1079    /// Admission-degrade-safe read verbs' audit rows refused before enqueue
1080    /// under transient audit-lane admission pressure (ADR-103 Amendment 3,
1081    /// ADR-133 Amendment 1) — a confirmed, terminal accounting loss. Disjoint
1082    /// from `degraded_rows` and from `admission_unresolved_obligations`.
1083    pub admission_refused_obligations: u64,
1084    /// Wall-clock ms at which `admission_refused_obligations` last moved;
1085    /// `None` until it moves. A total cannot say when it was last earned.
1086    pub admission_refused_obligations_last_at_ms: Option<u64>,
1087    /// Admission-degrade-safe read verbs' audit rows that were already
1088    /// enqueued but had not resolved when the caller's admission wait
1089    /// deadline elapsed (ADR-103 Amendment 3, ADR-133 Amendment 1). Not a
1090    /// confirmed loss — the row may still commit — so this is an upper bound
1091    /// on the eventual undercount. Disjoint from `degraded_rows` and from
1092    /// `admission_refused_obligations`.
1093    pub admission_unresolved_obligations: u64,
1094    /// Wall-clock ms at which `admission_unresolved_obligations` last moved;
1095    /// `None` until it moves.
1096    pub admission_unresolved_obligations_last_at_ms: Option<u64>,
1097}
1098
1099/// Live graph-edge rows compared with the durable list-cursor ledger.
1100///
1101/// `duplicate_edge_id_groups > 0` is the exact pre-V14 state in which two
1102/// namespaces share an edge UUID and a multi-namespace cursor walk can drop
1103/// one row during ID-based deduplication. The two raw row counts are reported
1104/// separately because sequence rows intentionally survive hard deletion;
1105/// count inequality by itself is therefore not proof of corruption.
1106///
1107/// `live_entities_carrying_merged_into` counts entity rows that are live
1108/// (`deleted_at IS NULL`) while still carrying merge provenance. A merge
1109/// tombstones its source; a restore that ran before restore refused merge
1110/// tombstones cleared the tombstone and left the provenance, so such rows
1111/// read as merged by `get` and as live by `list` and `search`. Restore now
1112/// names them as `live_merged_entity`; this count is where an operator finds
1113/// them across all namespaces.
1114#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1115pub struct GraphEdgeIntegrity {
1116    pub duplicate_edge_id_groups: i64,
1117    pub graph_edges_rows: i64,
1118    pub graph_edges_seq_rows: i64,
1119    pub pre_v14_duplicate_edge_state_detected: bool,
1120    pub live_entities_carrying_merged_into: i64,
1121}
1122
1123fn graph_edge_integrity(conn: &Connection) -> rusqlite::Result<GraphEdgeIntegrity> {
1124    conn.query_row(
1125        "SELECT
1126             (SELECT COUNT(*) FROM (
1127                 SELECT id FROM graph_edges GROUP BY id HAVING COUNT(*) > 1
1128             )),
1129             (SELECT COUNT(*) FROM graph_edges),
1130             (SELECT COUNT(*) FROM graph_edges_seq),
1131             (SELECT COUNT(*) FROM entities
1132              WHERE deleted_at IS NULL AND merged_into IS NOT NULL)",
1133        [],
1134        |row| {
1135            let duplicate_edge_id_groups = row.get(0)?;
1136            Ok(GraphEdgeIntegrity {
1137                duplicate_edge_id_groups,
1138                graph_edges_rows: row.get(1)?,
1139                graph_edges_seq_rows: row.get(2)?,
1140                pre_v14_duplicate_edge_state_detected: duplicate_edge_id_groups > 0,
1141                live_entities_carrying_merged_into: row.get(3)?,
1142            })
1143        },
1144    )
1145}
1146
1147const MAX_DATABASE_SIZE_OBJECTS: usize = 4_096;
1148
1149/// SQLite b-tree role reported by the size-composition diagnostic.
1150#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1151#[serde(rename_all = "snake_case")]
1152pub enum DatabaseObjectKind {
1153    Table,
1154    Index,
1155    Internal,
1156}
1157
1158/// Operational storage grouping. `mixed_row_and_embedding` is deliberately
1159/// separate: SQLite cannot attribute bytes within a table page to one column,
1160/// so counting all of `knowledge_sections` as pure vector bytes would be a
1161/// false precision claim.
1162#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1163#[serde(rename_all = "snake_case")]
1164pub enum DatabaseStorageClass {
1165    RowTable,
1166    Index,
1167    FullText,
1168    Vector,
1169    MixedRowAndEmbedding,
1170    Internal,
1171}
1172
1173#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
1174pub struct DatabaseObjectSize {
1175    pub name: String,
1176    pub owner_table: Option<String>,
1177    pub object_kind: DatabaseObjectKind,
1178    pub storage_class: DatabaseStorageClass,
1179    pub pages: u64,
1180    pub bytes: u64,
1181}
1182
1183/// Page-accounted file-size composition from SQLite's read-only `dbstat`
1184/// virtual table in aggregate mode.
1185#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
1186pub struct DatabaseSizeComposition {
1187    pub page_size_bytes: u64,
1188    pub page_count: u64,
1189    pub freelist_pages: u64,
1190    pub database_bytes: u64,
1191    pub freelist_bytes: u64,
1192    pub accounted_bytes: u64,
1193    pub unaccounted_bytes: u64,
1194    pub row_table_bytes: u64,
1195    pub index_bytes: u64,
1196    pub full_text_bytes: u64,
1197    pub vector_bytes: u64,
1198    pub mixed_embedding_bytes: u64,
1199    pub internal_bytes: u64,
1200    pub objects: Vec<DatabaseObjectSize>,
1201    pub objects_truncated: bool,
1202    pub objects_omitted: usize,
1203}
1204
1205fn nonnegative_sqlite_integer(column: usize, value: i64) -> rusqlite::Result<u64> {
1206    u64::try_from(value).map_err(|_| rusqlite::Error::IntegralValueOutOfRange(column, value))
1207}
1208
1209fn declares_embedding_blob(sql: Option<&str>) -> bool {
1210    let Some(sql) = sql else {
1211        return false;
1212    };
1213    let tokens: Vec<_> = sql
1214        .split(|character: char| !(character.is_ascii_alphanumeric() || character == '_'))
1215        .filter(|token| !token.is_empty())
1216        .collect();
1217    tokens.windows(2).any(|pair| {
1218        pair[0].eq_ignore_ascii_case("embedding") && pair[1].eq_ignore_ascii_case("blob")
1219    })
1220}
1221
1222fn classify_database_object(
1223    name: &str,
1224    sqlite_type: &str,
1225    sql: Option<&str>,
1226) -> (DatabaseObjectKind, DatabaseStorageClass) {
1227    let object_kind = match sqlite_type {
1228        "table" => DatabaseObjectKind::Table,
1229        "index" => DatabaseObjectKind::Index,
1230        _ => DatabaseObjectKind::Internal,
1231    };
1232    let lower_name = name.to_ascii_lowercase();
1233    let lower_sql = sql.unwrap_or_default().to_ascii_lowercase();
1234    let storage_class = if lower_name.starts_with("fts_") || lower_sql.contains("using fts5") {
1235        DatabaseStorageClass::FullText
1236    } else if lower_name.starts_with("vec_")
1237        || lower_name == "_embedding_models"
1238        || lower_sql.contains("using vec0")
1239    {
1240        DatabaseStorageClass::Vector
1241    } else if declares_embedding_blob(sql) {
1242        DatabaseStorageClass::MixedRowAndEmbedding
1243    } else if object_kind == DatabaseObjectKind::Index {
1244        DatabaseStorageClass::Index
1245    } else if name.starts_with("sqlite_") || object_kind == DatabaseObjectKind::Internal {
1246        DatabaseStorageClass::Internal
1247    } else {
1248        DatabaseStorageClass::RowTable
1249    };
1250    (object_kind, storage_class)
1251}
1252
1253fn database_size_composition(conn: &Connection) -> rusqlite::Result<DatabaseSizeComposition> {
1254    let page_size = nonnegative_sqlite_integer(
1255        0,
1256        conn.query_row("PRAGMA page_size", [], |row| row.get::<_, i64>(0))?,
1257    )?;
1258    let page_count = nonnegative_sqlite_integer(
1259        0,
1260        conn.query_row("PRAGMA page_count", [], |row| row.get::<_, i64>(0))?,
1261    )?;
1262    let freelist_pages = nonnegative_sqlite_integer(
1263        0,
1264        conn.query_row("PRAGMA freelist_count", [], |row| row.get::<_, i64>(0))?,
1265    )?;
1266
1267    let mut statement = conn.prepare(
1268        "SELECT d.name, COALESCE(s.type, 'internal'), s.tbl_name, s.sql, d.pageno, d.pgsize
1269         FROM dbstat AS d
1270         LEFT JOIN sqlite_schema AS s ON s.name = d.name
1271         WHERE d.aggregate = TRUE
1272         ORDER BY d.name",
1273    )?;
1274    let mut rows = statement.query([])?;
1275    let mut objects = Vec::new();
1276    let mut objects_omitted = 0usize;
1277    let mut accounted_bytes = 0u64;
1278    let mut row_table_bytes = 0u64;
1279    let mut index_bytes = 0u64;
1280    let mut full_text_bytes = 0u64;
1281    let mut vector_bytes = 0u64;
1282    let mut mixed_embedding_bytes = 0u64;
1283    let mut internal_bytes = 0u64;
1284
1285    while let Some(row) = rows.next()? {
1286        let name: String = row.get(0)?;
1287        let sqlite_type: String = row.get(1)?;
1288        let owner_table: Option<String> = row.get(2)?;
1289        let sql: Option<String> = row.get(3)?;
1290        let pages = nonnegative_sqlite_integer(4, row.get(4)?)?;
1291        let bytes = nonnegative_sqlite_integer(5, row.get(5)?)?;
1292        let (object_kind, storage_class) =
1293            classify_database_object(&name, &sqlite_type, sql.as_deref());
1294        accounted_bytes = accounted_bytes.saturating_add(bytes);
1295        let class_total = match storage_class {
1296            DatabaseStorageClass::RowTable => &mut row_table_bytes,
1297            DatabaseStorageClass::Index => &mut index_bytes,
1298            DatabaseStorageClass::FullText => &mut full_text_bytes,
1299            DatabaseStorageClass::Vector => &mut vector_bytes,
1300            DatabaseStorageClass::MixedRowAndEmbedding => &mut mixed_embedding_bytes,
1301            DatabaseStorageClass::Internal => &mut internal_bytes,
1302        };
1303        *class_total = class_total.saturating_add(bytes);
1304
1305        if objects.len() < MAX_DATABASE_SIZE_OBJECTS {
1306            objects.push(DatabaseObjectSize {
1307                name,
1308                owner_table,
1309                object_kind,
1310                storage_class,
1311                pages,
1312                bytes,
1313            });
1314        } else {
1315            objects_omitted = objects_omitted.saturating_add(1);
1316        }
1317    }
1318
1319    let database_bytes = page_count.saturating_mul(page_size);
1320    let freelist_bytes = freelist_pages.saturating_mul(page_size);
1321    let unaccounted_bytes = database_bytes
1322        .saturating_sub(freelist_bytes)
1323        .saturating_sub(accounted_bytes);
1324    Ok(DatabaseSizeComposition {
1325        page_size_bytes: page_size,
1326        page_count,
1327        freelist_pages,
1328        database_bytes,
1329        freelist_bytes,
1330        accounted_bytes,
1331        unaccounted_bytes,
1332        row_table_bytes,
1333        index_bytes,
1334        full_text_bytes,
1335        vector_bytes,
1336        mixed_embedding_bytes,
1337        internal_bytes,
1338        objects,
1339        objects_truncated: objects_omitted > 0,
1340        objects_omitted,
1341    })
1342}
1343
1344/// What this report cost to assemble, per section, in milliseconds.
1345///
1346/// `db_diagnostics` is the verb an operator reaches for when the system feels
1347/// slow, and its own cost was measured varying more than five-fold between
1348/// consecutive calls against one unchanged store on one box. A surface that
1349/// varies that much has to say which part varied; otherwise the reader is left
1350/// to guess, and the cheapest guess is "the store is stalled", which is the
1351/// conclusion this verb exists to test rather than to suggest.
1352///
1353/// Every field is wall-clock and measured in this process, so it includes time
1354/// spent waiting on a contended machine. That is deliberate: the contention is
1355/// the thing being reported.
1356#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1357pub struct CollectionCost {
1358    /// Whole-report assembly, from the first section to the last.
1359    pub total_ms: u64,
1360    /// SQLite-side inspection: the PASSIVE checkpoint probe, page composition,
1361    /// the graph-edge ledger counts and the FTS5 segment records.
1362    pub sqlite_ms: u64,
1363    /// `stat` of the WAL file. The cheap sibling of the census; reported so the
1364    /// comparison beside it is legible rather than asserted.
1365    pub wal_file_stat_ms: u64,
1366    /// The OS holder census: a walk of every process on the host.
1367    pub wal_pin_census_ms: u64,
1368    /// The sidecar directory enumeration the census is reconciled against.
1369    pub wal_pin_sidecar_ms: u64,
1370    /// The wall-clock budget the census was given, or `None` when it ran
1371    /// unbounded. An `Option` rather than a sentinel `0` because a producer
1372    /// that failed to record a budget would write `0`, which is exactly the
1373    /// value "unbounded" would have claimed.
1374    pub wal_pin_census_budget_ms: Option<u64>,
1375    /// Whether the census stopped because that budget was spent. When true the
1376    /// holder list is partial by design and `wal_pin.census.status` is
1377    /// `incomplete` for that reason and no other.
1378    pub wal_pin_census_budget_exhausted: bool,
1379}
1380
1381impl CollectionCost {
1382    /// A report whose file-backed sections never ran (an in-memory backend):
1383    /// the SQLite side still costs what it costs, and the census did not run
1384    /// at all rather than running fast.
1385    fn in_memory(total_ms: u64) -> Self {
1386        Self {
1387            total_ms,
1388            sqlite_ms: 0,
1389            wal_file_stat_ms: 0,
1390            wal_pin_census_ms: 0,
1391            wal_pin_sidecar_ms: 0,
1392            wal_pin_census_budget_ms: None,
1393            wal_pin_census_budget_exhausted: false,
1394        }
1395    }
1396}
1397
1398/// The default wall-clock budget for the OS holder census on the request path.
1399///
1400/// Sized from the measurements in the report that asked for the bound: five
1401/// consecutive samples on a 971-process box cost 3.4s to 18.6s, while every
1402/// contention counter the same report carried read zero. Two seconds keeps the
1403/// verb interactive on a busy machine and completes untruncated on an idle one.
1404const DEFAULT_CENSUS_BUDGET: Duration = Duration::from_millis(2000);
1405
1406/// Environment override for [`DEFAULT_CENSUS_BUDGET`], in milliseconds.
1407///
1408/// `0` disables the bound and restores the unbounded full-machine walk, for an
1409/// operator who would rather wait than read a partial holder list. An
1410/// unparseable value is ignored in favour of the default rather than failing
1411/// the request: a diagnostic verb that refuses to answer because its own
1412/// tuning knob is malformed is worse than one that answers on the default and
1413/// says which budget it used — which the report does, in
1414/// `collection_cost.wal_pin_census_budget_ms`.
1415const CENSUS_BUDGET_ENV: &str = "KHIVE_WALPIN_CENSUS_BUDGET_MS";
1416
1417fn request_census_budget() -> Option<Duration> {
1418    match std::env::var(CENSUS_BUDGET_ENV) {
1419        Ok(raw) => match raw.trim().parse::<u64>() {
1420            Ok(0) => None,
1421            Ok(ms) => Some(Duration::from_millis(ms)),
1422            Err(_) => Some(DEFAULT_CENSUS_BUDGET),
1423        },
1424        Err(_) => Some(DEFAULT_CENSUS_BUDGET),
1425    }
1426}
1427
1428/// The full database-integrity, reader/writer-contention, and WAL/checkpoint
1429/// payload.
1430#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1431pub struct WalCeilingDiagnostics {
1432    /// Resolved configuration value, retained even on a read-only backend.
1433    pub configured_bytes: u64,
1434    /// Active writer-policy limit; zero on read-only backends and when disabled.
1435    pub effective_bytes: u64,
1436    pub source: WalCeilingSource,
1437    pub enabled: bool,
1438    /// Why the configured policy is active or inactive for this backend.
1439    pub status: &'static str,
1440}
1441
1442impl WalCeilingDiagnostics {
1443    fn from_pool(pool: &ConnectionPool) -> Self {
1444        let policy = pool.config().wal_ceiling;
1445        let read_only = pool.config().read_only;
1446        let effective_bytes = policy.effective_bytes(read_only);
1447        Self {
1448            configured_bytes: policy.bytes,
1449            effective_bytes,
1450            source: policy.source,
1451            enabled: effective_bytes > 0,
1452            status: if effective_bytes > 0 {
1453                "enforced"
1454            } else if read_only && policy.bytes > 0 {
1455                "read_only_not_enforced"
1456            } else {
1457                "disabled"
1458            },
1459        }
1460    }
1461}
1462
1463#[derive(Debug, Clone, PartialEq, Serialize)]
1464pub struct DbDiagnostics {
1465    pub build: BuildIdentity,
1466    pub process: ProcessIdentity,
1467    /// Process-lifetime note-search vector route counts, independent of this
1468    /// report's database file and reset only with the serving process.
1469    pub note_search_ann_route_total: u64,
1470    pub note_search_fallback_route_total: u64,
1471    /// Pool-scoped, monotonic coordinator dispatch and note-candidate
1472    /// hydration counts used by ADR-166 G4/G5. Reset on pool reconstruction.
1473    pub search_mechanism: crate::pool::SearchMechanismSnapshot,
1474    /// `None` for an in-memory backend — the file-backed sections then carry
1475    /// their own unavailability reasons.
1476    pub db_path: Option<String>,
1477    /// Explicit WAL ceiling policy for this already-open database.
1478    pub wal_ceiling: WalCeilingDiagnostics,
1479    pub disk_guard: DiskGuardDiagnostics,
1480    pub wal_file: Option<WalFileState>,
1481    pub checkpoint_counters: CheckpointCounters,
1482    pub checkpoint_probe: Option<CheckpointProbe>,
1483    pub checkpoint_probe_error: Option<String>,
1484    #[serde(flatten)]
1485    /// Frame pin fields serialized beside `checkpoint_probe` at the report root.
1486    pub checkpoint_pin: CheckpointPinDiagnostics,
1487    /// Reader route, checkout saturation, and hold-lifecycle signals.
1488    pub reader_contention: ReaderContentionDiagnostics,
1489    /// Writer-pool and best-effort audit persistence signals.
1490    pub writer_contention: WriterContentionDiagnostics,
1491    pub size_composition: Option<DatabaseSizeComposition>,
1492    pub size_composition_error: Option<String>,
1493    pub graph_edge_integrity: Option<GraphEdgeIntegrity>,
1494    pub graph_edge_integrity_error: Option<String>,
1495    /// Current FTS5 segment structure, decoded from each index's documented
1496    /// one-row structure record (shadow-table row id 10). This does not scan
1497    /// either corpus or `%_idx` table.
1498    pub fts_segments: Option<crate::FtsSegmentDiagnostics>,
1499    pub fts_segments_error: Option<String>,
1500    /// Process-lifetime outcomes from the checkpoint task's bounded FTS5
1501    /// maintenance steps.
1502    pub fts_maintenance: crate::FtsMaintenanceCounters,
1503    pub wal_pin: WalPinAttribution,
1504    /// Per-section assembly cost for this report. See [`CollectionCost`].
1505    pub collection_cost: CollectionCost,
1506}
1507
1508/// Checkpoint-row evidence used to report a sustained backfill ceiling.
1509#[derive(Debug, Clone, PartialEq, Serialize)]
1510pub struct CheckpointPinDiagnostics {
1511    /// The probe row's checkpointed frame when it shows unbackfilled frames.
1512    pub backfill_ceiling: Option<i64>,
1513    /// Why the probe row did not establish a backfill ceiling.
1514    pub backfill_ceiling_unavailable_reason: Option<String>,
1515    /// The ceiling after a matching checkpoint run has lasted at least one second.
1516    pub oldest_pinned_frame: Option<i64>,
1517    /// Why the current probe did not establish an oldest pinned frame.
1518    pub oldest_pinned_frame_unavailable_reason: Option<String>,
1519    /// The first observation for the reported frame.
1520    pub oldest_pinned_frame_run: Option<crate::checkpoint::CheckpointRun>,
1521    /// Why the matching checkpoint run is unavailable.
1522    pub oldest_pinned_frame_run_unavailable_reason: Option<String>,
1523    /// Frames from the sustained pin to the current probe's log end.
1524    pub pin_depth: Option<i64>,
1525    /// Why this report cannot claim a pin depth.
1526    pub pin_depth_unavailable_reason: Option<String>,
1527}
1528
1529fn checkpoint_pin_diagnostics(
1530    pool: &ConnectionPool,
1531    probe: Option<&CheckpointProbe>,
1532    probe_error: Option<&str>,
1533) -> CheckpointPinDiagnostics {
1534    let (run_status, run_age) = checkpoint::checkpoint_run_snapshot(pool);
1535    checkpoint_pin_diagnostics_for_run(probe, probe_error, run_status, run_age)
1536}
1537
1538fn checkpoint_pin_diagnostics_for_run(
1539    probe: Option<&CheckpointProbe>,
1540    probe_error: Option<&str>,
1541    run_status: checkpoint::CheckpointRunStatus,
1542    run_age: Option<Duration>,
1543) -> CheckpointPinDiagnostics {
1544    let (backfill_ceiling, backfill_reason) = match probe {
1545        None => (
1546            None,
1547            Some(
1548                probe_error
1549                    .unwrap_or("checkpoint probe did not return a row")
1550                    .to_string(),
1551            ),
1552        ),
1553        Some(probe) if probe.busy != 0 => (
1554            None,
1555            Some(format!("PASSIVE checkpoint returned busy={}", probe.busy)),
1556        ),
1557        Some(probe) if probe.log_frames < 0 || probe.checkpointed_frames < 0 => (
1558            None,
1559            Some("PASSIVE checkpoint returned a negative frame count".to_string()),
1560        ),
1561        Some(probe) if probe.log_frames <= probe.checkpointed_frames => (
1562            None,
1563            Some("PASSIVE checkpoint found no frames beyond the backfill ceiling".to_string()),
1564        ),
1565        Some(probe) => (Some(probe.checkpointed_frames), None),
1566    };
1567
1568    let (oldest_pinned_frame, oldest_reason, oldest_pinned_frame_run, run_reason) =
1569        match backfill_ceiling {
1570            None => {
1571                let reason = backfill_reason
1572                    .clone()
1573                    .unwrap_or_else(|| "backfill ceiling is unavailable".to_string());
1574                (None, Some(reason.clone()), None, Some(reason))
1575            }
1576            Some(ceiling) => match run_status {
1577                checkpoint::CheckpointRunStatus::NoTask => {
1578                    let reason = "no checkpoint task in this process".to_string();
1579                    (None, Some(reason.clone()), None, Some(reason))
1580                }
1581                checkpoint::CheckpointRunStatus::NoObservation => {
1582                    let reason = "no checkpoint run has been observed for this backend".to_string();
1583                    (None, Some(reason.clone()), None, Some(reason))
1584                }
1585                checkpoint::CheckpointRunStatus::Observed(run) if run.frame != ceiling => {
1586                    let reason = "checkpoint run frame does not match the current backfill ceiling";
1587                    (
1588                        None,
1589                        Some(reason.to_string()),
1590                        None,
1591                        Some(reason.to_string()),
1592                    )
1593                }
1594                checkpoint::CheckpointRunStatus::Observed(_)
1595                    if run_age.is_none_or(|age| age < Duration::from_secs(1)) =>
1596                {
1597                    let reason = "checkpoint run has been observed for less than one second";
1598                    (
1599                        None,
1600                        Some(reason.to_string()),
1601                        None,
1602                        Some(reason.to_string()),
1603                    )
1604                }
1605                checkpoint::CheckpointRunStatus::Observed(run) => {
1606                    (Some(ceiling), None, Some(run), None)
1607                }
1608            },
1609        };
1610
1611    let pin_depth = oldest_pinned_frame.map(|frame| {
1612        probe
1613            .expect("a reported pin always has a probe row")
1614            .log_frames
1615            .saturating_sub(frame)
1616            .max(0)
1617    });
1618    let pin_depth_reason = if pin_depth.is_none() {
1619        oldest_reason.clone()
1620    } else {
1621        None
1622    };
1623
1624    CheckpointPinDiagnostics {
1625        backfill_ceiling,
1626        backfill_ceiling_unavailable_reason: backfill_reason,
1627        oldest_pinned_frame,
1628        oldest_pinned_frame_unavailable_reason: oldest_reason,
1629        oldest_pinned_frame_run,
1630        oldest_pinned_frame_run_unavailable_reason: run_reason,
1631        pin_depth,
1632        pin_depth_unavailable_reason: pin_depth_reason,
1633    }
1634}
1635
1636/// Assemble the report for `pool`'s database.
1637///
1638/// `db_path` in the returned report is the pool's own configured path, so
1639/// the report can never claim to describe a file the pool is not bound to.
1640/// Every operational probe — the WAL file, the sidecar, the OS holder census
1641/// — instead targets `pool.canonical_path()`, the same value `ConnectionPool`
1642/// and the checkpoint sidecar writers key off of; a symlinked or otherwise
1643/// aliased configured path would otherwise send those probes looking beside
1644/// the alias while the evidence sits beside the canonical file. An in-memory
1645/// pool has no path: its process- and pool-scoped counters are still real,
1646/// but every file-backed section degrades to an explicit "unavailable" with
1647/// a reason rather than being silently omitted.
1648///
1649/// The PASSIVE probe goes through the infrastructure-only untracked
1650/// standalone open, WITHOUT `SQLITE_OPEN_CREATE`, so a diagnostic request
1651/// against a missing file returns `checkpoint_probe: null` with a
1652/// `checkpoint_probe_error` instead of creating a database or incrementing
1653/// the write-traffic acquisition total. Running on a standalone connection
1654/// also keeps checkpoint I/O off the pooled writer mutex.
1655pub fn collect(
1656    pool: &ConnectionPool,
1657    build: BuildIdentity,
1658    sweep_interval: Duration,
1659) -> DbDiagnostics {
1660    collect_inner(pool, build, sweep_interval, None, None)
1661}
1662
1663/// Assemble the report with the runtime's process-wide count of swallowed
1664/// best-effort audit append failures.
1665pub fn collect_with_audit_append_failures(
1666    pool: &ConnectionPool,
1667    build: BuildIdentity,
1668    sweep_interval: Duration,
1669    audit_append_failures: u64,
1670) -> DbDiagnostics {
1671    collect_inner(
1672        pool,
1673        build,
1674        sweep_interval,
1675        Some(audit_append_failures),
1676        None,
1677    )
1678}
1679
1680/// Assemble diagnostics without allowing request abandonment to leave the
1681/// graph SELECT or OS holder census running detached.
1682///
1683/// The PASSIVE checkpoint is intentionally outside SQLite interruption: once
1684/// admitted it may backfill database pages and must reach its physical
1685/// completion. The following graph integrity SELECT registers the common
1686/// request progress callback on that exact connection. The OS census runs as
1687/// a second cooperative phase and polls a shared cancellation flag between
1688/// process and fd entries.
1689pub async fn collect_with_audit_append_failures_interruptibly(
1690    pool: Arc<ConnectionPool>,
1691    build: BuildIdentity,
1692    sweep_interval: Duration,
1693    audit_append_failures: u64,
1694) -> StorageResult<DbDiagnostics> {
1695    collect_with_runtime_audit_metrics_interruptibly(
1696        pool,
1697        build,
1698        sweep_interval,
1699        audit_append_failures,
1700        None,
1701    )
1702    .await
1703}
1704
1705/// Like [`collect_with_audit_append_failures_interruptibly`], additionally
1706/// threading through the runtime's audit-batch health counters. `None` when
1707/// no audit-batch control is registered with the calling runtime instance —
1708/// the corresponding `writer_contention` fields then report unavailable with
1709/// a reason, exactly like `audit_append_failures` does for a direct
1710/// `khive-db` caller.
1711pub async fn collect_with_runtime_audit_metrics_interruptibly(
1712    pool: Arc<ConnectionPool>,
1713    build: BuildIdentity,
1714    sweep_interval: Duration,
1715    audit_append_failures: u64,
1716    runtime_audit_batch_metrics: Option<RuntimeAuditBatchMetrics>,
1717) -> StorageResult<DbDiagnostics> {
1718    let process = ProcessIdentity::current(&pool);
1719    collect_with_runtime_audit_metrics_for_process_interruptibly(
1720        pool,
1721        build,
1722        process,
1723        sweep_interval,
1724        audit_append_failures,
1725        runtime_audit_batch_metrics,
1726    )
1727    .await
1728}
1729
1730/// Collect one already-open pool while retaining the main pool's process and
1731/// generation identity. A secondary pool must not claim or advance the main
1732/// pool generation when its pool-scoped counters are inspected.
1733pub async fn collect_with_runtime_audit_metrics_for_process_interruptibly(
1734    pool: Arc<ConnectionPool>,
1735    build: BuildIdentity,
1736    process: ProcessIdentity,
1737    sweep_interval: Duration,
1738    audit_append_failures: u64,
1739    runtime_audit_batch_metrics: Option<RuntimeAuditBatchMetrics>,
1740) -> StorageResult<DbDiagnostics> {
1741    crate::ensure_request_read_active("db_diagnostics")?;
1742    let started = Instant::now();
1743    let counters = checkpoint_counters();
1744    let reader_contention = ReaderContentionDiagnostics::snapshot(&pool);
1745    let search_mechanism = pool.search_mechanism_snapshot();
1746    let writer_contention = WriterContentionDiagnostics::snapshot(
1747        &pool,
1748        Some(audit_append_failures),
1749        runtime_audit_batch_metrics,
1750    );
1751
1752    let Some(path) = pool.config().path.clone() else {
1753        crate::ensure_request_read_active("db_diagnostics")?;
1754        return Ok(DbDiagnostics {
1755            build,
1756            process,
1757            note_search_ann_route_total: 0,
1758            note_search_fallback_route_total: 0,
1759            search_mechanism,
1760            db_path: None,
1761            wal_ceiling: WalCeilingDiagnostics::from_pool(&pool),
1762            disk_guard: DiskGuardDiagnostics::snapshot(&pool),
1763            wal_file: None,
1764            checkpoint_counters: counters,
1765            checkpoint_probe: None,
1766            checkpoint_probe_error: Some(
1767                "in-memory database: no WAL file and no checkpoint to probe".to_string(),
1768            ),
1769            checkpoint_pin: checkpoint_pin_diagnostics(
1770                &pool,
1771                None,
1772                Some("in-memory database: no WAL file and no checkpoint to probe"),
1773            ),
1774            reader_contention,
1775            writer_contention,
1776            size_composition: None,
1777            size_composition_error: Some(
1778                "in-memory database: no file-backed page composition to inspect".to_string(),
1779            ),
1780            graph_edge_integrity: None,
1781            graph_edge_integrity_error: Some(
1782                "in-memory database: no durable graph-edge ledger to inspect".to_string(),
1783            ),
1784            fts_segments: None,
1785            fts_segments_error: Some(
1786                "in-memory database: no durable FTS5 indexes to inspect".to_string(),
1787            ),
1788            fts_maintenance: crate::fts_maintenance_counters(),
1789            wal_pin: WalPinAttribution::unavailable(
1790                "in-memory database: no file for the OS holder census",
1791            ),
1792            collection_cost: CollectionCost::in_memory(elapsed_ms(started)),
1793        });
1794    };
1795
1796    let inspection_pool = Arc::clone(&pool);
1797    let sqlite_started = Instant::now();
1798    let inspection = crate::read_cancellation::run_interruptible_read(
1799        StorageCapability::Sql,
1800        "db_diagnostics.sqlite",
1801        move |scope| inspect_pool_interruptibly(&inspection_pool, scope),
1802    )
1803    .await?;
1804    let sqlite_ms = elapsed_ms(sqlite_started);
1805    crate::ensure_request_read_active("db_diagnostics")?;
1806    let canonical = operational_db_path(&pool, &path);
1807    let budget = request_census_budget();
1808    let (wal_file, wal_pin, file_state_cost) =
1809        inspect_file_state_interruptibly(canonical, sweep_interval, budget).await?;
1810    crate::ensure_request_read_active("db_diagnostics")?;
1811
1812    Ok(DbDiagnostics {
1813        build,
1814        process,
1815        note_search_ann_route_total: 0,
1816        note_search_fallback_route_total: 0,
1817        search_mechanism,
1818        db_path: Some(path.display().to_string()),
1819        wal_ceiling: WalCeilingDiagnostics::from_pool(&pool),
1820        disk_guard: DiskGuardDiagnostics::snapshot(&pool),
1821        wal_file: Some(wal_file),
1822        checkpoint_counters: counters,
1823        checkpoint_probe: inspection.checkpoint_probe,
1824        checkpoint_probe_error: inspection.checkpoint_probe_error,
1825        checkpoint_pin: inspection.checkpoint_pin,
1826        reader_contention,
1827        writer_contention,
1828        size_composition: inspection.size_composition,
1829        size_composition_error: inspection.size_composition_error,
1830        graph_edge_integrity: inspection.graph_edge_integrity,
1831        graph_edge_integrity_error: inspection.graph_edge_integrity_error,
1832        fts_segments: inspection.fts_segments,
1833        fts_segments_error: inspection.fts_segments_error,
1834        fts_maintenance: crate::fts_maintenance_counters(),
1835        wal_pin,
1836        collection_cost: CollectionCost {
1837            total_ms: elapsed_ms(started),
1838            sqlite_ms,
1839            wal_file_stat_ms: file_state_cost.wal_file_stat_ms,
1840            wal_pin_census_ms: file_state_cost.census_ms,
1841            wal_pin_sidecar_ms: file_state_cost.sidecar_ms,
1842            wal_pin_census_budget_ms: budget.map(|b| b.as_millis() as u64),
1843            wal_pin_census_budget_exhausted: file_state_cost.census_budget_exhausted,
1844        },
1845    })
1846}
1847
1848/// The path every operational probe (WAL file, sidecar directory, OS holder
1849/// census) resolves against for `pool`'s database. The checkpoint sidecar
1850/// writers and `ConnectionPool` itself key the WAL and sidecar directory off
1851/// `canonical_path()`, not the raw `configured` path — a symlinked or
1852/// otherwise aliased configured path would otherwise send every probe below
1853/// looking beside the alias while the evidence sits beside the canonical
1854/// file. `configured` remains the presentation value the sync and async
1855/// collectors put in `db_path`, since that is what the caller configured.
1856/// Both `collect_inner` and `collect_with_runtime_audit_metrics_interruptibly`
1857/// resolve through this one function so their aliasing behavior cannot drift
1858/// apart.
1859fn operational_db_path(pool: &ConnectionPool, configured: &Path) -> PathBuf {
1860    pool.canonical_path()
1861        .map(Path::to_path_buf)
1862        .unwrap_or_else(|| configured.to_path_buf())
1863}
1864
1865fn collect_inner(
1866    pool: &ConnectionPool,
1867    build: BuildIdentity,
1868    sweep_interval: Duration,
1869    audit_append_failures: Option<u64>,
1870    runtime_audit_batch_metrics: Option<RuntimeAuditBatchMetrics>,
1871) -> DbDiagnostics {
1872    let started = Instant::now();
1873    let process = ProcessIdentity::current(pool);
1874    let counters = checkpoint_counters();
1875    let reader_contention = ReaderContentionDiagnostics::snapshot(pool);
1876    let search_mechanism = pool.search_mechanism_snapshot();
1877    let writer_contention = WriterContentionDiagnostics::snapshot(
1878        pool,
1879        audit_append_failures,
1880        runtime_audit_batch_metrics,
1881    );
1882
1883    let Some(path) = pool.config().path.clone() else {
1884        return DbDiagnostics {
1885            build,
1886            process,
1887            note_search_ann_route_total: 0,
1888            note_search_fallback_route_total: 0,
1889            search_mechanism,
1890            db_path: None,
1891            wal_ceiling: WalCeilingDiagnostics::from_pool(pool),
1892            disk_guard: DiskGuardDiagnostics::snapshot(pool),
1893            wal_file: None,
1894            checkpoint_counters: counters,
1895            checkpoint_probe: None,
1896            checkpoint_probe_error: Some(
1897                "in-memory database: no WAL file and no checkpoint to probe".to_string(),
1898            ),
1899            checkpoint_pin: checkpoint_pin_diagnostics(
1900                pool,
1901                None,
1902                Some("in-memory database: no WAL file and no checkpoint to probe"),
1903            ),
1904            reader_contention,
1905            writer_contention,
1906            size_composition: None,
1907            size_composition_error: Some(
1908                "in-memory database: no file-backed page composition to inspect".to_string(),
1909            ),
1910            graph_edge_integrity: None,
1911            graph_edge_integrity_error: Some(
1912                "in-memory database: no durable graph-edge ledger to inspect".to_string(),
1913            ),
1914            fts_segments: None,
1915            fts_segments_error: Some(
1916                "in-memory database: no durable FTS5 indexes to inspect".to_string(),
1917            ),
1918            fts_maintenance: crate::fts_maintenance_counters(),
1919            wal_pin: WalPinAttribution::unavailable(
1920                "in-memory database: no file for the OS holder census",
1921            ),
1922            collection_cost: CollectionCost::in_memory(elapsed_ms(started)),
1923        };
1924    };
1925
1926    let sqlite_started = Instant::now();
1927    let inspection = inspect_pool(pool);
1928    let sqlite_ms = elapsed_ms(sqlite_started);
1929    let canonical = operational_db_path(pool, &path);
1930    let wal_file_started = Instant::now();
1931    let wal_file = wal_file_state(&canonical);
1932    let wal_file_stat_ms = elapsed_ms(wal_file_started);
1933    let census_started = Instant::now();
1934    let wal_pin = wal_pin_attribution(&canonical, sweep_interval);
1935    let wal_pin_ms = elapsed_ms(census_started);
1936
1937    DbDiagnostics {
1938        build,
1939        process,
1940        note_search_ann_route_total: 0,
1941        note_search_fallback_route_total: 0,
1942        search_mechanism,
1943        db_path: Some(path.display().to_string()),
1944        wal_ceiling: WalCeilingDiagnostics::from_pool(pool),
1945        disk_guard: DiskGuardDiagnostics::snapshot(pool),
1946        wal_file: Some(wal_file),
1947        checkpoint_counters: counters,
1948        checkpoint_probe: inspection.checkpoint_probe,
1949        checkpoint_probe_error: inspection.checkpoint_probe_error,
1950        checkpoint_pin: inspection.checkpoint_pin,
1951        reader_contention,
1952        writer_contention,
1953        size_composition: inspection.size_composition,
1954        size_composition_error: inspection.size_composition_error,
1955        graph_edge_integrity: inspection.graph_edge_integrity,
1956        graph_edge_integrity_error: inspection.graph_edge_integrity_error,
1957        fts_segments: inspection.fts_segments,
1958        fts_segments_error: inspection.fts_segments_error,
1959        fts_maintenance: crate::fts_maintenance_counters(),
1960        wal_pin,
1961        collection_cost: CollectionCost {
1962            total_ms: elapsed_ms(started),
1963            sqlite_ms,
1964            wal_file_stat_ms,
1965            // This path does not separate the census from the sidecar read it
1966            // is reconciled against: it calls the combined `wal_pin_attribution`,
1967            // so splitting them here would mean inventing a number. The
1968            // request path, which is the one the operator waits on, does split
1969            // them.
1970            wal_pin_census_ms: wal_pin_ms,
1971            wal_pin_sidecar_ms: 0,
1972            wal_pin_census_budget_ms: None,
1973            wal_pin_census_budget_exhausted: false,
1974        },
1975    }
1976}
1977
1978struct PoolInspection {
1979    checkpoint_probe: Option<CheckpointProbe>,
1980    checkpoint_probe_error: Option<String>,
1981    checkpoint_pin: CheckpointPinDiagnostics,
1982    size_composition: Option<DatabaseSizeComposition>,
1983    size_composition_error: Option<String>,
1984    graph_edge_integrity: Option<GraphEdgeIntegrity>,
1985    graph_edge_integrity_error: Option<String>,
1986    fts_segments: Option<crate::FtsSegmentDiagnostics>,
1987    fts_segments_error: Option<String>,
1988}
1989
1990/// Split an FTS segment inspection result into the `(value, error)` pair
1991/// shape every section of [`PoolInspection`]/[`DbDiagnostics`] uses, shared
1992/// by the interruptible and plain inspection paths below.
1993fn split_fts_segments_result(
1994    result: Result<crate::FtsSegmentDiagnostics, String>,
1995) -> (Option<crate::FtsSegmentDiagnostics>, Option<String>) {
1996    match result {
1997        Ok(segments) => (Some(segments), None),
1998        Err(error) => (None, Some(error)),
1999    }
2000}
2001
2002fn record_diagnostic_checkpoint_probe(
2003    pool: &ConnectionPool,
2004    result: &rusqlite::Result<CheckpointProbe>,
2005) -> (checkpoint::CheckpointRunStatus, Option<Duration>) {
2006    checkpoint::record_checkpoint_run_result(
2007        pool,
2008        result
2009            .as_ref()
2010            .ok()
2011            .map(|probe| (probe.busy, probe.log_frames, probe.checkpointed_frames)),
2012    );
2013    checkpoint::checkpoint_run_snapshot(pool)
2014}
2015
2016fn inspect_pool_interruptibly(
2017    pool: &ConnectionPool,
2018    scope: &crate::read_cancellation::InterruptibleReadScope,
2019) -> StorageResult<PoolInspection> {
2020    scope.ensure_active()?;
2021    let conn = match pool.open_standalone_writer_untracked() {
2022        Ok(conn) => conn,
2023        Err(e) => {
2024            scope.ensure_active()?;
2025            let reason = format!("guarded standalone open refused: {e}");
2026            return Ok(PoolInspection {
2027                checkpoint_probe: None,
2028                checkpoint_probe_error: Some(reason.clone()),
2029                checkpoint_pin: checkpoint_pin_diagnostics(pool, None, Some(&reason)),
2030                size_composition: None,
2031                size_composition_error: Some(reason.clone()),
2032                graph_edge_integrity: None,
2033                graph_edge_integrity_error: Some(reason.clone()),
2034                fts_segments: None,
2035                fts_segments_error: Some(reason),
2036            });
2037        }
2038    };
2039    // Opening is read-only filesystem work but can block. Cancellation that
2040    // arrived while it was in flight must prevent admission of the following
2041    // PASSIVE checkpoint, whose backfill I/O is intentionally noninterruptible
2042    // once started.
2043    scope.ensure_active()?;
2044
2045    // PASSIVE can perform write I/O. Never install sqlite3_interrupt for it.
2046    let probe_result = checkpoint_probe(&conn);
2047    let (run_status, run_age) = record_diagnostic_checkpoint_probe(pool, &probe_result);
2048    let (checkpoint_probe, checkpoint_probe_error) = match probe_result {
2049        Ok(probe) => (Some(probe), None),
2050        Err(e) => (
2051            None,
2052            Some(format!("PRAGMA wal_checkpoint(PASSIVE) failed: {e}")),
2053        ),
2054    };
2055    let checkpoint_pin = checkpoint_pin_diagnostics_for_run(
2056        checkpoint_probe.as_ref(),
2057        checkpoint_probe_error.as_deref(),
2058        run_status,
2059        run_age,
2060    );
2061    #[cfg(test)]
2062    if TEST_PAUSE_AFTER_PASSIVE.load(Ordering::SeqCst) {
2063        TEST_REACHED_AFTER_PASSIVE.store(true, Ordering::SeqCst);
2064        while TEST_PAUSE_AFTER_PASSIVE.load(Ordering::SeqCst) && !scope.should_stop() {
2065            std::thread::yield_now();
2066        }
2067    }
2068    scope.ensure_active()?;
2069
2070    // Install one progress/interrupt guard for all logical reads on this
2071    // connection. The scope deliberately refuses double registration, so
2072    // keep each query's ordinary SQLite result nested inside one guarded
2073    // execution while allowing the outer cancellation cause to remain typed.
2074    let (integrity, size_composition, fts_segments) = scope.run(&conn, || {
2075        Ok((
2076            graph_edge_integrity(&conn),
2077            database_size_composition(&conn),
2078            crate::fts_maintenance::inspect_fts_segments(&conn),
2079        ))
2080    })?;
2081    let (graph_edge_integrity, graph_edge_integrity_error) = match integrity {
2082        Ok(integrity) => (Some(integrity), None),
2083        Err(e) => (
2084            None,
2085            Some(format!("graph-edge integrity query failed: {e}")),
2086        ),
2087    };
2088    let (fts_segments, fts_segments_error) = split_fts_segments_result(fts_segments);
2089
2090    let (size_composition, size_composition_error) = match size_composition {
2091        Ok(composition) => (Some(composition), None),
2092        Err(error) => (
2093            None,
2094            Some(format!("database size composition query failed: {error}")),
2095        ),
2096    };
2097
2098    Ok(PoolInspection {
2099        checkpoint_probe,
2100        checkpoint_probe_error,
2101        checkpoint_pin,
2102        size_composition,
2103        size_composition_error,
2104        graph_edge_integrity,
2105        graph_edge_integrity_error,
2106        fts_segments,
2107        fts_segments_error,
2108    })
2109}
2110
2111#[cfg(test)]
2112static TEST_PAUSE_AFTER_PASSIVE: AtomicBool = AtomicBool::new(false);
2113#[cfg(test)]
2114static TEST_REACHED_AFTER_PASSIVE: AtomicBool = AtomicBool::new(false);
2115
2116struct StopCensusOnDrop {
2117    stopped: Arc<AtomicBool>,
2118    armed: bool,
2119}
2120
2121impl Drop for StopCensusOnDrop {
2122    fn drop(&mut self) {
2123        if self.armed {
2124            self.stopped.store(true, Ordering::SeqCst);
2125        }
2126    }
2127}
2128
2129/// Wall-clock cost of the file-backed half of a report, split so the census
2130/// can be told apart from the two cheap reads it sits between.
2131#[derive(Debug, Clone, Copy, Default)]
2132struct FileStateCost {
2133    wal_file_stat_ms: u64,
2134    census_ms: u64,
2135    sidecar_ms: u64,
2136    census_budget_exhausted: bool,
2137}
2138
2139fn elapsed_ms(since: Instant) -> u64 {
2140    since.elapsed().as_millis() as u64
2141}
2142
2143async fn inspect_file_state_interruptibly(
2144    path: PathBuf,
2145    sweep_interval: Duration,
2146    census_budget: Option<Duration>,
2147) -> StorageResult<(WalFileState, WalPinAttribution, FileStateCost)> {
2148    const OPERATION: &str = "db_diagnostics.wal_holder_census";
2149    crate::ensure_request_read_active(OPERATION)?;
2150    let stopped = Arc::new(AtomicBool::new(false));
2151    let worker_stopped = Arc::clone(&stopped);
2152    let mut stop_on_drop = StopCensusOnDrop {
2153        stopped: Arc::clone(&stopped),
2154        armed: true,
2155    };
2156    let mut worker = tokio::task::spawn_blocking(move || {
2157        let mut cost = FileStateCost::default();
2158        let wal_file_started = Instant::now();
2159        let wal_file = wal_file_state(&path);
2160        cost.wal_file_stat_ms = elapsed_ms(wal_file_started);
2161        if worker_stopped.load(Ordering::SeqCst) {
2162            return Err(std::io::Error::new(
2163                std::io::ErrorKind::Interrupted,
2164                "WAL holder census cancelled",
2165            ));
2166        }
2167        #[cfg(unix)]
2168        let census_started = Instant::now();
2169        #[cfg(unix)]
2170        let census_result = match census_budget {
2171            Some(budget) => crate::walpin::census_holders_until_within(
2172                &path,
2173                || worker_stopped.load(Ordering::SeqCst),
2174                budget,
2175            ),
2176            None => {
2177                crate::walpin::census_holders_until(&path, || worker_stopped.load(Ordering::SeqCst))
2178            }
2179        };
2180        #[cfg(unix)]
2181        {
2182            cost.census_ms = elapsed_ms(census_started);
2183        }
2184        #[cfg(unix)]
2185        let attribution = match census_result {
2186            Ok(census) => {
2187                cost.census_budget_exhausted = census.budget_exhausted;
2188                if worker_stopped.load(Ordering::SeqCst) {
2189                    return Err(std::io::Error::new(
2190                        std::io::ErrorKind::Interrupted,
2191                        "WAL sidecar inspection cancelled",
2192                    ));
2193                }
2194                if !crate::walpin::sidecar_enabled(true) {
2195                    wal_pin_attribution_without_sidecar(census, SIDECAR_DISABLED_REASON.to_string())
2196                } else {
2197                    let sidecar_started = Instant::now();
2198                    let sidecar = crate::walpin::inspect_live(
2199                        &crate::walpin::sidecar_dir_for(&path),
2200                        sweep_interval,
2201                    );
2202                    cost.sidecar_ms = elapsed_ms(sidecar_started);
2203                    if worker_stopped.load(Ordering::SeqCst) {
2204                        return Err(std::io::Error::new(
2205                            std::io::ErrorKind::Interrupted,
2206                            "WAL sidecar inspection cancelled",
2207                        ));
2208                    }
2209                    match sidecar {
2210                        Ok(sidecar) => wal_pin_attribution_from_evidence(census, sidecar),
2211                        Err(error) => wal_pin_attribution_without_sidecar(
2212                            census,
2213                            format!("read-only sidecar enumeration failed: {error}"),
2214                        ),
2215                    }
2216                }
2217            }
2218            Err(error) if error.kind() == std::io::ErrorKind::Interrupted => return Err(error),
2219            Err(error) => WalPinAttribution::unavailable(format!("census_holders failed: {error}")),
2220        };
2221        #[cfg(not(unix))]
2222        let attribution = {
2223            let _ = census_budget;
2224            let census_started = Instant::now();
2225            let attribution = wal_pin_attribution(&path, sweep_interval);
2226            cost.census_ms = elapsed_ms(census_started);
2227            attribution
2228        };
2229        Ok((wal_file, attribution, cost))
2230    });
2231
2232    tokio::select! {
2233        joined = &mut worker => {
2234            stop_on_drop.armed = false;
2235            let result = joined
2236                .map_err(|error| StorageError::driver(StorageCapability::Sql, OPERATION, error))?
2237                .map_err(|error| StorageError::driver(StorageCapability::Sql, OPERATION, error))?;
2238            crate::ensure_request_read_active(OPERATION)?;
2239            Ok(result)
2240        }
2241        _ = crate::wait_for_request_read_cancellation() => {
2242            stopped.store(true, Ordering::SeqCst);
2243            if tokio::time::timeout(crate::sqlite_interrupt_grace_from_env(), &mut worker)
2244                .await
2245                .is_err()
2246            {
2247                worker.abort();
2248            }
2249            stop_on_drop.armed = false;
2250            Err(StorageError::Timeout { operation: OPERATION.into() })
2251        }
2252    }
2253}
2254
2255/// Run the PASSIVE probe and graph-ledger reads on one guarded standalone
2256/// connection.
2257///
2258/// Every failure — missing file, read-only pool, in-memory pool, or the
2259/// pragma/query itself — comes back in its section's error field. Nothing
2260/// here can create a file.
2261fn inspect_pool(pool: &ConnectionPool) -> PoolInspection {
2262    let conn = match pool.open_standalone_writer_untracked() {
2263        Ok(conn) => conn,
2264        Err(e) => {
2265            let reason = format!("guarded standalone open refused: {e}");
2266            return PoolInspection {
2267                checkpoint_probe: None,
2268                checkpoint_probe_error: Some(reason.clone()),
2269                checkpoint_pin: checkpoint_pin_diagnostics(pool, None, Some(&reason)),
2270                size_composition: None,
2271                size_composition_error: Some(reason.clone()),
2272                graph_edge_integrity: None,
2273                graph_edge_integrity_error: Some(reason.clone()),
2274                fts_segments: None,
2275                fts_segments_error: Some(reason),
2276            };
2277        }
2278    };
2279
2280    let probe_result = checkpoint_probe(&conn);
2281    let (run_status, run_age) = record_diagnostic_checkpoint_probe(pool, &probe_result);
2282    let (checkpoint_probe, checkpoint_probe_error) = match probe_result {
2283        Ok(probe) => (Some(probe), None),
2284        Err(e) => (
2285            None,
2286            Some(format!("PRAGMA wal_checkpoint(PASSIVE) failed: {e}")),
2287        ),
2288    };
2289    let checkpoint_pin = checkpoint_pin_diagnostics_for_run(
2290        checkpoint_probe.as_ref(),
2291        checkpoint_probe_error.as_deref(),
2292        run_status,
2293        run_age,
2294    );
2295    let (graph_edge_integrity, graph_edge_integrity_error) = match graph_edge_integrity(&conn) {
2296        Ok(integrity) => (Some(integrity), None),
2297        Err(e) => (
2298            None,
2299            Some(format!("graph-edge integrity query failed: {e}")),
2300        ),
2301    };
2302    let (size_composition, size_composition_error) = match database_size_composition(&conn) {
2303        Ok(composition) => (Some(composition), None),
2304        Err(error) => (
2305            None,
2306            Some(format!("database size composition query failed: {error}")),
2307        ),
2308    };
2309    let (fts_segments, fts_segments_error) =
2310        split_fts_segments_result(crate::fts_maintenance::inspect_fts_segments(&conn));
2311
2312    PoolInspection {
2313        checkpoint_probe,
2314        checkpoint_probe_error,
2315        checkpoint_pin,
2316        size_composition,
2317        size_composition_error,
2318        graph_edge_integrity,
2319        graph_edge_integrity_error,
2320        fts_segments,
2321        fts_segments_error,
2322    }
2323}
2324
2325#[cfg(test)]
2326mod tests {
2327    use serial_test::serial;
2328
2329    use super::*;
2330    use crate::pool::{ConnectionPool, PoolConfig, WalCeilingPolicy, WalCeilingSource};
2331
2332    include!("diagnostics/wal_ceiling_tests.rs");
2333    include!("diagnostics/environment_tests.rs");
2334    include!("diagnostics_census_evidence_tests.rs");
2335
2336    fn seeded_pool(dir: &tempfile::TempDir) -> (ConnectionPool, PathBuf) {
2337        let path = dir.path().join("diag.db");
2338        let pool = ConnectionPool::new(PoolConfig {
2339            path: Some(path.clone()),
2340            ..PoolConfig::for_test()
2341        })
2342        .expect("pool open");
2343        {
2344            let writer = pool.try_writer().expect("writer");
2345            writer
2346                .conn()
2347                .execute_batch(
2348                    "CREATE TABLE t (x INTEGER); \
2349                     CREATE TABLE entities (id TEXT PRIMARY KEY, deleted_at INTEGER, merged_into TEXT); \
2350                     CREATE TABLE graph_edges (
2351                         namespace TEXT NOT NULL,
2352                         id TEXT NOT NULL,
2353                         PRIMARY KEY (namespace, id)
2354                     ); \
2355                     CREATE TABLE graph_edges_seq (
2356                         seq INTEGER PRIMARY KEY AUTOINCREMENT,
2357                         edge_id TEXT NOT NULL UNIQUE
2358                     ); \
2359                     CREATE VIRTUAL TABLE fts_entities USING fts5(
2360                         namespace UNINDEXED, subject_id UNINDEXED, title, body,
2361                         tokenize='trigram'
2362                     ); \
2363                     CREATE VIRTUAL TABLE fts_notes USING fts5(
2364                         namespace UNINDEXED, subject_id UNINDEXED, title, body,
2365                         tokenize='trigram'
2366                     ); \
2367                     INSERT INTO t VALUES (1), (2), (3); \
2368                     INSERT INTO fts_entities(namespace, subject_id, title, body)
2369                         VALUES('local', 'entity-1', 'entity title', 'entity diagnostic body'); \
2370                     INSERT INTO fts_notes(namespace, subject_id, title, body)
2371                         VALUES('local', 'note-1', 'note title', 'note diagnostic body');",
2372                )
2373                .expect("seed writes");
2374        }
2375        (pool, path)
2376    }
2377
2378    #[test]
2379    fn dropping_census_future_guard_requests_cooperative_stop() {
2380        let stopped = Arc::new(AtomicBool::new(false));
2381        let guard = StopCensusOnDrop {
2382            stopped: Arc::clone(&stopped),
2383            armed: true,
2384        };
2385
2386        drop(guard);
2387
2388        assert!(
2389            stopped.load(Ordering::SeqCst),
2390            "dropping diagnostics while its census worker is live must stop the PID/fd walk"
2391        );
2392    }
2393
2394    #[tokio::test]
2395    async fn runtime_audit_batch_fields_are_additive_and_unavailable_without_a_control() {
2396        let dir = tempfile::tempdir().expect("tempdir");
2397        let (pool, _path) = seeded_pool(&dir);
2398        let pool = Arc::new(pool);
2399
2400        let without_control = collect_with_audit_append_failures_interruptibly(
2401            Arc::clone(&pool),
2402            BuildIdentity::from_env("test", None),
2403            Duration::from_secs(30),
2404            0,
2405        )
2406        .await
2407        .expect("diagnostics succeed");
2408        assert!(without_control
2409            .writer_contention
2410            .audit_batch_flush_failures
2411            .is_none());
2412        assert!(
2413            without_control
2414                .writer_contention
2415                .audit_batch_flush_failures_unavailable_reason
2416                .is_some(),
2417            "no audit-batch control was supplied, so the field must carry a reason, not a \
2418             fabricated zero"
2419        );
2420        assert!(without_control
2421            .writer_contention
2422            .audit_degraded_rows
2423            .is_none());
2424        assert!(without_control.writer_contention.audit_degraded.is_none());
2425
2426        let with_control = collect_with_runtime_audit_metrics_interruptibly(
2427            Arc::clone(&pool),
2428            BuildIdentity::from_env("test", None),
2429            Duration::from_secs(30),
2430            0,
2431            Some(RuntimeAuditBatchMetrics {
2432                flush_failures: 3,
2433                degraded_rows: 7,
2434                degraded: true,
2435                admission_refused_obligations: 5,
2436                admission_refused_obligations_last_at_ms: Some(1_700_000_000_123),
2437                admission_unresolved_obligations: 2,
2438                admission_unresolved_obligations_last_at_ms: Some(1_700_000_000_456),
2439            }),
2440        )
2441        .await
2442        .expect("diagnostics succeed");
2443        assert_eq!(
2444            with_control.writer_contention.audit_batch_flush_failures,
2445            Some(3)
2446        );
2447        assert!(with_control
2448            .writer_contention
2449            .audit_batch_flush_failures_unavailable_reason
2450            .is_none());
2451        assert_eq!(with_control.writer_contention.audit_degraded_rows, Some(7));
2452        assert_eq!(with_control.writer_contention.audit_degraded, Some(true));
2453        assert_eq!(
2454            with_control
2455                .writer_contention
2456                .audit_admission_refused_obligations,
2457            Some(5),
2458            "an operator must be able to read the admission-refused obligation count from \
2459             db_diagnostics without a test-only feature gate (ADR-103 Amendment 3)"
2460        );
2461        assert!(with_control
2462            .writer_contention
2463            .audit_admission_refused_obligations_unavailable_reason
2464            .is_none());
2465        assert_eq!(
2466            with_control
2467                .writer_contention
2468                .audit_admission_unresolved_obligations,
2469            Some(2),
2470            "an operator must be able to distinguish enqueued-but-unresolved rows from \
2471             confirmed-refused rows (ADR-103 Amendment 3)"
2472        );
2473        assert!(with_control
2474            .writer_contention
2475            .audit_admission_unresolved_obligations_unavailable_reason
2476            .is_none());
2477        assert!(without_control
2478            .writer_contention
2479            .audit_admission_refused_obligations
2480            .is_none());
2481        assert!(without_control
2482            .writer_contention
2483            .audit_admission_refused_obligations_unavailable_reason
2484            .is_some());
2485        assert!(without_control
2486            .writer_contention
2487            .audit_admission_unresolved_obligations
2488            .is_none());
2489        assert!(without_control
2490            .writer_contention
2491            .audit_admission_unresolved_obligations_unavailable_reason
2492            .is_some());
2493
2494        // Existing fields must be unaffected by the new ones — additive, not
2495        // a reshuffle.
2496        assert_eq!(
2497            with_control.writer_contention.writer_acquisitions,
2498            without_control.writer_contention.writer_acquisitions
2499        );
2500    }
2501
2502    #[test]
2503    fn writer_task_pool_sourced_counters_are_always_populated_directly() {
2504        let dir = tempfile::tempdir().expect("tempdir");
2505        let (pool, _path) = seeded_pool(&dir);
2506
2507        let report = collect(
2508            &pool,
2509            BuildIdentity::from_env("9.9.9", None),
2510            Duration::from_secs(30),
2511        );
2512
2513        // Unlike the runtime-supplied audit-batch fields, these two come
2514        // straight from the pool's own counters and are never `Option`.
2515        assert_eq!(report.writer_contention.writer_task_request_failures, 0);
2516        assert_eq!(report.writer_contention.writer_task_side_effects_unknown, 0);
2517    }
2518
2519    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2520    #[serial]
2521    async fn request_cancellation_after_passive_stops_before_graph_and_census() {
2522        let dir = tempfile::tempdir().expect("tempdir");
2523        let (pool, _) = seeded_pool(&dir);
2524        let pool = Arc::new(pool);
2525        TEST_REACHED_AFTER_PASSIVE.store(false, Ordering::SeqCst);
2526        TEST_PAUSE_AFTER_PASSIVE.store(true, Ordering::SeqCst);
2527        let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
2528        let diagnostic_pool = Arc::clone(&pool);
2529        let task = tokio::spawn(crate::scope_request_read_cancellation(
2530            cancel_rx,
2531            async move {
2532                collect_with_audit_append_failures_interruptibly(
2533                    diagnostic_pool,
2534                    BuildIdentity::from_env("test", None),
2535                    Duration::from_secs(30),
2536                    0,
2537                )
2538                .await
2539            },
2540        ));
2541
2542        tokio::time::timeout(Duration::from_secs(1), async {
2543            while !TEST_REACHED_AFTER_PASSIVE.load(Ordering::SeqCst) {
2544                tokio::task::yield_now().await;
2545            }
2546        })
2547        .await
2548        .expect("diagnostics never completed its admitted PASSIVE phase");
2549        cancel_tx.send(true).unwrap();
2550        let result = tokio::time::timeout(Duration::from_secs(1), task)
2551            .await
2552            .expect("cancelled diagnostics did not stop promptly")
2553            .expect("diagnostics task panicked");
2554        TEST_PAUSE_AFTER_PASSIVE.store(false, Ordering::SeqCst);
2555        assert!(matches!(result, Err(StorageError::Timeout { .. })));
2556
2557        let one: i64 = pool
2558            .reader()
2559            .expect("diagnostics returned its connection")
2560            .conn()
2561            .query_row("SELECT 1", [], |row| row.get(0))
2562            .unwrap();
2563        assert_eq!(one, 1);
2564    }
2565
2566    #[test]
2567    fn checkpoint_probe_returns_a_well_formed_triple_on_a_file_backed_db() {
2568        let dir = tempfile::tempdir().expect("tempdir");
2569        let (pool, _path) = seeded_pool(&dir);
2570        let conn = pool
2571            .open_standalone_writer_untracked()
2572            .expect("standalone probe connection");
2573
2574        let probe = checkpoint_probe(&conn).expect("probe must succeed on a WAL database");
2575
2576        assert!(
2577            probe.busy == 0 || probe.busy == 1,
2578            "busy is a 0/1 flag, got {}",
2579            probe.busy
2580        );
2581        assert!(
2582            probe.log_frames >= 0,
2583            "a WAL database must report a non-negative frame count, got {}",
2584            probe.log_frames
2585        );
2586        assert!(
2587            probe.checkpointed_frames >= 0,
2588            "checkpointed frames must be non-negative, got {}",
2589            probe.checkpointed_frames
2590        );
2591        assert!(
2592            probe.checkpointed_frames <= probe.log_frames,
2593            "a PASSIVE pass cannot checkpoint more frames than the WAL holds: {probe:?}"
2594        );
2595        assert!(
2596            probe.backfill_gap_frames() >= 0,
2597            "the one-row backfill gap clamps at 0: {probe:?}"
2598        );
2599    }
2600
2601    #[cfg(unix)]
2602    #[test]
2603    fn checkpoint_probe_reports_replaced_pool_file_instead_of_probing_it() {
2604        let dir = tempfile::tempdir().expect("tempdir");
2605        let (pool, path) = seeded_pool(&dir);
2606        let replacement = dir.path().join("replacement.db");
2607        let replacement_conn = Connection::open(&replacement).expect("replacement database");
2608        replacement_conn
2609            .execute_batch("CREATE TABLE replacement_marker (value INTEGER)")
2610            .expect("initialize replacement database");
2611        drop(replacement_conn);
2612        std::fs::rename(&replacement, &path).expect("replace the pool's path");
2613
2614        let inspection = inspect_pool(&pool);
2615        assert!(inspection.checkpoint_probe.is_none());
2616        assert!(
2617            inspection
2618                .checkpoint_probe_error
2619                .as_deref()
2620                .is_some_and(|error| error.contains("file identity changed")),
2621            "replacement must be reported as a probe failure"
2622        );
2623    }
2624
2625    #[test]
2626    fn checkpoint_probe_backfill_gap_is_a_row_difference_not_a_pin_claim() {
2627        let probe = CheckpointProbe {
2628            busy: 0,
2629            log_frames: 12,
2630            checkpointed_frames: 7,
2631        };
2632        assert_eq!(probe.backfill_gap_frames(), 5);
2633        let no_wal = CheckpointProbe {
2634            busy: 0,
2635            log_frames: -1,
2636            checkpointed_frames: -1,
2637        };
2638        assert_eq!(no_wal.backfill_gap_frames(), 0);
2639        let busy = CheckpointProbe { busy: 1, ..probe };
2640        assert_eq!(busy.backfill_gap_frames(), 5);
2641    }
2642
2643    #[test]
2644    fn pin_age_does_not_depend_on_the_wall_clock_epoch() {
2645        let probe = CheckpointProbe {
2646            busy: 0,
2647            log_frames: 12,
2648            checkpointed_frames: 7,
2649        };
2650        let future_epoch = checkpoint::CheckpointRunStatus::Observed(checkpoint::CheckpointRun {
2651            frame: 7,
2652            first_observed_at_unix_ms: 10_000,
2653        });
2654        let past_epoch = checkpoint::CheckpointRunStatus::Observed(checkpoint::CheckpointRun {
2655            frame: 7,
2656            first_observed_at_unix_ms: 1,
2657        });
2658        let young = checkpoint_pin_diagnostics_for_run(
2659            Some(&probe),
2660            None,
2661            future_epoch,
2662            Some(Duration::from_millis(10)),
2663        );
2664        let forward = checkpoint_pin_diagnostics_for_run(
2665            Some(&probe),
2666            None,
2667            future_epoch,
2668            Some(Duration::from_millis(1010)),
2669        );
2670        let backward = checkpoint_pin_diagnostics_for_run(
2671            Some(&probe),
2672            None,
2673            past_epoch,
2674            Some(Duration::from_millis(1010)),
2675        );
2676        assert_eq!(young.oldest_pinned_frame, None);
2677        assert_eq!(forward.oldest_pinned_frame, Some(7));
2678        assert_eq!(forward.pin_depth, Some(5));
2679        assert_eq!(backward.oldest_pinned_frame, Some(7));
2680        let busy = CheckpointProbe { busy: 1, ..probe };
2681        let unavailable = checkpoint_pin_diagnostics_for_run(
2682            Some(&busy),
2683            None,
2684            future_epoch,
2685            Some(Duration::from_millis(1010)),
2686        );
2687        assert_eq!(unavailable.oldest_pinned_frame, None);
2688        assert_eq!(unavailable.pin_depth, None);
2689    }
2690
2691    #[test]
2692    fn checkpoint_pin_report_waits_one_second_for_a_matching_run() {
2693        let probe = CheckpointProbe {
2694            busy: 0,
2695            log_frames: 12,
2696            checkpointed_frames: 7,
2697        };
2698        let run = checkpoint::CheckpointRunStatus::Observed(checkpoint::CheckpointRun {
2699            frame: 7,
2700            first_observed_at_unix_ms: 1_000,
2701        });
2702
2703        let young = checkpoint_pin_diagnostics_for_run(
2704            Some(&probe),
2705            None,
2706            run,
2707            Some(Duration::from_millis(999)),
2708        );
2709        assert_eq!(young.backfill_ceiling, Some(7));
2710        assert_eq!(young.oldest_pinned_frame, None);
2711        assert!(young
2712            .oldest_pinned_frame_unavailable_reason
2713            .as_deref()
2714            .is_some_and(|reason| reason.contains("less than one second")));
2715        assert_eq!(young.oldest_pinned_frame_run, None);
2716        assert_eq!(young.pin_depth, None);
2717        assert!(young.pin_depth_unavailable_reason.is_some());
2718
2719        let aged = checkpoint_pin_diagnostics_for_run(
2720            Some(&probe),
2721            None,
2722            run,
2723            Some(Duration::from_secs(1)),
2724        );
2725        assert_eq!(aged.backfill_ceiling, Some(7));
2726        assert_eq!(aged.oldest_pinned_frame, Some(7));
2727        assert_eq!(aged.pin_depth, Some(5));
2728        assert_eq!(aged.pin_depth_unavailable_reason, None);
2729        assert_eq!(
2730            aged.oldest_pinned_frame_run,
2731            Some(match run {
2732                checkpoint::CheckpointRunStatus::Observed(value) => value,
2733                _ => unreachable!(),
2734            })
2735        );
2736    }
2737
2738    #[test]
2739    fn checkpoint_pin_report_keeps_busy_probe_fields_null_even_with_an_aged_run() {
2740        let run = checkpoint::CheckpointRunStatus::Observed(checkpoint::CheckpointRun {
2741            frame: 7,
2742            first_observed_at_unix_ms: 1,
2743        });
2744        for probe in [
2745            Some(CheckpointProbe {
2746                busy: 1,
2747                log_frames: 12,
2748                checkpointed_frames: 7,
2749            }),
2750            Some(CheckpointProbe {
2751                busy: 0,
2752                log_frames: -1,
2753                checkpointed_frames: -1,
2754            }),
2755            Some(CheckpointProbe {
2756                busy: 0,
2757                log_frames: 12,
2758                checkpointed_frames: 12,
2759            }),
2760        ] {
2761            let result = checkpoint_pin_diagnostics_for_run(
2762                probe.as_ref(),
2763                None,
2764                run,
2765                Some(Duration::from_secs(1)),
2766            );
2767            assert_eq!(result.backfill_ceiling, None);
2768            assert!(result.backfill_ceiling_unavailable_reason.is_some());
2769            assert_eq!(result.oldest_pinned_frame, None);
2770            assert!(result.oldest_pinned_frame_unavailable_reason.is_some());
2771            assert_eq!(result.oldest_pinned_frame_run, None);
2772            assert!(result.oldest_pinned_frame_run_unavailable_reason.is_some());
2773            assert_eq!(result.pin_depth, None);
2774            assert!(result.pin_depth_unavailable_reason.is_some());
2775        }
2776
2777        let error = checkpoint_pin_diagnostics_for_run(
2778            None,
2779            Some("probe failed"),
2780            run,
2781            Some(Duration::from_secs(1)),
2782        );
2783        assert_eq!(error.backfill_ceiling, None);
2784        assert_eq!(
2785            error.backfill_ceiling_unavailable_reason.as_deref(),
2786            Some("probe failed")
2787        );
2788        assert_eq!(error.oldest_pinned_frame, None);
2789        assert!(error.oldest_pinned_frame_run_unavailable_reason.is_some());
2790        assert_eq!(error.pin_depth, None);
2791        assert!(error.pin_depth_unavailable_reason.is_some());
2792    }
2793
2794    #[cfg(all(unix, any(target_os = "linux", target_os = "macos")))]
2795    #[test]
2796    fn db_diagnostics_reports_start_times_for_all_holders_and_identifies_reporter() {
2797        let dir = tempfile::tempdir().expect("tempdir");
2798        let (pool, _) = seeded_pool(&dir);
2799
2800        let report = collect(
2801            &pool,
2802            BuildIdentity::from_env("test", None),
2803            Duration::from_secs(30),
2804        );
2805        let reporter = crate::walpin::reporting_pid();
2806        let process = report
2807            .wal_pin
2808            .census_process_start_times
2809            .iter()
2810            .find(|process| process.pid == reporter)
2811            .expect("the reporting process must remain in the census list");
2812
2813        assert_eq!(report.wal_pin.reporting_pid, reporter);
2814        assert_eq!(report.wal_pin.reporting_process_is_holder, Some(true));
2815        assert_eq!(
2816            report.wal_pin.census_process_start_times.len(),
2817            report.wal_pin.census_holder_pids.len(),
2818            "every confirmed holder gets raw start-time data"
2819        );
2820        assert_eq!(
2821            process.process_start_time_secs,
2822            crate::walpin::process_start_time_secs(reporter)
2823        );
2824        assert_eq!(process.process_start_time_unavailable_reason, None);
2825        #[cfg(target_os = "linux")]
2826        assert_eq!(report.wal_pin.start_time_resolution_secs, Some(2));
2827        #[cfg(target_os = "macos")]
2828        assert_eq!(report.wal_pin.start_time_resolution_secs, Some(1));
2829        assert_eq!(
2830            serde_json::to_value(&report).unwrap()["wal_pin"]["census_process_start_times"]
2831                .as_array()
2832                .unwrap()
2833                .iter()
2834                .filter_map(|entry| entry["pid"].as_u64())
2835                .collect::<Vec<_>>(),
2836            report
2837                .wal_pin
2838                .census_holder_pids
2839                .iter()
2840                .map(|pid| u64::from(*pid))
2841                .collect::<Vec<_>>(),
2842            "start-time reporting does not filter census holders"
2843        );
2844    }
2845
2846    /// The verb must not perturb the state it reports: the probe touches none
2847    /// of the ADR-091 counters.
2848    #[test]
2849    #[serial(checkpoint_skip_metrics)]
2850    fn checkpoint_probe_does_not_perturb_the_adr091_counters() {
2851        crate::checkpoint::reset_checkpoint_metrics_for_tests();
2852        let dir = tempfile::tempdir().expect("tempdir");
2853        let (pool, _path) = seeded_pool(&dir);
2854        let conn = pool
2855            .open_standalone_writer_untracked()
2856            .expect("standalone probe connection");
2857
2858        let before = checkpoint_counters();
2859        for _ in 0..3 {
2860            checkpoint_probe(&conn).expect("probe must succeed");
2861        }
2862        let after = checkpoint_counters();
2863
2864        assert_eq!(
2865            before, after,
2866            "checkpoint_probe must leave every ADR-091 counter untouched"
2867        );
2868    }
2869
2870    #[test]
2871    fn wal_file_state_reports_the_sidecar_size_for_a_live_db() {
2872        let dir = tempfile::tempdir().expect("tempdir");
2873        let (_pool, path) = seeded_pool(&dir);
2874
2875        let state = wal_file_state(&path);
2876        assert!(
2877            state.wal_path.ends_with("diag.db-wal"),
2878            "WAL path is the db path plus a -wal suffix, got {}",
2879            state.wal_path
2880        );
2881        assert!(
2882            state.wal_size_bytes.is_some(),
2883            "a seeded WAL database must have a stat-able -wal file: {state:?}"
2884        );
2885        assert!(state.unavailable_reason.is_none(), "{state:?}");
2886    }
2887
2888    #[test]
2889    fn wal_file_state_degrades_with_a_reason_when_the_sidecar_is_absent() {
2890        let dir = tempfile::tempdir().expect("tempdir");
2891        let state = wal_file_state(&dir.path().join("never-created.db"));
2892        assert!(state.wal_size_bytes.is_none());
2893        assert!(
2894            state.unavailable_reason.is_some(),
2895            "an absent WAL file must carry a reason, not a silent zero: {state:?}"
2896        );
2897    }
2898
2899    #[test]
2900    fn collect_on_a_file_backed_db_carries_build_identity_and_every_counter() {
2901        let dir = tempfile::tempdir().expect("tempdir");
2902        let (pool, _path) = seeded_pool(&dir);
2903        let reader_admission_capacity = pool.max_readers().max(1);
2904
2905        let report = collect(
2906            &pool,
2907            BuildIdentity::from_env("9.9.9", Some("deadbeef")),
2908            Duration::from_secs(30),
2909        );
2910
2911        assert_eq!(report.build.version, "9.9.9");
2912        assert_eq!(report.build.build_hash.as_deref(), Some("deadbeef"));
2913        assert!(report.db_path.is_some());
2914        assert!(
2915            report.checkpoint_probe.is_some(),
2916            "file-backed collect must land a probe; error was {:?}",
2917            report.checkpoint_probe_error
2918        );
2919        assert!(
2920            report.wal_file.as_ref().and_then(|w| w.wal_size_bytes) >= Some(0),
2921            "wal_size_bytes must be a non-negative byte count when present"
2922        );
2923
2924        let json = serde_json::to_value(&report).expect("report serializes");
2925        let counters = json
2926            .get("checkpoint_counters")
2927            .expect("counters section present");
2928        for key in [
2929            "last_observed_wal_pages",
2930            "truncate_attempts",
2931            "truncate_consecutive_failures",
2932            "checkpoint_skipped_ticks",
2933            "checkpoint_consecutive_skips",
2934            "checkpoint_last_skip_wal_pages",
2935            "checkpoint_pressure_elevated_ticks",
2936            "checkpoint_pressure_episodes_started",
2937            "checkpoint_pressure_episodes_recovered",
2938            "checkpoint_lifecycle_append_attempts",
2939            "checkpoint_lifecycle_append_failures",
2940            "checkpoint_lifecycle_enqueue_drops",
2941            "read_tx_max_age_evictions",
2942        ] {
2943            assert!(counters.get(key).is_some(), "counter {key} must be present");
2944        }
2945        assert_eq!(
2946            report.writer_contention.writer_acquisitions, 1,
2947            "the seed write checked the finite-wait pooled writer out once"
2948        );
2949        assert_eq!(report.writer_contention.pooled_writer_acquisitions, 1);
2950        assert_eq!(report.writer_contention.standalone_writer_acquisitions, 0);
2951        assert_eq!(report.writer_contention.writer_task_acquisitions, 0);
2952        assert_eq!(report.writer_contention.writer_acquisition_timeouts, 0);
2953        assert_eq!(
2954            report.reader_contention,
2955            ReaderContentionDiagnostics {
2956                configured_reader_cap: pool.config().max_readers,
2957                configured_checkout_timeout_ms: u64::try_from(
2958                    pool.config().checkout_timeout.as_millis(),
2959                )
2960                .unwrap_or(u64::MAX),
2961                configured_busy_timeout_ms: u64::try_from(pool.config().busy_timeout.as_millis())
2962                    .unwrap_or(u64::MAX),
2963                reader_admission_capacity,
2964                available_reader_admission_slots: reader_admission_capacity,
2965                reader_acquisitions: 0,
2966                pooled_reader_checkouts: 0,
2967                standalone_reader_opens: 0,
2968                infrastructure_standalone_reader_opens: 0,
2969                reader_checkout_timeouts: 0,
2970                reader_busy_timeouts: 0,
2971                active_pooled_reader_checkouts: 0,
2972                peak_active_pooled_reader_checkouts: 0,
2973                completed_pooled_reader_checkouts: 0,
2974                max_completed_reader_hold_micros: 0,
2975                max_completed_reader_hold_operation: None,
2976                reader_replacement_open_failures: 0,
2977            },
2978            "the diagnostics probe itself must not masquerade as request reader traffic"
2979        );
2980        assert_eq!(
2981            report.graph_edge_integrity,
2982            Some(GraphEdgeIntegrity {
2983                duplicate_edge_id_groups: 0,
2984                graph_edges_rows: 0,
2985                graph_edges_seq_rows: 0,
2986                pre_v14_duplicate_edge_state_detected: false,
2987                live_entities_carrying_merged_into: 0,
2988            })
2989        );
2990        assert!(report.graph_edge_integrity_error.is_none());
2991        let fts_segments = report
2992            .fts_segments
2993            .as_ref()
2994            .expect("file-backed diagnostics must decode both FTS structure rows");
2995        assert_eq!(fts_segments.entities.segment_count, 1);
2996        assert_eq!(fts_segments.notes.segment_count, 1);
2997        assert_eq!(fts_segments.total_segments, 2);
2998        assert!(report.fts_segments_error.is_none());
2999        let fts_json = json
3000            .get("fts_segments")
3001            .expect("FTS segment diagnostics serialize");
3002        assert_eq!(fts_json["entities"]["segment_count"], 1);
3003        assert!(json.get("fts_maintenance").is_some());
3004        assert!(
3005            report.size_composition.is_some(),
3006            "file-backed diagnostics must include page composition; error was {:?}",
3007            report.size_composition_error
3008        );
3009        assert!(report.size_composition_error.is_none());
3010        assert!(report.writer_contention.audit_append_failures.is_none());
3011        assert!(report
3012            .writer_contention
3013            .audit_obligation_append_failures
3014            .is_none());
3015        assert!(report
3016            .writer_contention
3017            .audit_obligation_append_failures_unavailable_reason
3018            .is_some());
3019        assert!(json["writer_contention"]["audit_obligation_append_failures"].is_null());
3020        assert!(
3021            report
3022                .writer_contention
3023                .audit_append_failures_unavailable_reason
3024                .is_some(),
3025            "a direct khive-db snapshot must not fabricate a runtime audit count"
3026        );
3027    }
3028
3029    #[test]
3030    fn diagnostics_exposes_reader_saturation_and_completed_hold_evidence() {
3031        let dir = tempfile::tempdir().expect("tempdir");
3032        // A pre-opened reader makes setup independent of the exhaustion timeout.
3033        let pool = ConnectionPool::new(PoolConfig {
3034            path: Some(dir.path().join("reader_saturation.db")),
3035            max_readers: 1,
3036            checkout_timeout: Duration::from_millis(2),
3037            ..PoolConfig::default()
3038        })
3039        .expect("one-reader file-backed pool");
3040        let held = pool.reader().expect("first reader checkout");
3041        assert!(
3042            pool.reader().is_err(),
3043            "the live checkout must exhaust the one-slot reader budget"
3044        );
3045        drop(held);
3046
3047        let report = collect(
3048            &pool,
3049            BuildIdentity::from_env("9.9.9", None),
3050            Duration::from_secs(30),
3051        );
3052        let reader = report.reader_contention;
3053        assert_eq!(reader.reader_admission_capacity, 1);
3054        assert_eq!(reader.available_reader_admission_slots, 1);
3055        assert_eq!(reader.reader_acquisitions, 1);
3056        assert_eq!(reader.pooled_reader_checkouts, 1);
3057        assert_eq!(reader.standalone_reader_opens, 0);
3058        assert_eq!(reader.infrastructure_standalone_reader_opens, 0);
3059        assert_eq!(reader.reader_checkout_timeouts, 1);
3060        assert_eq!(reader.active_pooled_reader_checkouts, 0);
3061        assert_eq!(reader.peak_active_pooled_reader_checkouts, 1);
3062        assert_eq!(reader.completed_pooled_reader_checkouts, 1);
3063        assert!(reader.max_completed_reader_hold_micros > 0);
3064
3065        let json = serde_json::to_value(&report).expect("report serializes");
3066        assert_eq!(
3067            json.pointer("/reader_contention/reader_admission_capacity"),
3068            Some(&serde_json::json!(1)),
3069            "the operator wire payload must expose the reader admission budget"
3070        );
3071        assert_eq!(
3072            json.pointer("/reader_contention/reader_checkout_timeouts"),
3073            Some(&serde_json::json!(1)),
3074            "the operator wire payload must expose the reader timeout phase"
3075        );
3076        assert!(
3077            json.pointer("/reader_contention/max_completed_reader_hold_micros")
3078                .is_some(),
3079            "the operator wire payload must expose completed hold-time evidence"
3080        );
3081    }
3082
3083    /// A query refused with SQLITE_BUSY after the busy handler gives up shows up
3084    /// in `reader_busy_timeouts`, apart from `reader_checkout_timeouts`. WAL
3085    /// readers are never blocked by a writer, so the fixture uses a
3086    /// rollback-journal database, where a connection holding an exclusive lock
3087    /// refuses every other reader.
3088    #[test]
3089    fn diagnostics_counts_reader_busy_handler_timeouts_apart_from_checkout_timeouts() {
3090        let dir = tempfile::tempdir().expect("tempdir");
3091        let path = dir.path().join("reader_busy_timeouts.db");
3092        let pool = ConnectionPool::new(PoolConfig {
3093            path: Some(path.clone()),
3094            wal_mode: false,
3095            write_queue_enabled: Some(false),
3096            busy_timeout: Duration::from_millis(50),
3097            ..PoolConfig::default()
3098        })
3099        .expect("rollback-journal file-backed pool");
3100        pool.writer()
3101            .expect("writer")
3102            .conn()
3103            .execute_batch("CREATE TABLE busy_fixture (id INTEGER PRIMARY KEY)")
3104            .expect("fixture table");
3105
3106        // Control: with no lock held the read succeeds and nothing is counted.
3107        let reader = pool.reader().expect("reader checkout");
3108        let rows = reader
3109            .query_row("SELECT count(*) FROM busy_fixture", [], |row| {
3110                row.get::<_, i64>(0)
3111            })
3112            .expect("unlocked read");
3113        assert_eq!(rows, 0);
3114        drop(reader);
3115        assert_eq!(
3116            ReaderContentionDiagnostics::snapshot(&pool).reader_busy_timeouts,
3117            0
3118        );
3119
3120        let holder = Connection::open(&path).expect("second connection");
3121        holder
3122            .execute_batch("BEGIN EXCLUSIVE")
3123            .expect("exclusive lock");
3124        let reader = pool.reader().expect("reader checkout");
3125        let refused = reader
3126            .query_row("SELECT count(*) FROM busy_fixture", [], |row| {
3127                row.get::<_, i64>(0)
3128            })
3129            .expect_err("a read behind an exclusive lock must be refused");
3130        assert!(
3131            matches!(
3132                &refused,
3133                crate::SqliteError::Rusqlite(rusqlite::Error::SqliteFailure(code, _))
3134                    if code.code == rusqlite::ErrorCode::DatabaseBusy
3135            ),
3136            "the refusal must be SQLITE_BUSY: {refused}"
3137        );
3138        holder.execute_batch("ROLLBACK").expect("release lock");
3139        drop(reader);
3140
3141        let snapshot = ReaderContentionDiagnostics::snapshot(&pool);
3142        assert_eq!(snapshot.reader_busy_timeouts, 1);
3143        assert_eq!(
3144            snapshot.reader_checkout_timeouts, 0,
3145            "a busy-handler refusal after checkout is not a checkout timeout"
3146        );
3147        let json = serde_json::to_value(snapshot).expect("snapshot serializes");
3148        assert_eq!(
3149            json.pointer("/reader_busy_timeouts"),
3150            Some(&serde_json::json!(1)),
3151            "the operator wire payload must expose the busy-handler count"
3152        );
3153    }
3154
3155    #[test]
3156    fn diagnostics_reports_configured_reader_budget_and_both_deadlines() {
3157        let pool = ConnectionPool::new(PoolConfig {
3158            max_readers: 6,
3159            checkout_timeout: Duration::from_millis(17),
3160            busy_timeout: Duration::from_millis(31),
3161            ..PoolConfig::default()
3162        })
3163        .expect("in-memory pool");
3164        let report = collect(
3165            &pool,
3166            BuildIdentity::from_env("9.9.9", None),
3167            Duration::from_secs(30),
3168        );
3169        let reader = report.reader_contention;
3170        assert_eq!(reader.reader_admission_capacity, 1);
3171
3172        let json = serde_json::to_value(&report).expect("report serializes");
3173        assert_eq!(
3174            json.pointer("/reader_contention/configured_reader_cap"),
3175            Some(&serde_json::json!(6))
3176        );
3177        assert_eq!(
3178            json.pointer("/reader_contention/configured_checkout_timeout_ms"),
3179            Some(&serde_json::json!(17))
3180        );
3181        assert_eq!(
3182            json.pointer("/reader_contention/configured_busy_timeout_ms"),
3183            Some(&serde_json::json!(31))
3184        );
3185    }
3186
3187    #[test]
3188    fn diagnostics_composes_file_backed_standalone_acquisitions_without_counting_its_probe() {
3189        let dir = tempfile::tempdir().expect("tempdir");
3190        let (pool, _path) = seeded_pool(&dir);
3191
3192        drop(
3193            pool.open_standalone_writer()
3194                .expect("write-traffic standalone connection"),
3195        );
3196
3197        let report = collect_with_audit_append_failures(
3198            &pool,
3199            BuildIdentity::from_env("9.9.9", None),
3200            Duration::from_secs(30),
3201            0,
3202        );
3203        assert_eq!(report.writer_contention.writer_acquisitions, 2);
3204        assert_eq!(report.writer_contention.pooled_writer_acquisitions, 1);
3205        assert_eq!(report.writer_contention.standalone_writer_acquisitions, 1);
3206        assert_eq!(report.writer_contention.writer_task_acquisitions, 0);
3207        assert_eq!(report.writer_contention.writer_acquisition_timeouts, 0);
3208
3209        let second = collect_with_audit_append_failures(
3210            &pool,
3211            BuildIdentity::from_env("9.9.9", None),
3212            Duration::from_secs(30),
3213            0,
3214        );
3215        assert_eq!(
3216            second.writer_contention, report.writer_contention,
3217            "the diagnostics PASSIVE probe must not inflate write-traffic counters"
3218        );
3219    }
3220
3221    #[test]
3222    fn runtime_aware_collect_exposes_the_supplied_audit_failure_counter() {
3223        let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool");
3224
3225        let report = collect_with_audit_append_failures(
3226            &pool,
3227            BuildIdentity::from_env("9.9.9", None),
3228            Duration::from_secs(30),
3229            17,
3230        );
3231
3232        assert_eq!(report.writer_contention.audit_append_failures, Some(17));
3233        assert!(
3234            report
3235                .writer_contention
3236                .audit_append_failures_unavailable_reason
3237                .is_none(),
3238            "a supplied runtime counter must not carry an unavailable reason"
3239        );
3240    }
3241
3242    #[test]
3243    fn diagnostics_exposes_an_induced_writer_checkout_timeout() {
3244        let pool = ConnectionPool::new(PoolConfig {
3245            checkout_timeout: Duration::from_millis(1),
3246            ..PoolConfig::default()
3247        })
3248        .expect("in-memory pool");
3249
3250        let held = pool.writer().expect("first writer checkout succeeds");
3251        assert!(
3252            matches!(
3253                pool.writer(),
3254                Err(crate::SqliteError::WriterPoolCheckoutTimeout { .. })
3255            ),
3256            "holding the sole writer must exercise the typed timeout path"
3257        );
3258        drop(held);
3259
3260        let report = collect_with_audit_append_failures(
3261            &pool,
3262            BuildIdentity::from_env("9.9.9", None),
3263            Duration::from_secs(30),
3264            0,
3265        );
3266        assert_eq!(report.writer_contention.writer_acquisitions, 1);
3267        assert_eq!(report.writer_contention.pooled_writer_acquisitions, 1);
3268        assert_eq!(report.writer_contention.standalone_writer_acquisitions, 0);
3269        assert_eq!(report.writer_contention.writer_task_acquisitions, 0);
3270        assert_eq!(report.writer_contention.writer_acquisition_timeouts, 1);
3271    }
3272
3273    /// The `u64::MAX` never-observed sentinel must serialize as `null`, never
3274    /// as a huge number an operator would read as a real page count.
3275    #[test]
3276    fn never_observed_sentinels_serialize_as_null() {
3277        let counters = CheckpointCounters {
3278            last_observed_wal_pages: None,
3279            truncate_attempts: 0,
3280            truncate_consecutive_failures: 0,
3281            checkpoint_skipped_ticks: 0,
3282            checkpoint_consecutive_skips: 0,
3283            checkpoint_last_skip_wal_pages: None,
3284            checkpoint_pressure_elevated_ticks: 0,
3285            checkpoint_pressure_episodes_started: 0,
3286            checkpoint_pressure_episodes_recovered: 0,
3287            checkpoint_lifecycle_append_attempts: 0,
3288            checkpoint_lifecycle_append_failures: 0,
3289            checkpoint_lifecycle_enqueue_drops: 0,
3290            read_tx_max_age_evictions: 0,
3291        };
3292        let json = serde_json::to_value(counters).expect("serializes");
3293        assert!(json["last_observed_wal_pages"].is_null());
3294        assert!(json["checkpoint_last_skip_wal_pages"].is_null());
3295    }
3296
3297    #[test]
3298    fn graph_edge_integrity_detects_the_pre_v14_duplicate_state() {
3299        let conn = Connection::open_in_memory().expect("in-memory sqlite");
3300        conn.execute_batch(
3301            "CREATE TABLE graph_edges (
3302                 namespace TEXT NOT NULL,
3303                 id TEXT NOT NULL,
3304                 PRIMARY KEY (namespace, id)
3305             );
3306             CREATE TABLE graph_edges_seq (
3307                 seq INTEGER PRIMARY KEY AUTOINCREMENT,
3308                 edge_id TEXT NOT NULL UNIQUE
3309             );
3310             CREATE TABLE entities (id TEXT PRIMARY KEY, deleted_at INTEGER, merged_into TEXT);
3311             INSERT INTO graph_edges(namespace, id)
3312             VALUES ('alpha', 'shared-edge'), ('beta', 'shared-edge');
3313             INSERT INTO graph_edges_seq(edge_id) VALUES ('shared-edge');",
3314        )
3315        .expect("seed the state possible before the V14 uniqueness guard");
3316
3317        let integrity = graph_edge_integrity(&conn).expect("integrity query succeeds");
3318
3319        assert_eq!(integrity.duplicate_edge_id_groups, 1);
3320        assert_eq!(integrity.graph_edges_rows, 2);
3321        assert_eq!(integrity.graph_edges_seq_rows, 1);
3322        assert!(integrity.pre_v14_duplicate_edge_state_detected);
3323    }
3324
3325    #[test]
3326    fn graph_edge_integrity_counts_only_live_rows_that_still_carry_merge_provenance() {
3327        let conn = Connection::open_in_memory().expect("in-memory sqlite");
3328        conn.execute_batch(
3329            "CREATE TABLE graph_edges (
3330                 namespace TEXT NOT NULL,
3331                 id TEXT NOT NULL,
3332                 PRIMARY KEY (namespace, id)
3333             );
3334             CREATE TABLE graph_edges_seq (
3335                 seq INTEGER PRIMARY KEY AUTOINCREMENT,
3336                 edge_id TEXT NOT NULL UNIQUE
3337             );
3338             CREATE TABLE entities (id TEXT PRIMARY KEY, deleted_at INTEGER, merged_into TEXT);
3339             INSERT INTO entities(id, deleted_at, merged_into) VALUES
3340                 ('kept', NULL, NULL),
3341                 ('tombstoned-source', 1, 'kept'),
3342                 ('left-live-by-an-old-restore', NULL, 'kept'),
3343                 ('plain-soft-delete', 1, NULL);",
3344        )
3345        .expect("seed one row of each shape");
3346
3347        let integrity = graph_edge_integrity(&conn).expect("integrity query succeeds");
3348
3349        assert_eq!(
3350            integrity.live_entities_carrying_merged_into, 1,
3351            "a merge tombstone and a plain live row are both in order; only the live row \
3352             carrying merged_into is the invariant violation"
3353        );
3354        assert_eq!(integrity.duplicate_edge_id_groups, 0);
3355    }
3356
3357    #[test]
3358    fn graph_edge_integrity_does_not_mislabel_retained_delete_history_as_a_duplicate() {
3359        let conn = Connection::open_in_memory().expect("in-memory sqlite");
3360        conn.execute_batch(
3361            "CREATE TABLE graph_edges (
3362                 namespace TEXT NOT NULL,
3363                 id TEXT NOT NULL,
3364                 PRIMARY KEY (namespace, id)
3365             );
3366             CREATE TABLE graph_edges_seq (
3367                 seq INTEGER PRIMARY KEY AUTOINCREMENT,
3368                 edge_id TEXT NOT NULL UNIQUE
3369             );
3370             CREATE TABLE entities (id TEXT PRIMARY KEY, deleted_at INTEGER, merged_into TEXT);
3371             INSERT INTO graph_edges(namespace, id) VALUES ('local', 'live-edge');
3372             INSERT INTO graph_edges_seq(edge_id)
3373             VALUES ('deleted-edge'), ('live-edge');",
3374        )
3375        .expect("seed a retained sequence row for a hard-deleted edge");
3376
3377        let integrity = graph_edge_integrity(&conn).expect("integrity query succeeds");
3378
3379        assert_eq!(integrity.duplicate_edge_id_groups, 0);
3380        assert_eq!(integrity.graph_edges_rows, 1);
3381        assert_eq!(integrity.graph_edges_seq_rows, 2);
3382        assert!(
3383            !integrity.pre_v14_duplicate_edge_state_detected,
3384            "ledger rows intentionally survive hard deletion; count mismatch alone is not the \
3385             pre-V14 duplicate state"
3386        );
3387    }
3388
3389    #[cfg(unix)]
3390    #[test]
3391    fn unmeasured_sidecar_cleanup_fields_are_absent_from_the_wire_payload() {
3392        let pin = wal_pin_attribution_from_census(crate::walpin::CensusResult {
3393            holders: std::collections::HashSet::new(),
3394            uninspectable_pids: Vec::new(),
3395            truncated: false,
3396            budget_exhausted: false,
3397        });
3398
3399        assert_eq!(pin.sidecar_listing_truncated, None);
3400        assert_eq!(pin.sidecar_entries_cleanup_would_reap, None);
3401        let json = serde_json::to_value(pin).expect("attribution serializes");
3402        assert!(
3403            json.get("sidecar_listing_truncated").is_none(),
3404            "a skipped enumeration must omit sidecar_listing_truncated, not fabricate false"
3405        );
3406        assert!(
3407            json.get("sidecar_entries_cleanup_would_reap").is_none(),
3408            "a skipped enumeration must omit sidecar_entries_cleanup_would_reap, not fabricate 0"
3409        );
3410    }
3411
3412    #[cfg(unix)]
3413    #[test]
3414    fn wal_pin_census_serializes_only_under_the_nested_carrier() {
3415        let pin = wal_pin_attribution_from_census(crate::walpin::CensusResult {
3416            holders: std::collections::HashSet::from([41, 7]),
3417            uninspectable_pids: vec![99],
3418            truncated: true,
3419            budget_exhausted: false,
3420        });
3421
3422        let json = serde_json::to_value(pin).expect("attribution serializes");
3423        assert_eq!(json["census"]["holder_pids"], serde_json::json!([7, 41]));
3424        assert_eq!(
3425            json["census"]["uninspectable_pids"],
3426            serde_json::json!([99])
3427        );
3428        assert_eq!(json["census"]["truncated"], true);
3429
3430        for duplicate in [
3431            "census_holder_pids",
3432            "census_uninspectable_pids",
3433            "census_truncated",
3434            "census_is_complete",
3435        ] {
3436            assert!(
3437                json.get(duplicate).is_none(),
3438                "wal_pin.{duplicate} must not duplicate wal_pin.census: {json}"
3439            );
3440        }
3441    }
3442
3443    /// A missing configured path must never be created by a diagnostic
3444    /// request. The untracked standalone open omits `SQLITE_OPEN_CREATE`, so
3445    /// the probe degrades to an error and the file stays absent.
3446    #[test]
3447    fn probe_refuses_a_missing_configured_path_without_creating_it() {
3448        let dir = tempfile::tempdir().expect("tempdir");
3449        let (pool, path) = seeded_pool(&dir);
3450
3451        for suffix in ["", "-wal", "-shm"] {
3452            let mut p = path.as_os_str().to_os_string();
3453            p.push(suffix);
3454            let _ = std::fs::remove_file(PathBuf::from(p));
3455        }
3456        assert!(!path.exists(), "precondition: the database file is gone");
3457
3458        let report = collect(
3459            &pool,
3460            BuildIdentity::from_env("0.0.0", None),
3461            Duration::from_secs(30),
3462        );
3463
3464        assert!(
3465            report.checkpoint_probe.is_none(),
3466            "a missing database must not yield a probe result: {report:?}"
3467        );
3468        assert!(
3469            report.checkpoint_probe_error.is_some(),
3470            "a missing database must say why there is no probe: {report:?}"
3471        );
3472        assert!(
3473            !path.exists(),
3474            "a diagnostics request must never create the database it was asked about"
3475        );
3476    }
3477
3478    /// In-memory backends have no WAL file and no census target: the report
3479    /// still returns, with explicit reasons rather than missing sections.
3480    #[test]
3481    fn collect_degrades_gracefully_for_an_in_memory_backend() {
3482        let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool");
3483        let report = collect(
3484            &pool,
3485            BuildIdentity::from_env("0.0.0", None),
3486            Duration::from_secs(30),
3487        );
3488
3489        assert!(report.db_path.is_none());
3490        assert!(report.wal_file.is_none());
3491        assert!(report.checkpoint_probe.is_none());
3492        assert!(report.fts_segments.is_none());
3493        assert!(report.fts_segments_error.is_some());
3494        assert!(
3495            report.checkpoint_probe_error.is_some(),
3496            "an in-memory report must say WHY there is no probe"
3497        );
3498        assert!(!report.wal_pin.available);
3499        assert!(report.wal_pin.unavailable_reason.is_some());
3500        assert_eq!(report.wal_pin.status, WalPinAttributionStatus::Unavailable);
3501        assert!(matches!(
3502            report.wal_pin.census,
3503            WalPinCensus::Unavailable { .. }
3504        ));
3505        assert!(report.size_composition.is_none());
3506        assert!(report
3507            .size_composition_error
3508            .as_deref()
3509            .is_some_and(|reason| reason.contains("no file-backed page composition")));
3510    }
3511
3512    #[cfg(unix)]
3513    #[test]
3514    fn incomplete_holder_census_is_a_tagged_degraded_result() {
3515        let census = crate::walpin::CensusResult {
3516            holders: std::collections::HashSet::from([41, 7]),
3517            uninspectable_pids: vec![99, 99],
3518            truncated: true,
3519            budget_exhausted: false,
3520        };
3521
3522        let pin = wal_pin_attribution_from_census(census);
3523
3524        assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
3525        assert!(!pin.available);
3526        assert!(!pin.census_is_complete);
3527        assert_eq!(pin.census_holder_pids, vec![7, 41]);
3528        assert_eq!(pin.census_uninspectable_pids, vec![99]);
3529        assert!(
3530            pin.unavailable_reason.as_deref().is_some_and(
3531                |reason| reason.contains("additional database holders cannot be ruled out")
3532            ),
3533            "the legacy reason must also fail loud for old consumers: {pin:?}"
3534        );
3535        match &pin.census {
3536            WalPinCensus::Incomplete {
3537                holder_pids,
3538                uninspectable_pids,
3539                truncated,
3540                reason,
3541            } => {
3542                assert_eq!(holder_pids, &vec![7, 41]);
3543                assert_eq!(uninspectable_pids, &vec![99]);
3544                assert!(*truncated);
3545                assert!(reason.contains("additional database holders cannot be ruled out"));
3546            }
3547            other => panic!("incomplete scan must serialize as incomplete, got {other:?}"),
3548        }
3549
3550        let json = serde_json::to_value(&pin).expect("serializes");
3551        assert_eq!(json["status"], "degraded");
3552        assert_eq!(json["census"]["status"], "incomplete");
3553    }
3554
3555    #[cfg(unix)]
3556    #[test]
3557    fn complete_holder_census_stays_explicit_while_attribution_is_degraded() {
3558        let census = crate::walpin::CensusResult {
3559            holders: std::collections::HashSet::from([7]),
3560            uninspectable_pids: Vec::new(),
3561            truncated: false,
3562            budget_exhausted: false,
3563        };
3564
3565        let pin = wal_pin_attribution_from_census(census);
3566
3567        assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
3568        assert!(pin.census_is_complete);
3569        assert!(matches!(
3570            pin.census,
3571            WalPinCensus::Complete { ref holder_pids } if holder_pids == &vec![7]
3572        ));
3573        assert_eq!(
3574            pin.status_reasons.len(),
3575            1,
3576            "only missing sidecar reconciliation degrades a complete OS census"
3577        );
3578    }
3579
3580    #[cfg(unix)]
3581    #[test]
3582    fn complete_holder_and_read_only_sidecar_evidence_reconcile_to_complete() {
3583        let census = crate::walpin::CensusResult {
3584            holders: std::collections::HashSet::from([7]),
3585            uninspectable_pids: Vec::new(),
3586            truncated: false,
3587            budget_exhausted: false,
3588        };
3589        let sidecar = crate::walpin::WalpinReport {
3590            entries: vec![crate::walpin::WalpinPidHealth::RegisteredSilent { pid: 7 }],
3591            sidecar_listing_truncated: false,
3592            cleanup_would_reap: 0,
3593            orphan_temps_reaped: 0,
3594        };
3595
3596        let pin = wal_pin_attribution_from_evidence(census, sidecar);
3597
3598        assert_eq!(pin.status, WalPinAttributionStatus::Complete);
3599        assert!(pin.available);
3600        assert!(pin.fully_attributed);
3601        assert!(pin.status_reasons.is_empty());
3602        assert_eq!(pin.registered_silent_pids, vec![7]);
3603        assert!(pin.census_pids_without_attribution.is_empty());
3604        assert_eq!(pin.sidecar_listing_truncated, Some(false));
3605        assert_eq!(pin.sidecar_entries_cleanup_would_reap, Some(0));
3606    }
3607
3608    #[cfg(unix)]
3609    #[test]
3610    fn complete_census_with_an_unregistered_holder_is_degraded_not_exonerated() {
3611        let census = crate::walpin::CensusResult {
3612            holders: std::collections::HashSet::from([7, 41]),
3613            uninspectable_pids: Vec::new(),
3614            truncated: false,
3615            budget_exhausted: false,
3616        };
3617        let sidecar = crate::walpin::WalpinReport {
3618            entries: vec![crate::walpin::WalpinPidHealth::RegisteredSilent { pid: 7 }],
3619            sidecar_listing_truncated: false,
3620            cleanup_would_reap: 0,
3621            orphan_temps_reaped: 0,
3622        };
3623
3624        let pin = wal_pin_attribution_from_evidence(census, sidecar);
3625
3626        assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
3627        assert!(!pin.available);
3628        assert!(!pin.fully_attributed);
3629        assert_eq!(pin.census_pids_without_attribution, vec![41]);
3630        assert!(pin
3631            .status_reasons
3632            .iter()
3633            .any(|reason| reason.contains("holder(s) have no sidecar attribution")));
3634    }
3635
3636    #[test]
3637    fn database_size_composition_reports_tables_indexes_fts_and_vectors_separately() {
3638        let conn = Connection::open_in_memory().expect("in-memory sqlite");
3639        conn.execute_batch(
3640            "CREATE TABLE docs(id INTEGER PRIMARY KEY, body TEXT NOT NULL); \
3641             CREATE INDEX idx_docs_body ON docs(body); \
3642             CREATE TABLE fts_demo_data(id INTEGER PRIMARY KEY, block BLOB); \
3643             CREATE TABLE vec_demo_chunks(id INTEGER PRIMARY KEY, vectors BLOB); \
3644             CREATE TABLE knowledge_sections(id INTEGER PRIMARY KEY, embedding BLOB); \
3645             INSERT INTO docs(body) VALUES (zeroblob(8192)); \
3646             INSERT INTO fts_demo_data(block) VALUES (zeroblob(8192)); \
3647             INSERT INTO vec_demo_chunks(vectors) VALUES (zeroblob(8192)); \
3648             INSERT INTO knowledge_sections(embedding) VALUES (zeroblob(8192));",
3649        )
3650        .expect("seed size classes");
3651
3652        let composition = database_size_composition(&conn).expect("dbstat composition");
3653        let class_for = |name: &str| {
3654            composition
3655                .objects
3656                .iter()
3657                .find(|object| object.name == name)
3658                .map(|object| object.storage_class)
3659        };
3660
3661        assert_eq!(class_for("docs"), Some(DatabaseStorageClass::RowTable));
3662        assert_eq!(
3663            class_for("idx_docs_body"),
3664            Some(DatabaseStorageClass::Index)
3665        );
3666        assert_eq!(
3667            class_for("fts_demo_data"),
3668            Some(DatabaseStorageClass::FullText)
3669        );
3670        assert_eq!(
3671            class_for("vec_demo_chunks"),
3672            Some(DatabaseStorageClass::Vector)
3673        );
3674        assert_eq!(
3675            class_for("knowledge_sections"),
3676            Some(DatabaseStorageClass::MixedRowAndEmbedding)
3677        );
3678        assert!(composition.vector_bytes > 0);
3679        assert!(composition.full_text_bytes > 0);
3680        assert!(composition.mixed_embedding_bytes > 0);
3681        assert_eq!(
3682            composition
3683                .accounted_bytes
3684                .saturating_add(composition.freelist_bytes)
3685                .saturating_add(composition.unaccounted_bytes),
3686            composition.database_bytes
3687        );
3688    }
3689
3690    /// A holder with no sidecar registration remains degraded even though the
3691    /// read-only sidecar pass itself completed successfully.
3692    #[cfg(unix)]
3693    #[test]
3694    #[serial(khive_walpin_sidecar_env)]
3695    fn wal_pin_attribution_degrades_when_a_holder_has_no_sidecar_registration() {
3696        let dir = tempfile::tempdir().expect("tempdir");
3697        let (pool, path) = seeded_pool(&dir);
3698        let _ = &pool;
3699
3700        let pin = wal_pin_attribution(&path, Duration::from_secs(30));
3701
3702        assert!(
3703            !pin.fully_attributed,
3704            "the pool's OS holder has no test sidecar registration"
3705        );
3706        assert!(
3707            pin.unavailable_reason.is_some(),
3708            "the missing holder attribution must be explained: {pin:?}"
3709        );
3710        assert!(pin.sidecar_entries.is_empty());
3711        assert!(pin.reporting.is_empty());
3712        assert_eq!(pin.sidecar_listing_truncated, Some(false));
3713        assert_eq!(pin.sidecar_entries_cleanup_would_reap, Some(0));
3714        assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
3715        assert!(matches!(
3716            pin.census,
3717            WalPinCensus::Complete { .. } | WalPinCensus::Incomplete { .. }
3718        ));
3719    }
3720
3721    /// `sidecar_dir_for` is a purely lexical derivation from whatever path it
3722    /// is handed (`pool.rs`'s `sidecar_dir_for_alias_convergence` proves
3723    /// this at the primitive level). Diagnostics must feed it
3724    /// `pool.canonical_path()` — the same value the checkpoint sidecar
3725    /// writers use — never the pool's raw configured path, or a symlinked
3726    /// database misses its own sidecar evidence entirely.
3727    #[cfg(unix)]
3728    #[test]
3729    #[serial(khive_walpin_sidecar_env)]
3730    fn diagnostics_finds_sidecar_evidence_through_an_aliased_database_path() {
3731        let dir = tempfile::tempdir().expect("tempdir");
3732        let real_dir = dir.path().join("real");
3733        std::fs::create_dir(&real_dir).expect("mkdir real dir");
3734        let real_path = real_dir.join("diag.db");
3735        std::fs::write(&real_path, b"").expect("create real file");
3736        let alias_path = dir.path().join("alias.db");
3737        std::os::unix::fs::symlink(&real_path, &alias_path).expect("symlink alias");
3738
3739        let pool = ConnectionPool::new(PoolConfig {
3740            path: Some(alias_path.clone()),
3741            ..PoolConfig::for_test()
3742        })
3743        .expect("pool open through symlinked path");
3744        {
3745            let writer = pool.try_writer().expect("writer");
3746            writer
3747                .conn()
3748                .execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
3749                .expect("seed a write so the WAL file exists");
3750        }
3751
3752        let canonical = pool
3753            .canonical_path()
3754            .expect("file-backed pool")
3755            .to_path_buf();
3756        assert_ne!(
3757            canonical, alias_path,
3758            "the alias must actually differ from the canonical path for this test to mean \
3759             anything"
3760        );
3761
3762        let pid = std::process::id();
3763        let sidecar_dir = crate::walpin::sidecar_dir_for(&canonical);
3764        let beacon = crate::walpin::WalpinBeacon {
3765            pid,
3766            process_role: "session".to_string(),
3767            started_at: crate::walpin::process_start_time_secs(pid).unwrap_or(0),
3768            sweep_interval_ms: 5_000,
3769        };
3770        crate::walpin::write_beacon(&sidecar_dir, &beacon).expect("seed this process's beacon");
3771
3772        let report = collect(
3773            &pool,
3774            BuildIdentity::from_env("test", None),
3775            Duration::from_secs(30),
3776        );
3777
3778        // The wider OS holder census can be legitimately incomplete in a
3779        // sandboxed test environment (other, unrelated processes this test
3780        // has no permission to inspect) — that variability is orthogonal to
3781        // what this test checks. What must hold regardless is that the
3782        // sidecar enumeration itself, which is keyed off the canonical
3783        // path, actually ran, completed untruncated, and found this
3784        // process's own beacon rather than missing it beside the wrong
3785        // (aliased) directory.
3786        assert_eq!(
3787            report.wal_pin.sidecar_listing_truncated,
3788            Some(false),
3789            "the sidecar enumeration must run to completion: {:?}",
3790            report.wal_pin
3791        );
3792        assert!(
3793            report.wal_pin.registered_silent_pids.contains(&pid),
3794            "the beacon written beside the canonical path must be found: {:?}",
3795            report.wal_pin
3796        );
3797        assert!(
3798            report.wal_pin.census_holder_pids.contains(&pid),
3799            "the OS census must find this process holding its own database open: {:?}",
3800            report.wal_pin
3801        );
3802        assert!(
3803            !report
3804                .wal_pin
3805                .census_pids_without_attribution
3806                .contains(&pid),
3807            "this process's own holder entry must be attributed by its own sidecar evidence, \
3808             not left unexplained: {:?}",
3809            report.wal_pin
3810        );
3811    }
3812
3813    /// The async collector (`collect_with_audit_append_failures_interruptibly`)
3814    /// resolves its operational path independently of the sync collector
3815    /// (`collect`) — see `operational_db_path`. A regression that reintroduced
3816    /// the raw configured path on only the async side would leave the sync
3817    /// alias test above green while the async path silently missed its own
3818    /// sidecar evidence; this exercises the same aliasing scenario through
3819    /// the async entry point.
3820    #[cfg(unix)]
3821    #[tokio::test]
3822    #[serial(khive_walpin_sidecar_env)]
3823    async fn diagnostics_finds_sidecar_evidence_through_an_aliased_database_path_async() {
3824        let dir = tempfile::tempdir().expect("tempdir");
3825        let real_dir = dir.path().join("real");
3826        std::fs::create_dir(&real_dir).expect("mkdir real dir");
3827        let real_path = real_dir.join("diag.db");
3828        std::fs::write(&real_path, b"").expect("create real file");
3829        let alias_path = dir.path().join("alias.db");
3830        std::os::unix::fs::symlink(&real_path, &alias_path).expect("symlink alias");
3831
3832        let pool = ConnectionPool::new(PoolConfig {
3833            path: Some(alias_path.clone()),
3834            ..PoolConfig::for_test()
3835        })
3836        .expect("pool open through symlinked path");
3837        {
3838            let writer = pool.try_writer().expect("writer");
3839            writer
3840                .conn()
3841                .execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
3842                .expect("seed a write so the WAL file exists");
3843        }
3844
3845        let canonical = pool
3846            .canonical_path()
3847            .expect("file-backed pool")
3848            .to_path_buf();
3849        assert_ne!(
3850            canonical, alias_path,
3851            "the alias must actually differ from the canonical path for this test to mean \
3852             anything"
3853        );
3854
3855        let pid = std::process::id();
3856        let sidecar_dir = crate::walpin::sidecar_dir_for(&canonical);
3857        let beacon = crate::walpin::WalpinBeacon {
3858            pid,
3859            process_role: "session".to_string(),
3860            started_at: crate::walpin::process_start_time_secs(pid).unwrap_or(0),
3861            sweep_interval_ms: 5_000,
3862        };
3863        crate::walpin::write_beacon(&sidecar_dir, &beacon).expect("seed this process's beacon");
3864
3865        let pool = Arc::new(pool);
3866        let report = collect_with_audit_append_failures_interruptibly(
3867            Arc::clone(&pool),
3868            BuildIdentity::from_env("test", None),
3869            Duration::from_secs(30),
3870            0,
3871        )
3872        .await
3873        .expect("diagnostics succeed");
3874
3875        // Same rationale as the sync test: the wider OS holder census can be
3876        // legitimately incomplete in a sandboxed environment, but the sidecar
3877        // enumeration itself — keyed off the canonical path — must run to
3878        // completion and find this process's own beacon.
3879        assert_eq!(
3880            report.wal_pin.sidecar_listing_truncated,
3881            Some(false),
3882            "the sidecar enumeration must run to completion: {:?}",
3883            report.wal_pin
3884        );
3885        assert!(
3886            report.wal_pin.registered_silent_pids.contains(&pid),
3887            "the beacon written beside the canonical path must be found: {:?}",
3888            report.wal_pin
3889        );
3890        assert!(
3891            report.wal_pin.census_holder_pids.contains(&pid),
3892            "the OS census must find this process holding its own database open: {:?}",
3893            report.wal_pin
3894        );
3895        assert!(
3896            !report
3897                .wal_pin
3898                .census_pids_without_attribution
3899                .contains(&pid),
3900            "this process's own holder entry must be attributed by its own sidecar evidence, \
3901             not left unexplained: {:?}",
3902            report.wal_pin
3903        );
3904    }
3905}