use std::path::{Path, PathBuf};
use std::time::Duration;
use rusqlite::Connection;
use serde::Serialize;
use crate::checkpoint;
use crate::pool::ConnectionPool;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
pub struct CheckpointProbe {
pub busy: i64,
pub log_frames: i64,
pub checkpointed_frames: i64,
}
impl CheckpointProbe {
pub fn pin_depth(&self) -> i64 {
(self.log_frames - self.checkpointed_frames).max(0)
}
}
pub fn checkpoint_probe(conn: &Connection) -> rusqlite::Result<CheckpointProbe> {
conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
Ok(CheckpointProbe {
busy: row.get(0)?,
log_frames: row.get(1)?,
checkpointed_frames: row.get(2)?,
})
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
pub struct CheckpointCounters {
pub last_observed_wal_pages: Option<u64>,
pub truncate_attempts: u64,
pub truncate_consecutive_failures: u64,
pub checkpoint_skipped_ticks: u64,
pub checkpoint_consecutive_skips: u64,
pub checkpoint_last_skip_wal_pages: Option<u64>,
}
pub fn checkpoint_counters() -> CheckpointCounters {
CheckpointCounters {
last_observed_wal_pages: checkpoint::last_observed_wal_pages(),
truncate_attempts: checkpoint::truncate_attempts(),
truncate_consecutive_failures: checkpoint::truncate_consecutive_failures(),
checkpoint_skipped_ticks: checkpoint::checkpoint_skipped_ticks(),
checkpoint_consecutive_skips: checkpoint::checkpoint_consecutive_skips(),
checkpoint_last_skip_wal_pages: checkpoint::checkpoint_last_skip_wal_pages(),
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct BuildIdentity {
pub version: String,
pub build_hash: Option<String>,
}
impl BuildIdentity {
pub fn from_env(version: &str, build_hash: Option<&str>) -> Self {
Self {
version: version.to_string(),
build_hash: build_hash.map(str::to_string),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct WalFileState {
pub wal_path: String,
pub wal_size_bytes: Option<u64>,
pub unavailable_reason: Option<String>,
}
pub fn wal_file_state(db_path: &Path) -> WalFileState {
let wal_path = wal_sidecar_path(db_path);
match std::fs::metadata(&wal_path) {
Ok(md) => WalFileState {
wal_path: wal_path.display().to_string(),
wal_size_bytes: Some(md.len()),
unavailable_reason: None,
},
Err(e) => WalFileState {
wal_path: wal_path.display().to_string(),
wal_size_bytes: None,
unavailable_reason: Some(e.to_string()),
},
}
}
fn wal_sidecar_path(db_path: &Path) -> PathBuf {
let mut s = db_path.as_os_str().to_os_string();
s.push("-wal");
PathBuf::from(s)
}
#[derive(Debug, Clone, PartialEq, Serialize)]
pub struct WalPinAttribution {
pub available: bool,
pub unavailable_reason: Option<String>,
pub census_holder_pids: Vec<u32>,
pub census_uninspectable_pids: Vec<u32>,
pub census_truncated: bool,
pub census_is_complete: bool,
pub reporting: Vec<WalPinHolder>,
pub registered_silent_pids: Vec<u32>,
pub unknown_pids: Vec<u32>,
pub census_pids_without_attribution: Vec<u32>,
pub fully_attributed: bool,
pub sidecar_entries: Vec<serde_json::Value>,
pub sidecar_listing_truncated: bool,
pub sidecar_entries_cleanup_would_reap: usize,
}
#[derive(Debug, Clone, PartialEq, Serialize)]
pub struct WalPinHolder {
pub pid: u32,
pub process_role: String,
pub current_oldest_tx_age_secs: f64,
pub oldest_tx_label: Option<String>,
pub attribution_is_evidence_backed: bool,
}
impl WalPinAttribution {
fn unavailable(reason: impl Into<String>) -> Self {
Self {
available: false,
unavailable_reason: Some(reason.into()),
census_holder_pids: Vec::new(),
census_uninspectable_pids: Vec::new(),
census_truncated: false,
census_is_complete: false,
reporting: Vec::new(),
registered_silent_pids: Vec::new(),
unknown_pids: Vec::new(),
census_pids_without_attribution: Vec::new(),
fully_attributed: false,
sidecar_entries: Vec::new(),
sidecar_listing_truncated: false,
sidecar_entries_cleanup_would_reap: 0,
}
}
}
#[cfg(unix)]
pub fn wal_pin_attribution(db_path: &Path, _sweep_interval: Duration) -> WalPinAttribution {
use crate::walpin;
let census = match walpin::census_holders(db_path) {
Ok(c) => c,
Err(e) => return WalPinAttribution::unavailable(format!("census_holders failed: {e}")),
};
let mut census_holder_pids: Vec<u32> = census.holders.iter().copied().collect();
census_holder_pids.sort_unstable();
WalPinAttribution {
available: false,
unavailable_reason: Some(
"sidecar-to-holder reconciliation not available: this tree's khive-db exposes \
sidecar enumeration only via walpin::enumerate_live, which deletes stale/malformed \
entries as part of its cleanup pass; a diagnostics probe must not delete forensic \
sidecar evidence, so only the OS holder census below was collected"
.to_string(),
),
census_holder_pids,
census_uninspectable_pids: census.uninspectable_pids.clone(),
census_truncated: census.truncated,
census_is_complete: census.is_complete(),
reporting: Vec::new(),
registered_silent_pids: Vec::new(),
unknown_pids: Vec::new(),
census_pids_without_attribution: Vec::new(),
fully_attributed: false,
sidecar_entries: Vec::new(),
sidecar_listing_truncated: false,
sidecar_entries_cleanup_would_reap: 0,
}
}
#[cfg(not(unix))]
pub fn wal_pin_attribution(_db_path: &Path, _sweep_interval: Duration) -> WalPinAttribution {
WalPinAttribution::unavailable("WAL-pin attribution requires a Unix platform")
}
#[derive(Debug, Clone, PartialEq, Serialize)]
pub struct DbDiagnostics {
pub build: BuildIdentity,
pub db_path: Option<String>,
pub wal_file: Option<WalFileState>,
pub checkpoint_counters: CheckpointCounters,
pub checkpoint_probe: Option<CheckpointProbe>,
pub checkpoint_probe_error: Option<String>,
pub wal_pin: WalPinAttribution,
}
pub fn collect(
pool: &ConnectionPool,
build: BuildIdentity,
sweep_interval: Duration,
) -> DbDiagnostics {
let counters = checkpoint_counters();
let Some(path) = pool.config().path.clone() else {
return DbDiagnostics {
build,
db_path: None,
wal_file: None,
checkpoint_counters: counters,
checkpoint_probe: None,
checkpoint_probe_error: Some(
"in-memory database: no WAL file and no checkpoint to probe".to_string(),
),
wal_pin: WalPinAttribution::unavailable(
"in-memory database: no file for the OS holder census",
),
};
};
let (probe, probe_error) = match probe_pool(pool) {
Ok(p) => (Some(p), None),
Err(e) => (None, Some(e)),
};
DbDiagnostics {
build,
db_path: Some(path.display().to_string()),
wal_file: Some(wal_file_state(&path)),
checkpoint_counters: counters,
checkpoint_probe: probe,
checkpoint_probe_error: probe_error,
wal_pin: wal_pin_attribution(&path, sweep_interval),
}
}
fn probe_pool(pool: &ConnectionPool) -> Result<CheckpointProbe, String> {
let conn = pool
.open_standalone_writer()
.map_err(|e| format!("guarded standalone open refused: {e}"))?;
checkpoint_probe(&conn).map_err(|e| format!("PRAGMA wal_checkpoint(PASSIVE) failed: {e}"))
}
#[cfg(test)]
mod tests {
use serial_test::serial;
use super::*;
use crate::pool::{ConnectionPool, PoolConfig};
fn seeded_pool(dir: &tempfile::TempDir) -> (ConnectionPool, PathBuf) {
let path = dir.path().join("diag.db");
let pool = ConnectionPool::new(PoolConfig {
path: Some(path.clone()),
..PoolConfig::default()
})
.expect("pool open");
{
let writer = pool.try_writer().expect("writer");
writer
.conn()
.execute_batch(
"CREATE TABLE t (x INTEGER); \
INSERT INTO t VALUES (1), (2), (3);",
)
.expect("seed writes");
}
(pool, path)
}
#[test]
fn checkpoint_probe_returns_a_well_formed_triple_on_a_file_backed_db() {
let dir = tempfile::tempdir().expect("tempdir");
let (pool, _path) = seeded_pool(&dir);
let conn = pool.open_standalone_writer().expect("standalone");
let probe = checkpoint_probe(&conn).expect("probe must succeed on a WAL database");
assert!(
probe.busy == 0 || probe.busy == 1,
"busy is a 0/1 flag, got {}",
probe.busy
);
assert!(
probe.log_frames >= 0,
"a WAL database must report a non-negative frame count, got {}",
probe.log_frames
);
assert!(
probe.checkpointed_frames >= 0,
"checkpointed frames must be non-negative, got {}",
probe.checkpointed_frames
);
assert!(
probe.checkpointed_frames <= probe.log_frames,
"a PASSIVE pass cannot checkpoint more frames than the WAL holds: {probe:?}"
);
assert!(probe.pin_depth() >= 0, "pin depth clamps at 0: {probe:?}");
}
#[test]
#[serial(checkpoint_skip_metrics)]
fn checkpoint_probe_does_not_perturb_the_adr091_counters() {
crate::checkpoint::reset_checkpoint_metrics_for_tests();
let dir = tempfile::tempdir().expect("tempdir");
let (pool, _path) = seeded_pool(&dir);
let conn = pool.open_standalone_writer().expect("standalone");
let before = checkpoint_counters();
for _ in 0..3 {
checkpoint_probe(&conn).expect("probe must succeed");
}
let after = checkpoint_counters();
assert_eq!(
before, after,
"checkpoint_probe must leave every ADR-091 counter untouched"
);
}
#[test]
fn wal_file_state_reports_the_sidecar_size_for_a_live_db() {
let dir = tempfile::tempdir().expect("tempdir");
let (_pool, path) = seeded_pool(&dir);
let state = wal_file_state(&path);
assert!(
state.wal_path.ends_with("diag.db-wal"),
"WAL path is the db path plus a -wal suffix, got {}",
state.wal_path
);
assert!(
state.wal_size_bytes.is_some(),
"a seeded WAL database must have a stat-able -wal file: {state:?}"
);
assert!(state.unavailable_reason.is_none(), "{state:?}");
}
#[test]
fn wal_file_state_degrades_with_a_reason_when_the_sidecar_is_absent() {
let dir = tempfile::tempdir().expect("tempdir");
let state = wal_file_state(&dir.path().join("never-created.db"));
assert!(state.wal_size_bytes.is_none());
assert!(
state.unavailable_reason.is_some(),
"an absent WAL file must carry a reason, not a silent zero: {state:?}"
);
}
#[test]
fn collect_on_a_file_backed_db_carries_build_identity_and_every_counter() {
let dir = tempfile::tempdir().expect("tempdir");
let (pool, _path) = seeded_pool(&dir);
let report = collect(
&pool,
BuildIdentity::from_env("9.9.9", Some("deadbeef")),
Duration::from_secs(30),
);
assert_eq!(report.build.version, "9.9.9");
assert_eq!(report.build.build_hash.as_deref(), Some("deadbeef"));
assert!(report.db_path.is_some());
assert!(
report.checkpoint_probe.is_some(),
"file-backed collect must land a probe; error was {:?}",
report.checkpoint_probe_error
);
assert!(
report.wal_file.as_ref().and_then(|w| w.wal_size_bytes) >= Some(0),
"wal_size_bytes must be a non-negative byte count when present"
);
let json = serde_json::to_value(&report).expect("report serializes");
let counters = json
.get("checkpoint_counters")
.expect("counters section present");
for key in [
"last_observed_wal_pages",
"truncate_attempts",
"truncate_consecutive_failures",
"checkpoint_skipped_ticks",
"checkpoint_consecutive_skips",
"checkpoint_last_skip_wal_pages",
] {
assert!(counters.get(key).is_some(), "counter {key} must be present");
}
}
#[test]
fn never_observed_sentinels_serialize_as_null() {
let counters = CheckpointCounters {
last_observed_wal_pages: None,
truncate_attempts: 0,
truncate_consecutive_failures: 0,
checkpoint_skipped_ticks: 0,
checkpoint_consecutive_skips: 0,
checkpoint_last_skip_wal_pages: None,
};
let json = serde_json::to_value(counters).expect("serializes");
assert!(json["last_observed_wal_pages"].is_null());
assert!(json["checkpoint_last_skip_wal_pages"].is_null());
}
#[test]
fn probe_refuses_a_missing_configured_path_without_creating_it() {
let dir = tempfile::tempdir().expect("tempdir");
let (pool, path) = seeded_pool(&dir);
for suffix in ["", "-wal", "-shm"] {
let mut p = path.as_os_str().to_os_string();
p.push(suffix);
let _ = std::fs::remove_file(PathBuf::from(p));
}
assert!(!path.exists(), "precondition: the database file is gone");
let report = collect(
&pool,
BuildIdentity::from_env("0.0.0", None),
Duration::from_secs(30),
);
assert!(
report.checkpoint_probe.is_none(),
"a missing database must not yield a probe result: {report:?}"
);
assert!(
report.checkpoint_probe_error.is_some(),
"a missing database must say why there is no probe: {report:?}"
);
assert!(
!path.exists(),
"a diagnostics request must never create the database it was asked about"
);
}
#[test]
fn collect_degrades_gracefully_for_an_in_memory_backend() {
let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool");
let report = collect(
&pool,
BuildIdentity::from_env("0.0.0", None),
Duration::from_secs(30),
);
assert!(report.db_path.is_none());
assert!(report.wal_file.is_none());
assert!(report.checkpoint_probe.is_none());
assert!(
report.checkpoint_probe_error.is_some(),
"an in-memory report must say WHY there is no probe"
);
assert!(!report.wal_pin.available);
assert!(report.wal_pin.unavailable_reason.is_some());
}
#[cfg(unix)]
#[test]
fn wal_pin_attribution_reports_census_but_never_claims_full_attribution() {
let dir = tempfile::tempdir().expect("tempdir");
let (pool, path) = seeded_pool(&dir);
let _ = &pool;
let pin = wal_pin_attribution(&path, Duration::from_secs(30));
assert!(
!pin.fully_attributed,
"sidecar reconciliation is not ported: this must never claim completeness"
);
assert!(
pin.unavailable_reason.is_some(),
"the gap must be explained, not silent: {pin:?}"
);
assert!(pin.sidecar_entries.is_empty());
assert!(pin.reporting.is_empty());
}
}