1use std::path::{Path, PathBuf};
62use std::sync::atomic::{AtomicBool, Ordering};
63use std::sync::Arc;
64use std::time::{Duration, Instant};
65#[cfg(unix)]
66use std::time::{SystemTime, UNIX_EPOCH};
67
68use khive_storage::error::StorageError;
69use khive_storage::types::StorageResult;
70use khive_storage::StorageCapability;
71use rusqlite::Connection;
72use serde::Serialize;
73
74use crate::checkpoint;
75use crate::pool::ConnectionPool;
76
77#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
84pub struct CheckpointProbe {
85 pub busy: i64,
86 pub log_frames: i64,
87 pub checkpointed_frames: i64,
88}
89
90impl CheckpointProbe {
91 pub fn pin_depth(&self) -> i64 {
94 (self.log_frames - self.checkpointed_frames).max(0)
95 }
96}
97
98pub fn checkpoint_probe(conn: &Connection) -> rusqlite::Result<CheckpointProbe> {
112 conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
113 Ok(CheckpointProbe {
114 busy: row.get(0)?,
115 log_frames: row.get(1)?,
116 checkpointed_frames: row.get(2)?,
117 })
118 })
119}
120
121#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
127pub struct CheckpointCounters {
128 pub last_observed_wal_pages: Option<u64>,
129 pub truncate_attempts: u64,
130 pub truncate_consecutive_failures: u64,
131 pub checkpoint_skipped_ticks: u64,
132 pub checkpoint_consecutive_skips: u64,
133 pub checkpoint_last_skip_wal_pages: Option<u64>,
134 pub checkpoint_pressure_elevated_ticks: u64,
135 pub checkpoint_pressure_episodes_started: u64,
136 pub checkpoint_pressure_episodes_recovered: u64,
137 pub checkpoint_lifecycle_append_attempts: u64,
138 pub checkpoint_lifecycle_append_failures: u64,
139 pub checkpoint_lifecycle_enqueue_drops: u64,
140 pub read_tx_max_age_evictions: u64,
144}
145
146pub fn checkpoint_counters() -> CheckpointCounters {
148 CheckpointCounters {
149 last_observed_wal_pages: checkpoint::last_observed_wal_pages(),
150 truncate_attempts: checkpoint::truncate_attempts(),
151 truncate_consecutive_failures: checkpoint::truncate_consecutive_failures(),
152 checkpoint_skipped_ticks: checkpoint::checkpoint_skipped_ticks(),
153 checkpoint_consecutive_skips: checkpoint::checkpoint_consecutive_skips(),
154 checkpoint_last_skip_wal_pages: checkpoint::checkpoint_last_skip_wal_pages(),
155 checkpoint_pressure_elevated_ticks: checkpoint::checkpoint_pressure_elevated_ticks(),
156 checkpoint_pressure_episodes_started: checkpoint::checkpoint_pressure_episodes_started(),
157 checkpoint_pressure_episodes_recovered: checkpoint::checkpoint_pressure_episodes_recovered(
158 ),
159 checkpoint_lifecycle_append_attempts: checkpoint::checkpoint_lifecycle_append_attempts(),
160 checkpoint_lifecycle_append_failures: checkpoint::checkpoint_lifecycle_append_failures(),
161 checkpoint_lifecycle_enqueue_drops: checkpoint::checkpoint_lifecycle_enqueue_drops(),
162 read_tx_max_age_evictions: checkpoint::read_tx_max_age_evictions(),
163 }
164}
165
166#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
173pub struct BuildIdentity {
174 pub version: String,
175 pub build_hash: Option<String>,
176}
177
178impl BuildIdentity {
179 pub fn from_env(version: &str, build_hash: Option<&str>) -> Self {
181 Self {
182 version: version.to_string(),
183 build_hash: build_hash.map(str::to_string),
184 }
185 }
186}
187
188#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
190pub struct ProcessIdentity {
191 pub pid: u32,
192 pub started_at: Option<i64>,
195 pub started_at_unavailable_reason: Option<String>,
198 pub pool_generation: u64,
199}
200
201impl ProcessIdentity {
202 pub fn current(pool: &ConnectionPool) -> Self {
203 let pid = std::process::id();
204 Self::from_start_time(
205 pid,
206 crate::walpin::process_start_time_secs(pid),
207 pool.main_pool_generation(),
208 )
209 }
210
211 fn from_start_time(pid: u32, started_at: Option<i64>, pool_generation: u64) -> Self {
212 Self {
213 pid,
214 started_at,
215 started_at_unavailable_reason: started_at.is_none().then(|| {
216 format!(
217 "OS process start time is unsupported or unavailable on {}",
218 std::env::consts::OS
219 )
220 }),
221 pool_generation,
222 }
223 }
224}
225
226#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
228pub struct WalFileState {
229 pub wal_path: String,
230 pub wal_size_bytes: Option<u64>,
233 pub unavailable_reason: Option<String>,
234}
235
236pub fn wal_file_state(db_path: &Path) -> WalFileState {
238 let wal_path = wal_sidecar_path(db_path);
239 match std::fs::metadata(&wal_path) {
240 Ok(md) => WalFileState {
241 wal_path: wal_path.display().to_string(),
242 wal_size_bytes: Some(md.len()),
243 unavailable_reason: None,
244 },
245 Err(e) => WalFileState {
246 wal_path: wal_path.display().to_string(),
247 wal_size_bytes: None,
248 unavailable_reason: Some(e.to_string()),
249 },
250 }
251}
252
253fn wal_sidecar_path(db_path: &Path) -> PathBuf {
256 let mut s = db_path.as_os_str().to_os_string();
257 s.push("-wal");
258 PathBuf::from(s)
259}
260
261#[derive(Debug, Clone, PartialEq, Serialize)]
270pub struct WalPinAttribution {
271 pub status: WalPinAttributionStatus,
273 pub status_reasons: Vec<String>,
275 pub census: WalPinCensus,
277 pub available: bool,
280 pub unavailable_reason: Option<String>,
281 #[serde(skip_serializing)]
283 pub census_holder_pids: Vec<u32>,
284 #[serde(skip_serializing)]
285 pub census_uninspectable_pids: Vec<u32>,
286 #[serde(skip_serializing)]
287 pub census_truncated: bool,
288 #[serde(skip_serializing)]
289 pub census_is_complete: bool,
290 pub reporting: Vec<WalPinHolder>,
292 pub registered_silent_pids: Vec<u32>,
294 pub unknown_pids: Vec<u32>,
296 pub census_pids_without_attribution: Vec<u32>,
298 pub fully_attributed: bool,
301 pub sidecar_entries: Vec<serde_json::Value>,
303 #[serde(skip_serializing_if = "Option::is_none")]
305 pub sidecar_listing_truncated: Option<bool>,
306 #[serde(skip_serializing_if = "Option::is_none")]
309 pub sidecar_entries_cleanup_would_reap: Option<usize>,
310}
311
312#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
317#[serde(rename_all = "snake_case")]
318pub enum WalPinAttributionStatus {
319 Complete,
321 Degraded,
323 Unavailable,
325}
326
327#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
335#[serde(tag = "status", rename_all = "snake_case")]
336pub enum WalPinCensus {
337 Complete {
339 holder_pids: Vec<u32>,
341 },
342 Incomplete {
344 holder_pids: Vec<u32>,
346 uninspectable_pids: Vec<u32>,
348 truncated: bool,
350 reason: String,
352 },
353 Unavailable {
355 reason: String,
357 },
358}
359
360fn census_truncation_cause(budget_exhausted: bool) -> String {
370 if budget_exhausted {
371 "the OS process walk stopped at its wall-clock budget (see collection_cost.wal_pin_census_budget_ms)"
372 .to_string()
373 } else {
374 "the OS process walk was truncated".to_string()
375 }
376}
377
378#[derive(Debug, Clone, PartialEq, Serialize)]
380pub struct WalPinHolder {
381 pub pid: u32,
382 pub process_role: String,
383 pub current_oldest_tx_age_secs: f64,
384 pub oldest_tx_label: Option<String>,
385 pub attribution_is_evidence_backed: bool,
386}
387
388impl WalPinAttribution {
389 fn unavailable(reason: impl Into<String>) -> Self {
390 let reason = reason.into();
391 Self {
392 status: WalPinAttributionStatus::Unavailable,
393 status_reasons: vec![reason.clone()],
394 census: WalPinCensus::Unavailable {
395 reason: reason.clone(),
396 },
397 available: false,
398 unavailable_reason: Some(reason),
399 census_holder_pids: Vec::new(),
400 census_uninspectable_pids: Vec::new(),
401 census_truncated: false,
402 census_is_complete: false,
403 reporting: Vec::new(),
404 registered_silent_pids: Vec::new(),
405 unknown_pids: Vec::new(),
406 census_pids_without_attribution: Vec::new(),
407 fully_attributed: false,
408 sidecar_entries: Vec::new(),
409 sidecar_listing_truncated: None,
410 sidecar_entries_cleanup_would_reap: None,
411 }
412 }
413}
414
415#[cfg(all(unix, test))]
416fn wal_pin_attribution_from_census(census: crate::walpin::CensusResult) -> WalPinAttribution {
417 wal_pin_attribution_without_sidecar(
418 census,
419 "read-only sidecar enumeration did not run for this attribution snapshot".to_string(),
420 )
421}
422
423#[cfg(unix)]
424fn wal_pin_attribution_without_sidecar(
425 census: crate::walpin::CensusResult,
426 sidecar_reason: String,
427) -> WalPinAttribution {
428 let census_is_complete = census.is_complete();
429 let mut census_holder_pids: Vec<u32> = census.holders.iter().copied().collect();
430 census_holder_pids.sort_unstable();
431 let mut census_uninspectable_pids = census.uninspectable_pids;
432 census_uninspectable_pids.sort_unstable();
433 census_uninspectable_pids.dedup();
434 let census_truncated = census.truncated;
435 let census_budget_exhausted = census.budget_exhausted;
436
437 let mut status_reasons = vec![sidecar_reason];
438 let census = if census_is_complete {
439 WalPinCensus::Complete {
440 holder_pids: census_holder_pids.clone(),
441 }
442 } else {
443 let mut causes = Vec::new();
444 if census_truncated {
445 causes.push(census_truncation_cause(census_budget_exhausted));
446 }
447 if !census_uninspectable_pids.is_empty() {
448 causes.push(format!(
449 "{} PID(s) could not be inspected",
450 census_uninspectable_pids.len()
451 ));
452 }
453 let reason = format!(
454 "OS holder census is incomplete: {}; additional database holders cannot be ruled out",
455 causes.join("; ")
456 );
457 status_reasons.push(reason.clone());
458 WalPinCensus::Incomplete {
459 holder_pids: census_holder_pids.clone(),
460 uninspectable_pids: census_uninspectable_pids.clone(),
461 truncated: census_truncated,
462 reason,
463 }
464 };
465
466 WalPinAttribution {
467 status: WalPinAttributionStatus::Degraded,
468 unavailable_reason: Some(status_reasons.join("; ")),
469 status_reasons,
470 census,
471 available: false,
472 census_holder_pids,
473 census_uninspectable_pids,
474 census_truncated,
475 census_is_complete,
476 reporting: Vec::new(),
477 registered_silent_pids: Vec::new(),
478 unknown_pids: Vec::new(),
479 census_pids_without_attribution: Vec::new(),
480 fully_attributed: false,
481 sidecar_entries: Vec::new(),
482 sidecar_listing_truncated: None,
483 sidecar_entries_cleanup_would_reap: None,
484 }
485}
486
487#[cfg(unix)]
488fn wal_pin_attribution_from_evidence(
489 census: crate::walpin::CensusResult,
490 sidecar: crate::walpin::WalpinReport,
491) -> WalPinAttribution {
492 use std::collections::BTreeSet;
493
494 let census_is_complete = census.is_complete();
495 let mut census_holder_pids: Vec<u32> = census.holders.iter().copied().collect();
496 census_holder_pids.sort_unstable();
497 let mut census_uninspectable_pids = census.uninspectable_pids;
498 census_uninspectable_pids.sort_unstable();
499 census_uninspectable_pids.dedup();
500 let census_truncated = census.truncated;
501 let census_budget_exhausted = census.budget_exhausted;
502 let census_carrier = if census_is_complete {
503 WalPinCensus::Complete {
504 holder_pids: census_holder_pids.clone(),
505 }
506 } else {
507 let mut causes = Vec::new();
508 if census_truncated {
509 causes.push(census_truncation_cause(census_budget_exhausted));
510 }
511 if !census_uninspectable_pids.is_empty() {
512 causes.push(format!(
513 "{} PID(s) could not be inspected",
514 census_uninspectable_pids.len()
515 ));
516 }
517 let reason = format!(
518 "OS holder census is incomplete: {}; additional database holders cannot be ruled out",
519 causes.join("; ")
520 );
521 WalPinCensus::Incomplete {
522 holder_pids: census_holder_pids.clone(),
523 uninspectable_pids: census_uninspectable_pids.clone(),
524 truncated: census_truncated,
525 reason,
526 }
527 };
528
529 let now_epoch_secs = SystemTime::now()
530 .duration_since(UNIX_EPOCH)
531 .map(|duration| duration.as_secs() as i64)
532 .unwrap_or(0);
533 let sidecar_listing_truncated = sidecar.sidecar_listing_truncated;
534 let sidecar_entries_cleanup_would_reap = sidecar.cleanup_would_reap;
535 let mut reporting = Vec::new();
536 let mut registered_silent_pids = Vec::new();
537 let mut unknown_pids = Vec::new();
538 let mut sidecar_entries = Vec::new();
539 let mut sidecar_known_pids = BTreeSet::new();
540
541 for entry in sidecar.entries {
542 match entry {
543 crate::walpin::WalpinPidHealth::Reporting(heartbeat) => {
544 let current_oldest_tx_age_secs =
545 heartbeat.current_oldest_tx_age_secs(now_epoch_secs);
546 let attribution_is_evidence_backed = heartbeat.attribution_is_evidence_backed();
547 sidecar_known_pids.insert(heartbeat.pid);
548 reporting.push(WalPinHolder {
549 pid: heartbeat.pid,
550 process_role: heartbeat.process_role.clone(),
551 current_oldest_tx_age_secs,
552 oldest_tx_label: heartbeat.oldest_tx_label.clone(),
553 attribution_is_evidence_backed,
554 });
555 sidecar_entries.push((
556 heartbeat.pid,
557 0u8,
558 serde_json::json!({
559 "pid": heartbeat.pid,
560 "status": "reporting",
561 "process_role": heartbeat.process_role,
562 "current_oldest_tx_age_secs": current_oldest_tx_age_secs,
563 "oldest_tx_label": heartbeat.oldest_tx_label,
564 "attribution_is_evidence_backed": attribution_is_evidence_backed,
565 }),
566 ));
567 }
568 crate::walpin::WalpinPidHealth::RegisteredSilent { pid } => {
569 sidecar_known_pids.insert(pid);
570 registered_silent_pids.push(pid);
571 sidecar_entries.push((
572 pid,
573 1u8,
574 serde_json::json!({"pid": pid, "status": "registered_silent"}),
575 ));
576 }
577 crate::walpin::WalpinPidHealth::Unknown { pid, reason } => {
578 sidecar_known_pids.insert(pid);
579 unknown_pids.push(pid);
580 sidecar_entries.push((
581 pid,
582 2u8,
583 serde_json::json!({"pid": pid, "status": "unknown", "reason": reason}),
584 ));
585 }
586 }
587 }
588
589 reporting.sort_by_key(|holder| holder.pid);
590 reporting.dedup_by_key(|holder| holder.pid);
591 registered_silent_pids.sort_unstable();
592 registered_silent_pids.dedup();
593 unknown_pids.sort_unstable();
594 unknown_pids.dedup();
595 sidecar_entries.sort_by_key(|(pid, status_rank, _)| (*pid, *status_rank));
596 let sidecar_entries = sidecar_entries
597 .into_iter()
598 .map(|(_, _, entry)| entry)
599 .collect();
600 let census_pids_without_attribution: Vec<u32> = census_holder_pids
601 .iter()
602 .copied()
603 .filter(|pid| !sidecar_known_pids.contains(pid))
604 .collect();
605
606 let mut status_reasons = Vec::new();
607 if let WalPinCensus::Incomplete { reason, .. } = &census_carrier {
608 status_reasons.push(reason.clone());
609 }
610 if sidecar_listing_truncated {
611 status_reasons.push(
612 "read-only sidecar enumeration reached its entry cap; additional entries may exist"
613 .to_string(),
614 );
615 }
616 if !unknown_pids.is_empty() {
617 status_reasons.push(format!(
618 "{} sidecar PID(s) could not be classified conclusively",
619 unknown_pids.len()
620 ));
621 }
622 if !census_pids_without_attribution.is_empty() {
623 status_reasons.push(format!(
624 "{} OS-confirmed holder(s) have no sidecar attribution",
625 census_pids_without_attribution.len()
626 ));
627 }
628
629 let fully_attributed = census_is_complete
630 && !sidecar_listing_truncated
631 && unknown_pids.is_empty()
632 && census_pids_without_attribution.is_empty();
633 let status = if fully_attributed {
634 WalPinAttributionStatus::Complete
635 } else {
636 WalPinAttributionStatus::Degraded
637 };
638 let unavailable_reason = (!fully_attributed).then(|| status_reasons.join("; "));
639
640 WalPinAttribution {
641 status,
642 status_reasons,
643 census: census_carrier,
644 available: fully_attributed,
645 unavailable_reason,
646 census_holder_pids,
647 census_uninspectable_pids,
648 census_truncated,
649 census_is_complete,
650 reporting,
651 registered_silent_pids,
652 unknown_pids,
653 census_pids_without_attribution,
654 fully_attributed,
655 sidecar_entries,
656 sidecar_listing_truncated: Some(sidecar_listing_truncated),
657 sidecar_entries_cleanup_would_reap: Some(sidecar_entries_cleanup_would_reap),
658 }
659}
660
661#[cfg(unix)]
670pub fn wal_pin_attribution(db_path: &Path, sweep_interval: Duration) -> WalPinAttribution {
671 use crate::walpin;
672
673 let census = match walpin::census_holders(db_path) {
674 Ok(c) => c,
675 Err(e) => return WalPinAttribution::unavailable(format!("census_holders failed: {e}")),
676 };
677 if !walpin::sidecar_enabled(true) {
678 return wal_pin_attribution_without_sidecar(census, SIDECAR_DISABLED_REASON.to_string());
679 }
680 match walpin::inspect_live(&walpin::sidecar_dir_for(db_path), sweep_interval) {
681 Ok(sidecar) => wal_pin_attribution_from_evidence(census, sidecar),
682 Err(error) => wal_pin_attribution_without_sidecar(
683 census,
684 format!("read-only sidecar enumeration failed: {error}"),
685 ),
686 }
687}
688
689#[cfg(unix)]
691const SIDECAR_DISABLED_REASON: &str =
692 "walpin sidecar is explicitly disabled (KHIVE_WALPIN_SIDECAR); attribution has no sidecar \
693 evidence to reconcile against the OS holder census";
694
695#[cfg(not(unix))]
696pub fn wal_pin_attribution(_db_path: &Path, _sweep_interval: Duration) -> WalPinAttribution {
697 WalPinAttribution::unavailable("WAL-pin attribution requires a Unix platform")
698}
699
700#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
707pub struct ReaderContentionDiagnostics {
708 pub configured_reader_cap: usize,
711 pub configured_checkout_timeout_ms: u64,
713 pub configured_busy_timeout_ms: u64,
716 pub reader_admission_capacity: usize,
719 pub available_reader_admission_slots: usize,
721 pub reader_acquisitions: u64,
724 pub pooled_reader_checkouts: u64,
726 pub standalone_reader_opens: u64,
729 pub infrastructure_standalone_reader_opens: u64,
732 pub reader_checkout_timeouts: u64,
735 pub active_pooled_reader_checkouts: u64,
737 pub peak_active_pooled_reader_checkouts: u64,
739 pub completed_pooled_reader_checkouts: u64,
741 pub max_completed_reader_hold_micros: u64,
743 pub max_completed_reader_hold_operation: Option<&'static str>,
748 pub reader_replacement_open_failures: u64,
753}
754
755impl ReaderContentionDiagnostics {
756 fn snapshot(pool: &ConnectionPool) -> Self {
757 let reader = pool.reader_acquisition_snapshot();
758 Self {
759 configured_reader_cap: pool.config().max_readers,
760 configured_checkout_timeout_ms: u64::try_from(
761 pool.config().checkout_timeout.as_millis(),
762 )
763 .unwrap_or(u64::MAX),
764 configured_busy_timeout_ms: u64::try_from(pool.config().busy_timeout.as_millis())
765 .unwrap_or(u64::MAX),
766 reader_admission_capacity: reader.reader_admission_capacity,
767 available_reader_admission_slots: reader.available_reader_admission_slots,
768 reader_acquisitions: reader.acquisitions,
769 pooled_reader_checkouts: reader.pooled_checkouts,
770 standalone_reader_opens: reader.standalone_opens,
771 infrastructure_standalone_reader_opens: reader.infrastructure_standalone_opens,
772 reader_checkout_timeouts: reader.checkout_timeouts,
773 active_pooled_reader_checkouts: reader.active_pooled_checkouts,
774 peak_active_pooled_reader_checkouts: reader.peak_active_pooled_checkouts,
775 completed_pooled_reader_checkouts: reader.completed_pooled_checkouts,
776 max_completed_reader_hold_micros: reader.max_completed_hold_micros,
777 max_completed_reader_hold_operation: reader.max_completed_hold_operation,
778 reader_replacement_open_failures: reader.reader_replacement_open_failures,
779 }
780 }
781}
782
783#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
798pub struct WriterContentionDiagnostics {
799 pub writer_acquisitions: u64,
802 pub pooled_writer_acquisitions: u64,
804 pub standalone_writer_acquisitions: u64,
806 pub writer_task_acquisitions: u64,
808 pub writer_acquisition_timeouts: u64,
810 pub writer_task_begin_busy: u64,
815 pub writer_task_begin_busy_absorbed: u64,
820 pub writer_task_begin_errors: u64,
823 pub writer_task_request_failures: u64,
828 pub writer_task_side_effects_unknown: u64,
831 pub audit_append_failures: Option<u64>,
843 pub audit_append_failures_unavailable_reason: Option<String>,
845 pub audit_obligation_append_failures: Option<u64>,
850 pub audit_obligation_append_failures_unavailable_reason: Option<String>,
852 pub audit_batch_flush_failures: Option<u64>,
857 pub audit_batch_flush_failures_unavailable_reason: Option<String>,
859 pub audit_degraded_rows: Option<u64>,
862 pub audit_degraded_rows_unavailable_reason: Option<String>,
864 pub audit_degraded: Option<bool>,
870 pub audit_degraded_unavailable_reason: Option<String>,
872 pub audit_admission_refused_obligations: Option<u64>,
892 pub audit_admission_refused_obligations_last_at_ms: Option<u64>,
898 pub audit_admission_refused_obligations_unavailable_reason: Option<String>,
900 pub audit_admission_unresolved_obligations: Option<u64>,
921 pub audit_admission_unresolved_obligations_last_at_ms: Option<u64>,
925 pub audit_admission_unresolved_obligations_unavailable_reason: Option<String>,
928}
929
930#[derive(Debug, Clone, Copy, PartialEq, Eq)]
935pub struct RuntimeAuditBatchMetrics {
936 pub flush_failures: u64,
939 pub degraded_rows: u64,
941 pub degraded: bool,
943 pub admission_refused_obligations: u64,
948 pub admission_refused_obligations_last_at_ms: Option<u64>,
951 pub admission_unresolved_obligations: u64,
958 pub admission_unresolved_obligations_last_at_ms: Option<u64>,
961}
962
963impl WriterContentionDiagnostics {
964 fn snapshot(
965 pool: &ConnectionPool,
966 audit_append_failures: Option<u64>,
967 runtime_audit_batch_metrics: Option<RuntimeAuditBatchMetrics>,
968 ) -> Self {
969 let writer = pool.writer_acquisition_snapshot();
970 let unavailable_reason =
971 || Some("no audit-batch control is registered with this runtime instance".to_string());
972 Self {
973 writer_acquisitions: writer.acquisitions,
974 pooled_writer_acquisitions: writer.pooled_acquisitions,
975 standalone_writer_acquisitions: writer.standalone_acquisitions,
976 writer_task_acquisitions: writer.writer_task_acquisitions,
977 writer_acquisition_timeouts: writer.timeouts,
978 writer_task_begin_busy: writer.writer_task_begin_busy,
979 writer_task_begin_busy_absorbed: writer.writer_task_begin_busy_absorbed,
980 writer_task_begin_errors: writer.writer_task_begin_errors,
981 writer_task_request_failures: writer.writer_task_request_failures,
982 writer_task_side_effects_unknown: writer.writer_task_side_effects_unknown,
983 audit_append_failures,
984 audit_append_failures_unavailable_reason: audit_append_failures.is_none().then(|| {
985 "runtime audit instrumentation was not supplied to khive-db diagnostics".to_string()
986 }),
987 audit_obligation_append_failures: None,
988 audit_obligation_append_failures_unavailable_reason: Some(
989 "runtime obligation audit instrumentation was not supplied to khive-db diagnostics"
990 .to_string(),
991 ),
992 audit_batch_flush_failures: runtime_audit_batch_metrics.map(|m| m.flush_failures),
993 audit_batch_flush_failures_unavailable_reason: runtime_audit_batch_metrics
994 .is_none()
995 .then(unavailable_reason)
996 .flatten(),
997 audit_degraded_rows: runtime_audit_batch_metrics.map(|m| m.degraded_rows),
998 audit_degraded_rows_unavailable_reason: runtime_audit_batch_metrics
999 .is_none()
1000 .then(unavailable_reason)
1001 .flatten(),
1002 audit_degraded: runtime_audit_batch_metrics.map(|m| m.degraded),
1003 audit_degraded_unavailable_reason: runtime_audit_batch_metrics
1004 .is_none()
1005 .then(unavailable_reason)
1006 .flatten(),
1007 audit_admission_refused_obligations: runtime_audit_batch_metrics
1008 .map(|m| m.admission_refused_obligations),
1009 audit_admission_refused_obligations_last_at_ms: runtime_audit_batch_metrics
1010 .and_then(|m| m.admission_refused_obligations_last_at_ms),
1011 audit_admission_refused_obligations_unavailable_reason: runtime_audit_batch_metrics
1012 .is_none()
1013 .then(unavailable_reason)
1014 .flatten(),
1015 audit_admission_unresolved_obligations: runtime_audit_batch_metrics
1016 .map(|m| m.admission_unresolved_obligations),
1017 audit_admission_unresolved_obligations_last_at_ms: runtime_audit_batch_metrics
1018 .and_then(|m| m.admission_unresolved_obligations_last_at_ms),
1019 audit_admission_unresolved_obligations_unavailable_reason: runtime_audit_batch_metrics
1020 .is_none()
1021 .then(unavailable_reason)
1022 .flatten(),
1023 }
1024 }
1025}
1026
1027#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1043pub struct GraphEdgeIntegrity {
1044 pub duplicate_edge_id_groups: i64,
1045 pub graph_edges_rows: i64,
1046 pub graph_edges_seq_rows: i64,
1047 pub pre_v14_duplicate_edge_state_detected: bool,
1048 pub live_entities_carrying_merged_into: i64,
1049}
1050
1051fn graph_edge_integrity(conn: &Connection) -> rusqlite::Result<GraphEdgeIntegrity> {
1052 conn.query_row(
1053 "SELECT
1054 (SELECT COUNT(*) FROM (
1055 SELECT id FROM graph_edges GROUP BY id HAVING COUNT(*) > 1
1056 )),
1057 (SELECT COUNT(*) FROM graph_edges),
1058 (SELECT COUNT(*) FROM graph_edges_seq),
1059 (SELECT COUNT(*) FROM entities
1060 WHERE deleted_at IS NULL AND merged_into IS NOT NULL)",
1061 [],
1062 |row| {
1063 let duplicate_edge_id_groups = row.get(0)?;
1064 Ok(GraphEdgeIntegrity {
1065 duplicate_edge_id_groups,
1066 graph_edges_rows: row.get(1)?,
1067 graph_edges_seq_rows: row.get(2)?,
1068 pre_v14_duplicate_edge_state_detected: duplicate_edge_id_groups > 0,
1069 live_entities_carrying_merged_into: row.get(3)?,
1070 })
1071 },
1072 )
1073}
1074
1075const MAX_DATABASE_SIZE_OBJECTS: usize = 4_096;
1076
1077#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1079#[serde(rename_all = "snake_case")]
1080pub enum DatabaseObjectKind {
1081 Table,
1082 Index,
1083 Internal,
1084}
1085
1086#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1091#[serde(rename_all = "snake_case")]
1092pub enum DatabaseStorageClass {
1093 RowTable,
1094 Index,
1095 FullText,
1096 Vector,
1097 MixedRowAndEmbedding,
1098 Internal,
1099}
1100
1101#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
1102pub struct DatabaseObjectSize {
1103 pub name: String,
1104 pub owner_table: Option<String>,
1105 pub object_kind: DatabaseObjectKind,
1106 pub storage_class: DatabaseStorageClass,
1107 pub pages: u64,
1108 pub bytes: u64,
1109}
1110
1111#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
1114pub struct DatabaseSizeComposition {
1115 pub page_size_bytes: u64,
1116 pub page_count: u64,
1117 pub freelist_pages: u64,
1118 pub database_bytes: u64,
1119 pub freelist_bytes: u64,
1120 pub accounted_bytes: u64,
1121 pub unaccounted_bytes: u64,
1122 pub row_table_bytes: u64,
1123 pub index_bytes: u64,
1124 pub full_text_bytes: u64,
1125 pub vector_bytes: u64,
1126 pub mixed_embedding_bytes: u64,
1127 pub internal_bytes: u64,
1128 pub objects: Vec<DatabaseObjectSize>,
1129 pub objects_truncated: bool,
1130 pub objects_omitted: usize,
1131}
1132
1133fn nonnegative_sqlite_integer(column: usize, value: i64) -> rusqlite::Result<u64> {
1134 u64::try_from(value).map_err(|_| rusqlite::Error::IntegralValueOutOfRange(column, value))
1135}
1136
1137fn declares_embedding_blob(sql: Option<&str>) -> bool {
1138 let Some(sql) = sql else {
1139 return false;
1140 };
1141 let tokens: Vec<_> = sql
1142 .split(|character: char| !(character.is_ascii_alphanumeric() || character == '_'))
1143 .filter(|token| !token.is_empty())
1144 .collect();
1145 tokens.windows(2).any(|pair| {
1146 pair[0].eq_ignore_ascii_case("embedding") && pair[1].eq_ignore_ascii_case("blob")
1147 })
1148}
1149
1150fn classify_database_object(
1151 name: &str,
1152 sqlite_type: &str,
1153 sql: Option<&str>,
1154) -> (DatabaseObjectKind, DatabaseStorageClass) {
1155 let object_kind = match sqlite_type {
1156 "table" => DatabaseObjectKind::Table,
1157 "index" => DatabaseObjectKind::Index,
1158 _ => DatabaseObjectKind::Internal,
1159 };
1160 let lower_name = name.to_ascii_lowercase();
1161 let lower_sql = sql.unwrap_or_default().to_ascii_lowercase();
1162 let storage_class = if lower_name.starts_with("fts_") || lower_sql.contains("using fts5") {
1163 DatabaseStorageClass::FullText
1164 } else if lower_name.starts_with("vec_")
1165 || lower_name == "_embedding_models"
1166 || lower_sql.contains("using vec0")
1167 {
1168 DatabaseStorageClass::Vector
1169 } else if declares_embedding_blob(sql) {
1170 DatabaseStorageClass::MixedRowAndEmbedding
1171 } else if object_kind == DatabaseObjectKind::Index {
1172 DatabaseStorageClass::Index
1173 } else if name.starts_with("sqlite_") || object_kind == DatabaseObjectKind::Internal {
1174 DatabaseStorageClass::Internal
1175 } else {
1176 DatabaseStorageClass::RowTable
1177 };
1178 (object_kind, storage_class)
1179}
1180
1181fn database_size_composition(conn: &Connection) -> rusqlite::Result<DatabaseSizeComposition> {
1182 let page_size = nonnegative_sqlite_integer(
1183 0,
1184 conn.query_row("PRAGMA page_size", [], |row| row.get::<_, i64>(0))?,
1185 )?;
1186 let page_count = nonnegative_sqlite_integer(
1187 0,
1188 conn.query_row("PRAGMA page_count", [], |row| row.get::<_, i64>(0))?,
1189 )?;
1190 let freelist_pages = nonnegative_sqlite_integer(
1191 0,
1192 conn.query_row("PRAGMA freelist_count", [], |row| row.get::<_, i64>(0))?,
1193 )?;
1194
1195 let mut statement = conn.prepare(
1196 "SELECT d.name, COALESCE(s.type, 'internal'), s.tbl_name, s.sql, d.pageno, d.pgsize
1197 FROM dbstat AS d
1198 LEFT JOIN sqlite_schema AS s ON s.name = d.name
1199 WHERE d.aggregate = TRUE
1200 ORDER BY d.name",
1201 )?;
1202 let mut rows = statement.query([])?;
1203 let mut objects = Vec::new();
1204 let mut objects_omitted = 0usize;
1205 let mut accounted_bytes = 0u64;
1206 let mut row_table_bytes = 0u64;
1207 let mut index_bytes = 0u64;
1208 let mut full_text_bytes = 0u64;
1209 let mut vector_bytes = 0u64;
1210 let mut mixed_embedding_bytes = 0u64;
1211 let mut internal_bytes = 0u64;
1212
1213 while let Some(row) = rows.next()? {
1214 let name: String = row.get(0)?;
1215 let sqlite_type: String = row.get(1)?;
1216 let owner_table: Option<String> = row.get(2)?;
1217 let sql: Option<String> = row.get(3)?;
1218 let pages = nonnegative_sqlite_integer(4, row.get(4)?)?;
1219 let bytes = nonnegative_sqlite_integer(5, row.get(5)?)?;
1220 let (object_kind, storage_class) =
1221 classify_database_object(&name, &sqlite_type, sql.as_deref());
1222 accounted_bytes = accounted_bytes.saturating_add(bytes);
1223 let class_total = match storage_class {
1224 DatabaseStorageClass::RowTable => &mut row_table_bytes,
1225 DatabaseStorageClass::Index => &mut index_bytes,
1226 DatabaseStorageClass::FullText => &mut full_text_bytes,
1227 DatabaseStorageClass::Vector => &mut vector_bytes,
1228 DatabaseStorageClass::MixedRowAndEmbedding => &mut mixed_embedding_bytes,
1229 DatabaseStorageClass::Internal => &mut internal_bytes,
1230 };
1231 *class_total = class_total.saturating_add(bytes);
1232
1233 if objects.len() < MAX_DATABASE_SIZE_OBJECTS {
1234 objects.push(DatabaseObjectSize {
1235 name,
1236 owner_table,
1237 object_kind,
1238 storage_class,
1239 pages,
1240 bytes,
1241 });
1242 } else {
1243 objects_omitted = objects_omitted.saturating_add(1);
1244 }
1245 }
1246
1247 let database_bytes = page_count.saturating_mul(page_size);
1248 let freelist_bytes = freelist_pages.saturating_mul(page_size);
1249 let unaccounted_bytes = database_bytes
1250 .saturating_sub(freelist_bytes)
1251 .saturating_sub(accounted_bytes);
1252 Ok(DatabaseSizeComposition {
1253 page_size_bytes: page_size,
1254 page_count,
1255 freelist_pages,
1256 database_bytes,
1257 freelist_bytes,
1258 accounted_bytes,
1259 unaccounted_bytes,
1260 row_table_bytes,
1261 index_bytes,
1262 full_text_bytes,
1263 vector_bytes,
1264 mixed_embedding_bytes,
1265 internal_bytes,
1266 objects,
1267 objects_truncated: objects_omitted > 0,
1268 objects_omitted,
1269 })
1270}
1271
1272#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1285pub struct CollectionCost {
1286 pub total_ms: u64,
1288 pub sqlite_ms: u64,
1291 pub wal_file_stat_ms: u64,
1294 pub wal_pin_census_ms: u64,
1296 pub wal_pin_sidecar_ms: u64,
1298 pub wal_pin_census_budget_ms: Option<u64>,
1303 pub wal_pin_census_budget_exhausted: bool,
1307}
1308
1309impl CollectionCost {
1310 fn in_memory(total_ms: u64) -> Self {
1314 Self {
1315 total_ms,
1316 sqlite_ms: 0,
1317 wal_file_stat_ms: 0,
1318 wal_pin_census_ms: 0,
1319 wal_pin_sidecar_ms: 0,
1320 wal_pin_census_budget_ms: None,
1321 wal_pin_census_budget_exhausted: false,
1322 }
1323 }
1324}
1325
1326const DEFAULT_CENSUS_BUDGET: Duration = Duration::from_millis(2000);
1333
1334const CENSUS_BUDGET_ENV: &str = "KHIVE_WALPIN_CENSUS_BUDGET_MS";
1344
1345fn request_census_budget() -> Option<Duration> {
1346 match std::env::var(CENSUS_BUDGET_ENV) {
1347 Ok(raw) => match raw.trim().parse::<u64>() {
1348 Ok(0) => None,
1349 Ok(ms) => Some(Duration::from_millis(ms)),
1350 Err(_) => Some(DEFAULT_CENSUS_BUDGET),
1351 },
1352 Err(_) => Some(DEFAULT_CENSUS_BUDGET),
1353 }
1354}
1355
1356#[derive(Debug, Clone, PartialEq, Serialize)]
1359pub struct DbDiagnostics {
1360 pub build: BuildIdentity,
1361 pub process: ProcessIdentity,
1362 pub db_path: Option<String>,
1365 pub wal_file: Option<WalFileState>,
1366 pub checkpoint_counters: CheckpointCounters,
1367 pub checkpoint_probe: Option<CheckpointProbe>,
1368 pub checkpoint_probe_error: Option<String>,
1369 pub reader_contention: ReaderContentionDiagnostics,
1371 pub writer_contention: WriterContentionDiagnostics,
1373 pub size_composition: Option<DatabaseSizeComposition>,
1374 pub size_composition_error: Option<String>,
1375 pub graph_edge_integrity: Option<GraphEdgeIntegrity>,
1376 pub graph_edge_integrity_error: Option<String>,
1377 pub fts_segments: Option<crate::FtsSegmentDiagnostics>,
1381 pub fts_segments_error: Option<String>,
1382 pub fts_maintenance: crate::FtsMaintenanceCounters,
1385 pub wal_pin: WalPinAttribution,
1386 pub collection_cost: CollectionCost,
1388}
1389
1390pub fn collect(
1410 pool: &ConnectionPool,
1411 build: BuildIdentity,
1412 sweep_interval: Duration,
1413) -> DbDiagnostics {
1414 collect_inner(pool, build, sweep_interval, None, None)
1415}
1416
1417pub fn collect_with_audit_append_failures(
1420 pool: &ConnectionPool,
1421 build: BuildIdentity,
1422 sweep_interval: Duration,
1423 audit_append_failures: u64,
1424) -> DbDiagnostics {
1425 collect_inner(
1426 pool,
1427 build,
1428 sweep_interval,
1429 Some(audit_append_failures),
1430 None,
1431 )
1432}
1433
1434pub async fn collect_with_audit_append_failures_interruptibly(
1444 pool: Arc<ConnectionPool>,
1445 build: BuildIdentity,
1446 sweep_interval: Duration,
1447 audit_append_failures: u64,
1448) -> StorageResult<DbDiagnostics> {
1449 collect_with_runtime_audit_metrics_interruptibly(
1450 pool,
1451 build,
1452 sweep_interval,
1453 audit_append_failures,
1454 None,
1455 )
1456 .await
1457}
1458
1459pub async fn collect_with_runtime_audit_metrics_interruptibly(
1466 pool: Arc<ConnectionPool>,
1467 build: BuildIdentity,
1468 sweep_interval: Duration,
1469 audit_append_failures: u64,
1470 runtime_audit_batch_metrics: Option<RuntimeAuditBatchMetrics>,
1471) -> StorageResult<DbDiagnostics> {
1472 crate::ensure_request_read_active("db_diagnostics")?;
1473 let started = Instant::now();
1474 let process = ProcessIdentity::current(&pool);
1475 let counters = checkpoint_counters();
1476 let reader_contention = ReaderContentionDiagnostics::snapshot(&pool);
1477 let writer_contention = WriterContentionDiagnostics::snapshot(
1478 &pool,
1479 Some(audit_append_failures),
1480 runtime_audit_batch_metrics,
1481 );
1482
1483 let Some(path) = pool.config().path.clone() else {
1484 crate::ensure_request_read_active("db_diagnostics")?;
1485 return Ok(DbDiagnostics {
1486 build,
1487 process,
1488 db_path: None,
1489 wal_file: None,
1490 checkpoint_counters: counters,
1491 checkpoint_probe: None,
1492 checkpoint_probe_error: Some(
1493 "in-memory database: no WAL file and no checkpoint to probe".to_string(),
1494 ),
1495 reader_contention,
1496 writer_contention,
1497 size_composition: None,
1498 size_composition_error: Some(
1499 "in-memory database: no file-backed page composition to inspect".to_string(),
1500 ),
1501 graph_edge_integrity: None,
1502 graph_edge_integrity_error: Some(
1503 "in-memory database: no durable graph-edge ledger to inspect".to_string(),
1504 ),
1505 fts_segments: None,
1506 fts_segments_error: Some(
1507 "in-memory database: no durable FTS5 indexes to inspect".to_string(),
1508 ),
1509 fts_maintenance: crate::fts_maintenance_counters(),
1510 wal_pin: WalPinAttribution::unavailable(
1511 "in-memory database: no file for the OS holder census",
1512 ),
1513 collection_cost: CollectionCost::in_memory(elapsed_ms(started)),
1514 });
1515 };
1516
1517 let inspection_pool = Arc::clone(&pool);
1518 let sqlite_started = Instant::now();
1519 let inspection = crate::read_cancellation::run_interruptible_read(
1520 StorageCapability::Sql,
1521 "db_diagnostics.sqlite",
1522 move |scope| inspect_pool_interruptibly(&inspection_pool, scope),
1523 )
1524 .await?;
1525 let sqlite_ms = elapsed_ms(sqlite_started);
1526 crate::ensure_request_read_active("db_diagnostics")?;
1527 let canonical = operational_db_path(&pool, &path);
1528 let budget = request_census_budget();
1529 let (wal_file, wal_pin, file_state_cost) =
1530 inspect_file_state_interruptibly(canonical, sweep_interval, budget).await?;
1531 crate::ensure_request_read_active("db_diagnostics")?;
1532
1533 Ok(DbDiagnostics {
1534 build,
1535 process,
1536 db_path: Some(path.display().to_string()),
1537 wal_file: Some(wal_file),
1538 checkpoint_counters: counters,
1539 checkpoint_probe: inspection.checkpoint_probe,
1540 checkpoint_probe_error: inspection.checkpoint_probe_error,
1541 reader_contention,
1542 writer_contention,
1543 size_composition: inspection.size_composition,
1544 size_composition_error: inspection.size_composition_error,
1545 graph_edge_integrity: inspection.graph_edge_integrity,
1546 graph_edge_integrity_error: inspection.graph_edge_integrity_error,
1547 fts_segments: inspection.fts_segments,
1548 fts_segments_error: inspection.fts_segments_error,
1549 fts_maintenance: crate::fts_maintenance_counters(),
1550 wal_pin,
1551 collection_cost: CollectionCost {
1552 total_ms: elapsed_ms(started),
1553 sqlite_ms,
1554 wal_file_stat_ms: file_state_cost.wal_file_stat_ms,
1555 wal_pin_census_ms: file_state_cost.census_ms,
1556 wal_pin_sidecar_ms: file_state_cost.sidecar_ms,
1557 wal_pin_census_budget_ms: budget.map(|b| b.as_millis() as u64),
1558 wal_pin_census_budget_exhausted: file_state_cost.census_budget_exhausted,
1559 },
1560 })
1561}
1562
1563fn operational_db_path(pool: &ConnectionPool, configured: &Path) -> PathBuf {
1575 pool.canonical_path()
1576 .map(Path::to_path_buf)
1577 .unwrap_or_else(|| configured.to_path_buf())
1578}
1579
1580fn collect_inner(
1581 pool: &ConnectionPool,
1582 build: BuildIdentity,
1583 sweep_interval: Duration,
1584 audit_append_failures: Option<u64>,
1585 runtime_audit_batch_metrics: Option<RuntimeAuditBatchMetrics>,
1586) -> DbDiagnostics {
1587 let started = Instant::now();
1588 let process = ProcessIdentity::current(pool);
1589 let counters = checkpoint_counters();
1590 let reader_contention = ReaderContentionDiagnostics::snapshot(pool);
1591 let writer_contention = WriterContentionDiagnostics::snapshot(
1592 pool,
1593 audit_append_failures,
1594 runtime_audit_batch_metrics,
1595 );
1596
1597 let Some(path) = pool.config().path.clone() else {
1598 return DbDiagnostics {
1599 build,
1600 process,
1601 db_path: None,
1602 wal_file: None,
1603 checkpoint_counters: counters,
1604 checkpoint_probe: None,
1605 checkpoint_probe_error: Some(
1606 "in-memory database: no WAL file and no checkpoint to probe".to_string(),
1607 ),
1608 reader_contention,
1609 writer_contention,
1610 size_composition: None,
1611 size_composition_error: Some(
1612 "in-memory database: no file-backed page composition to inspect".to_string(),
1613 ),
1614 graph_edge_integrity: None,
1615 graph_edge_integrity_error: Some(
1616 "in-memory database: no durable graph-edge ledger to inspect".to_string(),
1617 ),
1618 fts_segments: None,
1619 fts_segments_error: Some(
1620 "in-memory database: no durable FTS5 indexes to inspect".to_string(),
1621 ),
1622 fts_maintenance: crate::fts_maintenance_counters(),
1623 wal_pin: WalPinAttribution::unavailable(
1624 "in-memory database: no file for the OS holder census",
1625 ),
1626 collection_cost: CollectionCost::in_memory(elapsed_ms(started)),
1627 };
1628 };
1629
1630 let sqlite_started = Instant::now();
1631 let inspection = inspect_pool(pool);
1632 let sqlite_ms = elapsed_ms(sqlite_started);
1633 let canonical = operational_db_path(pool, &path);
1634 let wal_file_started = Instant::now();
1635 let wal_file = wal_file_state(&canonical);
1636 let wal_file_stat_ms = elapsed_ms(wal_file_started);
1637 let census_started = Instant::now();
1638 let wal_pin = wal_pin_attribution(&canonical, sweep_interval);
1639 let wal_pin_ms = elapsed_ms(census_started);
1640
1641 DbDiagnostics {
1642 build,
1643 process,
1644 db_path: Some(path.display().to_string()),
1645 wal_file: Some(wal_file),
1646 checkpoint_counters: counters,
1647 checkpoint_probe: inspection.checkpoint_probe,
1648 checkpoint_probe_error: inspection.checkpoint_probe_error,
1649 reader_contention,
1650 writer_contention,
1651 size_composition: inspection.size_composition,
1652 size_composition_error: inspection.size_composition_error,
1653 graph_edge_integrity: inspection.graph_edge_integrity,
1654 graph_edge_integrity_error: inspection.graph_edge_integrity_error,
1655 fts_segments: inspection.fts_segments,
1656 fts_segments_error: inspection.fts_segments_error,
1657 fts_maintenance: crate::fts_maintenance_counters(),
1658 wal_pin,
1659 collection_cost: CollectionCost {
1660 total_ms: elapsed_ms(started),
1661 sqlite_ms,
1662 wal_file_stat_ms,
1663 wal_pin_census_ms: wal_pin_ms,
1669 wal_pin_sidecar_ms: 0,
1670 wal_pin_census_budget_ms: None,
1671 wal_pin_census_budget_exhausted: false,
1672 },
1673 }
1674}
1675
1676struct PoolInspection {
1677 checkpoint_probe: Option<CheckpointProbe>,
1678 checkpoint_probe_error: Option<String>,
1679 size_composition: Option<DatabaseSizeComposition>,
1680 size_composition_error: Option<String>,
1681 graph_edge_integrity: Option<GraphEdgeIntegrity>,
1682 graph_edge_integrity_error: Option<String>,
1683 fts_segments: Option<crate::FtsSegmentDiagnostics>,
1684 fts_segments_error: Option<String>,
1685}
1686
1687fn split_fts_segments_result(
1691 result: Result<crate::FtsSegmentDiagnostics, String>,
1692) -> (Option<crate::FtsSegmentDiagnostics>, Option<String>) {
1693 match result {
1694 Ok(segments) => (Some(segments), None),
1695 Err(error) => (None, Some(error)),
1696 }
1697}
1698
1699fn inspect_pool_interruptibly(
1700 pool: &ConnectionPool,
1701 scope: &crate::read_cancellation::InterruptibleReadScope,
1702) -> StorageResult<PoolInspection> {
1703 scope.ensure_active()?;
1704 let conn = match pool.open_standalone_writer_untracked() {
1705 Ok(conn) => conn,
1706 Err(e) => {
1707 scope.ensure_active()?;
1708 let reason = format!("guarded standalone open refused: {e}");
1709 return Ok(PoolInspection {
1710 checkpoint_probe: None,
1711 checkpoint_probe_error: Some(reason.clone()),
1712 size_composition: None,
1713 size_composition_error: Some(reason.clone()),
1714 graph_edge_integrity: None,
1715 graph_edge_integrity_error: Some(reason.clone()),
1716 fts_segments: None,
1717 fts_segments_error: Some(reason),
1718 });
1719 }
1720 };
1721 scope.ensure_active()?;
1726
1727 let (checkpoint_probe, checkpoint_probe_error) = match checkpoint_probe(&conn) {
1729 Ok(probe) => (Some(probe), None),
1730 Err(e) => (
1731 None,
1732 Some(format!("PRAGMA wal_checkpoint(PASSIVE) failed: {e}")),
1733 ),
1734 };
1735 #[cfg(test)]
1736 if TEST_PAUSE_AFTER_PASSIVE.load(Ordering::SeqCst) {
1737 TEST_REACHED_AFTER_PASSIVE.store(true, Ordering::SeqCst);
1738 while TEST_PAUSE_AFTER_PASSIVE.load(Ordering::SeqCst) && !scope.should_stop() {
1739 std::thread::yield_now();
1740 }
1741 }
1742 scope.ensure_active()?;
1743
1744 let (integrity, size_composition, fts_segments) = scope.run(&conn, || {
1749 Ok((
1750 graph_edge_integrity(&conn),
1751 database_size_composition(&conn),
1752 crate::fts_maintenance::inspect_fts_segments(&conn),
1753 ))
1754 })?;
1755 let (graph_edge_integrity, graph_edge_integrity_error) = match integrity {
1756 Ok(integrity) => (Some(integrity), None),
1757 Err(e) => (
1758 None,
1759 Some(format!("graph-edge integrity query failed: {e}")),
1760 ),
1761 };
1762 let (fts_segments, fts_segments_error) = split_fts_segments_result(fts_segments);
1763
1764 let (size_composition, size_composition_error) = match size_composition {
1765 Ok(composition) => (Some(composition), None),
1766 Err(error) => (
1767 None,
1768 Some(format!("database size composition query failed: {error}")),
1769 ),
1770 };
1771
1772 Ok(PoolInspection {
1773 checkpoint_probe,
1774 checkpoint_probe_error,
1775 size_composition,
1776 size_composition_error,
1777 graph_edge_integrity,
1778 graph_edge_integrity_error,
1779 fts_segments,
1780 fts_segments_error,
1781 })
1782}
1783
1784#[cfg(test)]
1785static TEST_PAUSE_AFTER_PASSIVE: AtomicBool = AtomicBool::new(false);
1786#[cfg(test)]
1787static TEST_REACHED_AFTER_PASSIVE: AtomicBool = AtomicBool::new(false);
1788
1789struct StopCensusOnDrop {
1790 stopped: Arc<AtomicBool>,
1791 armed: bool,
1792}
1793
1794impl Drop for StopCensusOnDrop {
1795 fn drop(&mut self) {
1796 if self.armed {
1797 self.stopped.store(true, Ordering::SeqCst);
1798 }
1799 }
1800}
1801
1802#[derive(Debug, Clone, Copy, Default)]
1805struct FileStateCost {
1806 wal_file_stat_ms: u64,
1807 census_ms: u64,
1808 sidecar_ms: u64,
1809 census_budget_exhausted: bool,
1810}
1811
1812fn elapsed_ms(since: Instant) -> u64 {
1813 since.elapsed().as_millis() as u64
1814}
1815
1816async fn inspect_file_state_interruptibly(
1817 path: PathBuf,
1818 sweep_interval: Duration,
1819 census_budget: Option<Duration>,
1820) -> StorageResult<(WalFileState, WalPinAttribution, FileStateCost)> {
1821 const OPERATION: &str = "db_diagnostics.wal_holder_census";
1822 crate::ensure_request_read_active(OPERATION)?;
1823 let stopped = Arc::new(AtomicBool::new(false));
1824 let worker_stopped = Arc::clone(&stopped);
1825 let mut stop_on_drop = StopCensusOnDrop {
1826 stopped: Arc::clone(&stopped),
1827 armed: true,
1828 };
1829 let mut worker = tokio::task::spawn_blocking(move || {
1830 let mut cost = FileStateCost::default();
1831 let wal_file_started = Instant::now();
1832 let wal_file = wal_file_state(&path);
1833 cost.wal_file_stat_ms = elapsed_ms(wal_file_started);
1834 if worker_stopped.load(Ordering::SeqCst) {
1835 return Err(std::io::Error::new(
1836 std::io::ErrorKind::Interrupted,
1837 "WAL holder census cancelled",
1838 ));
1839 }
1840 #[cfg(unix)]
1841 let census_started = Instant::now();
1842 #[cfg(unix)]
1843 let census_result = match census_budget {
1844 Some(budget) => crate::walpin::census_holders_until_within(
1845 &path,
1846 || worker_stopped.load(Ordering::SeqCst),
1847 budget,
1848 ),
1849 None => {
1850 crate::walpin::census_holders_until(&path, || worker_stopped.load(Ordering::SeqCst))
1851 }
1852 };
1853 #[cfg(unix)]
1854 {
1855 cost.census_ms = elapsed_ms(census_started);
1856 }
1857 #[cfg(unix)]
1858 let attribution = match census_result {
1859 Ok(census) => {
1860 cost.census_budget_exhausted = census.budget_exhausted;
1861 if worker_stopped.load(Ordering::SeqCst) {
1862 return Err(std::io::Error::new(
1863 std::io::ErrorKind::Interrupted,
1864 "WAL sidecar inspection cancelled",
1865 ));
1866 }
1867 if !crate::walpin::sidecar_enabled(true) {
1868 wal_pin_attribution_without_sidecar(census, SIDECAR_DISABLED_REASON.to_string())
1869 } else {
1870 let sidecar_started = Instant::now();
1871 let sidecar = crate::walpin::inspect_live(
1872 &crate::walpin::sidecar_dir_for(&path),
1873 sweep_interval,
1874 );
1875 cost.sidecar_ms = elapsed_ms(sidecar_started);
1876 if worker_stopped.load(Ordering::SeqCst) {
1877 return Err(std::io::Error::new(
1878 std::io::ErrorKind::Interrupted,
1879 "WAL sidecar inspection cancelled",
1880 ));
1881 }
1882 match sidecar {
1883 Ok(sidecar) => wal_pin_attribution_from_evidence(census, sidecar),
1884 Err(error) => wal_pin_attribution_without_sidecar(
1885 census,
1886 format!("read-only sidecar enumeration failed: {error}"),
1887 ),
1888 }
1889 }
1890 }
1891 Err(error) if error.kind() == std::io::ErrorKind::Interrupted => return Err(error),
1892 Err(error) => WalPinAttribution::unavailable(format!("census_holders failed: {error}")),
1893 };
1894 #[cfg(not(unix))]
1895 let attribution = {
1896 let _ = census_budget;
1897 let census_started = Instant::now();
1898 let attribution = wal_pin_attribution(&path, sweep_interval);
1899 cost.census_ms = elapsed_ms(census_started);
1900 attribution
1901 };
1902 Ok((wal_file, attribution, cost))
1903 });
1904
1905 tokio::select! {
1906 joined = &mut worker => {
1907 stop_on_drop.armed = false;
1908 let result = joined
1909 .map_err(|error| StorageError::driver(StorageCapability::Sql, OPERATION, error))?
1910 .map_err(|error| StorageError::driver(StorageCapability::Sql, OPERATION, error))?;
1911 crate::ensure_request_read_active(OPERATION)?;
1912 Ok(result)
1913 }
1914 _ = crate::wait_for_request_read_cancellation() => {
1915 stopped.store(true, Ordering::SeqCst);
1916 if tokio::time::timeout(crate::sqlite_interrupt_grace_from_env(), &mut worker)
1917 .await
1918 .is_err()
1919 {
1920 worker.abort();
1921 }
1922 stop_on_drop.armed = false;
1923 Err(StorageError::Timeout { operation: OPERATION.into() })
1924 }
1925 }
1926}
1927
1928fn inspect_pool(pool: &ConnectionPool) -> PoolInspection {
1935 let conn = match pool.open_standalone_writer_untracked() {
1936 Ok(conn) => conn,
1937 Err(e) => {
1938 let reason = format!("guarded standalone open refused: {e}");
1939 return PoolInspection {
1940 checkpoint_probe: None,
1941 checkpoint_probe_error: Some(reason.clone()),
1942 size_composition: None,
1943 size_composition_error: Some(reason.clone()),
1944 graph_edge_integrity: None,
1945 graph_edge_integrity_error: Some(reason.clone()),
1946 fts_segments: None,
1947 fts_segments_error: Some(reason),
1948 };
1949 }
1950 };
1951
1952 let (checkpoint_probe, checkpoint_probe_error) = match checkpoint_probe(&conn) {
1953 Ok(probe) => (Some(probe), None),
1954 Err(e) => (
1955 None,
1956 Some(format!("PRAGMA wal_checkpoint(PASSIVE) failed: {e}")),
1957 ),
1958 };
1959 let (graph_edge_integrity, graph_edge_integrity_error) = match graph_edge_integrity(&conn) {
1960 Ok(integrity) => (Some(integrity), None),
1961 Err(e) => (
1962 None,
1963 Some(format!("graph-edge integrity query failed: {e}")),
1964 ),
1965 };
1966 let (size_composition, size_composition_error) = match database_size_composition(&conn) {
1967 Ok(composition) => (Some(composition), None),
1968 Err(error) => (
1969 None,
1970 Some(format!("database size composition query failed: {error}")),
1971 ),
1972 };
1973 let (fts_segments, fts_segments_error) =
1974 split_fts_segments_result(crate::fts_maintenance::inspect_fts_segments(&conn));
1975
1976 PoolInspection {
1977 checkpoint_probe,
1978 checkpoint_probe_error,
1979 size_composition,
1980 size_composition_error,
1981 graph_edge_integrity,
1982 graph_edge_integrity_error,
1983 fts_segments,
1984 fts_segments_error,
1985 }
1986}
1987
1988#[cfg(test)]
1989mod tests {
1990 use serial_test::serial;
1991
1992 use super::*;
1993 use crate::pool::{ConnectionPool, PoolConfig};
1994
1995 #[test]
2000 #[serial_test::serial(khive_walpin_census_budget_env)]
2001 fn census_budget_reads_zero_as_unbounded_and_survives_a_malformed_value() {
2002 let _guard = crate::walpin::EnvVarGuard::capture(CENSUS_BUDGET_ENV);
2003
2004 std::env::remove_var(CENSUS_BUDGET_ENV);
2005 assert_eq!(
2006 request_census_budget(),
2007 Some(DEFAULT_CENSUS_BUDGET),
2008 "an unset variable takes the default bound"
2009 );
2010
2011 std::env::set_var(CENSUS_BUDGET_ENV, "0");
2012 assert_eq!(
2013 request_census_budget(),
2014 None,
2015 "0 restores the unbounded full-machine walk"
2016 );
2017
2018 std::env::set_var(CENSUS_BUDGET_ENV, " 750 ");
2019 assert_eq!(
2020 request_census_budget(),
2021 Some(Duration::from_millis(750)),
2022 "a surrounding-whitespace value is still a number"
2023 );
2024
2025 std::env::set_var(CENSUS_BUDGET_ENV, "soon");
2026 assert_eq!(
2027 request_census_budget(),
2028 Some(DEFAULT_CENSUS_BUDGET),
2029 "a malformed budget must not fail the request; the report states \
2030 which budget was actually used"
2031 );
2032 }
2033
2034 #[test]
2038 fn a_budget_stop_and_an_enumeration_failure_do_not_share_a_reason() {
2039 let budget = census_truncation_cause(true);
2040 let failure = census_truncation_cause(false);
2041 assert_ne!(budget, failure);
2042 assert!(
2043 budget.contains("budget"),
2044 "the budget reason must name the budget: {budget}"
2045 );
2046 assert!(
2047 budget.contains("wal_pin_census_budget_ms"),
2048 "and must point at the field carrying the value: {budget}"
2049 );
2050 assert!(
2051 !failure.contains("budget"),
2052 "an enumeration failure must not be described as a budget stop: {failure}"
2053 );
2054 }
2055
2056 #[test]
2060 fn collection_cost_distinguishes_an_unbounded_census_from_a_zero_cost_one() {
2061 let unbounded = CollectionCost {
2062 total_ms: 9,
2063 sqlite_ms: 4,
2064 wal_file_stat_ms: 0,
2065 wal_pin_census_ms: 5,
2066 wal_pin_sidecar_ms: 0,
2067 wal_pin_census_budget_ms: None,
2068 wal_pin_census_budget_exhausted: false,
2069 };
2070 let bounded = CollectionCost {
2071 wal_pin_census_budget_ms: Some(2000),
2072 wal_pin_census_budget_exhausted: true,
2073 ..unbounded
2074 };
2075
2076 let unbounded = serde_json::to_value(unbounded).unwrap();
2077 let bounded = serde_json::to_value(bounded).unwrap();
2078 assert_eq!(
2079 unbounded["wal_pin_census_budget_ms"],
2080 serde_json::Value::Null
2081 );
2082 assert_eq!(bounded["wal_pin_census_budget_ms"], 2000);
2083 assert_eq!(unbounded["wal_pin_census_budget_exhausted"], false);
2084 assert_eq!(bounded["wal_pin_census_budget_exhausted"], true);
2085 }
2086
2087 #[test]
2088 fn process_identity_serializes_os_start_time_or_explicit_unavailability() {
2089 let known = ProcessIdentity::from_start_time(42, Some(1_000_000_000), 3);
2090 assert_eq!(
2091 serde_json::to_value(known).unwrap(),
2092 serde_json::json!({
2093 "pid": 42,
2094 "started_at": 1_000_000_000,
2095 "started_at_unavailable_reason": null,
2096 "pool_generation": 3
2097 }),
2098 "preserve the OS timestamp, not request time or a derived uptime"
2099 );
2100
2101 let unknown = ProcessIdentity::from_start_time(42, None, 3);
2102 let json = serde_json::to_value(unknown).unwrap();
2103 assert_eq!(json["pid"], 42);
2104 assert!(json["started_at"].is_null(), "never invent a start time");
2105 let reason = json["started_at_unavailable_reason"].as_str().unwrap();
2106 assert!(!reason.is_empty());
2107 assert!(reason.contains(std::env::consts::OS));
2108 assert_eq!(json["pool_generation"], 3);
2109 }
2110
2111 #[tokio::test]
2112 async fn process_identity_is_present_in_every_collector_path() {
2113 let dir = tempfile::tempdir().expect("tempdir");
2114 let (file_pool, _) = seeded_pool(&dir);
2115 let memory_pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool");
2116 for pool in [file_pool, memory_pool] {
2117 let pool = Arc::new(pool);
2118 let expected = ProcessIdentity::current(&pool);
2119 let writer_before = pool.writer_acquisition_snapshot();
2120 let sync = collect(
2121 &pool,
2122 BuildIdentity::from_env("test", None),
2123 Duration::from_secs(30),
2124 );
2125 let asynchronous = collect_with_runtime_audit_metrics_interruptibly(
2126 Arc::clone(&pool),
2127 BuildIdentity::from_env("test", None),
2128 Duration::from_secs(30),
2129 0,
2130 None,
2131 )
2132 .await
2133 .expect("diagnostics");
2134 for report in [sync, asynchronous] {
2135 assert_eq!(report.process, expected);
2136 let json = serde_json::to_value(report).unwrap();
2137 assert_eq!(json["process"], serde_json::to_value(&expected).unwrap());
2138 assert_eq!(json["build"]["version"], "test");
2139 }
2140 assert_eq!(
2141 pool.writer_acquisition_snapshot(),
2142 writer_before,
2143 "diagnostics must not count its probes as write traffic"
2144 );
2145 }
2146 }
2147
2148 #[test]
2149 fn process_identity_survives_pool_reconstruction_with_reset_counters() {
2150 let pool = ConnectionPool::new(PoolConfig::default()).expect("first pool");
2151 drop(pool.try_writer().expect("writer acquisition"));
2152 drop(pool.reader().expect("reader acquisition"));
2153 let before = collect(
2154 &pool,
2155 BuildIdentity::from_env("test", None),
2156 Duration::from_secs(30),
2157 );
2158 assert!(before.writer_contention.writer_acquisitions > 0);
2159 assert!(before.reader_contention.reader_acquisitions > 0);
2160 drop(pool);
2161
2162 let replacement = ConnectionPool::new(PoolConfig::default()).expect("replacement pool");
2163 let after = collect(
2164 &replacement,
2165 BuildIdentity::from_env("test", None),
2166 Duration::from_secs(30),
2167 );
2168 assert_eq!(after.process.pid, before.process.pid);
2169 assert_eq!(after.process.started_at, before.process.started_at);
2170 assert_eq!(
2171 after.process.started_at_unavailable_reason,
2172 before.process.started_at_unavailable_reason
2173 );
2174 assert!(after.process.pool_generation > before.process.pool_generation);
2175 assert_eq!(after.writer_contention.writer_acquisitions, 0);
2176 assert_eq!(after.reader_contention.reader_acquisitions, 0);
2177 }
2178
2179 fn seeded_pool(dir: &tempfile::TempDir) -> (ConnectionPool, PathBuf) {
2180 let path = dir.path().join("diag.db");
2181 let pool = ConnectionPool::new(PoolConfig {
2182 path: Some(path.clone()),
2183 ..PoolConfig::for_test()
2184 })
2185 .expect("pool open");
2186 {
2187 let writer = pool.try_writer().expect("writer");
2188 writer
2189 .conn()
2190 .execute_batch(
2191 "CREATE TABLE t (x INTEGER); \
2192 CREATE TABLE entities (id TEXT PRIMARY KEY, deleted_at INTEGER, merged_into TEXT); \
2193 CREATE TABLE graph_edges (
2194 namespace TEXT NOT NULL,
2195 id TEXT NOT NULL,
2196 PRIMARY KEY (namespace, id)
2197 ); \
2198 CREATE TABLE graph_edges_seq (
2199 seq INTEGER PRIMARY KEY AUTOINCREMENT,
2200 edge_id TEXT NOT NULL UNIQUE
2201 ); \
2202 CREATE VIRTUAL TABLE fts_entities USING fts5(
2203 namespace UNINDEXED, subject_id UNINDEXED, title, body,
2204 tokenize='trigram'
2205 ); \
2206 CREATE VIRTUAL TABLE fts_notes USING fts5(
2207 namespace UNINDEXED, subject_id UNINDEXED, title, body,
2208 tokenize='trigram'
2209 ); \
2210 INSERT INTO t VALUES (1), (2), (3); \
2211 INSERT INTO fts_entities(namespace, subject_id, title, body)
2212 VALUES('local', 'entity-1', 'entity title', 'entity diagnostic body'); \
2213 INSERT INTO fts_notes(namespace, subject_id, title, body)
2214 VALUES('local', 'note-1', 'note title', 'note diagnostic body');",
2215 )
2216 .expect("seed writes");
2217 }
2218 (pool, path)
2219 }
2220
2221 #[test]
2222 fn dropping_census_future_guard_requests_cooperative_stop() {
2223 let stopped = Arc::new(AtomicBool::new(false));
2224 let guard = StopCensusOnDrop {
2225 stopped: Arc::clone(&stopped),
2226 armed: true,
2227 };
2228
2229 drop(guard);
2230
2231 assert!(
2232 stopped.load(Ordering::SeqCst),
2233 "dropping diagnostics while its census worker is live must stop the PID/fd walk"
2234 );
2235 }
2236
2237 #[tokio::test]
2238 async fn runtime_audit_batch_fields_are_additive_and_unavailable_without_a_control() {
2239 let dir = tempfile::tempdir().expect("tempdir");
2240 let (pool, _path) = seeded_pool(&dir);
2241 let pool = Arc::new(pool);
2242
2243 let without_control = collect_with_audit_append_failures_interruptibly(
2244 Arc::clone(&pool),
2245 BuildIdentity::from_env("test", None),
2246 Duration::from_secs(30),
2247 0,
2248 )
2249 .await
2250 .expect("diagnostics succeed");
2251 assert!(without_control
2252 .writer_contention
2253 .audit_batch_flush_failures
2254 .is_none());
2255 assert!(
2256 without_control
2257 .writer_contention
2258 .audit_batch_flush_failures_unavailable_reason
2259 .is_some(),
2260 "no audit-batch control was supplied, so the field must carry a reason, not a \
2261 fabricated zero"
2262 );
2263 assert!(without_control
2264 .writer_contention
2265 .audit_degraded_rows
2266 .is_none());
2267 assert!(without_control.writer_contention.audit_degraded.is_none());
2268
2269 let with_control = collect_with_runtime_audit_metrics_interruptibly(
2270 Arc::clone(&pool),
2271 BuildIdentity::from_env("test", None),
2272 Duration::from_secs(30),
2273 0,
2274 Some(RuntimeAuditBatchMetrics {
2275 flush_failures: 3,
2276 degraded_rows: 7,
2277 degraded: true,
2278 admission_refused_obligations: 5,
2279 admission_refused_obligations_last_at_ms: Some(1_700_000_000_123),
2280 admission_unresolved_obligations: 2,
2281 admission_unresolved_obligations_last_at_ms: Some(1_700_000_000_456),
2282 }),
2283 )
2284 .await
2285 .expect("diagnostics succeed");
2286 assert_eq!(
2287 with_control.writer_contention.audit_batch_flush_failures,
2288 Some(3)
2289 );
2290 assert!(with_control
2291 .writer_contention
2292 .audit_batch_flush_failures_unavailable_reason
2293 .is_none());
2294 assert_eq!(with_control.writer_contention.audit_degraded_rows, Some(7));
2295 assert_eq!(with_control.writer_contention.audit_degraded, Some(true));
2296 assert_eq!(
2297 with_control
2298 .writer_contention
2299 .audit_admission_refused_obligations,
2300 Some(5),
2301 "an operator must be able to read the admission-refused obligation count from \
2302 db_diagnostics without a test-only feature gate (ADR-103 Amendment 3)"
2303 );
2304 assert!(with_control
2305 .writer_contention
2306 .audit_admission_refused_obligations_unavailable_reason
2307 .is_none());
2308 assert_eq!(
2309 with_control
2310 .writer_contention
2311 .audit_admission_unresolved_obligations,
2312 Some(2),
2313 "an operator must be able to distinguish enqueued-but-unresolved rows from \
2314 confirmed-refused rows (ADR-103 Amendment 3)"
2315 );
2316 assert!(with_control
2317 .writer_contention
2318 .audit_admission_unresolved_obligations_unavailable_reason
2319 .is_none());
2320 assert!(without_control
2321 .writer_contention
2322 .audit_admission_refused_obligations
2323 .is_none());
2324 assert!(without_control
2325 .writer_contention
2326 .audit_admission_refused_obligations_unavailable_reason
2327 .is_some());
2328 assert!(without_control
2329 .writer_contention
2330 .audit_admission_unresolved_obligations
2331 .is_none());
2332 assert!(without_control
2333 .writer_contention
2334 .audit_admission_unresolved_obligations_unavailable_reason
2335 .is_some());
2336
2337 assert_eq!(
2340 with_control.writer_contention.writer_acquisitions,
2341 without_control.writer_contention.writer_acquisitions
2342 );
2343 }
2344
2345 #[test]
2346 fn writer_task_pool_sourced_counters_are_always_populated_directly() {
2347 let dir = tempfile::tempdir().expect("tempdir");
2348 let (pool, _path) = seeded_pool(&dir);
2349
2350 let report = collect(
2351 &pool,
2352 BuildIdentity::from_env("9.9.9", None),
2353 Duration::from_secs(30),
2354 );
2355
2356 assert_eq!(report.writer_contention.writer_task_request_failures, 0);
2359 assert_eq!(report.writer_contention.writer_task_side_effects_unknown, 0);
2360 }
2361
2362 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2363 #[serial]
2364 async fn request_cancellation_after_passive_stops_before_graph_and_census() {
2365 let dir = tempfile::tempdir().expect("tempdir");
2366 let (pool, _) = seeded_pool(&dir);
2367 let pool = Arc::new(pool);
2368 TEST_REACHED_AFTER_PASSIVE.store(false, Ordering::SeqCst);
2369 TEST_PAUSE_AFTER_PASSIVE.store(true, Ordering::SeqCst);
2370 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
2371 let diagnostic_pool = Arc::clone(&pool);
2372 let task = tokio::spawn(crate::scope_request_read_cancellation(
2373 cancel_rx,
2374 async move {
2375 collect_with_audit_append_failures_interruptibly(
2376 diagnostic_pool,
2377 BuildIdentity::from_env("test", None),
2378 Duration::from_secs(30),
2379 0,
2380 )
2381 .await
2382 },
2383 ));
2384
2385 tokio::time::timeout(Duration::from_secs(1), async {
2386 while !TEST_REACHED_AFTER_PASSIVE.load(Ordering::SeqCst) {
2387 tokio::task::yield_now().await;
2388 }
2389 })
2390 .await
2391 .expect("diagnostics never completed its admitted PASSIVE phase");
2392 cancel_tx.send(true).unwrap();
2393 let result = tokio::time::timeout(Duration::from_secs(1), task)
2394 .await
2395 .expect("cancelled diagnostics did not stop promptly")
2396 .expect("diagnostics task panicked");
2397 TEST_PAUSE_AFTER_PASSIVE.store(false, Ordering::SeqCst);
2398 assert!(matches!(result, Err(StorageError::Timeout { .. })));
2399
2400 let one: i64 = pool
2401 .reader()
2402 .expect("diagnostics returned its connection")
2403 .conn()
2404 .query_row("SELECT 1", [], |row| row.get(0))
2405 .unwrap();
2406 assert_eq!(one, 1);
2407 }
2408
2409 #[test]
2410 fn checkpoint_probe_returns_a_well_formed_triple_on_a_file_backed_db() {
2411 let dir = tempfile::tempdir().expect("tempdir");
2412 let (pool, _path) = seeded_pool(&dir);
2413 let conn = pool
2414 .open_standalone_writer_untracked()
2415 .expect("standalone probe connection");
2416
2417 let probe = checkpoint_probe(&conn).expect("probe must succeed on a WAL database");
2418
2419 assert!(
2420 probe.busy == 0 || probe.busy == 1,
2421 "busy is a 0/1 flag, got {}",
2422 probe.busy
2423 );
2424 assert!(
2425 probe.log_frames >= 0,
2426 "a WAL database must report a non-negative frame count, got {}",
2427 probe.log_frames
2428 );
2429 assert!(
2430 probe.checkpointed_frames >= 0,
2431 "checkpointed frames must be non-negative, got {}",
2432 probe.checkpointed_frames
2433 );
2434 assert!(
2435 probe.checkpointed_frames <= probe.log_frames,
2436 "a PASSIVE pass cannot checkpoint more frames than the WAL holds: {probe:?}"
2437 );
2438 assert!(probe.pin_depth() >= 0, "pin depth clamps at 0: {probe:?}");
2439 }
2440
2441 #[test]
2444 #[serial(checkpoint_skip_metrics)]
2445 fn checkpoint_probe_does_not_perturb_the_adr091_counters() {
2446 crate::checkpoint::reset_checkpoint_metrics_for_tests();
2447 let dir = tempfile::tempdir().expect("tempdir");
2448 let (pool, _path) = seeded_pool(&dir);
2449 let conn = pool
2450 .open_standalone_writer_untracked()
2451 .expect("standalone probe connection");
2452
2453 let before = checkpoint_counters();
2454 for _ in 0..3 {
2455 checkpoint_probe(&conn).expect("probe must succeed");
2456 }
2457 let after = checkpoint_counters();
2458
2459 assert_eq!(
2460 before, after,
2461 "checkpoint_probe must leave every ADR-091 counter untouched"
2462 );
2463 }
2464
2465 #[test]
2466 fn wal_file_state_reports_the_sidecar_size_for_a_live_db() {
2467 let dir = tempfile::tempdir().expect("tempdir");
2468 let (_pool, path) = seeded_pool(&dir);
2469
2470 let state = wal_file_state(&path);
2471 assert!(
2472 state.wal_path.ends_with("diag.db-wal"),
2473 "WAL path is the db path plus a -wal suffix, got {}",
2474 state.wal_path
2475 );
2476 assert!(
2477 state.wal_size_bytes.is_some(),
2478 "a seeded WAL database must have a stat-able -wal file: {state:?}"
2479 );
2480 assert!(state.unavailable_reason.is_none(), "{state:?}");
2481 }
2482
2483 #[test]
2484 fn wal_file_state_degrades_with_a_reason_when_the_sidecar_is_absent() {
2485 let dir = tempfile::tempdir().expect("tempdir");
2486 let state = wal_file_state(&dir.path().join("never-created.db"));
2487 assert!(state.wal_size_bytes.is_none());
2488 assert!(
2489 state.unavailable_reason.is_some(),
2490 "an absent WAL file must carry a reason, not a silent zero: {state:?}"
2491 );
2492 }
2493
2494 #[test]
2495 fn collect_on_a_file_backed_db_carries_build_identity_and_every_counter() {
2496 let dir = tempfile::tempdir().expect("tempdir");
2497 let (pool, _path) = seeded_pool(&dir);
2498 let reader_admission_capacity = pool.max_readers().max(1);
2499
2500 let report = collect(
2501 &pool,
2502 BuildIdentity::from_env("9.9.9", Some("deadbeef")),
2503 Duration::from_secs(30),
2504 );
2505
2506 assert_eq!(report.build.version, "9.9.9");
2507 assert_eq!(report.build.build_hash.as_deref(), Some("deadbeef"));
2508 assert!(report.db_path.is_some());
2509 assert!(
2510 report.checkpoint_probe.is_some(),
2511 "file-backed collect must land a probe; error was {:?}",
2512 report.checkpoint_probe_error
2513 );
2514 assert!(
2515 report.wal_file.as_ref().and_then(|w| w.wal_size_bytes) >= Some(0),
2516 "wal_size_bytes must be a non-negative byte count when present"
2517 );
2518
2519 let json = serde_json::to_value(&report).expect("report serializes");
2520 let counters = json
2521 .get("checkpoint_counters")
2522 .expect("counters section present");
2523 for key in [
2524 "last_observed_wal_pages",
2525 "truncate_attempts",
2526 "truncate_consecutive_failures",
2527 "checkpoint_skipped_ticks",
2528 "checkpoint_consecutive_skips",
2529 "checkpoint_last_skip_wal_pages",
2530 "checkpoint_pressure_elevated_ticks",
2531 "checkpoint_pressure_episodes_started",
2532 "checkpoint_pressure_episodes_recovered",
2533 "checkpoint_lifecycle_append_attempts",
2534 "checkpoint_lifecycle_append_failures",
2535 "checkpoint_lifecycle_enqueue_drops",
2536 "read_tx_max_age_evictions",
2537 ] {
2538 assert!(counters.get(key).is_some(), "counter {key} must be present");
2539 }
2540 assert_eq!(
2541 report.writer_contention.writer_acquisitions, 1,
2542 "the seed write checked the finite-wait pooled writer out once"
2543 );
2544 assert_eq!(report.writer_contention.pooled_writer_acquisitions, 1);
2545 assert_eq!(report.writer_contention.standalone_writer_acquisitions, 0);
2546 assert_eq!(report.writer_contention.writer_task_acquisitions, 0);
2547 assert_eq!(report.writer_contention.writer_acquisition_timeouts, 0);
2548 assert_eq!(
2549 report.reader_contention,
2550 ReaderContentionDiagnostics {
2551 configured_reader_cap: pool.config().max_readers,
2552 configured_checkout_timeout_ms: u64::try_from(
2553 pool.config().checkout_timeout.as_millis(),
2554 )
2555 .unwrap_or(u64::MAX),
2556 configured_busy_timeout_ms: u64::try_from(pool.config().busy_timeout.as_millis())
2557 .unwrap_or(u64::MAX),
2558 reader_admission_capacity,
2559 available_reader_admission_slots: reader_admission_capacity,
2560 reader_acquisitions: 0,
2561 pooled_reader_checkouts: 0,
2562 standalone_reader_opens: 0,
2563 infrastructure_standalone_reader_opens: 0,
2564 reader_checkout_timeouts: 0,
2565 active_pooled_reader_checkouts: 0,
2566 peak_active_pooled_reader_checkouts: 0,
2567 completed_pooled_reader_checkouts: 0,
2568 max_completed_reader_hold_micros: 0,
2569 max_completed_reader_hold_operation: None,
2570 reader_replacement_open_failures: 0,
2571 },
2572 "the diagnostics probe itself must not masquerade as request reader traffic"
2573 );
2574 assert_eq!(
2575 report.graph_edge_integrity,
2576 Some(GraphEdgeIntegrity {
2577 duplicate_edge_id_groups: 0,
2578 graph_edges_rows: 0,
2579 graph_edges_seq_rows: 0,
2580 pre_v14_duplicate_edge_state_detected: false,
2581 live_entities_carrying_merged_into: 0,
2582 })
2583 );
2584 assert!(report.graph_edge_integrity_error.is_none());
2585 let fts_segments = report
2586 .fts_segments
2587 .as_ref()
2588 .expect("file-backed diagnostics must decode both FTS structure rows");
2589 assert_eq!(fts_segments.entities.segment_count, 1);
2590 assert_eq!(fts_segments.notes.segment_count, 1);
2591 assert_eq!(fts_segments.total_segments, 2);
2592 assert!(report.fts_segments_error.is_none());
2593 let fts_json = json
2594 .get("fts_segments")
2595 .expect("FTS segment diagnostics serialize");
2596 assert_eq!(fts_json["entities"]["segment_count"], 1);
2597 assert!(json.get("fts_maintenance").is_some());
2598 assert!(
2599 report.size_composition.is_some(),
2600 "file-backed diagnostics must include page composition; error was {:?}",
2601 report.size_composition_error
2602 );
2603 assert!(report.size_composition_error.is_none());
2604 assert!(report.writer_contention.audit_append_failures.is_none());
2605 assert!(report
2606 .writer_contention
2607 .audit_obligation_append_failures
2608 .is_none());
2609 assert!(report
2610 .writer_contention
2611 .audit_obligation_append_failures_unavailable_reason
2612 .is_some());
2613 assert!(json["writer_contention"]["audit_obligation_append_failures"].is_null());
2614 assert!(
2615 report
2616 .writer_contention
2617 .audit_append_failures_unavailable_reason
2618 .is_some(),
2619 "a direct khive-db snapshot must not fabricate a runtime audit count"
2620 );
2621 }
2622
2623 #[test]
2624 fn diagnostics_exposes_reader_saturation_and_completed_hold_evidence() {
2625 let dir = tempfile::tempdir().expect("tempdir");
2626 let pool = ConnectionPool::new(PoolConfig {
2628 path: Some(dir.path().join("reader_saturation.db")),
2629 max_readers: 1,
2630 checkout_timeout: Duration::from_millis(2),
2631 ..PoolConfig::default()
2632 })
2633 .expect("one-reader file-backed pool");
2634 let held = pool.reader().expect("first reader checkout");
2635 assert!(
2636 pool.reader().is_err(),
2637 "the live checkout must exhaust the one-slot reader budget"
2638 );
2639 drop(held);
2640
2641 let report = collect(
2642 &pool,
2643 BuildIdentity::from_env("9.9.9", None),
2644 Duration::from_secs(30),
2645 );
2646 let reader = report.reader_contention;
2647 assert_eq!(reader.reader_admission_capacity, 1);
2648 assert_eq!(reader.available_reader_admission_slots, 1);
2649 assert_eq!(reader.reader_acquisitions, 1);
2650 assert_eq!(reader.pooled_reader_checkouts, 1);
2651 assert_eq!(reader.standalone_reader_opens, 0);
2652 assert_eq!(reader.infrastructure_standalone_reader_opens, 0);
2653 assert_eq!(reader.reader_checkout_timeouts, 1);
2654 assert_eq!(reader.active_pooled_reader_checkouts, 0);
2655 assert_eq!(reader.peak_active_pooled_reader_checkouts, 1);
2656 assert_eq!(reader.completed_pooled_reader_checkouts, 1);
2657 assert!(reader.max_completed_reader_hold_micros > 0);
2658
2659 let json = serde_json::to_value(&report).expect("report serializes");
2660 assert_eq!(
2661 json.pointer("/reader_contention/reader_admission_capacity"),
2662 Some(&serde_json::json!(1)),
2663 "the operator wire payload must expose the reader admission budget"
2664 );
2665 assert_eq!(
2666 json.pointer("/reader_contention/reader_checkout_timeouts"),
2667 Some(&serde_json::json!(1)),
2668 "the operator wire payload must expose the reader timeout phase"
2669 );
2670 assert!(
2671 json.pointer("/reader_contention/max_completed_reader_hold_micros")
2672 .is_some(),
2673 "the operator wire payload must expose completed hold-time evidence"
2674 );
2675 }
2676
2677 #[test]
2678 fn diagnostics_reports_configured_reader_budget_and_both_deadlines() {
2679 let pool = ConnectionPool::new(PoolConfig {
2680 max_readers: 6,
2681 checkout_timeout: Duration::from_millis(17),
2682 busy_timeout: Duration::from_millis(31),
2683 ..PoolConfig::default()
2684 })
2685 .expect("in-memory pool");
2686 let report = collect(
2687 &pool,
2688 BuildIdentity::from_env("9.9.9", None),
2689 Duration::from_secs(30),
2690 );
2691 let reader = report.reader_contention;
2692 assert_eq!(reader.reader_admission_capacity, 1);
2693
2694 let json = serde_json::to_value(&report).expect("report serializes");
2695 assert_eq!(
2696 json.pointer("/reader_contention/configured_reader_cap"),
2697 Some(&serde_json::json!(6))
2698 );
2699 assert_eq!(
2700 json.pointer("/reader_contention/configured_checkout_timeout_ms"),
2701 Some(&serde_json::json!(17))
2702 );
2703 assert_eq!(
2704 json.pointer("/reader_contention/configured_busy_timeout_ms"),
2705 Some(&serde_json::json!(31))
2706 );
2707 }
2708
2709 #[test]
2710 fn diagnostics_composes_file_backed_standalone_acquisitions_without_counting_its_probe() {
2711 let dir = tempfile::tempdir().expect("tempdir");
2712 let (pool, _path) = seeded_pool(&dir);
2713
2714 drop(
2715 pool.open_standalone_writer()
2716 .expect("write-traffic standalone connection"),
2717 );
2718
2719 let report = collect_with_audit_append_failures(
2720 &pool,
2721 BuildIdentity::from_env("9.9.9", None),
2722 Duration::from_secs(30),
2723 0,
2724 );
2725 assert_eq!(report.writer_contention.writer_acquisitions, 2);
2726 assert_eq!(report.writer_contention.pooled_writer_acquisitions, 1);
2727 assert_eq!(report.writer_contention.standalone_writer_acquisitions, 1);
2728 assert_eq!(report.writer_contention.writer_task_acquisitions, 0);
2729 assert_eq!(report.writer_contention.writer_acquisition_timeouts, 0);
2730
2731 let second = collect_with_audit_append_failures(
2732 &pool,
2733 BuildIdentity::from_env("9.9.9", None),
2734 Duration::from_secs(30),
2735 0,
2736 );
2737 assert_eq!(
2738 second.writer_contention, report.writer_contention,
2739 "the diagnostics PASSIVE probe must not inflate write-traffic counters"
2740 );
2741 }
2742
2743 #[test]
2744 fn runtime_aware_collect_exposes_the_supplied_audit_failure_counter() {
2745 let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool");
2746
2747 let report = collect_with_audit_append_failures(
2748 &pool,
2749 BuildIdentity::from_env("9.9.9", None),
2750 Duration::from_secs(30),
2751 17,
2752 );
2753
2754 assert_eq!(report.writer_contention.audit_append_failures, Some(17));
2755 assert!(
2756 report
2757 .writer_contention
2758 .audit_append_failures_unavailable_reason
2759 .is_none(),
2760 "a supplied runtime counter must not carry an unavailable reason"
2761 );
2762 }
2763
2764 #[test]
2765 fn diagnostics_exposes_an_induced_writer_checkout_timeout() {
2766 let pool = ConnectionPool::new(PoolConfig {
2767 checkout_timeout: Duration::from_millis(1),
2768 ..PoolConfig::default()
2769 })
2770 .expect("in-memory pool");
2771
2772 let held = pool.writer().expect("first writer checkout succeeds");
2773 assert!(
2774 matches!(
2775 pool.writer(),
2776 Err(crate::SqliteError::WriterPoolCheckoutTimeout { .. })
2777 ),
2778 "holding the sole writer must exercise the typed timeout path"
2779 );
2780 drop(held);
2781
2782 let report = collect_with_audit_append_failures(
2783 &pool,
2784 BuildIdentity::from_env("9.9.9", None),
2785 Duration::from_secs(30),
2786 0,
2787 );
2788 assert_eq!(report.writer_contention.writer_acquisitions, 1);
2789 assert_eq!(report.writer_contention.pooled_writer_acquisitions, 1);
2790 assert_eq!(report.writer_contention.standalone_writer_acquisitions, 0);
2791 assert_eq!(report.writer_contention.writer_task_acquisitions, 0);
2792 assert_eq!(report.writer_contention.writer_acquisition_timeouts, 1);
2793 }
2794
2795 #[test]
2798 fn never_observed_sentinels_serialize_as_null() {
2799 let counters = CheckpointCounters {
2800 last_observed_wal_pages: None,
2801 truncate_attempts: 0,
2802 truncate_consecutive_failures: 0,
2803 checkpoint_skipped_ticks: 0,
2804 checkpoint_consecutive_skips: 0,
2805 checkpoint_last_skip_wal_pages: None,
2806 checkpoint_pressure_elevated_ticks: 0,
2807 checkpoint_pressure_episodes_started: 0,
2808 checkpoint_pressure_episodes_recovered: 0,
2809 checkpoint_lifecycle_append_attempts: 0,
2810 checkpoint_lifecycle_append_failures: 0,
2811 checkpoint_lifecycle_enqueue_drops: 0,
2812 read_tx_max_age_evictions: 0,
2813 };
2814 let json = serde_json::to_value(counters).expect("serializes");
2815 assert!(json["last_observed_wal_pages"].is_null());
2816 assert!(json["checkpoint_last_skip_wal_pages"].is_null());
2817 }
2818
2819 #[test]
2820 fn graph_edge_integrity_detects_the_pre_v14_duplicate_state() {
2821 let conn = Connection::open_in_memory().expect("in-memory sqlite");
2822 conn.execute_batch(
2823 "CREATE TABLE graph_edges (
2824 namespace TEXT NOT NULL,
2825 id TEXT NOT NULL,
2826 PRIMARY KEY (namespace, id)
2827 );
2828 CREATE TABLE graph_edges_seq (
2829 seq INTEGER PRIMARY KEY AUTOINCREMENT,
2830 edge_id TEXT NOT NULL UNIQUE
2831 );
2832 CREATE TABLE entities (id TEXT PRIMARY KEY, deleted_at INTEGER, merged_into TEXT);
2833 INSERT INTO graph_edges(namespace, id)
2834 VALUES ('alpha', 'shared-edge'), ('beta', 'shared-edge');
2835 INSERT INTO graph_edges_seq(edge_id) VALUES ('shared-edge');",
2836 )
2837 .expect("seed the state possible before the V14 uniqueness guard");
2838
2839 let integrity = graph_edge_integrity(&conn).expect("integrity query succeeds");
2840
2841 assert_eq!(integrity.duplicate_edge_id_groups, 1);
2842 assert_eq!(integrity.graph_edges_rows, 2);
2843 assert_eq!(integrity.graph_edges_seq_rows, 1);
2844 assert!(integrity.pre_v14_duplicate_edge_state_detected);
2845 }
2846
2847 #[test]
2848 fn graph_edge_integrity_counts_only_live_rows_that_still_carry_merge_provenance() {
2849 let conn = Connection::open_in_memory().expect("in-memory sqlite");
2850 conn.execute_batch(
2851 "CREATE TABLE graph_edges (
2852 namespace TEXT NOT NULL,
2853 id TEXT NOT NULL,
2854 PRIMARY KEY (namespace, id)
2855 );
2856 CREATE TABLE graph_edges_seq (
2857 seq INTEGER PRIMARY KEY AUTOINCREMENT,
2858 edge_id TEXT NOT NULL UNIQUE
2859 );
2860 CREATE TABLE entities (id TEXT PRIMARY KEY, deleted_at INTEGER, merged_into TEXT);
2861 INSERT INTO entities(id, deleted_at, merged_into) VALUES
2862 ('kept', NULL, NULL),
2863 ('tombstoned-source', 1, 'kept'),
2864 ('left-live-by-an-old-restore', NULL, 'kept'),
2865 ('plain-soft-delete', 1, NULL);",
2866 )
2867 .expect("seed one row of each shape");
2868
2869 let integrity = graph_edge_integrity(&conn).expect("integrity query succeeds");
2870
2871 assert_eq!(
2872 integrity.live_entities_carrying_merged_into, 1,
2873 "a merge tombstone and a plain live row are both in order; only the live row \
2874 carrying merged_into is the invariant violation"
2875 );
2876 assert_eq!(integrity.duplicate_edge_id_groups, 0);
2877 }
2878
2879 #[test]
2880 fn graph_edge_integrity_does_not_mislabel_retained_delete_history_as_a_duplicate() {
2881 let conn = Connection::open_in_memory().expect("in-memory sqlite");
2882 conn.execute_batch(
2883 "CREATE TABLE graph_edges (
2884 namespace TEXT NOT NULL,
2885 id TEXT NOT NULL,
2886 PRIMARY KEY (namespace, id)
2887 );
2888 CREATE TABLE graph_edges_seq (
2889 seq INTEGER PRIMARY KEY AUTOINCREMENT,
2890 edge_id TEXT NOT NULL UNIQUE
2891 );
2892 CREATE TABLE entities (id TEXT PRIMARY KEY, deleted_at INTEGER, merged_into TEXT);
2893 INSERT INTO graph_edges(namespace, id) VALUES ('local', 'live-edge');
2894 INSERT INTO graph_edges_seq(edge_id)
2895 VALUES ('deleted-edge'), ('live-edge');",
2896 )
2897 .expect("seed a retained sequence row for a hard-deleted edge");
2898
2899 let integrity = graph_edge_integrity(&conn).expect("integrity query succeeds");
2900
2901 assert_eq!(integrity.duplicate_edge_id_groups, 0);
2902 assert_eq!(integrity.graph_edges_rows, 1);
2903 assert_eq!(integrity.graph_edges_seq_rows, 2);
2904 assert!(
2905 !integrity.pre_v14_duplicate_edge_state_detected,
2906 "ledger rows intentionally survive hard deletion; count mismatch alone is not the \
2907 pre-V14 duplicate state"
2908 );
2909 }
2910
2911 #[cfg(unix)]
2912 #[test]
2913 fn unmeasured_sidecar_cleanup_fields_are_absent_from_the_wire_payload() {
2914 let pin = wal_pin_attribution_from_census(crate::walpin::CensusResult {
2915 holders: std::collections::HashSet::new(),
2916 uninspectable_pids: Vec::new(),
2917 truncated: false,
2918 budget_exhausted: false,
2919 });
2920
2921 assert_eq!(pin.sidecar_listing_truncated, None);
2922 assert_eq!(pin.sidecar_entries_cleanup_would_reap, None);
2923 let json = serde_json::to_value(pin).expect("attribution serializes");
2924 assert!(
2925 json.get("sidecar_listing_truncated").is_none(),
2926 "a skipped enumeration must omit sidecar_listing_truncated, not fabricate false"
2927 );
2928 assert!(
2929 json.get("sidecar_entries_cleanup_would_reap").is_none(),
2930 "a skipped enumeration must omit sidecar_entries_cleanup_would_reap, not fabricate 0"
2931 );
2932 }
2933
2934 #[cfg(unix)]
2935 #[test]
2936 fn wal_pin_census_serializes_only_under_the_nested_carrier() {
2937 let pin = wal_pin_attribution_from_census(crate::walpin::CensusResult {
2938 holders: std::collections::HashSet::from([41, 7]),
2939 uninspectable_pids: vec![99],
2940 truncated: true,
2941 budget_exhausted: false,
2942 });
2943
2944 let json = serde_json::to_value(pin).expect("attribution serializes");
2945 assert_eq!(json["census"]["holder_pids"], serde_json::json!([7, 41]));
2946 assert_eq!(
2947 json["census"]["uninspectable_pids"],
2948 serde_json::json!([99])
2949 );
2950 assert_eq!(json["census"]["truncated"], true);
2951
2952 for duplicate in [
2953 "census_holder_pids",
2954 "census_uninspectable_pids",
2955 "census_truncated",
2956 "census_is_complete",
2957 ] {
2958 assert!(
2959 json.get(duplicate).is_none(),
2960 "wal_pin.{duplicate} must not duplicate wal_pin.census: {json}"
2961 );
2962 }
2963 }
2964
2965 #[test]
2969 fn probe_refuses_a_missing_configured_path_without_creating_it() {
2970 let dir = tempfile::tempdir().expect("tempdir");
2971 let (pool, path) = seeded_pool(&dir);
2972
2973 for suffix in ["", "-wal", "-shm"] {
2974 let mut p = path.as_os_str().to_os_string();
2975 p.push(suffix);
2976 let _ = std::fs::remove_file(PathBuf::from(p));
2977 }
2978 assert!(!path.exists(), "precondition: the database file is gone");
2979
2980 let report = collect(
2981 &pool,
2982 BuildIdentity::from_env("0.0.0", None),
2983 Duration::from_secs(30),
2984 );
2985
2986 assert!(
2987 report.checkpoint_probe.is_none(),
2988 "a missing database must not yield a probe result: {report:?}"
2989 );
2990 assert!(
2991 report.checkpoint_probe_error.is_some(),
2992 "a missing database must say why there is no probe: {report:?}"
2993 );
2994 assert!(
2995 !path.exists(),
2996 "a diagnostics request must never create the database it was asked about"
2997 );
2998 }
2999
3000 #[test]
3003 fn collect_degrades_gracefully_for_an_in_memory_backend() {
3004 let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool");
3005 let report = collect(
3006 &pool,
3007 BuildIdentity::from_env("0.0.0", None),
3008 Duration::from_secs(30),
3009 );
3010
3011 assert!(report.db_path.is_none());
3012 assert!(report.wal_file.is_none());
3013 assert!(report.checkpoint_probe.is_none());
3014 assert!(report.fts_segments.is_none());
3015 assert!(report.fts_segments_error.is_some());
3016 assert!(
3017 report.checkpoint_probe_error.is_some(),
3018 "an in-memory report must say WHY there is no probe"
3019 );
3020 assert!(!report.wal_pin.available);
3021 assert!(report.wal_pin.unavailable_reason.is_some());
3022 assert_eq!(report.wal_pin.status, WalPinAttributionStatus::Unavailable);
3023 assert!(matches!(
3024 report.wal_pin.census,
3025 WalPinCensus::Unavailable { .. }
3026 ));
3027 assert!(report.size_composition.is_none());
3028 assert!(report
3029 .size_composition_error
3030 .as_deref()
3031 .is_some_and(|reason| reason.contains("no file-backed page composition")));
3032 }
3033
3034 #[cfg(unix)]
3035 #[test]
3036 fn incomplete_holder_census_is_a_tagged_degraded_result() {
3037 let census = crate::walpin::CensusResult {
3038 holders: std::collections::HashSet::from([41, 7]),
3039 uninspectable_pids: vec![99, 99],
3040 truncated: true,
3041 budget_exhausted: false,
3042 };
3043
3044 let pin = wal_pin_attribution_from_census(census);
3045
3046 assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
3047 assert!(!pin.available);
3048 assert!(!pin.census_is_complete);
3049 assert_eq!(pin.census_holder_pids, vec![7, 41]);
3050 assert_eq!(pin.census_uninspectable_pids, vec![99]);
3051 assert!(
3052 pin.unavailable_reason.as_deref().is_some_and(
3053 |reason| reason.contains("additional database holders cannot be ruled out")
3054 ),
3055 "the legacy reason must also fail loud for old consumers: {pin:?}"
3056 );
3057 match &pin.census {
3058 WalPinCensus::Incomplete {
3059 holder_pids,
3060 uninspectable_pids,
3061 truncated,
3062 reason,
3063 } => {
3064 assert_eq!(holder_pids, &vec![7, 41]);
3065 assert_eq!(uninspectable_pids, &vec![99]);
3066 assert!(*truncated);
3067 assert!(reason.contains("additional database holders cannot be ruled out"));
3068 }
3069 other => panic!("incomplete scan must serialize as incomplete, got {other:?}"),
3070 }
3071
3072 let json = serde_json::to_value(&pin).expect("serializes");
3073 assert_eq!(json["status"], "degraded");
3074 assert_eq!(json["census"]["status"], "incomplete");
3075 }
3076
3077 #[cfg(unix)]
3078 #[test]
3079 fn complete_holder_census_stays_explicit_while_attribution_is_degraded() {
3080 let census = crate::walpin::CensusResult {
3081 holders: std::collections::HashSet::from([7]),
3082 uninspectable_pids: Vec::new(),
3083 truncated: false,
3084 budget_exhausted: false,
3085 };
3086
3087 let pin = wal_pin_attribution_from_census(census);
3088
3089 assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
3090 assert!(pin.census_is_complete);
3091 assert!(matches!(
3092 pin.census,
3093 WalPinCensus::Complete { ref holder_pids } if holder_pids == &vec![7]
3094 ));
3095 assert_eq!(
3096 pin.status_reasons.len(),
3097 1,
3098 "only missing sidecar reconciliation degrades a complete OS census"
3099 );
3100 }
3101
3102 #[cfg(unix)]
3103 #[test]
3104 fn complete_holder_and_read_only_sidecar_evidence_reconcile_to_complete() {
3105 let census = crate::walpin::CensusResult {
3106 holders: std::collections::HashSet::from([7]),
3107 uninspectable_pids: Vec::new(),
3108 truncated: false,
3109 budget_exhausted: false,
3110 };
3111 let sidecar = crate::walpin::WalpinReport {
3112 entries: vec![crate::walpin::WalpinPidHealth::RegisteredSilent { pid: 7 }],
3113 sidecar_listing_truncated: false,
3114 cleanup_would_reap: 0,
3115 orphan_temps_reaped: 0,
3116 };
3117
3118 let pin = wal_pin_attribution_from_evidence(census, sidecar);
3119
3120 assert_eq!(pin.status, WalPinAttributionStatus::Complete);
3121 assert!(pin.available);
3122 assert!(pin.fully_attributed);
3123 assert!(pin.status_reasons.is_empty());
3124 assert_eq!(pin.registered_silent_pids, vec![7]);
3125 assert!(pin.census_pids_without_attribution.is_empty());
3126 assert_eq!(pin.sidecar_listing_truncated, Some(false));
3127 assert_eq!(pin.sidecar_entries_cleanup_would_reap, Some(0));
3128 }
3129
3130 #[cfg(unix)]
3131 #[test]
3132 fn complete_census_with_an_unregistered_holder_is_degraded_not_exonerated() {
3133 let census = crate::walpin::CensusResult {
3134 holders: std::collections::HashSet::from([7, 41]),
3135 uninspectable_pids: Vec::new(),
3136 truncated: false,
3137 budget_exhausted: false,
3138 };
3139 let sidecar = crate::walpin::WalpinReport {
3140 entries: vec![crate::walpin::WalpinPidHealth::RegisteredSilent { pid: 7 }],
3141 sidecar_listing_truncated: false,
3142 cleanup_would_reap: 0,
3143 orphan_temps_reaped: 0,
3144 };
3145
3146 let pin = wal_pin_attribution_from_evidence(census, sidecar);
3147
3148 assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
3149 assert!(!pin.available);
3150 assert!(!pin.fully_attributed);
3151 assert_eq!(pin.census_pids_without_attribution, vec![41]);
3152 assert!(pin
3153 .status_reasons
3154 .iter()
3155 .any(|reason| reason.contains("holder(s) have no sidecar attribution")));
3156 }
3157
3158 #[test]
3159 fn database_size_composition_reports_tables_indexes_fts_and_vectors_separately() {
3160 let conn = Connection::open_in_memory().expect("in-memory sqlite");
3161 conn.execute_batch(
3162 "CREATE TABLE docs(id INTEGER PRIMARY KEY, body TEXT NOT NULL); \
3163 CREATE INDEX idx_docs_body ON docs(body); \
3164 CREATE TABLE fts_demo_data(id INTEGER PRIMARY KEY, block BLOB); \
3165 CREATE TABLE vec_demo_chunks(id INTEGER PRIMARY KEY, vectors BLOB); \
3166 CREATE TABLE knowledge_sections(id INTEGER PRIMARY KEY, embedding BLOB); \
3167 INSERT INTO docs(body) VALUES (zeroblob(8192)); \
3168 INSERT INTO fts_demo_data(block) VALUES (zeroblob(8192)); \
3169 INSERT INTO vec_demo_chunks(vectors) VALUES (zeroblob(8192)); \
3170 INSERT INTO knowledge_sections(embedding) VALUES (zeroblob(8192));",
3171 )
3172 .expect("seed size classes");
3173
3174 let composition = database_size_composition(&conn).expect("dbstat composition");
3175 let class_for = |name: &str| {
3176 composition
3177 .objects
3178 .iter()
3179 .find(|object| object.name == name)
3180 .map(|object| object.storage_class)
3181 };
3182
3183 assert_eq!(class_for("docs"), Some(DatabaseStorageClass::RowTable));
3184 assert_eq!(
3185 class_for("idx_docs_body"),
3186 Some(DatabaseStorageClass::Index)
3187 );
3188 assert_eq!(
3189 class_for("fts_demo_data"),
3190 Some(DatabaseStorageClass::FullText)
3191 );
3192 assert_eq!(
3193 class_for("vec_demo_chunks"),
3194 Some(DatabaseStorageClass::Vector)
3195 );
3196 assert_eq!(
3197 class_for("knowledge_sections"),
3198 Some(DatabaseStorageClass::MixedRowAndEmbedding)
3199 );
3200 assert!(composition.vector_bytes > 0);
3201 assert!(composition.full_text_bytes > 0);
3202 assert!(composition.mixed_embedding_bytes > 0);
3203 assert_eq!(
3204 composition
3205 .accounted_bytes
3206 .saturating_add(composition.freelist_bytes)
3207 .saturating_add(composition.unaccounted_bytes),
3208 composition.database_bytes
3209 );
3210 }
3211
3212 #[cfg(unix)]
3215 #[test]
3216 #[serial(khive_walpin_sidecar_env)]
3217 fn wal_pin_attribution_degrades_when_a_holder_has_no_sidecar_registration() {
3218 let dir = tempfile::tempdir().expect("tempdir");
3219 let (pool, path) = seeded_pool(&dir);
3220 let _ = &pool;
3221
3222 let pin = wal_pin_attribution(&path, Duration::from_secs(30));
3223
3224 assert!(
3225 !pin.fully_attributed,
3226 "the pool's OS holder has no test sidecar registration"
3227 );
3228 assert!(
3229 pin.unavailable_reason.is_some(),
3230 "the missing holder attribution must be explained: {pin:?}"
3231 );
3232 assert!(pin.sidecar_entries.is_empty());
3233 assert!(pin.reporting.is_empty());
3234 assert_eq!(pin.sidecar_listing_truncated, Some(false));
3235 assert_eq!(pin.sidecar_entries_cleanup_would_reap, Some(0));
3236 assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
3237 assert!(matches!(
3238 pin.census,
3239 WalPinCensus::Complete { .. } | WalPinCensus::Incomplete { .. }
3240 ));
3241 }
3242
3243 #[cfg(unix)]
3250 #[test]
3251 #[serial(khive_walpin_sidecar_env)]
3252 fn diagnostics_finds_sidecar_evidence_through_an_aliased_database_path() {
3253 let dir = tempfile::tempdir().expect("tempdir");
3254 let real_dir = dir.path().join("real");
3255 std::fs::create_dir(&real_dir).expect("mkdir real dir");
3256 let real_path = real_dir.join("diag.db");
3257 std::fs::write(&real_path, b"").expect("create real file");
3258 let alias_path = dir.path().join("alias.db");
3259 std::os::unix::fs::symlink(&real_path, &alias_path).expect("symlink alias");
3260
3261 let pool = ConnectionPool::new(PoolConfig {
3262 path: Some(alias_path.clone()),
3263 ..PoolConfig::for_test()
3264 })
3265 .expect("pool open through symlinked path");
3266 {
3267 let writer = pool.try_writer().expect("writer");
3268 writer
3269 .conn()
3270 .execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
3271 .expect("seed a write so the WAL file exists");
3272 }
3273
3274 let canonical = pool
3275 .canonical_path()
3276 .expect("file-backed pool")
3277 .to_path_buf();
3278 assert_ne!(
3279 canonical, alias_path,
3280 "the alias must actually differ from the canonical path for this test to mean \
3281 anything"
3282 );
3283
3284 let pid = std::process::id();
3285 let sidecar_dir = crate::walpin::sidecar_dir_for(&canonical);
3286 let beacon = crate::walpin::WalpinBeacon {
3287 pid,
3288 process_role: "session".to_string(),
3289 started_at: crate::walpin::process_start_time_secs(pid).unwrap_or(0),
3290 sweep_interval_ms: 5_000,
3291 };
3292 crate::walpin::write_beacon(&sidecar_dir, &beacon).expect("seed this process's beacon");
3293
3294 let report = collect(
3295 &pool,
3296 BuildIdentity::from_env("test", None),
3297 Duration::from_secs(30),
3298 );
3299
3300 assert_eq!(
3309 report.wal_pin.sidecar_listing_truncated,
3310 Some(false),
3311 "the sidecar enumeration must run to completion: {:?}",
3312 report.wal_pin
3313 );
3314 assert!(
3315 report.wal_pin.registered_silent_pids.contains(&pid),
3316 "the beacon written beside the canonical path must be found: {:?}",
3317 report.wal_pin
3318 );
3319 assert!(
3320 report.wal_pin.census_holder_pids.contains(&pid),
3321 "the OS census must find this process holding its own database open: {:?}",
3322 report.wal_pin
3323 );
3324 assert!(
3325 !report
3326 .wal_pin
3327 .census_pids_without_attribution
3328 .contains(&pid),
3329 "this process's own holder entry must be attributed by its own sidecar evidence, \
3330 not left unexplained: {:?}",
3331 report.wal_pin
3332 );
3333 }
3334
3335 #[cfg(unix)]
3343 #[tokio::test]
3344 #[serial(khive_walpin_sidecar_env)]
3345 async fn diagnostics_finds_sidecar_evidence_through_an_aliased_database_path_async() {
3346 let dir = tempfile::tempdir().expect("tempdir");
3347 let real_dir = dir.path().join("real");
3348 std::fs::create_dir(&real_dir).expect("mkdir real dir");
3349 let real_path = real_dir.join("diag.db");
3350 std::fs::write(&real_path, b"").expect("create real file");
3351 let alias_path = dir.path().join("alias.db");
3352 std::os::unix::fs::symlink(&real_path, &alias_path).expect("symlink alias");
3353
3354 let pool = ConnectionPool::new(PoolConfig {
3355 path: Some(alias_path.clone()),
3356 ..PoolConfig::for_test()
3357 })
3358 .expect("pool open through symlinked path");
3359 {
3360 let writer = pool.try_writer().expect("writer");
3361 writer
3362 .conn()
3363 .execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
3364 .expect("seed a write so the WAL file exists");
3365 }
3366
3367 let canonical = pool
3368 .canonical_path()
3369 .expect("file-backed pool")
3370 .to_path_buf();
3371 assert_ne!(
3372 canonical, alias_path,
3373 "the alias must actually differ from the canonical path for this test to mean \
3374 anything"
3375 );
3376
3377 let pid = std::process::id();
3378 let sidecar_dir = crate::walpin::sidecar_dir_for(&canonical);
3379 let beacon = crate::walpin::WalpinBeacon {
3380 pid,
3381 process_role: "session".to_string(),
3382 started_at: crate::walpin::process_start_time_secs(pid).unwrap_or(0),
3383 sweep_interval_ms: 5_000,
3384 };
3385 crate::walpin::write_beacon(&sidecar_dir, &beacon).expect("seed this process's beacon");
3386
3387 let pool = Arc::new(pool);
3388 let report = collect_with_audit_append_failures_interruptibly(
3389 Arc::clone(&pool),
3390 BuildIdentity::from_env("test", None),
3391 Duration::from_secs(30),
3392 0,
3393 )
3394 .await
3395 .expect("diagnostics succeed");
3396
3397 assert_eq!(
3402 report.wal_pin.sidecar_listing_truncated,
3403 Some(false),
3404 "the sidecar enumeration must run to completion: {:?}",
3405 report.wal_pin
3406 );
3407 assert!(
3408 report.wal_pin.registered_silent_pids.contains(&pid),
3409 "the beacon written beside the canonical path must be found: {:?}",
3410 report.wal_pin
3411 );
3412 assert!(
3413 report.wal_pin.census_holder_pids.contains(&pid),
3414 "the OS census must find this process holding its own database open: {:?}",
3415 report.wal_pin
3416 );
3417 assert!(
3418 !report
3419 .wal_pin
3420 .census_pids_without_attribution
3421 .contains(&pid),
3422 "this process's own holder entry must be attributed by its own sidecar evidence, \
3423 not left unexplained: {:?}",
3424 report.wal_pin
3425 );
3426 }
3427
3428 #[cfg(unix)]
3432 #[test]
3433 #[serial(khive_walpin_sidecar_env)]
3434 fn wal_pin_attribution_reports_disabled_when_the_sidecar_is_explicitly_off() {
3435 let dir = tempfile::tempdir().expect("tempdir");
3436 let (pool, path) = seeded_pool(&dir);
3437 let _ = &pool;
3438 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
3439 std::env::set_var("KHIVE_WALPIN_SIDECAR", "0");
3440
3441 let pin = wal_pin_attribution(&path, Duration::from_secs(30));
3442
3443 assert!(
3444 !pin.available,
3445 "an explicitly disabled sidecar can never produce a reconciled answer"
3446 );
3447 assert!(
3448 pin.unavailable_reason
3449 .as_deref()
3450 .is_some_and(|reason| reason.contains("disabled")),
3451 "the reason must name the disabled sidecar, not a generic enumeration failure: \
3452 {pin:?}"
3453 );
3454 assert!(pin.sidecar_entries.is_empty());
3455 assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
3456 }
3457}