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