use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
#[cfg(unix)]
use std::time::{SystemTime, UNIX_EPOCH};
use khive_storage::error::StorageError;
use khive_storage::types::StorageResult;
use khive_storage::StorageCapability;
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 checkpoint_pressure_elevated_ticks: u64,
pub checkpoint_pressure_episodes_started: u64,
pub checkpoint_pressure_episodes_recovered: u64,
pub checkpoint_lifecycle_append_attempts: u64,
pub checkpoint_lifecycle_append_failures: u64,
pub checkpoint_lifecycle_enqueue_drops: u64,
pub read_tx_max_age_evictions: 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(),
checkpoint_pressure_elevated_ticks: checkpoint::checkpoint_pressure_elevated_ticks(),
checkpoint_pressure_episodes_started: checkpoint::checkpoint_pressure_episodes_started(),
checkpoint_pressure_episodes_recovered: checkpoint::checkpoint_pressure_episodes_recovered(
),
checkpoint_lifecycle_append_attempts: checkpoint::checkpoint_lifecycle_append_attempts(),
checkpoint_lifecycle_append_failures: checkpoint::checkpoint_lifecycle_append_failures(),
checkpoint_lifecycle_enqueue_drops: checkpoint::checkpoint_lifecycle_enqueue_drops(),
read_tx_max_age_evictions: checkpoint::read_tx_max_age_evictions(),
}
}
#[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 ProcessIdentity {
pub pid: u32,
pub started_at: Option<i64>,
pub started_at_unavailable_reason: Option<String>,
pub pool_generation: u64,
}
impl ProcessIdentity {
pub fn current(pool: &ConnectionPool) -> Self {
let pid = std::process::id();
Self::from_start_time(
pid,
crate::walpin::process_start_time_secs(pid),
pool.main_pool_generation(),
)
}
fn from_start_time(pid: u32, started_at: Option<i64>, pool_generation: u64) -> Self {
Self {
pid,
started_at,
started_at_unavailable_reason: started_at.is_none().then(|| {
format!(
"OS process start time is unsupported or unavailable on {}",
std::env::consts::OS
)
}),
pool_generation,
}
}
}
#[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 status: WalPinAttributionStatus,
pub status_reasons: Vec<String>,
pub census: WalPinCensus,
pub available: bool,
pub unavailable_reason: Option<String>,
#[serde(skip_serializing)]
pub census_holder_pids: Vec<u32>,
#[serde(skip_serializing)]
pub census_uninspectable_pids: Vec<u32>,
#[serde(skip_serializing)]
pub census_truncated: bool,
#[serde(skip_serializing)]
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>,
#[serde(skip_serializing_if = "Option::is_none")]
pub sidecar_listing_truncated: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub sidecar_entries_cleanup_would_reap: Option<usize>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum WalPinAttributionStatus {
Complete,
Degraded,
Unavailable,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(tag = "status", rename_all = "snake_case")]
pub enum WalPinCensus {
Complete {
holder_pids: Vec<u32>,
},
Incomplete {
holder_pids: Vec<u32>,
uninspectable_pids: Vec<u32>,
truncated: bool,
reason: String,
},
Unavailable {
reason: String,
},
}
fn census_truncation_cause(budget_exhausted: bool) -> String {
if budget_exhausted {
"the OS process walk stopped at its wall-clock budget (see collection_cost.wal_pin_census_budget_ms)"
.to_string()
} else {
"the OS process walk was truncated".to_string()
}
}
#[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 {
let reason = reason.into();
Self {
status: WalPinAttributionStatus::Unavailable,
status_reasons: vec![reason.clone()],
census: WalPinCensus::Unavailable {
reason: reason.clone(),
},
available: false,
unavailable_reason: Some(reason),
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: None,
sidecar_entries_cleanup_would_reap: None,
}
}
}
#[cfg(all(unix, test))]
fn wal_pin_attribution_from_census(census: crate::walpin::CensusResult) -> WalPinAttribution {
wal_pin_attribution_without_sidecar(
census,
"read-only sidecar enumeration did not run for this attribution snapshot".to_string(),
)
}
#[cfg(unix)]
fn wal_pin_attribution_without_sidecar(
census: crate::walpin::CensusResult,
sidecar_reason: String,
) -> WalPinAttribution {
let census_is_complete = census.is_complete();
let mut census_holder_pids: Vec<u32> = census.holders.iter().copied().collect();
census_holder_pids.sort_unstable();
let mut census_uninspectable_pids = census.uninspectable_pids;
census_uninspectable_pids.sort_unstable();
census_uninspectable_pids.dedup();
let census_truncated = census.truncated;
let census_budget_exhausted = census.budget_exhausted;
let mut status_reasons = vec![sidecar_reason];
let census = if census_is_complete {
WalPinCensus::Complete {
holder_pids: census_holder_pids.clone(),
}
} else {
let mut causes = Vec::new();
if census_truncated {
causes.push(census_truncation_cause(census_budget_exhausted));
}
if !census_uninspectable_pids.is_empty() {
causes.push(format!(
"{} PID(s) could not be inspected",
census_uninspectable_pids.len()
));
}
let reason = format!(
"OS holder census is incomplete: {}; additional database holders cannot be ruled out",
causes.join("; ")
);
status_reasons.push(reason.clone());
WalPinCensus::Incomplete {
holder_pids: census_holder_pids.clone(),
uninspectable_pids: census_uninspectable_pids.clone(),
truncated: census_truncated,
reason,
}
};
WalPinAttribution {
status: WalPinAttributionStatus::Degraded,
unavailable_reason: Some(status_reasons.join("; ")),
status_reasons,
census,
available: false,
census_holder_pids,
census_uninspectable_pids,
census_truncated,
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: None,
sidecar_entries_cleanup_would_reap: None,
}
}
#[cfg(unix)]
fn wal_pin_attribution_from_evidence(
census: crate::walpin::CensusResult,
sidecar: crate::walpin::WalpinReport,
) -> WalPinAttribution {
use std::collections::BTreeSet;
let census_is_complete = census.is_complete();
let mut census_holder_pids: Vec<u32> = census.holders.iter().copied().collect();
census_holder_pids.sort_unstable();
let mut census_uninspectable_pids = census.uninspectable_pids;
census_uninspectable_pids.sort_unstable();
census_uninspectable_pids.dedup();
let census_truncated = census.truncated;
let census_budget_exhausted = census.budget_exhausted;
let census_carrier = if census_is_complete {
WalPinCensus::Complete {
holder_pids: census_holder_pids.clone(),
}
} else {
let mut causes = Vec::new();
if census_truncated {
causes.push(census_truncation_cause(census_budget_exhausted));
}
if !census_uninspectable_pids.is_empty() {
causes.push(format!(
"{} PID(s) could not be inspected",
census_uninspectable_pids.len()
));
}
let reason = format!(
"OS holder census is incomplete: {}; additional database holders cannot be ruled out",
causes.join("; ")
);
WalPinCensus::Incomplete {
holder_pids: census_holder_pids.clone(),
uninspectable_pids: census_uninspectable_pids.clone(),
truncated: census_truncated,
reason,
}
};
let now_epoch_secs = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|duration| duration.as_secs() as i64)
.unwrap_or(0);
let sidecar_listing_truncated = sidecar.sidecar_listing_truncated;
let sidecar_entries_cleanup_would_reap = sidecar.cleanup_would_reap;
let mut reporting = Vec::new();
let mut registered_silent_pids = Vec::new();
let mut unknown_pids = Vec::new();
let mut sidecar_entries = Vec::new();
let mut sidecar_known_pids = BTreeSet::new();
for entry in sidecar.entries {
match entry {
crate::walpin::WalpinPidHealth::Reporting(heartbeat) => {
let current_oldest_tx_age_secs =
heartbeat.current_oldest_tx_age_secs(now_epoch_secs);
let attribution_is_evidence_backed = heartbeat.attribution_is_evidence_backed();
sidecar_known_pids.insert(heartbeat.pid);
reporting.push(WalPinHolder {
pid: heartbeat.pid,
process_role: heartbeat.process_role.clone(),
current_oldest_tx_age_secs,
oldest_tx_label: heartbeat.oldest_tx_label.clone(),
attribution_is_evidence_backed,
});
sidecar_entries.push((
heartbeat.pid,
0u8,
serde_json::json!({
"pid": heartbeat.pid,
"status": "reporting",
"process_role": heartbeat.process_role,
"current_oldest_tx_age_secs": current_oldest_tx_age_secs,
"oldest_tx_label": heartbeat.oldest_tx_label,
"attribution_is_evidence_backed": attribution_is_evidence_backed,
}),
));
}
crate::walpin::WalpinPidHealth::RegisteredSilent { pid } => {
sidecar_known_pids.insert(pid);
registered_silent_pids.push(pid);
sidecar_entries.push((
pid,
1u8,
serde_json::json!({"pid": pid, "status": "registered_silent"}),
));
}
crate::walpin::WalpinPidHealth::Unknown { pid, reason } => {
sidecar_known_pids.insert(pid);
unknown_pids.push(pid);
sidecar_entries.push((
pid,
2u8,
serde_json::json!({"pid": pid, "status": "unknown", "reason": reason}),
));
}
}
}
reporting.sort_by_key(|holder| holder.pid);
reporting.dedup_by_key(|holder| holder.pid);
registered_silent_pids.sort_unstable();
registered_silent_pids.dedup();
unknown_pids.sort_unstable();
unknown_pids.dedup();
sidecar_entries.sort_by_key(|(pid, status_rank, _)| (*pid, *status_rank));
let sidecar_entries = sidecar_entries
.into_iter()
.map(|(_, _, entry)| entry)
.collect();
let census_pids_without_attribution: Vec<u32> = census_holder_pids
.iter()
.copied()
.filter(|pid| !sidecar_known_pids.contains(pid))
.collect();
let mut status_reasons = Vec::new();
if let WalPinCensus::Incomplete { reason, .. } = &census_carrier {
status_reasons.push(reason.clone());
}
if sidecar_listing_truncated {
status_reasons.push(
"read-only sidecar enumeration reached its entry cap; additional entries may exist"
.to_string(),
);
}
if !unknown_pids.is_empty() {
status_reasons.push(format!(
"{} sidecar PID(s) could not be classified conclusively",
unknown_pids.len()
));
}
if !census_pids_without_attribution.is_empty() {
status_reasons.push(format!(
"{} OS-confirmed holder(s) have no sidecar attribution",
census_pids_without_attribution.len()
));
}
let fully_attributed = census_is_complete
&& !sidecar_listing_truncated
&& unknown_pids.is_empty()
&& census_pids_without_attribution.is_empty();
let status = if fully_attributed {
WalPinAttributionStatus::Complete
} else {
WalPinAttributionStatus::Degraded
};
let unavailable_reason = (!fully_attributed).then(|| status_reasons.join("; "));
WalPinAttribution {
status,
status_reasons,
census: census_carrier,
available: fully_attributed,
unavailable_reason,
census_holder_pids,
census_uninspectable_pids,
census_truncated,
census_is_complete,
reporting,
registered_silent_pids,
unknown_pids,
census_pids_without_attribution,
fully_attributed,
sidecar_entries,
sidecar_listing_truncated: Some(sidecar_listing_truncated),
sidecar_entries_cleanup_would_reap: Some(sidecar_entries_cleanup_would_reap),
}
}
#[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}")),
};
if !walpin::sidecar_enabled(true) {
return wal_pin_attribution_without_sidecar(census, SIDECAR_DISABLED_REASON.to_string());
}
match walpin::inspect_live(&walpin::sidecar_dir_for(db_path), sweep_interval) {
Ok(sidecar) => wal_pin_attribution_from_evidence(census, sidecar),
Err(error) => wal_pin_attribution_without_sidecar(
census,
format!("read-only sidecar enumeration failed: {error}"),
),
}
}
#[cfg(unix)]
const SIDECAR_DISABLED_REASON: &str =
"walpin sidecar is explicitly disabled (KHIVE_WALPIN_SIDECAR); attribution has no sidecar \
evidence to reconcile against the OS holder census";
#[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, Copy, PartialEq, Eq, Serialize)]
pub struct ReaderContentionDiagnostics {
pub configured_reader_cap: usize,
pub configured_checkout_timeout_ms: u64,
pub configured_busy_timeout_ms: u64,
pub reader_admission_capacity: usize,
pub available_reader_admission_slots: usize,
pub reader_acquisitions: u64,
pub pooled_reader_checkouts: u64,
pub standalone_reader_opens: u64,
pub infrastructure_standalone_reader_opens: u64,
pub reader_checkout_timeouts: u64,
pub active_pooled_reader_checkouts: u64,
pub peak_active_pooled_reader_checkouts: u64,
pub completed_pooled_reader_checkouts: u64,
pub max_completed_reader_hold_micros: u64,
pub max_completed_reader_hold_operation: Option<&'static str>,
pub reader_replacement_open_failures: u64,
}
impl ReaderContentionDiagnostics {
fn snapshot(pool: &ConnectionPool) -> Self {
let reader = pool.reader_acquisition_snapshot();
Self {
configured_reader_cap: pool.config().max_readers,
configured_checkout_timeout_ms: u64::try_from(
pool.config().checkout_timeout.as_millis(),
)
.unwrap_or(u64::MAX),
configured_busy_timeout_ms: u64::try_from(pool.config().busy_timeout.as_millis())
.unwrap_or(u64::MAX),
reader_admission_capacity: reader.reader_admission_capacity,
available_reader_admission_slots: reader.available_reader_admission_slots,
reader_acquisitions: reader.acquisitions,
pooled_reader_checkouts: reader.pooled_checkouts,
standalone_reader_opens: reader.standalone_opens,
infrastructure_standalone_reader_opens: reader.infrastructure_standalone_opens,
reader_checkout_timeouts: reader.checkout_timeouts,
active_pooled_reader_checkouts: reader.active_pooled_checkouts,
peak_active_pooled_reader_checkouts: reader.peak_active_pooled_checkouts,
completed_pooled_reader_checkouts: reader.completed_pooled_checkouts,
max_completed_reader_hold_micros: reader.max_completed_hold_micros,
max_completed_reader_hold_operation: reader.max_completed_hold_operation,
reader_replacement_open_failures: reader.reader_replacement_open_failures,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct WriterContentionDiagnostics {
pub writer_acquisitions: u64,
pub pooled_writer_acquisitions: u64,
pub standalone_writer_acquisitions: u64,
pub writer_task_acquisitions: u64,
pub writer_acquisition_timeouts: u64,
pub writer_task_begin_busy: u64,
pub writer_task_begin_busy_absorbed: u64,
pub writer_task_begin_errors: u64,
pub writer_task_request_failures: u64,
pub writer_task_side_effects_unknown: u64,
pub audit_append_failures: Option<u64>,
pub audit_append_failures_unavailable_reason: Option<String>,
pub audit_obligation_append_failures: Option<u64>,
pub audit_obligation_append_failures_unavailable_reason: Option<String>,
pub audit_batch_flush_failures: Option<u64>,
pub audit_batch_flush_failures_unavailable_reason: Option<String>,
pub audit_degraded_rows: Option<u64>,
pub audit_degraded_rows_unavailable_reason: Option<String>,
pub audit_degraded: Option<bool>,
pub audit_degraded_unavailable_reason: Option<String>,
pub audit_admission_refused_obligations: Option<u64>,
pub audit_admission_refused_obligations_last_at_ms: Option<u64>,
pub audit_admission_refused_obligations_unavailable_reason: Option<String>,
pub audit_admission_unresolved_obligations: Option<u64>,
pub audit_admission_unresolved_obligations_last_at_ms: Option<u64>,
pub audit_admission_unresolved_obligations_unavailable_reason: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct RuntimeAuditBatchMetrics {
pub flush_failures: u64,
pub degraded_rows: u64,
pub degraded: bool,
pub admission_refused_obligations: u64,
pub admission_refused_obligations_last_at_ms: Option<u64>,
pub admission_unresolved_obligations: u64,
pub admission_unresolved_obligations_last_at_ms: Option<u64>,
}
impl WriterContentionDiagnostics {
fn snapshot(
pool: &ConnectionPool,
audit_append_failures: Option<u64>,
runtime_audit_batch_metrics: Option<RuntimeAuditBatchMetrics>,
) -> Self {
let writer = pool.writer_acquisition_snapshot();
let unavailable_reason =
|| Some("no audit-batch control is registered with this runtime instance".to_string());
Self {
writer_acquisitions: writer.acquisitions,
pooled_writer_acquisitions: writer.pooled_acquisitions,
standalone_writer_acquisitions: writer.standalone_acquisitions,
writer_task_acquisitions: writer.writer_task_acquisitions,
writer_acquisition_timeouts: writer.timeouts,
writer_task_begin_busy: writer.writer_task_begin_busy,
writer_task_begin_busy_absorbed: writer.writer_task_begin_busy_absorbed,
writer_task_begin_errors: writer.writer_task_begin_errors,
writer_task_request_failures: writer.writer_task_request_failures,
writer_task_side_effects_unknown: writer.writer_task_side_effects_unknown,
audit_append_failures,
audit_append_failures_unavailable_reason: audit_append_failures.is_none().then(|| {
"runtime audit instrumentation was not supplied to khive-db diagnostics".to_string()
}),
audit_obligation_append_failures: None,
audit_obligation_append_failures_unavailable_reason: Some(
"runtime obligation audit instrumentation was not supplied to khive-db diagnostics"
.to_string(),
),
audit_batch_flush_failures: runtime_audit_batch_metrics.map(|m| m.flush_failures),
audit_batch_flush_failures_unavailable_reason: runtime_audit_batch_metrics
.is_none()
.then(unavailable_reason)
.flatten(),
audit_degraded_rows: runtime_audit_batch_metrics.map(|m| m.degraded_rows),
audit_degraded_rows_unavailable_reason: runtime_audit_batch_metrics
.is_none()
.then(unavailable_reason)
.flatten(),
audit_degraded: runtime_audit_batch_metrics.map(|m| m.degraded),
audit_degraded_unavailable_reason: runtime_audit_batch_metrics
.is_none()
.then(unavailable_reason)
.flatten(),
audit_admission_refused_obligations: runtime_audit_batch_metrics
.map(|m| m.admission_refused_obligations),
audit_admission_refused_obligations_last_at_ms: runtime_audit_batch_metrics
.and_then(|m| m.admission_refused_obligations_last_at_ms),
audit_admission_refused_obligations_unavailable_reason: runtime_audit_batch_metrics
.is_none()
.then(unavailable_reason)
.flatten(),
audit_admission_unresolved_obligations: runtime_audit_batch_metrics
.map(|m| m.admission_unresolved_obligations),
audit_admission_unresolved_obligations_last_at_ms: runtime_audit_batch_metrics
.and_then(|m| m.admission_unresolved_obligations_last_at_ms),
audit_admission_unresolved_obligations_unavailable_reason: runtime_audit_batch_metrics
.is_none()
.then(unavailable_reason)
.flatten(),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
pub struct GraphEdgeIntegrity {
pub duplicate_edge_id_groups: i64,
pub graph_edges_rows: i64,
pub graph_edges_seq_rows: i64,
pub pre_v14_duplicate_edge_state_detected: bool,
pub live_entities_carrying_merged_into: i64,
}
fn graph_edge_integrity(conn: &Connection) -> rusqlite::Result<GraphEdgeIntegrity> {
conn.query_row(
"SELECT
(SELECT COUNT(*) FROM (
SELECT id FROM graph_edges GROUP BY id HAVING COUNT(*) > 1
)),
(SELECT COUNT(*) FROM graph_edges),
(SELECT COUNT(*) FROM graph_edges_seq),
(SELECT COUNT(*) FROM entities
WHERE deleted_at IS NULL AND merged_into IS NOT NULL)",
[],
|row| {
let duplicate_edge_id_groups = row.get(0)?;
Ok(GraphEdgeIntegrity {
duplicate_edge_id_groups,
graph_edges_rows: row.get(1)?,
graph_edges_seq_rows: row.get(2)?,
pre_v14_duplicate_edge_state_detected: duplicate_edge_id_groups > 0,
live_entities_carrying_merged_into: row.get(3)?,
})
},
)
}
const MAX_DATABASE_SIZE_OBJECTS: usize = 4_096;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum DatabaseObjectKind {
Table,
Index,
Internal,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum DatabaseStorageClass {
RowTable,
Index,
FullText,
Vector,
MixedRowAndEmbedding,
Internal,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct DatabaseObjectSize {
pub name: String,
pub owner_table: Option<String>,
pub object_kind: DatabaseObjectKind,
pub storage_class: DatabaseStorageClass,
pub pages: u64,
pub bytes: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct DatabaseSizeComposition {
pub page_size_bytes: u64,
pub page_count: u64,
pub freelist_pages: u64,
pub database_bytes: u64,
pub freelist_bytes: u64,
pub accounted_bytes: u64,
pub unaccounted_bytes: u64,
pub row_table_bytes: u64,
pub index_bytes: u64,
pub full_text_bytes: u64,
pub vector_bytes: u64,
pub mixed_embedding_bytes: u64,
pub internal_bytes: u64,
pub objects: Vec<DatabaseObjectSize>,
pub objects_truncated: bool,
pub objects_omitted: usize,
}
fn nonnegative_sqlite_integer(column: usize, value: i64) -> rusqlite::Result<u64> {
u64::try_from(value).map_err(|_| rusqlite::Error::IntegralValueOutOfRange(column, value))
}
fn declares_embedding_blob(sql: Option<&str>) -> bool {
let Some(sql) = sql else {
return false;
};
let tokens: Vec<_> = sql
.split(|character: char| !(character.is_ascii_alphanumeric() || character == '_'))
.filter(|token| !token.is_empty())
.collect();
tokens.windows(2).any(|pair| {
pair[0].eq_ignore_ascii_case("embedding") && pair[1].eq_ignore_ascii_case("blob")
})
}
fn classify_database_object(
name: &str,
sqlite_type: &str,
sql: Option<&str>,
) -> (DatabaseObjectKind, DatabaseStorageClass) {
let object_kind = match sqlite_type {
"table" => DatabaseObjectKind::Table,
"index" => DatabaseObjectKind::Index,
_ => DatabaseObjectKind::Internal,
};
let lower_name = name.to_ascii_lowercase();
let lower_sql = sql.unwrap_or_default().to_ascii_lowercase();
let storage_class = if lower_name.starts_with("fts_") || lower_sql.contains("using fts5") {
DatabaseStorageClass::FullText
} else if lower_name.starts_with("vec_")
|| lower_name == "_embedding_models"
|| lower_sql.contains("using vec0")
{
DatabaseStorageClass::Vector
} else if declares_embedding_blob(sql) {
DatabaseStorageClass::MixedRowAndEmbedding
} else if object_kind == DatabaseObjectKind::Index {
DatabaseStorageClass::Index
} else if name.starts_with("sqlite_") || object_kind == DatabaseObjectKind::Internal {
DatabaseStorageClass::Internal
} else {
DatabaseStorageClass::RowTable
};
(object_kind, storage_class)
}
fn database_size_composition(conn: &Connection) -> rusqlite::Result<DatabaseSizeComposition> {
let page_size = nonnegative_sqlite_integer(
0,
conn.query_row("PRAGMA page_size", [], |row| row.get::<_, i64>(0))?,
)?;
let page_count = nonnegative_sqlite_integer(
0,
conn.query_row("PRAGMA page_count", [], |row| row.get::<_, i64>(0))?,
)?;
let freelist_pages = nonnegative_sqlite_integer(
0,
conn.query_row("PRAGMA freelist_count", [], |row| row.get::<_, i64>(0))?,
)?;
let mut statement = conn.prepare(
"SELECT d.name, COALESCE(s.type, 'internal'), s.tbl_name, s.sql, d.pageno, d.pgsize
FROM dbstat AS d
LEFT JOIN sqlite_schema AS s ON s.name = d.name
WHERE d.aggregate = TRUE
ORDER BY d.name",
)?;
let mut rows = statement.query([])?;
let mut objects = Vec::new();
let mut objects_omitted = 0usize;
let mut accounted_bytes = 0u64;
let mut row_table_bytes = 0u64;
let mut index_bytes = 0u64;
let mut full_text_bytes = 0u64;
let mut vector_bytes = 0u64;
let mut mixed_embedding_bytes = 0u64;
let mut internal_bytes = 0u64;
while let Some(row) = rows.next()? {
let name: String = row.get(0)?;
let sqlite_type: String = row.get(1)?;
let owner_table: Option<String> = row.get(2)?;
let sql: Option<String> = row.get(3)?;
let pages = nonnegative_sqlite_integer(4, row.get(4)?)?;
let bytes = nonnegative_sqlite_integer(5, row.get(5)?)?;
let (object_kind, storage_class) =
classify_database_object(&name, &sqlite_type, sql.as_deref());
accounted_bytes = accounted_bytes.saturating_add(bytes);
let class_total = match storage_class {
DatabaseStorageClass::RowTable => &mut row_table_bytes,
DatabaseStorageClass::Index => &mut index_bytes,
DatabaseStorageClass::FullText => &mut full_text_bytes,
DatabaseStorageClass::Vector => &mut vector_bytes,
DatabaseStorageClass::MixedRowAndEmbedding => &mut mixed_embedding_bytes,
DatabaseStorageClass::Internal => &mut internal_bytes,
};
*class_total = class_total.saturating_add(bytes);
if objects.len() < MAX_DATABASE_SIZE_OBJECTS {
objects.push(DatabaseObjectSize {
name,
owner_table,
object_kind,
storage_class,
pages,
bytes,
});
} else {
objects_omitted = objects_omitted.saturating_add(1);
}
}
let database_bytes = page_count.saturating_mul(page_size);
let freelist_bytes = freelist_pages.saturating_mul(page_size);
let unaccounted_bytes = database_bytes
.saturating_sub(freelist_bytes)
.saturating_sub(accounted_bytes);
Ok(DatabaseSizeComposition {
page_size_bytes: page_size,
page_count,
freelist_pages,
database_bytes,
freelist_bytes,
accounted_bytes,
unaccounted_bytes,
row_table_bytes,
index_bytes,
full_text_bytes,
vector_bytes,
mixed_embedding_bytes,
internal_bytes,
objects,
objects_truncated: objects_omitted > 0,
objects_omitted,
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
pub struct CollectionCost {
pub total_ms: u64,
pub sqlite_ms: u64,
pub wal_file_stat_ms: u64,
pub wal_pin_census_ms: u64,
pub wal_pin_sidecar_ms: u64,
pub wal_pin_census_budget_ms: Option<u64>,
pub wal_pin_census_budget_exhausted: bool,
}
impl CollectionCost {
fn in_memory(total_ms: u64) -> Self {
Self {
total_ms,
sqlite_ms: 0,
wal_file_stat_ms: 0,
wal_pin_census_ms: 0,
wal_pin_sidecar_ms: 0,
wal_pin_census_budget_ms: None,
wal_pin_census_budget_exhausted: false,
}
}
}
const DEFAULT_CENSUS_BUDGET: Duration = Duration::from_millis(2000);
const CENSUS_BUDGET_ENV: &str = "KHIVE_WALPIN_CENSUS_BUDGET_MS";
fn request_census_budget() -> Option<Duration> {
match std::env::var(CENSUS_BUDGET_ENV) {
Ok(raw) => match raw.trim().parse::<u64>() {
Ok(0) => None,
Ok(ms) => Some(Duration::from_millis(ms)),
Err(_) => Some(DEFAULT_CENSUS_BUDGET),
},
Err(_) => Some(DEFAULT_CENSUS_BUDGET),
}
}
#[derive(Debug, Clone, PartialEq, Serialize)]
pub struct DbDiagnostics {
pub build: BuildIdentity,
pub process: ProcessIdentity,
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 reader_contention: ReaderContentionDiagnostics,
pub writer_contention: WriterContentionDiagnostics,
pub size_composition: Option<DatabaseSizeComposition>,
pub size_composition_error: Option<String>,
pub graph_edge_integrity: Option<GraphEdgeIntegrity>,
pub graph_edge_integrity_error: Option<String>,
pub fts_segments: Option<crate::FtsSegmentDiagnostics>,
pub fts_segments_error: Option<String>,
pub fts_maintenance: crate::FtsMaintenanceCounters,
pub wal_pin: WalPinAttribution,
pub collection_cost: CollectionCost,
}
pub fn collect(
pool: &ConnectionPool,
build: BuildIdentity,
sweep_interval: Duration,
) -> DbDiagnostics {
collect_inner(pool, build, sweep_interval, None, None)
}
pub fn collect_with_audit_append_failures(
pool: &ConnectionPool,
build: BuildIdentity,
sweep_interval: Duration,
audit_append_failures: u64,
) -> DbDiagnostics {
collect_inner(
pool,
build,
sweep_interval,
Some(audit_append_failures),
None,
)
}
pub async fn collect_with_audit_append_failures_interruptibly(
pool: Arc<ConnectionPool>,
build: BuildIdentity,
sweep_interval: Duration,
audit_append_failures: u64,
) -> StorageResult<DbDiagnostics> {
collect_with_runtime_audit_metrics_interruptibly(
pool,
build,
sweep_interval,
audit_append_failures,
None,
)
.await
}
pub async fn collect_with_runtime_audit_metrics_interruptibly(
pool: Arc<ConnectionPool>,
build: BuildIdentity,
sweep_interval: Duration,
audit_append_failures: u64,
runtime_audit_batch_metrics: Option<RuntimeAuditBatchMetrics>,
) -> StorageResult<DbDiagnostics> {
crate::ensure_request_read_active("db_diagnostics")?;
let started = Instant::now();
let process = ProcessIdentity::current(&pool);
let counters = checkpoint_counters();
let reader_contention = ReaderContentionDiagnostics::snapshot(&pool);
let writer_contention = WriterContentionDiagnostics::snapshot(
&pool,
Some(audit_append_failures),
runtime_audit_batch_metrics,
);
let Some(path) = pool.config().path.clone() else {
crate::ensure_request_read_active("db_diagnostics")?;
return Ok(DbDiagnostics {
build,
process,
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(),
),
reader_contention,
writer_contention,
size_composition: None,
size_composition_error: Some(
"in-memory database: no file-backed page composition to inspect".to_string(),
),
graph_edge_integrity: None,
graph_edge_integrity_error: Some(
"in-memory database: no durable graph-edge ledger to inspect".to_string(),
),
fts_segments: None,
fts_segments_error: Some(
"in-memory database: no durable FTS5 indexes to inspect".to_string(),
),
fts_maintenance: crate::fts_maintenance_counters(),
wal_pin: WalPinAttribution::unavailable(
"in-memory database: no file for the OS holder census",
),
collection_cost: CollectionCost::in_memory(elapsed_ms(started)),
});
};
let inspection_pool = Arc::clone(&pool);
let sqlite_started = Instant::now();
let inspection = crate::read_cancellation::run_interruptible_read(
StorageCapability::Sql,
"db_diagnostics.sqlite",
move |scope| inspect_pool_interruptibly(&inspection_pool, scope),
)
.await?;
let sqlite_ms = elapsed_ms(sqlite_started);
crate::ensure_request_read_active("db_diagnostics")?;
let canonical = operational_db_path(&pool, &path);
let budget = request_census_budget();
let (wal_file, wal_pin, file_state_cost) =
inspect_file_state_interruptibly(canonical, sweep_interval, budget).await?;
crate::ensure_request_read_active("db_diagnostics")?;
Ok(DbDiagnostics {
build,
process,
db_path: Some(path.display().to_string()),
wal_file: Some(wal_file),
checkpoint_counters: counters,
checkpoint_probe: inspection.checkpoint_probe,
checkpoint_probe_error: inspection.checkpoint_probe_error,
reader_contention,
writer_contention,
size_composition: inspection.size_composition,
size_composition_error: inspection.size_composition_error,
graph_edge_integrity: inspection.graph_edge_integrity,
graph_edge_integrity_error: inspection.graph_edge_integrity_error,
fts_segments: inspection.fts_segments,
fts_segments_error: inspection.fts_segments_error,
fts_maintenance: crate::fts_maintenance_counters(),
wal_pin,
collection_cost: CollectionCost {
total_ms: elapsed_ms(started),
sqlite_ms,
wal_file_stat_ms: file_state_cost.wal_file_stat_ms,
wal_pin_census_ms: file_state_cost.census_ms,
wal_pin_sidecar_ms: file_state_cost.sidecar_ms,
wal_pin_census_budget_ms: budget.map(|b| b.as_millis() as u64),
wal_pin_census_budget_exhausted: file_state_cost.census_budget_exhausted,
},
})
}
fn operational_db_path(pool: &ConnectionPool, configured: &Path) -> PathBuf {
pool.canonical_path()
.map(Path::to_path_buf)
.unwrap_or_else(|| configured.to_path_buf())
}
fn collect_inner(
pool: &ConnectionPool,
build: BuildIdentity,
sweep_interval: Duration,
audit_append_failures: Option<u64>,
runtime_audit_batch_metrics: Option<RuntimeAuditBatchMetrics>,
) -> DbDiagnostics {
let started = Instant::now();
let process = ProcessIdentity::current(pool);
let counters = checkpoint_counters();
let reader_contention = ReaderContentionDiagnostics::snapshot(pool);
let writer_contention = WriterContentionDiagnostics::snapshot(
pool,
audit_append_failures,
runtime_audit_batch_metrics,
);
let Some(path) = pool.config().path.clone() else {
return DbDiagnostics {
build,
process,
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(),
),
reader_contention,
writer_contention,
size_composition: None,
size_composition_error: Some(
"in-memory database: no file-backed page composition to inspect".to_string(),
),
graph_edge_integrity: None,
graph_edge_integrity_error: Some(
"in-memory database: no durable graph-edge ledger to inspect".to_string(),
),
fts_segments: None,
fts_segments_error: Some(
"in-memory database: no durable FTS5 indexes to inspect".to_string(),
),
fts_maintenance: crate::fts_maintenance_counters(),
wal_pin: WalPinAttribution::unavailable(
"in-memory database: no file for the OS holder census",
),
collection_cost: CollectionCost::in_memory(elapsed_ms(started)),
};
};
let sqlite_started = Instant::now();
let inspection = inspect_pool(pool);
let sqlite_ms = elapsed_ms(sqlite_started);
let canonical = operational_db_path(pool, &path);
let wal_file_started = Instant::now();
let wal_file = wal_file_state(&canonical);
let wal_file_stat_ms = elapsed_ms(wal_file_started);
let census_started = Instant::now();
let wal_pin = wal_pin_attribution(&canonical, sweep_interval);
let wal_pin_ms = elapsed_ms(census_started);
DbDiagnostics {
build,
process,
db_path: Some(path.display().to_string()),
wal_file: Some(wal_file),
checkpoint_counters: counters,
checkpoint_probe: inspection.checkpoint_probe,
checkpoint_probe_error: inspection.checkpoint_probe_error,
reader_contention,
writer_contention,
size_composition: inspection.size_composition,
size_composition_error: inspection.size_composition_error,
graph_edge_integrity: inspection.graph_edge_integrity,
graph_edge_integrity_error: inspection.graph_edge_integrity_error,
fts_segments: inspection.fts_segments,
fts_segments_error: inspection.fts_segments_error,
fts_maintenance: crate::fts_maintenance_counters(),
wal_pin,
collection_cost: CollectionCost {
total_ms: elapsed_ms(started),
sqlite_ms,
wal_file_stat_ms,
wal_pin_census_ms: wal_pin_ms,
wal_pin_sidecar_ms: 0,
wal_pin_census_budget_ms: None,
wal_pin_census_budget_exhausted: false,
},
}
}
struct PoolInspection {
checkpoint_probe: Option<CheckpointProbe>,
checkpoint_probe_error: Option<String>,
size_composition: Option<DatabaseSizeComposition>,
size_composition_error: Option<String>,
graph_edge_integrity: Option<GraphEdgeIntegrity>,
graph_edge_integrity_error: Option<String>,
fts_segments: Option<crate::FtsSegmentDiagnostics>,
fts_segments_error: Option<String>,
}
fn split_fts_segments_result(
result: Result<crate::FtsSegmentDiagnostics, String>,
) -> (Option<crate::FtsSegmentDiagnostics>, Option<String>) {
match result {
Ok(segments) => (Some(segments), None),
Err(error) => (None, Some(error)),
}
}
fn inspect_pool_interruptibly(
pool: &ConnectionPool,
scope: &crate::read_cancellation::InterruptibleReadScope,
) -> StorageResult<PoolInspection> {
scope.ensure_active()?;
let conn = match pool.open_standalone_writer_untracked() {
Ok(conn) => conn,
Err(e) => {
scope.ensure_active()?;
let reason = format!("guarded standalone open refused: {e}");
return Ok(PoolInspection {
checkpoint_probe: None,
checkpoint_probe_error: Some(reason.clone()),
size_composition: None,
size_composition_error: Some(reason.clone()),
graph_edge_integrity: None,
graph_edge_integrity_error: Some(reason.clone()),
fts_segments: None,
fts_segments_error: Some(reason),
});
}
};
scope.ensure_active()?;
let (checkpoint_probe, checkpoint_probe_error) = match checkpoint_probe(&conn) {
Ok(probe) => (Some(probe), None),
Err(e) => (
None,
Some(format!("PRAGMA wal_checkpoint(PASSIVE) failed: {e}")),
),
};
#[cfg(test)]
if TEST_PAUSE_AFTER_PASSIVE.load(Ordering::SeqCst) {
TEST_REACHED_AFTER_PASSIVE.store(true, Ordering::SeqCst);
while TEST_PAUSE_AFTER_PASSIVE.load(Ordering::SeqCst) && !scope.should_stop() {
std::thread::yield_now();
}
}
scope.ensure_active()?;
let (integrity, size_composition, fts_segments) = scope.run(&conn, || {
Ok((
graph_edge_integrity(&conn),
database_size_composition(&conn),
crate::fts_maintenance::inspect_fts_segments(&conn),
))
})?;
let (graph_edge_integrity, graph_edge_integrity_error) = match integrity {
Ok(integrity) => (Some(integrity), None),
Err(e) => (
None,
Some(format!("graph-edge integrity query failed: {e}")),
),
};
let (fts_segments, fts_segments_error) = split_fts_segments_result(fts_segments);
let (size_composition, size_composition_error) = match size_composition {
Ok(composition) => (Some(composition), None),
Err(error) => (
None,
Some(format!("database size composition query failed: {error}")),
),
};
Ok(PoolInspection {
checkpoint_probe,
checkpoint_probe_error,
size_composition,
size_composition_error,
graph_edge_integrity,
graph_edge_integrity_error,
fts_segments,
fts_segments_error,
})
}
#[cfg(test)]
static TEST_PAUSE_AFTER_PASSIVE: AtomicBool = AtomicBool::new(false);
#[cfg(test)]
static TEST_REACHED_AFTER_PASSIVE: AtomicBool = AtomicBool::new(false);
struct StopCensusOnDrop {
stopped: Arc<AtomicBool>,
armed: bool,
}
impl Drop for StopCensusOnDrop {
fn drop(&mut self) {
if self.armed {
self.stopped.store(true, Ordering::SeqCst);
}
}
}
#[derive(Debug, Clone, Copy, Default)]
struct FileStateCost {
wal_file_stat_ms: u64,
census_ms: u64,
sidecar_ms: u64,
census_budget_exhausted: bool,
}
fn elapsed_ms(since: Instant) -> u64 {
since.elapsed().as_millis() as u64
}
async fn inspect_file_state_interruptibly(
path: PathBuf,
sweep_interval: Duration,
census_budget: Option<Duration>,
) -> StorageResult<(WalFileState, WalPinAttribution, FileStateCost)> {
const OPERATION: &str = "db_diagnostics.wal_holder_census";
crate::ensure_request_read_active(OPERATION)?;
let stopped = Arc::new(AtomicBool::new(false));
let worker_stopped = Arc::clone(&stopped);
let mut stop_on_drop = StopCensusOnDrop {
stopped: Arc::clone(&stopped),
armed: true,
};
let mut worker = tokio::task::spawn_blocking(move || {
let mut cost = FileStateCost::default();
let wal_file_started = Instant::now();
let wal_file = wal_file_state(&path);
cost.wal_file_stat_ms = elapsed_ms(wal_file_started);
if worker_stopped.load(Ordering::SeqCst) {
return Err(std::io::Error::new(
std::io::ErrorKind::Interrupted,
"WAL holder census cancelled",
));
}
#[cfg(unix)]
let census_started = Instant::now();
#[cfg(unix)]
let census_result = match census_budget {
Some(budget) => crate::walpin::census_holders_until_within(
&path,
|| worker_stopped.load(Ordering::SeqCst),
budget,
),
None => {
crate::walpin::census_holders_until(&path, || worker_stopped.load(Ordering::SeqCst))
}
};
#[cfg(unix)]
{
cost.census_ms = elapsed_ms(census_started);
}
#[cfg(unix)]
let attribution = match census_result {
Ok(census) => {
cost.census_budget_exhausted = census.budget_exhausted;
if worker_stopped.load(Ordering::SeqCst) {
return Err(std::io::Error::new(
std::io::ErrorKind::Interrupted,
"WAL sidecar inspection cancelled",
));
}
if !crate::walpin::sidecar_enabled(true) {
wal_pin_attribution_without_sidecar(census, SIDECAR_DISABLED_REASON.to_string())
} else {
let sidecar_started = Instant::now();
let sidecar = crate::walpin::inspect_live(
&crate::walpin::sidecar_dir_for(&path),
sweep_interval,
);
cost.sidecar_ms = elapsed_ms(sidecar_started);
if worker_stopped.load(Ordering::SeqCst) {
return Err(std::io::Error::new(
std::io::ErrorKind::Interrupted,
"WAL sidecar inspection cancelled",
));
}
match sidecar {
Ok(sidecar) => wal_pin_attribution_from_evidence(census, sidecar),
Err(error) => wal_pin_attribution_without_sidecar(
census,
format!("read-only sidecar enumeration failed: {error}"),
),
}
}
}
Err(error) if error.kind() == std::io::ErrorKind::Interrupted => return Err(error),
Err(error) => WalPinAttribution::unavailable(format!("census_holders failed: {error}")),
};
#[cfg(not(unix))]
let attribution = {
let _ = census_budget;
let census_started = Instant::now();
let attribution = wal_pin_attribution(&path, sweep_interval);
cost.census_ms = elapsed_ms(census_started);
attribution
};
Ok((wal_file, attribution, cost))
});
tokio::select! {
joined = &mut worker => {
stop_on_drop.armed = false;
let result = joined
.map_err(|error| StorageError::driver(StorageCapability::Sql, OPERATION, error))?
.map_err(|error| StorageError::driver(StorageCapability::Sql, OPERATION, error))?;
crate::ensure_request_read_active(OPERATION)?;
Ok(result)
}
_ = crate::wait_for_request_read_cancellation() => {
stopped.store(true, Ordering::SeqCst);
if tokio::time::timeout(crate::sqlite_interrupt_grace_from_env(), &mut worker)
.await
.is_err()
{
worker.abort();
}
stop_on_drop.armed = false;
Err(StorageError::Timeout { operation: OPERATION.into() })
}
}
}
fn inspect_pool(pool: &ConnectionPool) -> PoolInspection {
let conn = match pool.open_standalone_writer_untracked() {
Ok(conn) => conn,
Err(e) => {
let reason = format!("guarded standalone open refused: {e}");
return PoolInspection {
checkpoint_probe: None,
checkpoint_probe_error: Some(reason.clone()),
size_composition: None,
size_composition_error: Some(reason.clone()),
graph_edge_integrity: None,
graph_edge_integrity_error: Some(reason.clone()),
fts_segments: None,
fts_segments_error: Some(reason),
};
}
};
let (checkpoint_probe, checkpoint_probe_error) = match checkpoint_probe(&conn) {
Ok(probe) => (Some(probe), None),
Err(e) => (
None,
Some(format!("PRAGMA wal_checkpoint(PASSIVE) failed: {e}")),
),
};
let (graph_edge_integrity, graph_edge_integrity_error) = match graph_edge_integrity(&conn) {
Ok(integrity) => (Some(integrity), None),
Err(e) => (
None,
Some(format!("graph-edge integrity query failed: {e}")),
),
};
let (size_composition, size_composition_error) = match database_size_composition(&conn) {
Ok(composition) => (Some(composition), None),
Err(error) => (
None,
Some(format!("database size composition query failed: {error}")),
),
};
let (fts_segments, fts_segments_error) =
split_fts_segments_result(crate::fts_maintenance::inspect_fts_segments(&conn));
PoolInspection {
checkpoint_probe,
checkpoint_probe_error,
size_composition,
size_composition_error,
graph_edge_integrity,
graph_edge_integrity_error,
fts_segments,
fts_segments_error,
}
}
#[cfg(test)]
mod tests {
use serial_test::serial;
use super::*;
use crate::pool::{ConnectionPool, PoolConfig};
#[test]
#[serial_test::serial(khive_walpin_census_budget_env)]
fn census_budget_reads_zero_as_unbounded_and_survives_a_malformed_value() {
let _guard = crate::walpin::EnvVarGuard::capture(CENSUS_BUDGET_ENV);
std::env::remove_var(CENSUS_BUDGET_ENV);
assert_eq!(
request_census_budget(),
Some(DEFAULT_CENSUS_BUDGET),
"an unset variable takes the default bound"
);
std::env::set_var(CENSUS_BUDGET_ENV, "0");
assert_eq!(
request_census_budget(),
None,
"0 restores the unbounded full-machine walk"
);
std::env::set_var(CENSUS_BUDGET_ENV, " 750 ");
assert_eq!(
request_census_budget(),
Some(Duration::from_millis(750)),
"a surrounding-whitespace value is still a number"
);
std::env::set_var(CENSUS_BUDGET_ENV, "soon");
assert_eq!(
request_census_budget(),
Some(DEFAULT_CENSUS_BUDGET),
"a malformed budget must not fail the request; the report states \
which budget was actually used"
);
}
#[test]
fn a_budget_stop_and_an_enumeration_failure_do_not_share_a_reason() {
let budget = census_truncation_cause(true);
let failure = census_truncation_cause(false);
assert_ne!(budget, failure);
assert!(
budget.contains("budget"),
"the budget reason must name the budget: {budget}"
);
assert!(
budget.contains("wal_pin_census_budget_ms"),
"and must point at the field carrying the value: {budget}"
);
assert!(
!failure.contains("budget"),
"an enumeration failure must not be described as a budget stop: {failure}"
);
}
#[test]
fn collection_cost_distinguishes_an_unbounded_census_from_a_zero_cost_one() {
let unbounded = CollectionCost {
total_ms: 9,
sqlite_ms: 4,
wal_file_stat_ms: 0,
wal_pin_census_ms: 5,
wal_pin_sidecar_ms: 0,
wal_pin_census_budget_ms: None,
wal_pin_census_budget_exhausted: false,
};
let bounded = CollectionCost {
wal_pin_census_budget_ms: Some(2000),
wal_pin_census_budget_exhausted: true,
..unbounded
};
let unbounded = serde_json::to_value(unbounded).unwrap();
let bounded = serde_json::to_value(bounded).unwrap();
assert_eq!(
unbounded["wal_pin_census_budget_ms"],
serde_json::Value::Null
);
assert_eq!(bounded["wal_pin_census_budget_ms"], 2000);
assert_eq!(unbounded["wal_pin_census_budget_exhausted"], false);
assert_eq!(bounded["wal_pin_census_budget_exhausted"], true);
}
#[test]
fn process_identity_serializes_os_start_time_or_explicit_unavailability() {
let known = ProcessIdentity::from_start_time(42, Some(1_000_000_000), 3);
assert_eq!(
serde_json::to_value(known).unwrap(),
serde_json::json!({
"pid": 42,
"started_at": 1_000_000_000,
"started_at_unavailable_reason": null,
"pool_generation": 3
}),
"preserve the OS timestamp, not request time or a derived uptime"
);
let unknown = ProcessIdentity::from_start_time(42, None, 3);
let json = serde_json::to_value(unknown).unwrap();
assert_eq!(json["pid"], 42);
assert!(json["started_at"].is_null(), "never invent a start time");
let reason = json["started_at_unavailable_reason"].as_str().unwrap();
assert!(!reason.is_empty());
assert!(reason.contains(std::env::consts::OS));
assert_eq!(json["pool_generation"], 3);
}
#[tokio::test]
async fn process_identity_is_present_in_every_collector_path() {
let dir = tempfile::tempdir().expect("tempdir");
let (file_pool, _) = seeded_pool(&dir);
let memory_pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool");
for pool in [file_pool, memory_pool] {
let pool = Arc::new(pool);
let expected = ProcessIdentity::current(&pool);
let writer_before = pool.writer_acquisition_snapshot();
let sync = collect(
&pool,
BuildIdentity::from_env("test", None),
Duration::from_secs(30),
);
let asynchronous = collect_with_runtime_audit_metrics_interruptibly(
Arc::clone(&pool),
BuildIdentity::from_env("test", None),
Duration::from_secs(30),
0,
None,
)
.await
.expect("diagnostics");
for report in [sync, asynchronous] {
assert_eq!(report.process, expected);
let json = serde_json::to_value(report).unwrap();
assert_eq!(json["process"], serde_json::to_value(&expected).unwrap());
assert_eq!(json["build"]["version"], "test");
}
assert_eq!(
pool.writer_acquisition_snapshot(),
writer_before,
"diagnostics must not count its probes as write traffic"
);
}
}
#[test]
fn process_identity_survives_pool_reconstruction_with_reset_counters() {
let pool = ConnectionPool::new(PoolConfig::default()).expect("first pool");
drop(pool.try_writer().expect("writer acquisition"));
drop(pool.reader().expect("reader acquisition"));
let before = collect(
&pool,
BuildIdentity::from_env("test", None),
Duration::from_secs(30),
);
assert!(before.writer_contention.writer_acquisitions > 0);
assert!(before.reader_contention.reader_acquisitions > 0);
drop(pool);
let replacement = ConnectionPool::new(PoolConfig::default()).expect("replacement pool");
let after = collect(
&replacement,
BuildIdentity::from_env("test", None),
Duration::from_secs(30),
);
assert_eq!(after.process.pid, before.process.pid);
assert_eq!(after.process.started_at, before.process.started_at);
assert_eq!(
after.process.started_at_unavailable_reason,
before.process.started_at_unavailable_reason
);
assert!(after.process.pool_generation > before.process.pool_generation);
assert_eq!(after.writer_contention.writer_acquisitions, 0);
assert_eq!(after.reader_contention.reader_acquisitions, 0);
}
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::for_test()
})
.expect("pool open");
{
let writer = pool.try_writer().expect("writer");
writer
.conn()
.execute_batch(
"CREATE TABLE t (x INTEGER); \
CREATE TABLE entities (id TEXT PRIMARY KEY, deleted_at INTEGER, merged_into TEXT); \
CREATE TABLE graph_edges (
namespace TEXT NOT NULL,
id TEXT NOT NULL,
PRIMARY KEY (namespace, id)
); \
CREATE TABLE graph_edges_seq (
seq INTEGER PRIMARY KEY AUTOINCREMENT,
edge_id TEXT NOT NULL UNIQUE
); \
CREATE VIRTUAL TABLE fts_entities USING fts5(
namespace UNINDEXED, subject_id UNINDEXED, title, body,
tokenize='trigram'
); \
CREATE VIRTUAL TABLE fts_notes USING fts5(
namespace UNINDEXED, subject_id UNINDEXED, title, body,
tokenize='trigram'
); \
INSERT INTO t VALUES (1), (2), (3); \
INSERT INTO fts_entities(namespace, subject_id, title, body)
VALUES('local', 'entity-1', 'entity title', 'entity diagnostic body'); \
INSERT INTO fts_notes(namespace, subject_id, title, body)
VALUES('local', 'note-1', 'note title', 'note diagnostic body');",
)
.expect("seed writes");
}
(pool, path)
}
#[test]
fn dropping_census_future_guard_requests_cooperative_stop() {
let stopped = Arc::new(AtomicBool::new(false));
let guard = StopCensusOnDrop {
stopped: Arc::clone(&stopped),
armed: true,
};
drop(guard);
assert!(
stopped.load(Ordering::SeqCst),
"dropping diagnostics while its census worker is live must stop the PID/fd walk"
);
}
#[tokio::test]
async fn runtime_audit_batch_fields_are_additive_and_unavailable_without_a_control() {
let dir = tempfile::tempdir().expect("tempdir");
let (pool, _path) = seeded_pool(&dir);
let pool = Arc::new(pool);
let without_control = collect_with_audit_append_failures_interruptibly(
Arc::clone(&pool),
BuildIdentity::from_env("test", None),
Duration::from_secs(30),
0,
)
.await
.expect("diagnostics succeed");
assert!(without_control
.writer_contention
.audit_batch_flush_failures
.is_none());
assert!(
without_control
.writer_contention
.audit_batch_flush_failures_unavailable_reason
.is_some(),
"no audit-batch control was supplied, so the field must carry a reason, not a \
fabricated zero"
);
assert!(without_control
.writer_contention
.audit_degraded_rows
.is_none());
assert!(without_control.writer_contention.audit_degraded.is_none());
let with_control = collect_with_runtime_audit_metrics_interruptibly(
Arc::clone(&pool),
BuildIdentity::from_env("test", None),
Duration::from_secs(30),
0,
Some(RuntimeAuditBatchMetrics {
flush_failures: 3,
degraded_rows: 7,
degraded: true,
admission_refused_obligations: 5,
admission_refused_obligations_last_at_ms: Some(1_700_000_000_123),
admission_unresolved_obligations: 2,
admission_unresolved_obligations_last_at_ms: Some(1_700_000_000_456),
}),
)
.await
.expect("diagnostics succeed");
assert_eq!(
with_control.writer_contention.audit_batch_flush_failures,
Some(3)
);
assert!(with_control
.writer_contention
.audit_batch_flush_failures_unavailable_reason
.is_none());
assert_eq!(with_control.writer_contention.audit_degraded_rows, Some(7));
assert_eq!(with_control.writer_contention.audit_degraded, Some(true));
assert_eq!(
with_control
.writer_contention
.audit_admission_refused_obligations,
Some(5),
"an operator must be able to read the admission-refused obligation count from \
db_diagnostics without a test-only feature gate (ADR-103 Amendment 3)"
);
assert!(with_control
.writer_contention
.audit_admission_refused_obligations_unavailable_reason
.is_none());
assert_eq!(
with_control
.writer_contention
.audit_admission_unresolved_obligations,
Some(2),
"an operator must be able to distinguish enqueued-but-unresolved rows from \
confirmed-refused rows (ADR-103 Amendment 3)"
);
assert!(with_control
.writer_contention
.audit_admission_unresolved_obligations_unavailable_reason
.is_none());
assert!(without_control
.writer_contention
.audit_admission_refused_obligations
.is_none());
assert!(without_control
.writer_contention
.audit_admission_refused_obligations_unavailable_reason
.is_some());
assert!(without_control
.writer_contention
.audit_admission_unresolved_obligations
.is_none());
assert!(without_control
.writer_contention
.audit_admission_unresolved_obligations_unavailable_reason
.is_some());
assert_eq!(
with_control.writer_contention.writer_acquisitions,
without_control.writer_contention.writer_acquisitions
);
}
#[test]
fn writer_task_pool_sourced_counters_are_always_populated_directly() {
let dir = tempfile::tempdir().expect("tempdir");
let (pool, _path) = seeded_pool(&dir);
let report = collect(
&pool,
BuildIdentity::from_env("9.9.9", None),
Duration::from_secs(30),
);
assert_eq!(report.writer_contention.writer_task_request_failures, 0);
assert_eq!(report.writer_contention.writer_task_side_effects_unknown, 0);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial]
async fn request_cancellation_after_passive_stops_before_graph_and_census() {
let dir = tempfile::tempdir().expect("tempdir");
let (pool, _) = seeded_pool(&dir);
let pool = Arc::new(pool);
TEST_REACHED_AFTER_PASSIVE.store(false, Ordering::SeqCst);
TEST_PAUSE_AFTER_PASSIVE.store(true, Ordering::SeqCst);
let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
let diagnostic_pool = Arc::clone(&pool);
let task = tokio::spawn(crate::scope_request_read_cancellation(
cancel_rx,
async move {
collect_with_audit_append_failures_interruptibly(
diagnostic_pool,
BuildIdentity::from_env("test", None),
Duration::from_secs(30),
0,
)
.await
},
));
tokio::time::timeout(Duration::from_secs(1), async {
while !TEST_REACHED_AFTER_PASSIVE.load(Ordering::SeqCst) {
tokio::task::yield_now().await;
}
})
.await
.expect("diagnostics never completed its admitted PASSIVE phase");
cancel_tx.send(true).unwrap();
let result = tokio::time::timeout(Duration::from_secs(1), task)
.await
.expect("cancelled diagnostics did not stop promptly")
.expect("diagnostics task panicked");
TEST_PAUSE_AFTER_PASSIVE.store(false, Ordering::SeqCst);
assert!(matches!(result, Err(StorageError::Timeout { .. })));
let one: i64 = pool
.reader()
.expect("diagnostics returned its connection")
.conn()
.query_row("SELECT 1", [], |row| row.get(0))
.unwrap();
assert_eq!(one, 1);
}
#[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_untracked()
.expect("standalone probe connection");
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_untracked()
.expect("standalone probe connection");
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 reader_admission_capacity = pool.max_readers().max(1);
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",
"checkpoint_pressure_elevated_ticks",
"checkpoint_pressure_episodes_started",
"checkpoint_pressure_episodes_recovered",
"checkpoint_lifecycle_append_attempts",
"checkpoint_lifecycle_append_failures",
"checkpoint_lifecycle_enqueue_drops",
"read_tx_max_age_evictions",
] {
assert!(counters.get(key).is_some(), "counter {key} must be present");
}
assert_eq!(
report.writer_contention.writer_acquisitions, 1,
"the seed write checked the finite-wait pooled writer out once"
);
assert_eq!(report.writer_contention.pooled_writer_acquisitions, 1);
assert_eq!(report.writer_contention.standalone_writer_acquisitions, 0);
assert_eq!(report.writer_contention.writer_task_acquisitions, 0);
assert_eq!(report.writer_contention.writer_acquisition_timeouts, 0);
assert_eq!(
report.reader_contention,
ReaderContentionDiagnostics {
configured_reader_cap: pool.config().max_readers,
configured_checkout_timeout_ms: u64::try_from(
pool.config().checkout_timeout.as_millis(),
)
.unwrap_or(u64::MAX),
configured_busy_timeout_ms: u64::try_from(pool.config().busy_timeout.as_millis())
.unwrap_or(u64::MAX),
reader_admission_capacity,
available_reader_admission_slots: reader_admission_capacity,
reader_acquisitions: 0,
pooled_reader_checkouts: 0,
standalone_reader_opens: 0,
infrastructure_standalone_reader_opens: 0,
reader_checkout_timeouts: 0,
active_pooled_reader_checkouts: 0,
peak_active_pooled_reader_checkouts: 0,
completed_pooled_reader_checkouts: 0,
max_completed_reader_hold_micros: 0,
max_completed_reader_hold_operation: None,
reader_replacement_open_failures: 0,
},
"the diagnostics probe itself must not masquerade as request reader traffic"
);
assert_eq!(
report.graph_edge_integrity,
Some(GraphEdgeIntegrity {
duplicate_edge_id_groups: 0,
graph_edges_rows: 0,
graph_edges_seq_rows: 0,
pre_v14_duplicate_edge_state_detected: false,
live_entities_carrying_merged_into: 0,
})
);
assert!(report.graph_edge_integrity_error.is_none());
let fts_segments = report
.fts_segments
.as_ref()
.expect("file-backed diagnostics must decode both FTS structure rows");
assert_eq!(fts_segments.entities.segment_count, 1);
assert_eq!(fts_segments.notes.segment_count, 1);
assert_eq!(fts_segments.total_segments, 2);
assert!(report.fts_segments_error.is_none());
let fts_json = json
.get("fts_segments")
.expect("FTS segment diagnostics serialize");
assert_eq!(fts_json["entities"]["segment_count"], 1);
assert!(json.get("fts_maintenance").is_some());
assert!(
report.size_composition.is_some(),
"file-backed diagnostics must include page composition; error was {:?}",
report.size_composition_error
);
assert!(report.size_composition_error.is_none());
assert!(report.writer_contention.audit_append_failures.is_none());
assert!(report
.writer_contention
.audit_obligation_append_failures
.is_none());
assert!(report
.writer_contention
.audit_obligation_append_failures_unavailable_reason
.is_some());
assert!(json["writer_contention"]["audit_obligation_append_failures"].is_null());
assert!(
report
.writer_contention
.audit_append_failures_unavailable_reason
.is_some(),
"a direct khive-db snapshot must not fabricate a runtime audit count"
);
}
#[test]
fn diagnostics_exposes_reader_saturation_and_completed_hold_evidence() {
let dir = tempfile::tempdir().expect("tempdir");
let pool = ConnectionPool::new(PoolConfig {
path: Some(dir.path().join("reader_saturation.db")),
max_readers: 1,
checkout_timeout: Duration::from_millis(2),
..PoolConfig::default()
})
.expect("one-reader file-backed pool");
let held = pool.reader().expect("first reader checkout");
assert!(
pool.reader().is_err(),
"the live checkout must exhaust the one-slot reader budget"
);
drop(held);
let report = collect(
&pool,
BuildIdentity::from_env("9.9.9", None),
Duration::from_secs(30),
);
let reader = report.reader_contention;
assert_eq!(reader.reader_admission_capacity, 1);
assert_eq!(reader.available_reader_admission_slots, 1);
assert_eq!(reader.reader_acquisitions, 1);
assert_eq!(reader.pooled_reader_checkouts, 1);
assert_eq!(reader.standalone_reader_opens, 0);
assert_eq!(reader.infrastructure_standalone_reader_opens, 0);
assert_eq!(reader.reader_checkout_timeouts, 1);
assert_eq!(reader.active_pooled_reader_checkouts, 0);
assert_eq!(reader.peak_active_pooled_reader_checkouts, 1);
assert_eq!(reader.completed_pooled_reader_checkouts, 1);
assert!(reader.max_completed_reader_hold_micros > 0);
let json = serde_json::to_value(&report).expect("report serializes");
assert_eq!(
json.pointer("/reader_contention/reader_admission_capacity"),
Some(&serde_json::json!(1)),
"the operator wire payload must expose the reader admission budget"
);
assert_eq!(
json.pointer("/reader_contention/reader_checkout_timeouts"),
Some(&serde_json::json!(1)),
"the operator wire payload must expose the reader timeout phase"
);
assert!(
json.pointer("/reader_contention/max_completed_reader_hold_micros")
.is_some(),
"the operator wire payload must expose completed hold-time evidence"
);
}
#[test]
fn diagnostics_reports_configured_reader_budget_and_both_deadlines() {
let pool = ConnectionPool::new(PoolConfig {
max_readers: 6,
checkout_timeout: Duration::from_millis(17),
busy_timeout: Duration::from_millis(31),
..PoolConfig::default()
})
.expect("in-memory pool");
let report = collect(
&pool,
BuildIdentity::from_env("9.9.9", None),
Duration::from_secs(30),
);
let reader = report.reader_contention;
assert_eq!(reader.reader_admission_capacity, 1);
let json = serde_json::to_value(&report).expect("report serializes");
assert_eq!(
json.pointer("/reader_contention/configured_reader_cap"),
Some(&serde_json::json!(6))
);
assert_eq!(
json.pointer("/reader_contention/configured_checkout_timeout_ms"),
Some(&serde_json::json!(17))
);
assert_eq!(
json.pointer("/reader_contention/configured_busy_timeout_ms"),
Some(&serde_json::json!(31))
);
}
#[test]
fn diagnostics_composes_file_backed_standalone_acquisitions_without_counting_its_probe() {
let dir = tempfile::tempdir().expect("tempdir");
let (pool, _path) = seeded_pool(&dir);
drop(
pool.open_standalone_writer()
.expect("write-traffic standalone connection"),
);
let report = collect_with_audit_append_failures(
&pool,
BuildIdentity::from_env("9.9.9", None),
Duration::from_secs(30),
0,
);
assert_eq!(report.writer_contention.writer_acquisitions, 2);
assert_eq!(report.writer_contention.pooled_writer_acquisitions, 1);
assert_eq!(report.writer_contention.standalone_writer_acquisitions, 1);
assert_eq!(report.writer_contention.writer_task_acquisitions, 0);
assert_eq!(report.writer_contention.writer_acquisition_timeouts, 0);
let second = collect_with_audit_append_failures(
&pool,
BuildIdentity::from_env("9.9.9", None),
Duration::from_secs(30),
0,
);
assert_eq!(
second.writer_contention, report.writer_contention,
"the diagnostics PASSIVE probe must not inflate write-traffic counters"
);
}
#[test]
fn runtime_aware_collect_exposes_the_supplied_audit_failure_counter() {
let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool");
let report = collect_with_audit_append_failures(
&pool,
BuildIdentity::from_env("9.9.9", None),
Duration::from_secs(30),
17,
);
assert_eq!(report.writer_contention.audit_append_failures, Some(17));
assert!(
report
.writer_contention
.audit_append_failures_unavailable_reason
.is_none(),
"a supplied runtime counter must not carry an unavailable reason"
);
}
#[test]
fn diagnostics_exposes_an_induced_writer_checkout_timeout() {
let pool = ConnectionPool::new(PoolConfig {
checkout_timeout: Duration::from_millis(1),
..PoolConfig::default()
})
.expect("in-memory pool");
let held = pool.writer().expect("first writer checkout succeeds");
assert!(
matches!(
pool.writer(),
Err(crate::SqliteError::WriterPoolCheckoutTimeout { .. })
),
"holding the sole writer must exercise the typed timeout path"
);
drop(held);
let report = collect_with_audit_append_failures(
&pool,
BuildIdentity::from_env("9.9.9", None),
Duration::from_secs(30),
0,
);
assert_eq!(report.writer_contention.writer_acquisitions, 1);
assert_eq!(report.writer_contention.pooled_writer_acquisitions, 1);
assert_eq!(report.writer_contention.standalone_writer_acquisitions, 0);
assert_eq!(report.writer_contention.writer_task_acquisitions, 0);
assert_eq!(report.writer_contention.writer_acquisition_timeouts, 1);
}
#[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,
checkpoint_pressure_elevated_ticks: 0,
checkpoint_pressure_episodes_started: 0,
checkpoint_pressure_episodes_recovered: 0,
checkpoint_lifecycle_append_attempts: 0,
checkpoint_lifecycle_append_failures: 0,
checkpoint_lifecycle_enqueue_drops: 0,
read_tx_max_age_evictions: 0,
};
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 graph_edge_integrity_detects_the_pre_v14_duplicate_state() {
let conn = Connection::open_in_memory().expect("in-memory sqlite");
conn.execute_batch(
"CREATE TABLE graph_edges (
namespace TEXT NOT NULL,
id TEXT NOT NULL,
PRIMARY KEY (namespace, id)
);
CREATE TABLE graph_edges_seq (
seq INTEGER PRIMARY KEY AUTOINCREMENT,
edge_id TEXT NOT NULL UNIQUE
);
CREATE TABLE entities (id TEXT PRIMARY KEY, deleted_at INTEGER, merged_into TEXT);
INSERT INTO graph_edges(namespace, id)
VALUES ('alpha', 'shared-edge'), ('beta', 'shared-edge');
INSERT INTO graph_edges_seq(edge_id) VALUES ('shared-edge');",
)
.expect("seed the state possible before the V14 uniqueness guard");
let integrity = graph_edge_integrity(&conn).expect("integrity query succeeds");
assert_eq!(integrity.duplicate_edge_id_groups, 1);
assert_eq!(integrity.graph_edges_rows, 2);
assert_eq!(integrity.graph_edges_seq_rows, 1);
assert!(integrity.pre_v14_duplicate_edge_state_detected);
}
#[test]
fn graph_edge_integrity_counts_only_live_rows_that_still_carry_merge_provenance() {
let conn = Connection::open_in_memory().expect("in-memory sqlite");
conn.execute_batch(
"CREATE TABLE graph_edges (
namespace TEXT NOT NULL,
id TEXT NOT NULL,
PRIMARY KEY (namespace, id)
);
CREATE TABLE graph_edges_seq (
seq INTEGER PRIMARY KEY AUTOINCREMENT,
edge_id TEXT NOT NULL UNIQUE
);
CREATE TABLE entities (id TEXT PRIMARY KEY, deleted_at INTEGER, merged_into TEXT);
INSERT INTO entities(id, deleted_at, merged_into) VALUES
('kept', NULL, NULL),
('tombstoned-source', 1, 'kept'),
('left-live-by-an-old-restore', NULL, 'kept'),
('plain-soft-delete', 1, NULL);",
)
.expect("seed one row of each shape");
let integrity = graph_edge_integrity(&conn).expect("integrity query succeeds");
assert_eq!(
integrity.live_entities_carrying_merged_into, 1,
"a merge tombstone and a plain live row are both in order; only the live row \
carrying merged_into is the invariant violation"
);
assert_eq!(integrity.duplicate_edge_id_groups, 0);
}
#[test]
fn graph_edge_integrity_does_not_mislabel_retained_delete_history_as_a_duplicate() {
let conn = Connection::open_in_memory().expect("in-memory sqlite");
conn.execute_batch(
"CREATE TABLE graph_edges (
namespace TEXT NOT NULL,
id TEXT NOT NULL,
PRIMARY KEY (namespace, id)
);
CREATE TABLE graph_edges_seq (
seq INTEGER PRIMARY KEY AUTOINCREMENT,
edge_id TEXT NOT NULL UNIQUE
);
CREATE TABLE entities (id TEXT PRIMARY KEY, deleted_at INTEGER, merged_into TEXT);
INSERT INTO graph_edges(namespace, id) VALUES ('local', 'live-edge');
INSERT INTO graph_edges_seq(edge_id)
VALUES ('deleted-edge'), ('live-edge');",
)
.expect("seed a retained sequence row for a hard-deleted edge");
let integrity = graph_edge_integrity(&conn).expect("integrity query succeeds");
assert_eq!(integrity.duplicate_edge_id_groups, 0);
assert_eq!(integrity.graph_edges_rows, 1);
assert_eq!(integrity.graph_edges_seq_rows, 2);
assert!(
!integrity.pre_v14_duplicate_edge_state_detected,
"ledger rows intentionally survive hard deletion; count mismatch alone is not the \
pre-V14 duplicate state"
);
}
#[cfg(unix)]
#[test]
fn unmeasured_sidecar_cleanup_fields_are_absent_from_the_wire_payload() {
let pin = wal_pin_attribution_from_census(crate::walpin::CensusResult {
holders: std::collections::HashSet::new(),
uninspectable_pids: Vec::new(),
truncated: false,
budget_exhausted: false,
});
assert_eq!(pin.sidecar_listing_truncated, None);
assert_eq!(pin.sidecar_entries_cleanup_would_reap, None);
let json = serde_json::to_value(pin).expect("attribution serializes");
assert!(
json.get("sidecar_listing_truncated").is_none(),
"a skipped enumeration must omit sidecar_listing_truncated, not fabricate false"
);
assert!(
json.get("sidecar_entries_cleanup_would_reap").is_none(),
"a skipped enumeration must omit sidecar_entries_cleanup_would_reap, not fabricate 0"
);
}
#[cfg(unix)]
#[test]
fn wal_pin_census_serializes_only_under_the_nested_carrier() {
let pin = wal_pin_attribution_from_census(crate::walpin::CensusResult {
holders: std::collections::HashSet::from([41, 7]),
uninspectable_pids: vec![99],
truncated: true,
budget_exhausted: false,
});
let json = serde_json::to_value(pin).expect("attribution serializes");
assert_eq!(json["census"]["holder_pids"], serde_json::json!([7, 41]));
assert_eq!(
json["census"]["uninspectable_pids"],
serde_json::json!([99])
);
assert_eq!(json["census"]["truncated"], true);
for duplicate in [
"census_holder_pids",
"census_uninspectable_pids",
"census_truncated",
"census_is_complete",
] {
assert!(
json.get(duplicate).is_none(),
"wal_pin.{duplicate} must not duplicate wal_pin.census: {json}"
);
}
}
#[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.fts_segments.is_none());
assert!(report.fts_segments_error.is_some());
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());
assert_eq!(report.wal_pin.status, WalPinAttributionStatus::Unavailable);
assert!(matches!(
report.wal_pin.census,
WalPinCensus::Unavailable { .. }
));
assert!(report.size_composition.is_none());
assert!(report
.size_composition_error
.as_deref()
.is_some_and(|reason| reason.contains("no file-backed page composition")));
}
#[cfg(unix)]
#[test]
fn incomplete_holder_census_is_a_tagged_degraded_result() {
let census = crate::walpin::CensusResult {
holders: std::collections::HashSet::from([41, 7]),
uninspectable_pids: vec![99, 99],
truncated: true,
budget_exhausted: false,
};
let pin = wal_pin_attribution_from_census(census);
assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
assert!(!pin.available);
assert!(!pin.census_is_complete);
assert_eq!(pin.census_holder_pids, vec![7, 41]);
assert_eq!(pin.census_uninspectable_pids, vec![99]);
assert!(
pin.unavailable_reason.as_deref().is_some_and(
|reason| reason.contains("additional database holders cannot be ruled out")
),
"the legacy reason must also fail loud for old consumers: {pin:?}"
);
match &pin.census {
WalPinCensus::Incomplete {
holder_pids,
uninspectable_pids,
truncated,
reason,
} => {
assert_eq!(holder_pids, &vec![7, 41]);
assert_eq!(uninspectable_pids, &vec![99]);
assert!(*truncated);
assert!(reason.contains("additional database holders cannot be ruled out"));
}
other => panic!("incomplete scan must serialize as incomplete, got {other:?}"),
}
let json = serde_json::to_value(&pin).expect("serializes");
assert_eq!(json["status"], "degraded");
assert_eq!(json["census"]["status"], "incomplete");
}
#[cfg(unix)]
#[test]
fn complete_holder_census_stays_explicit_while_attribution_is_degraded() {
let census = crate::walpin::CensusResult {
holders: std::collections::HashSet::from([7]),
uninspectable_pids: Vec::new(),
truncated: false,
budget_exhausted: false,
};
let pin = wal_pin_attribution_from_census(census);
assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
assert!(pin.census_is_complete);
assert!(matches!(
pin.census,
WalPinCensus::Complete { ref holder_pids } if holder_pids == &vec![7]
));
assert_eq!(
pin.status_reasons.len(),
1,
"only missing sidecar reconciliation degrades a complete OS census"
);
}
#[cfg(unix)]
#[test]
fn complete_holder_and_read_only_sidecar_evidence_reconcile_to_complete() {
let census = crate::walpin::CensusResult {
holders: std::collections::HashSet::from([7]),
uninspectable_pids: Vec::new(),
truncated: false,
budget_exhausted: false,
};
let sidecar = crate::walpin::WalpinReport {
entries: vec![crate::walpin::WalpinPidHealth::RegisteredSilent { pid: 7 }],
sidecar_listing_truncated: false,
cleanup_would_reap: 0,
orphan_temps_reaped: 0,
};
let pin = wal_pin_attribution_from_evidence(census, sidecar);
assert_eq!(pin.status, WalPinAttributionStatus::Complete);
assert!(pin.available);
assert!(pin.fully_attributed);
assert!(pin.status_reasons.is_empty());
assert_eq!(pin.registered_silent_pids, vec![7]);
assert!(pin.census_pids_without_attribution.is_empty());
assert_eq!(pin.sidecar_listing_truncated, Some(false));
assert_eq!(pin.sidecar_entries_cleanup_would_reap, Some(0));
}
#[cfg(unix)]
#[test]
fn complete_census_with_an_unregistered_holder_is_degraded_not_exonerated() {
let census = crate::walpin::CensusResult {
holders: std::collections::HashSet::from([7, 41]),
uninspectable_pids: Vec::new(),
truncated: false,
budget_exhausted: false,
};
let sidecar = crate::walpin::WalpinReport {
entries: vec![crate::walpin::WalpinPidHealth::RegisteredSilent { pid: 7 }],
sidecar_listing_truncated: false,
cleanup_would_reap: 0,
orphan_temps_reaped: 0,
};
let pin = wal_pin_attribution_from_evidence(census, sidecar);
assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
assert!(!pin.available);
assert!(!pin.fully_attributed);
assert_eq!(pin.census_pids_without_attribution, vec![41]);
assert!(pin
.status_reasons
.iter()
.any(|reason| reason.contains("holder(s) have no sidecar attribution")));
}
#[test]
fn database_size_composition_reports_tables_indexes_fts_and_vectors_separately() {
let conn = Connection::open_in_memory().expect("in-memory sqlite");
conn.execute_batch(
"CREATE TABLE docs(id INTEGER PRIMARY KEY, body TEXT NOT NULL); \
CREATE INDEX idx_docs_body ON docs(body); \
CREATE TABLE fts_demo_data(id INTEGER PRIMARY KEY, block BLOB); \
CREATE TABLE vec_demo_chunks(id INTEGER PRIMARY KEY, vectors BLOB); \
CREATE TABLE knowledge_sections(id INTEGER PRIMARY KEY, embedding BLOB); \
INSERT INTO docs(body) VALUES (zeroblob(8192)); \
INSERT INTO fts_demo_data(block) VALUES (zeroblob(8192)); \
INSERT INTO vec_demo_chunks(vectors) VALUES (zeroblob(8192)); \
INSERT INTO knowledge_sections(embedding) VALUES (zeroblob(8192));",
)
.expect("seed size classes");
let composition = database_size_composition(&conn).expect("dbstat composition");
let class_for = |name: &str| {
composition
.objects
.iter()
.find(|object| object.name == name)
.map(|object| object.storage_class)
};
assert_eq!(class_for("docs"), Some(DatabaseStorageClass::RowTable));
assert_eq!(
class_for("idx_docs_body"),
Some(DatabaseStorageClass::Index)
);
assert_eq!(
class_for("fts_demo_data"),
Some(DatabaseStorageClass::FullText)
);
assert_eq!(
class_for("vec_demo_chunks"),
Some(DatabaseStorageClass::Vector)
);
assert_eq!(
class_for("knowledge_sections"),
Some(DatabaseStorageClass::MixedRowAndEmbedding)
);
assert!(composition.vector_bytes > 0);
assert!(composition.full_text_bytes > 0);
assert!(composition.mixed_embedding_bytes > 0);
assert_eq!(
composition
.accounted_bytes
.saturating_add(composition.freelist_bytes)
.saturating_add(composition.unaccounted_bytes),
composition.database_bytes
);
}
#[cfg(unix)]
#[test]
#[serial(khive_walpin_sidecar_env)]
fn wal_pin_attribution_degrades_when_a_holder_has_no_sidecar_registration() {
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,
"the pool's OS holder has no test sidecar registration"
);
assert!(
pin.unavailable_reason.is_some(),
"the missing holder attribution must be explained: {pin:?}"
);
assert!(pin.sidecar_entries.is_empty());
assert!(pin.reporting.is_empty());
assert_eq!(pin.sidecar_listing_truncated, Some(false));
assert_eq!(pin.sidecar_entries_cleanup_would_reap, Some(0));
assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
assert!(matches!(
pin.census,
WalPinCensus::Complete { .. } | WalPinCensus::Incomplete { .. }
));
}
#[cfg(unix)]
#[test]
#[serial(khive_walpin_sidecar_env)]
fn diagnostics_finds_sidecar_evidence_through_an_aliased_database_path() {
let dir = tempfile::tempdir().expect("tempdir");
let real_dir = dir.path().join("real");
std::fs::create_dir(&real_dir).expect("mkdir real dir");
let real_path = real_dir.join("diag.db");
std::fs::write(&real_path, b"").expect("create real file");
let alias_path = dir.path().join("alias.db");
std::os::unix::fs::symlink(&real_path, &alias_path).expect("symlink alias");
let pool = ConnectionPool::new(PoolConfig {
path: Some(alias_path.clone()),
..PoolConfig::for_test()
})
.expect("pool open through symlinked path");
{
let writer = pool.try_writer().expect("writer");
writer
.conn()
.execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
.expect("seed a write so the WAL file exists");
}
let canonical = pool
.canonical_path()
.expect("file-backed pool")
.to_path_buf();
assert_ne!(
canonical, alias_path,
"the alias must actually differ from the canonical path for this test to mean \
anything"
);
let pid = std::process::id();
let sidecar_dir = crate::walpin::sidecar_dir_for(&canonical);
let beacon = crate::walpin::WalpinBeacon {
pid,
process_role: "session".to_string(),
started_at: crate::walpin::process_start_time_secs(pid).unwrap_or(0),
sweep_interval_ms: 5_000,
};
crate::walpin::write_beacon(&sidecar_dir, &beacon).expect("seed this process's beacon");
let report = collect(
&pool,
BuildIdentity::from_env("test", None),
Duration::from_secs(30),
);
assert_eq!(
report.wal_pin.sidecar_listing_truncated,
Some(false),
"the sidecar enumeration must run to completion: {:?}",
report.wal_pin
);
assert!(
report.wal_pin.registered_silent_pids.contains(&pid),
"the beacon written beside the canonical path must be found: {:?}",
report.wal_pin
);
assert!(
report.wal_pin.census_holder_pids.contains(&pid),
"the OS census must find this process holding its own database open: {:?}",
report.wal_pin
);
assert!(
!report
.wal_pin
.census_pids_without_attribution
.contains(&pid),
"this process's own holder entry must be attributed by its own sidecar evidence, \
not left unexplained: {:?}",
report.wal_pin
);
}
#[cfg(unix)]
#[tokio::test]
#[serial(khive_walpin_sidecar_env)]
async fn diagnostics_finds_sidecar_evidence_through_an_aliased_database_path_async() {
let dir = tempfile::tempdir().expect("tempdir");
let real_dir = dir.path().join("real");
std::fs::create_dir(&real_dir).expect("mkdir real dir");
let real_path = real_dir.join("diag.db");
std::fs::write(&real_path, b"").expect("create real file");
let alias_path = dir.path().join("alias.db");
std::os::unix::fs::symlink(&real_path, &alias_path).expect("symlink alias");
let pool = ConnectionPool::new(PoolConfig {
path: Some(alias_path.clone()),
..PoolConfig::for_test()
})
.expect("pool open through symlinked path");
{
let writer = pool.try_writer().expect("writer");
writer
.conn()
.execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
.expect("seed a write so the WAL file exists");
}
let canonical = pool
.canonical_path()
.expect("file-backed pool")
.to_path_buf();
assert_ne!(
canonical, alias_path,
"the alias must actually differ from the canonical path for this test to mean \
anything"
);
let pid = std::process::id();
let sidecar_dir = crate::walpin::sidecar_dir_for(&canonical);
let beacon = crate::walpin::WalpinBeacon {
pid,
process_role: "session".to_string(),
started_at: crate::walpin::process_start_time_secs(pid).unwrap_or(0),
sweep_interval_ms: 5_000,
};
crate::walpin::write_beacon(&sidecar_dir, &beacon).expect("seed this process's beacon");
let pool = Arc::new(pool);
let report = collect_with_audit_append_failures_interruptibly(
Arc::clone(&pool),
BuildIdentity::from_env("test", None),
Duration::from_secs(30),
0,
)
.await
.expect("diagnostics succeed");
assert_eq!(
report.wal_pin.sidecar_listing_truncated,
Some(false),
"the sidecar enumeration must run to completion: {:?}",
report.wal_pin
);
assert!(
report.wal_pin.registered_silent_pids.contains(&pid),
"the beacon written beside the canonical path must be found: {:?}",
report.wal_pin
);
assert!(
report.wal_pin.census_holder_pids.contains(&pid),
"the OS census must find this process holding its own database open: {:?}",
report.wal_pin
);
assert!(
!report
.wal_pin
.census_pids_without_attribution
.contains(&pid),
"this process's own holder entry must be attributed by its own sidecar evidence, \
not left unexplained: {:?}",
report.wal_pin
);
}
#[cfg(unix)]
#[test]
#[serial(khive_walpin_sidecar_env)]
fn wal_pin_attribution_reports_disabled_when_the_sidecar_is_explicitly_off() {
let dir = tempfile::tempdir().expect("tempdir");
let (pool, path) = seeded_pool(&dir);
let _ = &pool;
let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
std::env::set_var("KHIVE_WALPIN_SIDECAR", "0");
let pin = wal_pin_attribution(&path, Duration::from_secs(30));
assert!(
!pin.available,
"an explicitly disabled sidecar can never produce a reconciled answer"
);
assert!(
pin.unavailable_reason
.as_deref()
.is_some_and(|reason| reason.contains("disabled")),
"the reason must name the disabled sidecar, not a generic enumeration failure: \
{pin:?}"
);
assert!(pin.sidecar_entries.is_empty());
assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
}
}