1#[path = "diagnostics/disk_guard.rs"]
64mod disk_guard;
65#[path = "diagnostics/writer_contention.rs"]
66mod writer_contention;
67
68use std::path::{Path, PathBuf};
69use std::sync::atomic::{AtomicBool, Ordering};
70use std::sync::Arc;
71use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
72
73use khive_storage::error::StorageError;
74use khive_storage::types::StorageResult;
75use khive_storage::StorageCapability;
76use rusqlite::Connection;
77use serde::Serialize;
78
79use crate::checkpoint;
80use crate::pool::{ConnectionPool, WalCeilingSource};
81
82pub use disk_guard::DiskGuardDiagnostics;
83
84#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
92pub struct CheckpointProbe {
93 pub busy: i64,
94 pub log_frames: i64,
95 pub checkpointed_frames: i64,
96}
97
98impl CheckpointProbe {
99 pub fn backfill_gap_frames(&self) -> i64 {
103 if self.log_frames < 0 || self.checkpointed_frames < 0 {
104 return 0;
105 }
106 self.log_frames
107 .saturating_sub(self.checkpointed_frames)
108 .max(0)
109 }
110}
111
112pub fn checkpoint_probe(conn: &Connection) -> rusqlite::Result<CheckpointProbe> {
126 conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
127 Ok(CheckpointProbe {
128 busy: row.get(0)?,
129 log_frames: row.get(1)?,
130 checkpointed_frames: row.get(2)?,
131 })
132 })
133}
134
135#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
141pub struct CheckpointCounters {
142 pub last_observed_wal_pages: Option<u64>,
143 pub truncate_attempts: u64,
144 pub truncate_consecutive_failures: u64,
145 pub checkpoint_skipped_ticks: u64,
146 pub checkpoint_consecutive_skips: u64,
147 pub checkpoint_last_skip_wal_pages: Option<u64>,
148 pub checkpoint_pressure_elevated_ticks: u64,
149 pub checkpoint_pressure_episodes_started: u64,
150 pub checkpoint_pressure_episodes_recovered: u64,
151 pub checkpoint_lifecycle_append_attempts: u64,
152 pub checkpoint_lifecycle_append_failures: u64,
153 pub checkpoint_lifecycle_enqueue_drops: u64,
154 pub read_tx_max_age_evictions: u64,
158}
159
160pub fn checkpoint_counters() -> CheckpointCounters {
162 CheckpointCounters {
163 last_observed_wal_pages: checkpoint::last_observed_wal_pages(),
164 truncate_attempts: checkpoint::truncate_attempts(),
165 truncate_consecutive_failures: checkpoint::truncate_consecutive_failures(),
166 checkpoint_skipped_ticks: checkpoint::checkpoint_skipped_ticks(),
167 checkpoint_consecutive_skips: checkpoint::checkpoint_consecutive_skips(),
168 checkpoint_last_skip_wal_pages: checkpoint::checkpoint_last_skip_wal_pages(),
169 checkpoint_pressure_elevated_ticks: checkpoint::checkpoint_pressure_elevated_ticks(),
170 checkpoint_pressure_episodes_started: checkpoint::checkpoint_pressure_episodes_started(),
171 checkpoint_pressure_episodes_recovered: checkpoint::checkpoint_pressure_episodes_recovered(
172 ),
173 checkpoint_lifecycle_append_attempts: checkpoint::checkpoint_lifecycle_append_attempts(),
174 checkpoint_lifecycle_append_failures: checkpoint::checkpoint_lifecycle_append_failures(),
175 checkpoint_lifecycle_enqueue_drops: checkpoint::checkpoint_lifecycle_enqueue_drops(),
176 read_tx_max_age_evictions: checkpoint::read_tx_max_age_evictions(),
177 }
178}
179
180#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
187pub struct BuildIdentity {
188 pub version: String,
189 pub build_hash: Option<String>,
190}
191
192impl BuildIdentity {
193 pub fn from_env(version: &str, build_hash: Option<&str>) -> Self {
195 Self {
196 version: version.to_string(),
197 build_hash: build_hash.map(str::to_string),
198 }
199 }
200}
201
202#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
204pub struct ProcessIdentity {
205 pub pid: u32,
206 pub started_at: Option<i64>,
209 pub started_at_unavailable_reason: Option<String>,
212 pub pool_generation: u64,
213}
214
215impl ProcessIdentity {
216 pub fn current(pool: &ConnectionPool) -> Self {
217 let pid = std::process::id();
218 Self::from_start_time(
219 pid,
220 crate::walpin::process_start_time_secs(pid),
221 pool.main_pool_generation(),
222 )
223 }
224
225 fn from_start_time(pid: u32, started_at: Option<i64>, pool_generation: u64) -> Self {
226 Self {
227 pid,
228 started_at,
229 started_at_unavailable_reason: started_at.is_none().then(|| {
230 format!(
231 "OS process start time is unsupported or unavailable on {}",
232 std::env::consts::OS
233 )
234 }),
235 pool_generation,
236 }
237 }
238}
239
240#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
242pub struct WalFileState {
243 pub wal_path: String,
244 pub wal_size_bytes: Option<u64>,
247 pub unavailable_reason: Option<String>,
248}
249
250pub fn wal_file_state(db_path: &Path) -> WalFileState {
252 let wal_path = wal_sidecar_path(db_path);
253 match std::fs::metadata(&wal_path) {
254 Ok(md) => WalFileState {
255 wal_path: wal_path.display().to_string(),
256 wal_size_bytes: Some(md.len()),
257 unavailable_reason: None,
258 },
259 Err(e) => WalFileState {
260 wal_path: wal_path.display().to_string(),
261 wal_size_bytes: None,
262 unavailable_reason: Some(e.to_string()),
263 },
264 }
265}
266
267fn wal_sidecar_path(db_path: &Path) -> PathBuf {
270 let mut s = db_path.as_os_str().to_os_string();
271 s.push("-wal");
272 PathBuf::from(s)
273}
274
275#[derive(Debug, Clone, PartialEq, Serialize)]
284pub struct WalPinAttribution {
285 pub status: WalPinAttributionStatus,
287 pub status_reasons: Vec<String>,
289 pub census: WalPinCensus,
291 pub reporting_pid: u32,
293 pub reporting_process_is_holder: Option<bool>,
295 pub reporting_process_is_holder_unavailable_reason: Option<String>,
297 pub census_process_start_times: Vec<WalPinCensusProcessStart>,
299 pub start_time_resolution_secs: Option<u64>,
301 pub start_time_resolution_unavailable_reason: Option<String>,
303 pub available: bool,
306 pub unavailable_reason: Option<String>,
307 #[serde(skip_serializing)]
309 pub census_holder_pids: Vec<u32>,
310 #[serde(skip_serializing)]
311 pub census_uninspectable_pids: Vec<u32>,
312 #[serde(skip_serializing)]
313 pub census_truncated: bool,
314 #[serde(skip_serializing)]
315 pub census_is_complete: bool,
316 pub reporting: Vec<WalPinHolder>,
318 pub registered_silent_pids: Vec<u32>,
320 pub unknown_pids: Vec<u32>,
322 pub census_pids_without_attribution: Vec<u32>,
324 pub fully_attributed: bool,
327 pub sidecar_entries: Vec<serde_json::Value>,
329 #[serde(skip_serializing_if = "Option::is_none")]
331 pub sidecar_listing_truncated: Option<bool>,
332 #[serde(skip_serializing_if = "Option::is_none")]
335 pub sidecar_entries_cleanup_would_reap: Option<usize>,
336}
337
338#[derive(Debug, Clone, PartialEq, Serialize)]
340pub struct WalPinCensusProcessStart {
341 pub pid: u32,
343 pub process_start_time_secs: Option<i64>,
345 pub process_start_time_unavailable_reason: Option<String>,
347}
348
349fn census_process_start_times(holder_pids: &[u32]) -> Vec<WalPinCensusProcessStart> {
350 holder_pids
351 .iter()
352 .map(|pid| {
353 let process_start_time_secs = crate::walpin::process_start_time_secs(*pid);
354 let process_start_time_unavailable_reason =
355 process_start_time_secs.is_none().then(|| {
356 if crate::walpin::start_time_resolution_secs().is_none() {
357 "process start time is unavailable on this platform".to_string()
358 } else {
359 "the operating system did not report a start time for this process"
360 .to_string()
361 }
362 });
363 WalPinCensusProcessStart {
364 pid: *pid,
365 process_start_time_secs,
366 process_start_time_unavailable_reason,
367 }
368 })
369 .collect()
370}
371
372#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
377#[serde(rename_all = "snake_case")]
378pub enum WalPinAttributionStatus {
379 Complete,
381 Degraded,
383 Unavailable,
385}
386
387#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
395#[serde(tag = "status", rename_all = "snake_case")]
396pub enum WalPinCensus {
397 Complete {
399 holder_pids: Vec<u32>,
401 },
402 Incomplete {
404 holder_pids: Vec<u32>,
406 uninspectable_pids: Vec<u32>,
408 truncated: bool,
410 reason: String,
412 },
413 Unavailable {
415 reason: String,
417 },
418}
419
420fn census_truncation_cause(budget_exhausted: bool) -> String {
430 if budget_exhausted {
431 "the OS process walk stopped at its wall-clock budget (see collection_cost.wal_pin_census_budget_ms)"
432 .to_string()
433 } else {
434 "the OS process walk was truncated".to_string()
435 }
436}
437
438#[derive(Debug, Clone, PartialEq, Serialize)]
440pub struct WalPinHolder {
441 pub pid: u32,
442 pub process_role: String,
443 pub current_oldest_tx_age_secs: f64,
444 pub oldest_tx_label: Option<String>,
445 pub attribution_is_evidence_backed: bool,
446}
447
448impl WalPinAttribution {
449 fn unavailable(reason: impl Into<String>) -> Self {
450 let reason = reason.into();
451 let start_time_resolution_secs = crate::walpin::start_time_resolution_secs();
452 Self {
453 status: WalPinAttributionStatus::Unavailable,
454 status_reasons: vec![reason.clone()],
455 census: WalPinCensus::Unavailable {
456 reason: reason.clone(),
457 },
458 reporting_pid: crate::walpin::reporting_pid(),
459 reporting_process_is_holder: None,
460 reporting_process_is_holder_unavailable_reason: Some(
461 "the OS holder census supplied no holder evidence".to_string(),
462 ),
463 census_process_start_times: Vec::new(),
464 start_time_resolution_secs,
465 start_time_resolution_unavailable_reason: start_time_resolution_secs
466 .is_none()
467 .then(|| "process start time is unavailable on this platform".to_string()),
468 available: false,
469 unavailable_reason: Some(reason),
470 census_holder_pids: Vec::new(),
471 census_uninspectable_pids: Vec::new(),
472 census_truncated: false,
473 census_is_complete: false,
474 reporting: Vec::new(),
475 registered_silent_pids: Vec::new(),
476 unknown_pids: Vec::new(),
477 census_pids_without_attribution: Vec::new(),
478 fully_attributed: false,
479 sidecar_entries: Vec::new(),
480 sidecar_listing_truncated: None,
481 sidecar_entries_cleanup_would_reap: None,
482 }
483 }
484}
485
486#[cfg(all(unix, test))]
487fn wal_pin_attribution_from_census(census: crate::walpin::CensusResult) -> WalPinAttribution {
488 wal_pin_attribution_without_sidecar(
489 census,
490 "read-only sidecar enumeration did not run for this attribution snapshot".to_string(),
491 )
492}
493
494#[cfg(unix)]
495fn wal_pin_attribution_without_sidecar(
496 census: crate::walpin::CensusResult,
497 sidecar_reason: String,
498) -> WalPinAttribution {
499 let census_is_complete = census.is_complete();
500 let mut census_holder_pids: Vec<u32> = census.holders.iter().copied().collect();
501 census_holder_pids.sort_unstable();
502 let reporting_pid = crate::walpin::reporting_pid();
503 let reporting_process_is_holder = census_holder_pids
504 .contains(&reporting_pid)
505 .then_some(true)
506 .or_else(|| census_is_complete.then_some(false));
507 let reporting_process_is_holder_unavailable_reason =
508 reporting_process_is_holder.is_none().then(|| {
509 "the reporting process was absent from the incomplete OS holder census".to_string()
510 });
511 let census_process_start_times = census_process_start_times(&census_holder_pids);
512 let start_time_resolution_secs = crate::walpin::start_time_resolution_secs();
513 let mut census_uninspectable_pids = census.uninspectable_pids;
514 census_uninspectable_pids.sort_unstable();
515 census_uninspectable_pids.dedup();
516 let census_truncated = census.truncated;
517 let census_budget_exhausted = census.budget_exhausted;
518
519 let mut status_reasons = vec![sidecar_reason];
520 let census = if census_is_complete {
521 WalPinCensus::Complete {
522 holder_pids: census_holder_pids.clone(),
523 }
524 } else {
525 let mut causes = Vec::new();
526 if census_truncated {
527 causes.push(census_truncation_cause(census_budget_exhausted));
528 }
529 if !census_uninspectable_pids.is_empty() {
530 causes.push(format!(
531 "{} PID(s) could not be inspected",
532 census_uninspectable_pids.len()
533 ));
534 }
535 let reason = format!(
536 "OS holder census is incomplete: {}; additional database holders cannot be ruled out",
537 causes.join("; ")
538 );
539 status_reasons.push(reason.clone());
540 WalPinCensus::Incomplete {
541 holder_pids: census_holder_pids.clone(),
542 uninspectable_pids: census_uninspectable_pids.clone(),
543 truncated: census_truncated,
544 reason,
545 }
546 };
547
548 WalPinAttribution {
549 status: WalPinAttributionStatus::Degraded,
550 unavailable_reason: Some(status_reasons.join("; ")),
551 status_reasons,
552 census,
553 reporting_pid,
554 reporting_process_is_holder,
555 reporting_process_is_holder_unavailable_reason,
556 census_process_start_times,
557 start_time_resolution_secs,
558 start_time_resolution_unavailable_reason: start_time_resolution_secs
559 .is_none()
560 .then(|| "process start time is unavailable on this platform".to_string()),
561 available: false,
562 census_holder_pids,
563 census_uninspectable_pids,
564 census_truncated,
565 census_is_complete,
566 reporting: Vec::new(),
567 registered_silent_pids: Vec::new(),
568 unknown_pids: Vec::new(),
569 census_pids_without_attribution: Vec::new(),
570 fully_attributed: false,
571 sidecar_entries: Vec::new(),
572 sidecar_listing_truncated: None,
573 sidecar_entries_cleanup_would_reap: None,
574 }
575}
576
577#[cfg(unix)]
578fn wal_pin_attribution_from_evidence(
579 census: crate::walpin::CensusResult,
580 sidecar: crate::walpin::WalpinReport,
581) -> WalPinAttribution {
582 use std::collections::BTreeSet;
583
584 let census_is_complete = census.is_complete();
585 let mut census_holder_pids: Vec<u32> = census.holders.iter().copied().collect();
586 census_holder_pids.sort_unstable();
587 let reporting_pid = crate::walpin::reporting_pid();
588 let reporting_process_is_holder = census_holder_pids
589 .contains(&reporting_pid)
590 .then_some(true)
591 .or_else(|| census_is_complete.then_some(false));
592 let reporting_process_is_holder_unavailable_reason =
593 reporting_process_is_holder.is_none().then(|| {
594 "the reporting process was absent from the incomplete OS holder census".to_string()
595 });
596 let census_process_start_times = census_process_start_times(&census_holder_pids);
597 let start_time_resolution_secs = crate::walpin::start_time_resolution_secs();
598 let mut census_uninspectable_pids = census.uninspectable_pids;
599 census_uninspectable_pids.sort_unstable();
600 census_uninspectable_pids.dedup();
601 let census_truncated = census.truncated;
602 let census_budget_exhausted = census.budget_exhausted;
603 let census_carrier = if census_is_complete {
604 WalPinCensus::Complete {
605 holder_pids: census_holder_pids.clone(),
606 }
607 } else {
608 let mut causes = Vec::new();
609 if census_truncated {
610 causes.push(census_truncation_cause(census_budget_exhausted));
611 }
612 if !census_uninspectable_pids.is_empty() {
613 causes.push(format!(
614 "{} PID(s) could not be inspected",
615 census_uninspectable_pids.len()
616 ));
617 }
618 let reason = format!(
619 "OS holder census is incomplete: {}; additional database holders cannot be ruled out",
620 causes.join("; ")
621 );
622 WalPinCensus::Incomplete {
623 holder_pids: census_holder_pids.clone(),
624 uninspectable_pids: census_uninspectable_pids.clone(),
625 truncated: census_truncated,
626 reason,
627 }
628 };
629
630 let now_epoch_secs = SystemTime::now()
631 .duration_since(UNIX_EPOCH)
632 .map(|duration| duration.as_secs() as i64)
633 .unwrap_or(0);
634 let sidecar_listing_truncated = sidecar.sidecar_listing_truncated;
635 let sidecar_entries_cleanup_would_reap = sidecar.cleanup_would_reap;
636 let mut reporting = Vec::new();
637 let mut registered_silent_pids = Vec::new();
638 let mut unknown_pids = Vec::new();
639 let mut sidecar_entries = Vec::new();
640 let mut sidecar_known_pids = BTreeSet::new();
641
642 for entry in sidecar.entries {
643 match entry {
644 crate::walpin::WalpinPidHealth::Reporting(heartbeat) => {
645 let current_oldest_tx_age_secs =
646 heartbeat.current_oldest_tx_age_secs(now_epoch_secs);
647 let attribution_is_evidence_backed = heartbeat.attribution_is_evidence_backed();
648 sidecar_known_pids.insert(heartbeat.pid);
649 reporting.push(WalPinHolder {
650 pid: heartbeat.pid,
651 process_role: heartbeat.process_role.clone(),
652 current_oldest_tx_age_secs,
653 oldest_tx_label: heartbeat.oldest_tx_label.clone(),
654 attribution_is_evidence_backed,
655 });
656 sidecar_entries.push((
657 heartbeat.pid,
658 0u8,
659 serde_json::json!({
660 "pid": heartbeat.pid,
661 "status": "reporting",
662 "process_role": heartbeat.process_role,
663 "current_oldest_tx_age_secs": current_oldest_tx_age_secs,
664 "oldest_tx_label": heartbeat.oldest_tx_label,
665 "attribution_is_evidence_backed": attribution_is_evidence_backed,
666 }),
667 ));
668 }
669 crate::walpin::WalpinPidHealth::RegisteredSilent { pid } => {
670 sidecar_known_pids.insert(pid);
671 registered_silent_pids.push(pid);
672 sidecar_entries.push((
673 pid,
674 1u8,
675 serde_json::json!({"pid": pid, "status": "registered_silent"}),
676 ));
677 }
678 crate::walpin::WalpinPidHealth::Unknown { pid, reason } => {
679 sidecar_known_pids.insert(pid);
680 unknown_pids.push(pid);
681 sidecar_entries.push((
682 pid,
683 2u8,
684 serde_json::json!({"pid": pid, "status": "unknown", "reason": reason}),
685 ));
686 }
687 }
688 }
689
690 reporting.sort_by_key(|holder| holder.pid);
691 reporting.dedup_by_key(|holder| holder.pid);
692 registered_silent_pids.sort_unstable();
693 registered_silent_pids.dedup();
694 unknown_pids.sort_unstable();
695 unknown_pids.dedup();
696 sidecar_entries.sort_by_key(|(pid, status_rank, _)| (*pid, *status_rank));
697 let sidecar_entries = sidecar_entries
698 .into_iter()
699 .map(|(_, _, entry)| entry)
700 .collect();
701 let census_pids_without_attribution: Vec<u32> = census_holder_pids
702 .iter()
703 .copied()
704 .filter(|pid| !sidecar_known_pids.contains(pid))
705 .collect();
706
707 let mut status_reasons = Vec::new();
708 if let WalPinCensus::Incomplete { reason, .. } = &census_carrier {
709 status_reasons.push(reason.clone());
710 }
711 if sidecar_listing_truncated {
712 status_reasons.push(
713 "read-only sidecar enumeration reached its entry cap; additional entries may exist"
714 .to_string(),
715 );
716 }
717 if !unknown_pids.is_empty() {
718 status_reasons.push(format!(
719 "{} sidecar PID(s) could not be classified conclusively",
720 unknown_pids.len()
721 ));
722 }
723 if !census_pids_without_attribution.is_empty() {
724 status_reasons.push(format!(
725 "{} OS-confirmed holder(s) have no sidecar attribution",
726 census_pids_without_attribution.len()
727 ));
728 }
729
730 let fully_attributed = census_is_complete
731 && !sidecar_listing_truncated
732 && unknown_pids.is_empty()
733 && census_pids_without_attribution.is_empty();
734 let status = if fully_attributed {
735 WalPinAttributionStatus::Complete
736 } else {
737 WalPinAttributionStatus::Degraded
738 };
739 let unavailable_reason = (!fully_attributed).then(|| status_reasons.join("; "));
740
741 WalPinAttribution {
742 status,
743 status_reasons,
744 census: census_carrier,
745 reporting_pid,
746 reporting_process_is_holder,
747 reporting_process_is_holder_unavailable_reason,
748 census_process_start_times,
749 start_time_resolution_secs,
750 start_time_resolution_unavailable_reason: start_time_resolution_secs
751 .is_none()
752 .then(|| "process start time is unavailable on this platform".to_string()),
753 available: fully_attributed,
754 unavailable_reason,
755 census_holder_pids,
756 census_uninspectable_pids,
757 census_truncated,
758 census_is_complete,
759 reporting,
760 registered_silent_pids,
761 unknown_pids,
762 census_pids_without_attribution,
763 fully_attributed,
764 sidecar_entries,
765 sidecar_listing_truncated: Some(sidecar_listing_truncated),
766 sidecar_entries_cleanup_would_reap: Some(sidecar_entries_cleanup_would_reap),
767 }
768}
769
770#[cfg(unix)]
779pub fn wal_pin_attribution(db_path: &Path, sweep_interval: Duration) -> WalPinAttribution {
780 use crate::walpin;
781
782 let census = match walpin::census_holders(db_path) {
783 Ok(c) => c,
784 Err(e) => return WalPinAttribution::unavailable(format!("census_holders failed: {e}")),
785 };
786 if !walpin::sidecar_enabled(true) {
787 return wal_pin_attribution_without_sidecar(census, SIDECAR_DISABLED_REASON.to_string());
788 }
789 match walpin::inspect_live(&walpin::sidecar_dir_for(db_path), sweep_interval) {
790 Ok(sidecar) => wal_pin_attribution_from_evidence(census, sidecar),
791 Err(error) => wal_pin_attribution_without_sidecar(
792 census,
793 format!("read-only sidecar enumeration failed: {error}"),
794 ),
795 }
796}
797
798#[cfg(unix)]
800const SIDECAR_DISABLED_REASON: &str =
801 "walpin sidecar is explicitly disabled (KHIVE_WALPIN_SIDECAR); attribution has no sidecar \
802 evidence to reconcile against the OS holder census";
803
804#[cfg(not(unix))]
805pub fn wal_pin_attribution(_db_path: &Path, _sweep_interval: Duration) -> WalPinAttribution {
806 WalPinAttribution::unavailable("WAL-pin attribution requires a Unix platform")
807}
808
809#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
816pub struct ReaderContentionDiagnostics {
817 pub configured_reader_cap: usize,
820 pub configured_checkout_timeout_ms: u64,
822 pub configured_busy_timeout_ms: u64,
825 pub reader_admission_capacity: usize,
828 pub available_reader_admission_slots: usize,
830 pub reader_acquisitions: u64,
833 pub pooled_reader_checkouts: u64,
835 pub standalone_reader_opens: u64,
838 pub infrastructure_standalone_reader_opens: u64,
841 pub reader_checkout_timeouts: u64,
844 pub reader_busy_timeouts: u64,
850 pub active_pooled_reader_checkouts: u64,
852 pub peak_active_pooled_reader_checkouts: u64,
854 pub completed_pooled_reader_checkouts: u64,
856 pub max_completed_reader_hold_micros: u64,
858 pub max_completed_reader_hold_operation: Option<&'static str>,
863 pub reader_replacement_open_failures: u64,
868}
869
870impl ReaderContentionDiagnostics {
871 fn snapshot(pool: &ConnectionPool) -> Self {
872 let reader = pool.reader_acquisition_snapshot();
873 Self {
874 configured_reader_cap: pool.config().max_readers,
875 configured_checkout_timeout_ms: u64::try_from(
876 pool.config().checkout_timeout.as_millis(),
877 )
878 .unwrap_or(u64::MAX),
879 configured_busy_timeout_ms: u64::try_from(pool.config().busy_timeout.as_millis())
880 .unwrap_or(u64::MAX),
881 reader_admission_capacity: reader.reader_admission_capacity,
882 available_reader_admission_slots: reader.available_reader_admission_slots,
883 reader_acquisitions: reader.acquisitions,
884 pooled_reader_checkouts: reader.pooled_checkouts,
885 standalone_reader_opens: reader.standalone_opens,
886 infrastructure_standalone_reader_opens: reader.infrastructure_standalone_opens,
887 reader_checkout_timeouts: reader.checkout_timeouts,
888 reader_busy_timeouts: reader.busy_timeouts,
889 active_pooled_reader_checkouts: reader.active_pooled_checkouts,
890 peak_active_pooled_reader_checkouts: reader.peak_active_pooled_checkouts,
891 completed_pooled_reader_checkouts: reader.completed_pooled_checkouts,
892 max_completed_reader_hold_micros: reader.max_completed_hold_micros,
893 max_completed_reader_hold_operation: reader.max_completed_hold_operation,
894 reader_replacement_open_failures: reader.reader_replacement_open_failures,
895 }
896 }
897}
898
899#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
914pub struct WriterContentionDiagnostics {
915 pub writer_acquisitions: u64,
918 pub pooled_writer_acquisitions: u64,
920 pub standalone_writer_acquisitions: u64,
922 pub writer_task_acquisitions: u64,
924 pub writer_acquisition_timeouts: u64,
926 pub writer_lease_timeouts: u64,
932 pub configured_guard_deadline_ms: Option<u64>,
935 pub configured_checkout_timeout_ms: u64,
938 pub effective_writer_wait_bound_ms: u64,
943 pub direct_writer_busy_refusals: u64,
946 pub writer_task_begin_busy: u64,
951 pub writer_task_begin_busy_absorbed: u64,
956 pub writer_task_begin_errors: u64,
959 pub writer_task_request_failures: u64,
964 pub writer_task_side_effects_unknown: u64,
967 pub audit_append_failures: Option<u64>,
979 pub audit_append_failures_unavailable_reason: Option<String>,
981 pub audit_obligation_append_failures: Option<u64>,
986 pub audit_obligation_append_failures_unavailable_reason: Option<String>,
988 pub audit_batch_flush_failures: Option<u64>,
993 pub audit_batch_flush_failures_unavailable_reason: Option<String>,
995 pub audit_degraded_rows: Option<u64>,
998 pub audit_degraded_rows_unavailable_reason: Option<String>,
1000 pub audit_degraded: Option<bool>,
1006 pub audit_degraded_unavailable_reason: Option<String>,
1008 pub audit_admission_refused_obligations: Option<u64>,
1028 pub audit_admission_refused_obligations_last_at_ms: Option<u64>,
1034 pub audit_admission_refused_obligations_unavailable_reason: Option<String>,
1036 pub audit_admission_unresolved_obligations: Option<u64>,
1057 pub audit_admission_unresolved_obligations_last_at_ms: Option<u64>,
1061 pub audit_admission_unresolved_obligations_unavailable_reason: Option<String>,
1064}
1065
1066#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1071pub struct RuntimeAuditBatchMetrics {
1072 pub flush_failures: u64,
1075 pub degraded_rows: u64,
1077 pub degraded: bool,
1079 pub admission_refused_obligations: u64,
1084 pub admission_refused_obligations_last_at_ms: Option<u64>,
1087 pub admission_unresolved_obligations: u64,
1094 pub admission_unresolved_obligations_last_at_ms: Option<u64>,
1097}
1098
1099#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1115pub struct GraphEdgeIntegrity {
1116 pub duplicate_edge_id_groups: i64,
1117 pub graph_edges_rows: i64,
1118 pub graph_edges_seq_rows: i64,
1119 pub pre_v14_duplicate_edge_state_detected: bool,
1120 pub live_entities_carrying_merged_into: i64,
1121}
1122
1123fn graph_edge_integrity(conn: &Connection) -> rusqlite::Result<GraphEdgeIntegrity> {
1124 conn.query_row(
1125 "SELECT
1126 (SELECT COUNT(*) FROM (
1127 SELECT id FROM graph_edges GROUP BY id HAVING COUNT(*) > 1
1128 )),
1129 (SELECT COUNT(*) FROM graph_edges),
1130 (SELECT COUNT(*) FROM graph_edges_seq),
1131 (SELECT COUNT(*) FROM entities
1132 WHERE deleted_at IS NULL AND merged_into IS NOT NULL)",
1133 [],
1134 |row| {
1135 let duplicate_edge_id_groups = row.get(0)?;
1136 Ok(GraphEdgeIntegrity {
1137 duplicate_edge_id_groups,
1138 graph_edges_rows: row.get(1)?,
1139 graph_edges_seq_rows: row.get(2)?,
1140 pre_v14_duplicate_edge_state_detected: duplicate_edge_id_groups > 0,
1141 live_entities_carrying_merged_into: row.get(3)?,
1142 })
1143 },
1144 )
1145}
1146
1147const MAX_DATABASE_SIZE_OBJECTS: usize = 4_096;
1148
1149#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1151#[serde(rename_all = "snake_case")]
1152pub enum DatabaseObjectKind {
1153 Table,
1154 Index,
1155 Internal,
1156}
1157
1158#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1163#[serde(rename_all = "snake_case")]
1164pub enum DatabaseStorageClass {
1165 RowTable,
1166 Index,
1167 FullText,
1168 Vector,
1169 MixedRowAndEmbedding,
1170 Internal,
1171}
1172
1173#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
1174pub struct DatabaseObjectSize {
1175 pub name: String,
1176 pub owner_table: Option<String>,
1177 pub object_kind: DatabaseObjectKind,
1178 pub storage_class: DatabaseStorageClass,
1179 pub pages: u64,
1180 pub bytes: u64,
1181}
1182
1183#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
1186pub struct DatabaseSizeComposition {
1187 pub page_size_bytes: u64,
1188 pub page_count: u64,
1189 pub freelist_pages: u64,
1190 pub database_bytes: u64,
1191 pub freelist_bytes: u64,
1192 pub accounted_bytes: u64,
1193 pub unaccounted_bytes: u64,
1194 pub row_table_bytes: u64,
1195 pub index_bytes: u64,
1196 pub full_text_bytes: u64,
1197 pub vector_bytes: u64,
1198 pub mixed_embedding_bytes: u64,
1199 pub internal_bytes: u64,
1200 pub objects: Vec<DatabaseObjectSize>,
1201 pub objects_truncated: bool,
1202 pub objects_omitted: usize,
1203}
1204
1205fn nonnegative_sqlite_integer(column: usize, value: i64) -> rusqlite::Result<u64> {
1206 u64::try_from(value).map_err(|_| rusqlite::Error::IntegralValueOutOfRange(column, value))
1207}
1208
1209fn declares_embedding_blob(sql: Option<&str>) -> bool {
1210 let Some(sql) = sql else {
1211 return false;
1212 };
1213 let tokens: Vec<_> = sql
1214 .split(|character: char| !(character.is_ascii_alphanumeric() || character == '_'))
1215 .filter(|token| !token.is_empty())
1216 .collect();
1217 tokens.windows(2).any(|pair| {
1218 pair[0].eq_ignore_ascii_case("embedding") && pair[1].eq_ignore_ascii_case("blob")
1219 })
1220}
1221
1222fn classify_database_object(
1223 name: &str,
1224 sqlite_type: &str,
1225 sql: Option<&str>,
1226) -> (DatabaseObjectKind, DatabaseStorageClass) {
1227 let object_kind = match sqlite_type {
1228 "table" => DatabaseObjectKind::Table,
1229 "index" => DatabaseObjectKind::Index,
1230 _ => DatabaseObjectKind::Internal,
1231 };
1232 let lower_name = name.to_ascii_lowercase();
1233 let lower_sql = sql.unwrap_or_default().to_ascii_lowercase();
1234 let storage_class = if lower_name.starts_with("fts_") || lower_sql.contains("using fts5") {
1235 DatabaseStorageClass::FullText
1236 } else if lower_name.starts_with("vec_")
1237 || lower_name == "_embedding_models"
1238 || lower_sql.contains("using vec0")
1239 {
1240 DatabaseStorageClass::Vector
1241 } else if declares_embedding_blob(sql) {
1242 DatabaseStorageClass::MixedRowAndEmbedding
1243 } else if object_kind == DatabaseObjectKind::Index {
1244 DatabaseStorageClass::Index
1245 } else if name.starts_with("sqlite_") || object_kind == DatabaseObjectKind::Internal {
1246 DatabaseStorageClass::Internal
1247 } else {
1248 DatabaseStorageClass::RowTable
1249 };
1250 (object_kind, storage_class)
1251}
1252
1253fn database_size_composition(conn: &Connection) -> rusqlite::Result<DatabaseSizeComposition> {
1254 let page_size = nonnegative_sqlite_integer(
1255 0,
1256 conn.query_row("PRAGMA page_size", [], |row| row.get::<_, i64>(0))?,
1257 )?;
1258 let page_count = nonnegative_sqlite_integer(
1259 0,
1260 conn.query_row("PRAGMA page_count", [], |row| row.get::<_, i64>(0))?,
1261 )?;
1262 let freelist_pages = nonnegative_sqlite_integer(
1263 0,
1264 conn.query_row("PRAGMA freelist_count", [], |row| row.get::<_, i64>(0))?,
1265 )?;
1266
1267 let mut statement = conn.prepare(
1268 "SELECT d.name, COALESCE(s.type, 'internal'), s.tbl_name, s.sql, d.pageno, d.pgsize
1269 FROM dbstat AS d
1270 LEFT JOIN sqlite_schema AS s ON s.name = d.name
1271 WHERE d.aggregate = TRUE
1272 ORDER BY d.name",
1273 )?;
1274 let mut rows = statement.query([])?;
1275 let mut objects = Vec::new();
1276 let mut objects_omitted = 0usize;
1277 let mut accounted_bytes = 0u64;
1278 let mut row_table_bytes = 0u64;
1279 let mut index_bytes = 0u64;
1280 let mut full_text_bytes = 0u64;
1281 let mut vector_bytes = 0u64;
1282 let mut mixed_embedding_bytes = 0u64;
1283 let mut internal_bytes = 0u64;
1284
1285 while let Some(row) = rows.next()? {
1286 let name: String = row.get(0)?;
1287 let sqlite_type: String = row.get(1)?;
1288 let owner_table: Option<String> = row.get(2)?;
1289 let sql: Option<String> = row.get(3)?;
1290 let pages = nonnegative_sqlite_integer(4, row.get(4)?)?;
1291 let bytes = nonnegative_sqlite_integer(5, row.get(5)?)?;
1292 let (object_kind, storage_class) =
1293 classify_database_object(&name, &sqlite_type, sql.as_deref());
1294 accounted_bytes = accounted_bytes.saturating_add(bytes);
1295 let class_total = match storage_class {
1296 DatabaseStorageClass::RowTable => &mut row_table_bytes,
1297 DatabaseStorageClass::Index => &mut index_bytes,
1298 DatabaseStorageClass::FullText => &mut full_text_bytes,
1299 DatabaseStorageClass::Vector => &mut vector_bytes,
1300 DatabaseStorageClass::MixedRowAndEmbedding => &mut mixed_embedding_bytes,
1301 DatabaseStorageClass::Internal => &mut internal_bytes,
1302 };
1303 *class_total = class_total.saturating_add(bytes);
1304
1305 if objects.len() < MAX_DATABASE_SIZE_OBJECTS {
1306 objects.push(DatabaseObjectSize {
1307 name,
1308 owner_table,
1309 object_kind,
1310 storage_class,
1311 pages,
1312 bytes,
1313 });
1314 } else {
1315 objects_omitted = objects_omitted.saturating_add(1);
1316 }
1317 }
1318
1319 let database_bytes = page_count.saturating_mul(page_size);
1320 let freelist_bytes = freelist_pages.saturating_mul(page_size);
1321 let unaccounted_bytes = database_bytes
1322 .saturating_sub(freelist_bytes)
1323 .saturating_sub(accounted_bytes);
1324 Ok(DatabaseSizeComposition {
1325 page_size_bytes: page_size,
1326 page_count,
1327 freelist_pages,
1328 database_bytes,
1329 freelist_bytes,
1330 accounted_bytes,
1331 unaccounted_bytes,
1332 row_table_bytes,
1333 index_bytes,
1334 full_text_bytes,
1335 vector_bytes,
1336 mixed_embedding_bytes,
1337 internal_bytes,
1338 objects,
1339 objects_truncated: objects_omitted > 0,
1340 objects_omitted,
1341 })
1342}
1343
1344#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1357pub struct CollectionCost {
1358 pub total_ms: u64,
1360 pub sqlite_ms: u64,
1363 pub wal_file_stat_ms: u64,
1366 pub wal_pin_census_ms: u64,
1368 pub wal_pin_sidecar_ms: u64,
1370 pub wal_pin_census_budget_ms: Option<u64>,
1375 pub wal_pin_census_budget_exhausted: bool,
1379}
1380
1381impl CollectionCost {
1382 fn in_memory(total_ms: u64) -> Self {
1386 Self {
1387 total_ms,
1388 sqlite_ms: 0,
1389 wal_file_stat_ms: 0,
1390 wal_pin_census_ms: 0,
1391 wal_pin_sidecar_ms: 0,
1392 wal_pin_census_budget_ms: None,
1393 wal_pin_census_budget_exhausted: false,
1394 }
1395 }
1396}
1397
1398const DEFAULT_CENSUS_BUDGET: Duration = Duration::from_millis(2000);
1405
1406const CENSUS_BUDGET_ENV: &str = "KHIVE_WALPIN_CENSUS_BUDGET_MS";
1416
1417fn request_census_budget() -> Option<Duration> {
1418 match std::env::var(CENSUS_BUDGET_ENV) {
1419 Ok(raw) => match raw.trim().parse::<u64>() {
1420 Ok(0) => None,
1421 Ok(ms) => Some(Duration::from_millis(ms)),
1422 Err(_) => Some(DEFAULT_CENSUS_BUDGET),
1423 },
1424 Err(_) => Some(DEFAULT_CENSUS_BUDGET),
1425 }
1426}
1427
1428#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1431pub struct WalCeilingDiagnostics {
1432 pub configured_bytes: u64,
1434 pub effective_bytes: u64,
1436 pub source: WalCeilingSource,
1437 pub enabled: bool,
1438 pub status: &'static str,
1440}
1441
1442impl WalCeilingDiagnostics {
1443 fn from_pool(pool: &ConnectionPool) -> Self {
1444 let policy = pool.config().wal_ceiling;
1445 let read_only = pool.config().read_only;
1446 let effective_bytes = policy.effective_bytes(read_only);
1447 Self {
1448 configured_bytes: policy.bytes,
1449 effective_bytes,
1450 source: policy.source,
1451 enabled: effective_bytes > 0,
1452 status: if effective_bytes > 0 {
1453 "enforced"
1454 } else if read_only && policy.bytes > 0 {
1455 "read_only_not_enforced"
1456 } else {
1457 "disabled"
1458 },
1459 }
1460 }
1461}
1462
1463#[derive(Debug, Clone, PartialEq, Serialize)]
1464pub struct DbDiagnostics {
1465 pub build: BuildIdentity,
1466 pub process: ProcessIdentity,
1467 pub note_search_ann_route_total: u64,
1470 pub note_search_fallback_route_total: u64,
1471 pub search_mechanism: crate::pool::SearchMechanismSnapshot,
1474 pub db_path: Option<String>,
1477 pub wal_ceiling: WalCeilingDiagnostics,
1479 pub disk_guard: DiskGuardDiagnostics,
1480 pub wal_file: Option<WalFileState>,
1481 pub checkpoint_counters: CheckpointCounters,
1482 pub checkpoint_probe: Option<CheckpointProbe>,
1483 pub checkpoint_probe_error: Option<String>,
1484 #[serde(flatten)]
1485 pub checkpoint_pin: CheckpointPinDiagnostics,
1487 pub reader_contention: ReaderContentionDiagnostics,
1489 pub writer_contention: WriterContentionDiagnostics,
1491 pub size_composition: Option<DatabaseSizeComposition>,
1492 pub size_composition_error: Option<String>,
1493 pub graph_edge_integrity: Option<GraphEdgeIntegrity>,
1494 pub graph_edge_integrity_error: Option<String>,
1495 pub fts_segments: Option<crate::FtsSegmentDiagnostics>,
1499 pub fts_segments_error: Option<String>,
1500 pub fts_maintenance: crate::FtsMaintenanceCounters,
1503 pub wal_pin: WalPinAttribution,
1504 pub collection_cost: CollectionCost,
1506}
1507
1508#[derive(Debug, Clone, PartialEq, Serialize)]
1510pub struct CheckpointPinDiagnostics {
1511 pub backfill_ceiling: Option<i64>,
1513 pub backfill_ceiling_unavailable_reason: Option<String>,
1515 pub oldest_pinned_frame: Option<i64>,
1517 pub oldest_pinned_frame_unavailable_reason: Option<String>,
1519 pub oldest_pinned_frame_run: Option<crate::checkpoint::CheckpointRun>,
1521 pub oldest_pinned_frame_run_unavailable_reason: Option<String>,
1523 pub pin_depth: Option<i64>,
1525 pub pin_depth_unavailable_reason: Option<String>,
1527}
1528
1529fn checkpoint_pin_diagnostics(
1530 pool: &ConnectionPool,
1531 probe: Option<&CheckpointProbe>,
1532 probe_error: Option<&str>,
1533) -> CheckpointPinDiagnostics {
1534 let (run_status, run_age) = checkpoint::checkpoint_run_snapshot(pool);
1535 checkpoint_pin_diagnostics_for_run(probe, probe_error, run_status, run_age)
1536}
1537
1538fn checkpoint_pin_diagnostics_for_run(
1539 probe: Option<&CheckpointProbe>,
1540 probe_error: Option<&str>,
1541 run_status: checkpoint::CheckpointRunStatus,
1542 run_age: Option<Duration>,
1543) -> CheckpointPinDiagnostics {
1544 let (backfill_ceiling, backfill_reason) = match probe {
1545 None => (
1546 None,
1547 Some(
1548 probe_error
1549 .unwrap_or("checkpoint probe did not return a row")
1550 .to_string(),
1551 ),
1552 ),
1553 Some(probe) if probe.busy != 0 => (
1554 None,
1555 Some(format!("PASSIVE checkpoint returned busy={}", probe.busy)),
1556 ),
1557 Some(probe) if probe.log_frames < 0 || probe.checkpointed_frames < 0 => (
1558 None,
1559 Some("PASSIVE checkpoint returned a negative frame count".to_string()),
1560 ),
1561 Some(probe) if probe.log_frames <= probe.checkpointed_frames => (
1562 None,
1563 Some("PASSIVE checkpoint found no frames beyond the backfill ceiling".to_string()),
1564 ),
1565 Some(probe) => (Some(probe.checkpointed_frames), None),
1566 };
1567
1568 let (oldest_pinned_frame, oldest_reason, oldest_pinned_frame_run, run_reason) =
1569 match backfill_ceiling {
1570 None => {
1571 let reason = backfill_reason
1572 .clone()
1573 .unwrap_or_else(|| "backfill ceiling is unavailable".to_string());
1574 (None, Some(reason.clone()), None, Some(reason))
1575 }
1576 Some(ceiling) => match run_status {
1577 checkpoint::CheckpointRunStatus::NoTask => {
1578 let reason = "no checkpoint task in this process".to_string();
1579 (None, Some(reason.clone()), None, Some(reason))
1580 }
1581 checkpoint::CheckpointRunStatus::NoObservation => {
1582 let reason = "no checkpoint run has been observed for this backend".to_string();
1583 (None, Some(reason.clone()), None, Some(reason))
1584 }
1585 checkpoint::CheckpointRunStatus::Observed(run) if run.frame != ceiling => {
1586 let reason = "checkpoint run frame does not match the current backfill ceiling";
1587 (
1588 None,
1589 Some(reason.to_string()),
1590 None,
1591 Some(reason.to_string()),
1592 )
1593 }
1594 checkpoint::CheckpointRunStatus::Observed(_)
1595 if run_age.is_none_or(|age| age < Duration::from_secs(1)) =>
1596 {
1597 let reason = "checkpoint run has been observed for less than one second";
1598 (
1599 None,
1600 Some(reason.to_string()),
1601 None,
1602 Some(reason.to_string()),
1603 )
1604 }
1605 checkpoint::CheckpointRunStatus::Observed(run) => {
1606 (Some(ceiling), None, Some(run), None)
1607 }
1608 },
1609 };
1610
1611 let pin_depth = oldest_pinned_frame.map(|frame| {
1612 probe
1613 .expect("a reported pin always has a probe row")
1614 .log_frames
1615 .saturating_sub(frame)
1616 .max(0)
1617 });
1618 let pin_depth_reason = if pin_depth.is_none() {
1619 oldest_reason.clone()
1620 } else {
1621 None
1622 };
1623
1624 CheckpointPinDiagnostics {
1625 backfill_ceiling,
1626 backfill_ceiling_unavailable_reason: backfill_reason,
1627 oldest_pinned_frame,
1628 oldest_pinned_frame_unavailable_reason: oldest_reason,
1629 oldest_pinned_frame_run,
1630 oldest_pinned_frame_run_unavailable_reason: run_reason,
1631 pin_depth,
1632 pin_depth_unavailable_reason: pin_depth_reason,
1633 }
1634}
1635
1636pub fn collect(
1656 pool: &ConnectionPool,
1657 build: BuildIdentity,
1658 sweep_interval: Duration,
1659) -> DbDiagnostics {
1660 collect_inner(pool, build, sweep_interval, None, None)
1661}
1662
1663pub fn collect_with_audit_append_failures(
1666 pool: &ConnectionPool,
1667 build: BuildIdentity,
1668 sweep_interval: Duration,
1669 audit_append_failures: u64,
1670) -> DbDiagnostics {
1671 collect_inner(
1672 pool,
1673 build,
1674 sweep_interval,
1675 Some(audit_append_failures),
1676 None,
1677 )
1678}
1679
1680pub async fn collect_with_audit_append_failures_interruptibly(
1690 pool: Arc<ConnectionPool>,
1691 build: BuildIdentity,
1692 sweep_interval: Duration,
1693 audit_append_failures: u64,
1694) -> StorageResult<DbDiagnostics> {
1695 collect_with_runtime_audit_metrics_interruptibly(
1696 pool,
1697 build,
1698 sweep_interval,
1699 audit_append_failures,
1700 None,
1701 )
1702 .await
1703}
1704
1705pub async fn collect_with_runtime_audit_metrics_interruptibly(
1712 pool: Arc<ConnectionPool>,
1713 build: BuildIdentity,
1714 sweep_interval: Duration,
1715 audit_append_failures: u64,
1716 runtime_audit_batch_metrics: Option<RuntimeAuditBatchMetrics>,
1717) -> StorageResult<DbDiagnostics> {
1718 let process = ProcessIdentity::current(&pool);
1719 collect_with_runtime_audit_metrics_for_process_interruptibly(
1720 pool,
1721 build,
1722 process,
1723 sweep_interval,
1724 audit_append_failures,
1725 runtime_audit_batch_metrics,
1726 )
1727 .await
1728}
1729
1730pub async fn collect_with_runtime_audit_metrics_for_process_interruptibly(
1734 pool: Arc<ConnectionPool>,
1735 build: BuildIdentity,
1736 process: ProcessIdentity,
1737 sweep_interval: Duration,
1738 audit_append_failures: u64,
1739 runtime_audit_batch_metrics: Option<RuntimeAuditBatchMetrics>,
1740) -> StorageResult<DbDiagnostics> {
1741 crate::ensure_request_read_active("db_diagnostics")?;
1742 let started = Instant::now();
1743 let counters = checkpoint_counters();
1744 let reader_contention = ReaderContentionDiagnostics::snapshot(&pool);
1745 let search_mechanism = pool.search_mechanism_snapshot();
1746 let writer_contention = WriterContentionDiagnostics::snapshot(
1747 &pool,
1748 Some(audit_append_failures),
1749 runtime_audit_batch_metrics,
1750 );
1751
1752 let Some(path) = pool.config().path.clone() else {
1753 crate::ensure_request_read_active("db_diagnostics")?;
1754 return Ok(DbDiagnostics {
1755 build,
1756 process,
1757 note_search_ann_route_total: 0,
1758 note_search_fallback_route_total: 0,
1759 search_mechanism,
1760 db_path: None,
1761 wal_ceiling: WalCeilingDiagnostics::from_pool(&pool),
1762 disk_guard: DiskGuardDiagnostics::snapshot(&pool),
1763 wal_file: None,
1764 checkpoint_counters: counters,
1765 checkpoint_probe: None,
1766 checkpoint_probe_error: Some(
1767 "in-memory database: no WAL file and no checkpoint to probe".to_string(),
1768 ),
1769 checkpoint_pin: checkpoint_pin_diagnostics(
1770 &pool,
1771 None,
1772 Some("in-memory database: no WAL file and no checkpoint to probe"),
1773 ),
1774 reader_contention,
1775 writer_contention,
1776 size_composition: None,
1777 size_composition_error: Some(
1778 "in-memory database: no file-backed page composition to inspect".to_string(),
1779 ),
1780 graph_edge_integrity: None,
1781 graph_edge_integrity_error: Some(
1782 "in-memory database: no durable graph-edge ledger to inspect".to_string(),
1783 ),
1784 fts_segments: None,
1785 fts_segments_error: Some(
1786 "in-memory database: no durable FTS5 indexes to inspect".to_string(),
1787 ),
1788 fts_maintenance: crate::fts_maintenance_counters(),
1789 wal_pin: WalPinAttribution::unavailable(
1790 "in-memory database: no file for the OS holder census",
1791 ),
1792 collection_cost: CollectionCost::in_memory(elapsed_ms(started)),
1793 });
1794 };
1795
1796 let inspection_pool = Arc::clone(&pool);
1797 let sqlite_started = Instant::now();
1798 let inspection = crate::read_cancellation::run_interruptible_read(
1799 StorageCapability::Sql,
1800 "db_diagnostics.sqlite",
1801 move |scope| inspect_pool_interruptibly(&inspection_pool, scope),
1802 )
1803 .await?;
1804 let sqlite_ms = elapsed_ms(sqlite_started);
1805 crate::ensure_request_read_active("db_diagnostics")?;
1806 let canonical = operational_db_path(&pool, &path);
1807 let budget = request_census_budget();
1808 let (wal_file, wal_pin, file_state_cost) =
1809 inspect_file_state_interruptibly(canonical, sweep_interval, budget).await?;
1810 crate::ensure_request_read_active("db_diagnostics")?;
1811
1812 Ok(DbDiagnostics {
1813 build,
1814 process,
1815 note_search_ann_route_total: 0,
1816 note_search_fallback_route_total: 0,
1817 search_mechanism,
1818 db_path: Some(path.display().to_string()),
1819 wal_ceiling: WalCeilingDiagnostics::from_pool(&pool),
1820 disk_guard: DiskGuardDiagnostics::snapshot(&pool),
1821 wal_file: Some(wal_file),
1822 checkpoint_counters: counters,
1823 checkpoint_probe: inspection.checkpoint_probe,
1824 checkpoint_probe_error: inspection.checkpoint_probe_error,
1825 checkpoint_pin: inspection.checkpoint_pin,
1826 reader_contention,
1827 writer_contention,
1828 size_composition: inspection.size_composition,
1829 size_composition_error: inspection.size_composition_error,
1830 graph_edge_integrity: inspection.graph_edge_integrity,
1831 graph_edge_integrity_error: inspection.graph_edge_integrity_error,
1832 fts_segments: inspection.fts_segments,
1833 fts_segments_error: inspection.fts_segments_error,
1834 fts_maintenance: crate::fts_maintenance_counters(),
1835 wal_pin,
1836 collection_cost: CollectionCost {
1837 total_ms: elapsed_ms(started),
1838 sqlite_ms,
1839 wal_file_stat_ms: file_state_cost.wal_file_stat_ms,
1840 wal_pin_census_ms: file_state_cost.census_ms,
1841 wal_pin_sidecar_ms: file_state_cost.sidecar_ms,
1842 wal_pin_census_budget_ms: budget.map(|b| b.as_millis() as u64),
1843 wal_pin_census_budget_exhausted: file_state_cost.census_budget_exhausted,
1844 },
1845 })
1846}
1847
1848fn operational_db_path(pool: &ConnectionPool, configured: &Path) -> PathBuf {
1860 pool.canonical_path()
1861 .map(Path::to_path_buf)
1862 .unwrap_or_else(|| configured.to_path_buf())
1863}
1864
1865fn collect_inner(
1866 pool: &ConnectionPool,
1867 build: BuildIdentity,
1868 sweep_interval: Duration,
1869 audit_append_failures: Option<u64>,
1870 runtime_audit_batch_metrics: Option<RuntimeAuditBatchMetrics>,
1871) -> DbDiagnostics {
1872 let started = Instant::now();
1873 let process = ProcessIdentity::current(pool);
1874 let counters = checkpoint_counters();
1875 let reader_contention = ReaderContentionDiagnostics::snapshot(pool);
1876 let search_mechanism = pool.search_mechanism_snapshot();
1877 let writer_contention = WriterContentionDiagnostics::snapshot(
1878 pool,
1879 audit_append_failures,
1880 runtime_audit_batch_metrics,
1881 );
1882
1883 let Some(path) = pool.config().path.clone() else {
1884 return DbDiagnostics {
1885 build,
1886 process,
1887 note_search_ann_route_total: 0,
1888 note_search_fallback_route_total: 0,
1889 search_mechanism,
1890 db_path: None,
1891 wal_ceiling: WalCeilingDiagnostics::from_pool(pool),
1892 disk_guard: DiskGuardDiagnostics::snapshot(pool),
1893 wal_file: None,
1894 checkpoint_counters: counters,
1895 checkpoint_probe: None,
1896 checkpoint_probe_error: Some(
1897 "in-memory database: no WAL file and no checkpoint to probe".to_string(),
1898 ),
1899 checkpoint_pin: checkpoint_pin_diagnostics(
1900 pool,
1901 None,
1902 Some("in-memory database: no WAL file and no checkpoint to probe"),
1903 ),
1904 reader_contention,
1905 writer_contention,
1906 size_composition: None,
1907 size_composition_error: Some(
1908 "in-memory database: no file-backed page composition to inspect".to_string(),
1909 ),
1910 graph_edge_integrity: None,
1911 graph_edge_integrity_error: Some(
1912 "in-memory database: no durable graph-edge ledger to inspect".to_string(),
1913 ),
1914 fts_segments: None,
1915 fts_segments_error: Some(
1916 "in-memory database: no durable FTS5 indexes to inspect".to_string(),
1917 ),
1918 fts_maintenance: crate::fts_maintenance_counters(),
1919 wal_pin: WalPinAttribution::unavailable(
1920 "in-memory database: no file for the OS holder census",
1921 ),
1922 collection_cost: CollectionCost::in_memory(elapsed_ms(started)),
1923 };
1924 };
1925
1926 let sqlite_started = Instant::now();
1927 let inspection = inspect_pool(pool);
1928 let sqlite_ms = elapsed_ms(sqlite_started);
1929 let canonical = operational_db_path(pool, &path);
1930 let wal_file_started = Instant::now();
1931 let wal_file = wal_file_state(&canonical);
1932 let wal_file_stat_ms = elapsed_ms(wal_file_started);
1933 let census_started = Instant::now();
1934 let wal_pin = wal_pin_attribution(&canonical, sweep_interval);
1935 let wal_pin_ms = elapsed_ms(census_started);
1936
1937 DbDiagnostics {
1938 build,
1939 process,
1940 note_search_ann_route_total: 0,
1941 note_search_fallback_route_total: 0,
1942 search_mechanism,
1943 db_path: Some(path.display().to_string()),
1944 wal_ceiling: WalCeilingDiagnostics::from_pool(pool),
1945 disk_guard: DiskGuardDiagnostics::snapshot(pool),
1946 wal_file: Some(wal_file),
1947 checkpoint_counters: counters,
1948 checkpoint_probe: inspection.checkpoint_probe,
1949 checkpoint_probe_error: inspection.checkpoint_probe_error,
1950 checkpoint_pin: inspection.checkpoint_pin,
1951 reader_contention,
1952 writer_contention,
1953 size_composition: inspection.size_composition,
1954 size_composition_error: inspection.size_composition_error,
1955 graph_edge_integrity: inspection.graph_edge_integrity,
1956 graph_edge_integrity_error: inspection.graph_edge_integrity_error,
1957 fts_segments: inspection.fts_segments,
1958 fts_segments_error: inspection.fts_segments_error,
1959 fts_maintenance: crate::fts_maintenance_counters(),
1960 wal_pin,
1961 collection_cost: CollectionCost {
1962 total_ms: elapsed_ms(started),
1963 sqlite_ms,
1964 wal_file_stat_ms,
1965 wal_pin_census_ms: wal_pin_ms,
1971 wal_pin_sidecar_ms: 0,
1972 wal_pin_census_budget_ms: None,
1973 wal_pin_census_budget_exhausted: false,
1974 },
1975 }
1976}
1977
1978struct PoolInspection {
1979 checkpoint_probe: Option<CheckpointProbe>,
1980 checkpoint_probe_error: Option<String>,
1981 checkpoint_pin: CheckpointPinDiagnostics,
1982 size_composition: Option<DatabaseSizeComposition>,
1983 size_composition_error: Option<String>,
1984 graph_edge_integrity: Option<GraphEdgeIntegrity>,
1985 graph_edge_integrity_error: Option<String>,
1986 fts_segments: Option<crate::FtsSegmentDiagnostics>,
1987 fts_segments_error: Option<String>,
1988}
1989
1990fn split_fts_segments_result(
1994 result: Result<crate::FtsSegmentDiagnostics, String>,
1995) -> (Option<crate::FtsSegmentDiagnostics>, Option<String>) {
1996 match result {
1997 Ok(segments) => (Some(segments), None),
1998 Err(error) => (None, Some(error)),
1999 }
2000}
2001
2002fn record_diagnostic_checkpoint_probe(
2003 pool: &ConnectionPool,
2004 result: &rusqlite::Result<CheckpointProbe>,
2005) -> (checkpoint::CheckpointRunStatus, Option<Duration>) {
2006 checkpoint::record_checkpoint_run_result(
2007 pool,
2008 result
2009 .as_ref()
2010 .ok()
2011 .map(|probe| (probe.busy, probe.log_frames, probe.checkpointed_frames)),
2012 );
2013 checkpoint::checkpoint_run_snapshot(pool)
2014}
2015
2016fn inspect_pool_interruptibly(
2017 pool: &ConnectionPool,
2018 scope: &crate::read_cancellation::InterruptibleReadScope,
2019) -> StorageResult<PoolInspection> {
2020 scope.ensure_active()?;
2021 let conn = match pool.open_standalone_writer_untracked() {
2022 Ok(conn) => conn,
2023 Err(e) => {
2024 scope.ensure_active()?;
2025 let reason = format!("guarded standalone open refused: {e}");
2026 return Ok(PoolInspection {
2027 checkpoint_probe: None,
2028 checkpoint_probe_error: Some(reason.clone()),
2029 checkpoint_pin: checkpoint_pin_diagnostics(pool, None, Some(&reason)),
2030 size_composition: None,
2031 size_composition_error: Some(reason.clone()),
2032 graph_edge_integrity: None,
2033 graph_edge_integrity_error: Some(reason.clone()),
2034 fts_segments: None,
2035 fts_segments_error: Some(reason),
2036 });
2037 }
2038 };
2039 scope.ensure_active()?;
2044
2045 let probe_result = checkpoint_probe(&conn);
2047 let (run_status, run_age) = record_diagnostic_checkpoint_probe(pool, &probe_result);
2048 let (checkpoint_probe, checkpoint_probe_error) = match probe_result {
2049 Ok(probe) => (Some(probe), None),
2050 Err(e) => (
2051 None,
2052 Some(format!("PRAGMA wal_checkpoint(PASSIVE) failed: {e}")),
2053 ),
2054 };
2055 let checkpoint_pin = checkpoint_pin_diagnostics_for_run(
2056 checkpoint_probe.as_ref(),
2057 checkpoint_probe_error.as_deref(),
2058 run_status,
2059 run_age,
2060 );
2061 #[cfg(test)]
2062 if TEST_PAUSE_AFTER_PASSIVE.load(Ordering::SeqCst) {
2063 TEST_REACHED_AFTER_PASSIVE.store(true, Ordering::SeqCst);
2064 while TEST_PAUSE_AFTER_PASSIVE.load(Ordering::SeqCst) && !scope.should_stop() {
2065 std::thread::yield_now();
2066 }
2067 }
2068 scope.ensure_active()?;
2069
2070 let (integrity, size_composition, fts_segments) = scope.run(&conn, || {
2075 Ok((
2076 graph_edge_integrity(&conn),
2077 database_size_composition(&conn),
2078 crate::fts_maintenance::inspect_fts_segments(&conn),
2079 ))
2080 })?;
2081 let (graph_edge_integrity, graph_edge_integrity_error) = match integrity {
2082 Ok(integrity) => (Some(integrity), None),
2083 Err(e) => (
2084 None,
2085 Some(format!("graph-edge integrity query failed: {e}")),
2086 ),
2087 };
2088 let (fts_segments, fts_segments_error) = split_fts_segments_result(fts_segments);
2089
2090 let (size_composition, size_composition_error) = match size_composition {
2091 Ok(composition) => (Some(composition), None),
2092 Err(error) => (
2093 None,
2094 Some(format!("database size composition query failed: {error}")),
2095 ),
2096 };
2097
2098 Ok(PoolInspection {
2099 checkpoint_probe,
2100 checkpoint_probe_error,
2101 checkpoint_pin,
2102 size_composition,
2103 size_composition_error,
2104 graph_edge_integrity,
2105 graph_edge_integrity_error,
2106 fts_segments,
2107 fts_segments_error,
2108 })
2109}
2110
2111#[cfg(test)]
2112static TEST_PAUSE_AFTER_PASSIVE: AtomicBool = AtomicBool::new(false);
2113#[cfg(test)]
2114static TEST_REACHED_AFTER_PASSIVE: AtomicBool = AtomicBool::new(false);
2115
2116struct StopCensusOnDrop {
2117 stopped: Arc<AtomicBool>,
2118 armed: bool,
2119}
2120
2121impl Drop for StopCensusOnDrop {
2122 fn drop(&mut self) {
2123 if self.armed {
2124 self.stopped.store(true, Ordering::SeqCst);
2125 }
2126 }
2127}
2128
2129#[derive(Debug, Clone, Copy, Default)]
2132struct FileStateCost {
2133 wal_file_stat_ms: u64,
2134 census_ms: u64,
2135 sidecar_ms: u64,
2136 census_budget_exhausted: bool,
2137}
2138
2139fn elapsed_ms(since: Instant) -> u64 {
2140 since.elapsed().as_millis() as u64
2141}
2142
2143async fn inspect_file_state_interruptibly(
2144 path: PathBuf,
2145 sweep_interval: Duration,
2146 census_budget: Option<Duration>,
2147) -> StorageResult<(WalFileState, WalPinAttribution, FileStateCost)> {
2148 const OPERATION: &str = "db_diagnostics.wal_holder_census";
2149 crate::ensure_request_read_active(OPERATION)?;
2150 let stopped = Arc::new(AtomicBool::new(false));
2151 let worker_stopped = Arc::clone(&stopped);
2152 let mut stop_on_drop = StopCensusOnDrop {
2153 stopped: Arc::clone(&stopped),
2154 armed: true,
2155 };
2156 let mut worker = tokio::task::spawn_blocking(move || {
2157 let mut cost = FileStateCost::default();
2158 let wal_file_started = Instant::now();
2159 let wal_file = wal_file_state(&path);
2160 cost.wal_file_stat_ms = elapsed_ms(wal_file_started);
2161 if worker_stopped.load(Ordering::SeqCst) {
2162 return Err(std::io::Error::new(
2163 std::io::ErrorKind::Interrupted,
2164 "WAL holder census cancelled",
2165 ));
2166 }
2167 #[cfg(unix)]
2168 let census_started = Instant::now();
2169 #[cfg(unix)]
2170 let census_result = match census_budget {
2171 Some(budget) => crate::walpin::census_holders_until_within(
2172 &path,
2173 || worker_stopped.load(Ordering::SeqCst),
2174 budget,
2175 ),
2176 None => {
2177 crate::walpin::census_holders_until(&path, || worker_stopped.load(Ordering::SeqCst))
2178 }
2179 };
2180 #[cfg(unix)]
2181 {
2182 cost.census_ms = elapsed_ms(census_started);
2183 }
2184 #[cfg(unix)]
2185 let attribution = match census_result {
2186 Ok(census) => {
2187 cost.census_budget_exhausted = census.budget_exhausted;
2188 if worker_stopped.load(Ordering::SeqCst) {
2189 return Err(std::io::Error::new(
2190 std::io::ErrorKind::Interrupted,
2191 "WAL sidecar inspection cancelled",
2192 ));
2193 }
2194 if !crate::walpin::sidecar_enabled(true) {
2195 wal_pin_attribution_without_sidecar(census, SIDECAR_DISABLED_REASON.to_string())
2196 } else {
2197 let sidecar_started = Instant::now();
2198 let sidecar = crate::walpin::inspect_live(
2199 &crate::walpin::sidecar_dir_for(&path),
2200 sweep_interval,
2201 );
2202 cost.sidecar_ms = elapsed_ms(sidecar_started);
2203 if worker_stopped.load(Ordering::SeqCst) {
2204 return Err(std::io::Error::new(
2205 std::io::ErrorKind::Interrupted,
2206 "WAL sidecar inspection cancelled",
2207 ));
2208 }
2209 match sidecar {
2210 Ok(sidecar) => wal_pin_attribution_from_evidence(census, sidecar),
2211 Err(error) => wal_pin_attribution_without_sidecar(
2212 census,
2213 format!("read-only sidecar enumeration failed: {error}"),
2214 ),
2215 }
2216 }
2217 }
2218 Err(error) if error.kind() == std::io::ErrorKind::Interrupted => return Err(error),
2219 Err(error) => WalPinAttribution::unavailable(format!("census_holders failed: {error}")),
2220 };
2221 #[cfg(not(unix))]
2222 let attribution = {
2223 let _ = census_budget;
2224 let census_started = Instant::now();
2225 let attribution = wal_pin_attribution(&path, sweep_interval);
2226 cost.census_ms = elapsed_ms(census_started);
2227 attribution
2228 };
2229 Ok((wal_file, attribution, cost))
2230 });
2231
2232 tokio::select! {
2233 joined = &mut worker => {
2234 stop_on_drop.armed = false;
2235 let result = joined
2236 .map_err(|error| StorageError::driver(StorageCapability::Sql, OPERATION, error))?
2237 .map_err(|error| StorageError::driver(StorageCapability::Sql, OPERATION, error))?;
2238 crate::ensure_request_read_active(OPERATION)?;
2239 Ok(result)
2240 }
2241 _ = crate::wait_for_request_read_cancellation() => {
2242 stopped.store(true, Ordering::SeqCst);
2243 if tokio::time::timeout(crate::sqlite_interrupt_grace_from_env(), &mut worker)
2244 .await
2245 .is_err()
2246 {
2247 worker.abort();
2248 }
2249 stop_on_drop.armed = false;
2250 Err(StorageError::Timeout { operation: OPERATION.into() })
2251 }
2252 }
2253}
2254
2255fn inspect_pool(pool: &ConnectionPool) -> PoolInspection {
2262 let conn = match pool.open_standalone_writer_untracked() {
2263 Ok(conn) => conn,
2264 Err(e) => {
2265 let reason = format!("guarded standalone open refused: {e}");
2266 return PoolInspection {
2267 checkpoint_probe: None,
2268 checkpoint_probe_error: Some(reason.clone()),
2269 checkpoint_pin: checkpoint_pin_diagnostics(pool, None, Some(&reason)),
2270 size_composition: None,
2271 size_composition_error: Some(reason.clone()),
2272 graph_edge_integrity: None,
2273 graph_edge_integrity_error: Some(reason.clone()),
2274 fts_segments: None,
2275 fts_segments_error: Some(reason),
2276 };
2277 }
2278 };
2279
2280 let probe_result = checkpoint_probe(&conn);
2281 let (run_status, run_age) = record_diagnostic_checkpoint_probe(pool, &probe_result);
2282 let (checkpoint_probe, checkpoint_probe_error) = match probe_result {
2283 Ok(probe) => (Some(probe), None),
2284 Err(e) => (
2285 None,
2286 Some(format!("PRAGMA wal_checkpoint(PASSIVE) failed: {e}")),
2287 ),
2288 };
2289 let checkpoint_pin = checkpoint_pin_diagnostics_for_run(
2290 checkpoint_probe.as_ref(),
2291 checkpoint_probe_error.as_deref(),
2292 run_status,
2293 run_age,
2294 );
2295 let (graph_edge_integrity, graph_edge_integrity_error) = match graph_edge_integrity(&conn) {
2296 Ok(integrity) => (Some(integrity), None),
2297 Err(e) => (
2298 None,
2299 Some(format!("graph-edge integrity query failed: {e}")),
2300 ),
2301 };
2302 let (size_composition, size_composition_error) = match database_size_composition(&conn) {
2303 Ok(composition) => (Some(composition), None),
2304 Err(error) => (
2305 None,
2306 Some(format!("database size composition query failed: {error}")),
2307 ),
2308 };
2309 let (fts_segments, fts_segments_error) =
2310 split_fts_segments_result(crate::fts_maintenance::inspect_fts_segments(&conn));
2311
2312 PoolInspection {
2313 checkpoint_probe,
2314 checkpoint_probe_error,
2315 checkpoint_pin,
2316 size_composition,
2317 size_composition_error,
2318 graph_edge_integrity,
2319 graph_edge_integrity_error,
2320 fts_segments,
2321 fts_segments_error,
2322 }
2323}
2324
2325#[cfg(test)]
2326mod tests {
2327 use serial_test::serial;
2328
2329 use super::*;
2330 use crate::pool::{ConnectionPool, PoolConfig, WalCeilingPolicy, WalCeilingSource};
2331
2332 include!("diagnostics/wal_ceiling_tests.rs");
2333 include!("diagnostics/environment_tests.rs");
2334 include!("diagnostics_census_evidence_tests.rs");
2335
2336 fn seeded_pool(dir: &tempfile::TempDir) -> (ConnectionPool, PathBuf) {
2337 let path = dir.path().join("diag.db");
2338 let pool = ConnectionPool::new(PoolConfig {
2339 path: Some(path.clone()),
2340 ..PoolConfig::for_test()
2341 })
2342 .expect("pool open");
2343 {
2344 let writer = pool.try_writer().expect("writer");
2345 writer
2346 .conn()
2347 .execute_batch(
2348 "CREATE TABLE t (x INTEGER); \
2349 CREATE TABLE entities (id TEXT PRIMARY KEY, deleted_at INTEGER, merged_into TEXT); \
2350 CREATE TABLE graph_edges (
2351 namespace TEXT NOT NULL,
2352 id TEXT NOT NULL,
2353 PRIMARY KEY (namespace, id)
2354 ); \
2355 CREATE TABLE graph_edges_seq (
2356 seq INTEGER PRIMARY KEY AUTOINCREMENT,
2357 edge_id TEXT NOT NULL UNIQUE
2358 ); \
2359 CREATE VIRTUAL TABLE fts_entities USING fts5(
2360 namespace UNINDEXED, subject_id UNINDEXED, title, body,
2361 tokenize='trigram'
2362 ); \
2363 CREATE VIRTUAL TABLE fts_notes USING fts5(
2364 namespace UNINDEXED, subject_id UNINDEXED, title, body,
2365 tokenize='trigram'
2366 ); \
2367 INSERT INTO t VALUES (1), (2), (3); \
2368 INSERT INTO fts_entities(namespace, subject_id, title, body)
2369 VALUES('local', 'entity-1', 'entity title', 'entity diagnostic body'); \
2370 INSERT INTO fts_notes(namespace, subject_id, title, body)
2371 VALUES('local', 'note-1', 'note title', 'note diagnostic body');",
2372 )
2373 .expect("seed writes");
2374 }
2375 (pool, path)
2376 }
2377
2378 #[test]
2379 fn dropping_census_future_guard_requests_cooperative_stop() {
2380 let stopped = Arc::new(AtomicBool::new(false));
2381 let guard = StopCensusOnDrop {
2382 stopped: Arc::clone(&stopped),
2383 armed: true,
2384 };
2385
2386 drop(guard);
2387
2388 assert!(
2389 stopped.load(Ordering::SeqCst),
2390 "dropping diagnostics while its census worker is live must stop the PID/fd walk"
2391 );
2392 }
2393
2394 #[tokio::test]
2395 async fn runtime_audit_batch_fields_are_additive_and_unavailable_without_a_control() {
2396 let dir = tempfile::tempdir().expect("tempdir");
2397 let (pool, _path) = seeded_pool(&dir);
2398 let pool = Arc::new(pool);
2399
2400 let without_control = collect_with_audit_append_failures_interruptibly(
2401 Arc::clone(&pool),
2402 BuildIdentity::from_env("test", None),
2403 Duration::from_secs(30),
2404 0,
2405 )
2406 .await
2407 .expect("diagnostics succeed");
2408 assert!(without_control
2409 .writer_contention
2410 .audit_batch_flush_failures
2411 .is_none());
2412 assert!(
2413 without_control
2414 .writer_contention
2415 .audit_batch_flush_failures_unavailable_reason
2416 .is_some(),
2417 "no audit-batch control was supplied, so the field must carry a reason, not a \
2418 fabricated zero"
2419 );
2420 assert!(without_control
2421 .writer_contention
2422 .audit_degraded_rows
2423 .is_none());
2424 assert!(without_control.writer_contention.audit_degraded.is_none());
2425
2426 let with_control = collect_with_runtime_audit_metrics_interruptibly(
2427 Arc::clone(&pool),
2428 BuildIdentity::from_env("test", None),
2429 Duration::from_secs(30),
2430 0,
2431 Some(RuntimeAuditBatchMetrics {
2432 flush_failures: 3,
2433 degraded_rows: 7,
2434 degraded: true,
2435 admission_refused_obligations: 5,
2436 admission_refused_obligations_last_at_ms: Some(1_700_000_000_123),
2437 admission_unresolved_obligations: 2,
2438 admission_unresolved_obligations_last_at_ms: Some(1_700_000_000_456),
2439 }),
2440 )
2441 .await
2442 .expect("diagnostics succeed");
2443 assert_eq!(
2444 with_control.writer_contention.audit_batch_flush_failures,
2445 Some(3)
2446 );
2447 assert!(with_control
2448 .writer_contention
2449 .audit_batch_flush_failures_unavailable_reason
2450 .is_none());
2451 assert_eq!(with_control.writer_contention.audit_degraded_rows, Some(7));
2452 assert_eq!(with_control.writer_contention.audit_degraded, Some(true));
2453 assert_eq!(
2454 with_control
2455 .writer_contention
2456 .audit_admission_refused_obligations,
2457 Some(5),
2458 "an operator must be able to read the admission-refused obligation count from \
2459 db_diagnostics without a test-only feature gate (ADR-103 Amendment 3)"
2460 );
2461 assert!(with_control
2462 .writer_contention
2463 .audit_admission_refused_obligations_unavailable_reason
2464 .is_none());
2465 assert_eq!(
2466 with_control
2467 .writer_contention
2468 .audit_admission_unresolved_obligations,
2469 Some(2),
2470 "an operator must be able to distinguish enqueued-but-unresolved rows from \
2471 confirmed-refused rows (ADR-103 Amendment 3)"
2472 );
2473 assert!(with_control
2474 .writer_contention
2475 .audit_admission_unresolved_obligations_unavailable_reason
2476 .is_none());
2477 assert!(without_control
2478 .writer_contention
2479 .audit_admission_refused_obligations
2480 .is_none());
2481 assert!(without_control
2482 .writer_contention
2483 .audit_admission_refused_obligations_unavailable_reason
2484 .is_some());
2485 assert!(without_control
2486 .writer_contention
2487 .audit_admission_unresolved_obligations
2488 .is_none());
2489 assert!(without_control
2490 .writer_contention
2491 .audit_admission_unresolved_obligations_unavailable_reason
2492 .is_some());
2493
2494 assert_eq!(
2497 with_control.writer_contention.writer_acquisitions,
2498 without_control.writer_contention.writer_acquisitions
2499 );
2500 }
2501
2502 #[test]
2503 fn writer_task_pool_sourced_counters_are_always_populated_directly() {
2504 let dir = tempfile::tempdir().expect("tempdir");
2505 let (pool, _path) = seeded_pool(&dir);
2506
2507 let report = collect(
2508 &pool,
2509 BuildIdentity::from_env("9.9.9", None),
2510 Duration::from_secs(30),
2511 );
2512
2513 assert_eq!(report.writer_contention.writer_task_request_failures, 0);
2516 assert_eq!(report.writer_contention.writer_task_side_effects_unknown, 0);
2517 }
2518
2519 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2520 #[serial]
2521 async fn request_cancellation_after_passive_stops_before_graph_and_census() {
2522 let dir = tempfile::tempdir().expect("tempdir");
2523 let (pool, _) = seeded_pool(&dir);
2524 let pool = Arc::new(pool);
2525 TEST_REACHED_AFTER_PASSIVE.store(false, Ordering::SeqCst);
2526 TEST_PAUSE_AFTER_PASSIVE.store(true, Ordering::SeqCst);
2527 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
2528 let diagnostic_pool = Arc::clone(&pool);
2529 let task = tokio::spawn(crate::scope_request_read_cancellation(
2530 cancel_rx,
2531 async move {
2532 collect_with_audit_append_failures_interruptibly(
2533 diagnostic_pool,
2534 BuildIdentity::from_env("test", None),
2535 Duration::from_secs(30),
2536 0,
2537 )
2538 .await
2539 },
2540 ));
2541
2542 tokio::time::timeout(Duration::from_secs(1), async {
2543 while !TEST_REACHED_AFTER_PASSIVE.load(Ordering::SeqCst) {
2544 tokio::task::yield_now().await;
2545 }
2546 })
2547 .await
2548 .expect("diagnostics never completed its admitted PASSIVE phase");
2549 cancel_tx.send(true).unwrap();
2550 let result = tokio::time::timeout(Duration::from_secs(1), task)
2551 .await
2552 .expect("cancelled diagnostics did not stop promptly")
2553 .expect("diagnostics task panicked");
2554 TEST_PAUSE_AFTER_PASSIVE.store(false, Ordering::SeqCst);
2555 assert!(matches!(result, Err(StorageError::Timeout { .. })));
2556
2557 let one: i64 = pool
2558 .reader()
2559 .expect("diagnostics returned its connection")
2560 .conn()
2561 .query_row("SELECT 1", [], |row| row.get(0))
2562 .unwrap();
2563 assert_eq!(one, 1);
2564 }
2565
2566 #[test]
2567 fn checkpoint_probe_returns_a_well_formed_triple_on_a_file_backed_db() {
2568 let dir = tempfile::tempdir().expect("tempdir");
2569 let (pool, _path) = seeded_pool(&dir);
2570 let conn = pool
2571 .open_standalone_writer_untracked()
2572 .expect("standalone probe connection");
2573
2574 let probe = checkpoint_probe(&conn).expect("probe must succeed on a WAL database");
2575
2576 assert!(
2577 probe.busy == 0 || probe.busy == 1,
2578 "busy is a 0/1 flag, got {}",
2579 probe.busy
2580 );
2581 assert!(
2582 probe.log_frames >= 0,
2583 "a WAL database must report a non-negative frame count, got {}",
2584 probe.log_frames
2585 );
2586 assert!(
2587 probe.checkpointed_frames >= 0,
2588 "checkpointed frames must be non-negative, got {}",
2589 probe.checkpointed_frames
2590 );
2591 assert!(
2592 probe.checkpointed_frames <= probe.log_frames,
2593 "a PASSIVE pass cannot checkpoint more frames than the WAL holds: {probe:?}"
2594 );
2595 assert!(
2596 probe.backfill_gap_frames() >= 0,
2597 "the one-row backfill gap clamps at 0: {probe:?}"
2598 );
2599 }
2600
2601 #[cfg(unix)]
2602 #[test]
2603 fn checkpoint_probe_reports_replaced_pool_file_instead_of_probing_it() {
2604 let dir = tempfile::tempdir().expect("tempdir");
2605 let (pool, path) = seeded_pool(&dir);
2606 let replacement = dir.path().join("replacement.db");
2607 let replacement_conn = Connection::open(&replacement).expect("replacement database");
2608 replacement_conn
2609 .execute_batch("CREATE TABLE replacement_marker (value INTEGER)")
2610 .expect("initialize replacement database");
2611 drop(replacement_conn);
2612 std::fs::rename(&replacement, &path).expect("replace the pool's path");
2613
2614 let inspection = inspect_pool(&pool);
2615 assert!(inspection.checkpoint_probe.is_none());
2616 assert!(
2617 inspection
2618 .checkpoint_probe_error
2619 .as_deref()
2620 .is_some_and(|error| error.contains("file identity changed")),
2621 "replacement must be reported as a probe failure"
2622 );
2623 }
2624
2625 #[test]
2626 fn checkpoint_probe_backfill_gap_is_a_row_difference_not_a_pin_claim() {
2627 let probe = CheckpointProbe {
2628 busy: 0,
2629 log_frames: 12,
2630 checkpointed_frames: 7,
2631 };
2632 assert_eq!(probe.backfill_gap_frames(), 5);
2633 let no_wal = CheckpointProbe {
2634 busy: 0,
2635 log_frames: -1,
2636 checkpointed_frames: -1,
2637 };
2638 assert_eq!(no_wal.backfill_gap_frames(), 0);
2639 let busy = CheckpointProbe { busy: 1, ..probe };
2640 assert_eq!(busy.backfill_gap_frames(), 5);
2641 }
2642
2643 #[test]
2644 fn pin_age_does_not_depend_on_the_wall_clock_epoch() {
2645 let probe = CheckpointProbe {
2646 busy: 0,
2647 log_frames: 12,
2648 checkpointed_frames: 7,
2649 };
2650 let future_epoch = checkpoint::CheckpointRunStatus::Observed(checkpoint::CheckpointRun {
2651 frame: 7,
2652 first_observed_at_unix_ms: 10_000,
2653 });
2654 let past_epoch = checkpoint::CheckpointRunStatus::Observed(checkpoint::CheckpointRun {
2655 frame: 7,
2656 first_observed_at_unix_ms: 1,
2657 });
2658 let young = checkpoint_pin_diagnostics_for_run(
2659 Some(&probe),
2660 None,
2661 future_epoch,
2662 Some(Duration::from_millis(10)),
2663 );
2664 let forward = checkpoint_pin_diagnostics_for_run(
2665 Some(&probe),
2666 None,
2667 future_epoch,
2668 Some(Duration::from_millis(1010)),
2669 );
2670 let backward = checkpoint_pin_diagnostics_for_run(
2671 Some(&probe),
2672 None,
2673 past_epoch,
2674 Some(Duration::from_millis(1010)),
2675 );
2676 assert_eq!(young.oldest_pinned_frame, None);
2677 assert_eq!(forward.oldest_pinned_frame, Some(7));
2678 assert_eq!(forward.pin_depth, Some(5));
2679 assert_eq!(backward.oldest_pinned_frame, Some(7));
2680 let busy = CheckpointProbe { busy: 1, ..probe };
2681 let unavailable = checkpoint_pin_diagnostics_for_run(
2682 Some(&busy),
2683 None,
2684 future_epoch,
2685 Some(Duration::from_millis(1010)),
2686 );
2687 assert_eq!(unavailable.oldest_pinned_frame, None);
2688 assert_eq!(unavailable.pin_depth, None);
2689 }
2690
2691 #[test]
2692 fn checkpoint_pin_report_waits_one_second_for_a_matching_run() {
2693 let probe = CheckpointProbe {
2694 busy: 0,
2695 log_frames: 12,
2696 checkpointed_frames: 7,
2697 };
2698 let run = checkpoint::CheckpointRunStatus::Observed(checkpoint::CheckpointRun {
2699 frame: 7,
2700 first_observed_at_unix_ms: 1_000,
2701 });
2702
2703 let young = checkpoint_pin_diagnostics_for_run(
2704 Some(&probe),
2705 None,
2706 run,
2707 Some(Duration::from_millis(999)),
2708 );
2709 assert_eq!(young.backfill_ceiling, Some(7));
2710 assert_eq!(young.oldest_pinned_frame, None);
2711 assert!(young
2712 .oldest_pinned_frame_unavailable_reason
2713 .as_deref()
2714 .is_some_and(|reason| reason.contains("less than one second")));
2715 assert_eq!(young.oldest_pinned_frame_run, None);
2716 assert_eq!(young.pin_depth, None);
2717 assert!(young.pin_depth_unavailable_reason.is_some());
2718
2719 let aged = checkpoint_pin_diagnostics_for_run(
2720 Some(&probe),
2721 None,
2722 run,
2723 Some(Duration::from_secs(1)),
2724 );
2725 assert_eq!(aged.backfill_ceiling, Some(7));
2726 assert_eq!(aged.oldest_pinned_frame, Some(7));
2727 assert_eq!(aged.pin_depth, Some(5));
2728 assert_eq!(aged.pin_depth_unavailable_reason, None);
2729 assert_eq!(
2730 aged.oldest_pinned_frame_run,
2731 Some(match run {
2732 checkpoint::CheckpointRunStatus::Observed(value) => value,
2733 _ => unreachable!(),
2734 })
2735 );
2736 }
2737
2738 #[test]
2739 fn checkpoint_pin_report_keeps_busy_probe_fields_null_even_with_an_aged_run() {
2740 let run = checkpoint::CheckpointRunStatus::Observed(checkpoint::CheckpointRun {
2741 frame: 7,
2742 first_observed_at_unix_ms: 1,
2743 });
2744 for probe in [
2745 Some(CheckpointProbe {
2746 busy: 1,
2747 log_frames: 12,
2748 checkpointed_frames: 7,
2749 }),
2750 Some(CheckpointProbe {
2751 busy: 0,
2752 log_frames: -1,
2753 checkpointed_frames: -1,
2754 }),
2755 Some(CheckpointProbe {
2756 busy: 0,
2757 log_frames: 12,
2758 checkpointed_frames: 12,
2759 }),
2760 ] {
2761 let result = checkpoint_pin_diagnostics_for_run(
2762 probe.as_ref(),
2763 None,
2764 run,
2765 Some(Duration::from_secs(1)),
2766 );
2767 assert_eq!(result.backfill_ceiling, None);
2768 assert!(result.backfill_ceiling_unavailable_reason.is_some());
2769 assert_eq!(result.oldest_pinned_frame, None);
2770 assert!(result.oldest_pinned_frame_unavailable_reason.is_some());
2771 assert_eq!(result.oldest_pinned_frame_run, None);
2772 assert!(result.oldest_pinned_frame_run_unavailable_reason.is_some());
2773 assert_eq!(result.pin_depth, None);
2774 assert!(result.pin_depth_unavailable_reason.is_some());
2775 }
2776
2777 let error = checkpoint_pin_diagnostics_for_run(
2778 None,
2779 Some("probe failed"),
2780 run,
2781 Some(Duration::from_secs(1)),
2782 );
2783 assert_eq!(error.backfill_ceiling, None);
2784 assert_eq!(
2785 error.backfill_ceiling_unavailable_reason.as_deref(),
2786 Some("probe failed")
2787 );
2788 assert_eq!(error.oldest_pinned_frame, None);
2789 assert!(error.oldest_pinned_frame_run_unavailable_reason.is_some());
2790 assert_eq!(error.pin_depth, None);
2791 assert!(error.pin_depth_unavailable_reason.is_some());
2792 }
2793
2794 #[cfg(all(unix, any(target_os = "linux", target_os = "macos")))]
2795 #[test]
2796 fn db_diagnostics_reports_start_times_for_all_holders_and_identifies_reporter() {
2797 let dir = tempfile::tempdir().expect("tempdir");
2798 let (pool, _) = seeded_pool(&dir);
2799
2800 let report = collect(
2801 &pool,
2802 BuildIdentity::from_env("test", None),
2803 Duration::from_secs(30),
2804 );
2805 let reporter = crate::walpin::reporting_pid();
2806 let process = report
2807 .wal_pin
2808 .census_process_start_times
2809 .iter()
2810 .find(|process| process.pid == reporter)
2811 .expect("the reporting process must remain in the census list");
2812
2813 assert_eq!(report.wal_pin.reporting_pid, reporter);
2814 assert_eq!(report.wal_pin.reporting_process_is_holder, Some(true));
2815 assert_eq!(
2816 report.wal_pin.census_process_start_times.len(),
2817 report.wal_pin.census_holder_pids.len(),
2818 "every confirmed holder gets raw start-time data"
2819 );
2820 assert_eq!(
2821 process.process_start_time_secs,
2822 crate::walpin::process_start_time_secs(reporter)
2823 );
2824 assert_eq!(process.process_start_time_unavailable_reason, None);
2825 #[cfg(target_os = "linux")]
2826 assert_eq!(report.wal_pin.start_time_resolution_secs, Some(2));
2827 #[cfg(target_os = "macos")]
2828 assert_eq!(report.wal_pin.start_time_resolution_secs, Some(1));
2829 assert_eq!(
2830 serde_json::to_value(&report).unwrap()["wal_pin"]["census_process_start_times"]
2831 .as_array()
2832 .unwrap()
2833 .iter()
2834 .filter_map(|entry| entry["pid"].as_u64())
2835 .collect::<Vec<_>>(),
2836 report
2837 .wal_pin
2838 .census_holder_pids
2839 .iter()
2840 .map(|pid| u64::from(*pid))
2841 .collect::<Vec<_>>(),
2842 "start-time reporting does not filter census holders"
2843 );
2844 }
2845
2846 #[test]
2849 #[serial(checkpoint_skip_metrics)]
2850 fn checkpoint_probe_does_not_perturb_the_adr091_counters() {
2851 crate::checkpoint::reset_checkpoint_metrics_for_tests();
2852 let dir = tempfile::tempdir().expect("tempdir");
2853 let (pool, _path) = seeded_pool(&dir);
2854 let conn = pool
2855 .open_standalone_writer_untracked()
2856 .expect("standalone probe connection");
2857
2858 let before = checkpoint_counters();
2859 for _ in 0..3 {
2860 checkpoint_probe(&conn).expect("probe must succeed");
2861 }
2862 let after = checkpoint_counters();
2863
2864 assert_eq!(
2865 before, after,
2866 "checkpoint_probe must leave every ADR-091 counter untouched"
2867 );
2868 }
2869
2870 #[test]
2871 fn wal_file_state_reports_the_sidecar_size_for_a_live_db() {
2872 let dir = tempfile::tempdir().expect("tempdir");
2873 let (_pool, path) = seeded_pool(&dir);
2874
2875 let state = wal_file_state(&path);
2876 assert!(
2877 state.wal_path.ends_with("diag.db-wal"),
2878 "WAL path is the db path plus a -wal suffix, got {}",
2879 state.wal_path
2880 );
2881 assert!(
2882 state.wal_size_bytes.is_some(),
2883 "a seeded WAL database must have a stat-able -wal file: {state:?}"
2884 );
2885 assert!(state.unavailable_reason.is_none(), "{state:?}");
2886 }
2887
2888 #[test]
2889 fn wal_file_state_degrades_with_a_reason_when_the_sidecar_is_absent() {
2890 let dir = tempfile::tempdir().expect("tempdir");
2891 let state = wal_file_state(&dir.path().join("never-created.db"));
2892 assert!(state.wal_size_bytes.is_none());
2893 assert!(
2894 state.unavailable_reason.is_some(),
2895 "an absent WAL file must carry a reason, not a silent zero: {state:?}"
2896 );
2897 }
2898
2899 #[test]
2900 fn collect_on_a_file_backed_db_carries_build_identity_and_every_counter() {
2901 let dir = tempfile::tempdir().expect("tempdir");
2902 let (pool, _path) = seeded_pool(&dir);
2903 let reader_admission_capacity = pool.max_readers().max(1);
2904
2905 let report = collect(
2906 &pool,
2907 BuildIdentity::from_env("9.9.9", Some("deadbeef")),
2908 Duration::from_secs(30),
2909 );
2910
2911 assert_eq!(report.build.version, "9.9.9");
2912 assert_eq!(report.build.build_hash.as_deref(), Some("deadbeef"));
2913 assert!(report.db_path.is_some());
2914 assert!(
2915 report.checkpoint_probe.is_some(),
2916 "file-backed collect must land a probe; error was {:?}",
2917 report.checkpoint_probe_error
2918 );
2919 assert!(
2920 report.wal_file.as_ref().and_then(|w| w.wal_size_bytes) >= Some(0),
2921 "wal_size_bytes must be a non-negative byte count when present"
2922 );
2923
2924 let json = serde_json::to_value(&report).expect("report serializes");
2925 let counters = json
2926 .get("checkpoint_counters")
2927 .expect("counters section present");
2928 for key in [
2929 "last_observed_wal_pages",
2930 "truncate_attempts",
2931 "truncate_consecutive_failures",
2932 "checkpoint_skipped_ticks",
2933 "checkpoint_consecutive_skips",
2934 "checkpoint_last_skip_wal_pages",
2935 "checkpoint_pressure_elevated_ticks",
2936 "checkpoint_pressure_episodes_started",
2937 "checkpoint_pressure_episodes_recovered",
2938 "checkpoint_lifecycle_append_attempts",
2939 "checkpoint_lifecycle_append_failures",
2940 "checkpoint_lifecycle_enqueue_drops",
2941 "read_tx_max_age_evictions",
2942 ] {
2943 assert!(counters.get(key).is_some(), "counter {key} must be present");
2944 }
2945 assert_eq!(
2946 report.writer_contention.writer_acquisitions, 1,
2947 "the seed write checked the finite-wait pooled writer out once"
2948 );
2949 assert_eq!(report.writer_contention.pooled_writer_acquisitions, 1);
2950 assert_eq!(report.writer_contention.standalone_writer_acquisitions, 0);
2951 assert_eq!(report.writer_contention.writer_task_acquisitions, 0);
2952 assert_eq!(report.writer_contention.writer_acquisition_timeouts, 0);
2953 assert_eq!(
2954 report.reader_contention,
2955 ReaderContentionDiagnostics {
2956 configured_reader_cap: pool.config().max_readers,
2957 configured_checkout_timeout_ms: u64::try_from(
2958 pool.config().checkout_timeout.as_millis(),
2959 )
2960 .unwrap_or(u64::MAX),
2961 configured_busy_timeout_ms: u64::try_from(pool.config().busy_timeout.as_millis())
2962 .unwrap_or(u64::MAX),
2963 reader_admission_capacity,
2964 available_reader_admission_slots: reader_admission_capacity,
2965 reader_acquisitions: 0,
2966 pooled_reader_checkouts: 0,
2967 standalone_reader_opens: 0,
2968 infrastructure_standalone_reader_opens: 0,
2969 reader_checkout_timeouts: 0,
2970 reader_busy_timeouts: 0,
2971 active_pooled_reader_checkouts: 0,
2972 peak_active_pooled_reader_checkouts: 0,
2973 completed_pooled_reader_checkouts: 0,
2974 max_completed_reader_hold_micros: 0,
2975 max_completed_reader_hold_operation: None,
2976 reader_replacement_open_failures: 0,
2977 },
2978 "the diagnostics probe itself must not masquerade as request reader traffic"
2979 );
2980 assert_eq!(
2981 report.graph_edge_integrity,
2982 Some(GraphEdgeIntegrity {
2983 duplicate_edge_id_groups: 0,
2984 graph_edges_rows: 0,
2985 graph_edges_seq_rows: 0,
2986 pre_v14_duplicate_edge_state_detected: false,
2987 live_entities_carrying_merged_into: 0,
2988 })
2989 );
2990 assert!(report.graph_edge_integrity_error.is_none());
2991 let fts_segments = report
2992 .fts_segments
2993 .as_ref()
2994 .expect("file-backed diagnostics must decode both FTS structure rows");
2995 assert_eq!(fts_segments.entities.segment_count, 1);
2996 assert_eq!(fts_segments.notes.segment_count, 1);
2997 assert_eq!(fts_segments.total_segments, 2);
2998 assert!(report.fts_segments_error.is_none());
2999 let fts_json = json
3000 .get("fts_segments")
3001 .expect("FTS segment diagnostics serialize");
3002 assert_eq!(fts_json["entities"]["segment_count"], 1);
3003 assert!(json.get("fts_maintenance").is_some());
3004 assert!(
3005 report.size_composition.is_some(),
3006 "file-backed diagnostics must include page composition; error was {:?}",
3007 report.size_composition_error
3008 );
3009 assert!(report.size_composition_error.is_none());
3010 assert!(report.writer_contention.audit_append_failures.is_none());
3011 assert!(report
3012 .writer_contention
3013 .audit_obligation_append_failures
3014 .is_none());
3015 assert!(report
3016 .writer_contention
3017 .audit_obligation_append_failures_unavailable_reason
3018 .is_some());
3019 assert!(json["writer_contention"]["audit_obligation_append_failures"].is_null());
3020 assert!(
3021 report
3022 .writer_contention
3023 .audit_append_failures_unavailable_reason
3024 .is_some(),
3025 "a direct khive-db snapshot must not fabricate a runtime audit count"
3026 );
3027 }
3028
3029 #[test]
3030 fn diagnostics_exposes_reader_saturation_and_completed_hold_evidence() {
3031 let dir = tempfile::tempdir().expect("tempdir");
3032 let pool = ConnectionPool::new(PoolConfig {
3034 path: Some(dir.path().join("reader_saturation.db")),
3035 max_readers: 1,
3036 checkout_timeout: Duration::from_millis(2),
3037 ..PoolConfig::default()
3038 })
3039 .expect("one-reader file-backed pool");
3040 let held = pool.reader().expect("first reader checkout");
3041 assert!(
3042 pool.reader().is_err(),
3043 "the live checkout must exhaust the one-slot reader budget"
3044 );
3045 drop(held);
3046
3047 let report = collect(
3048 &pool,
3049 BuildIdentity::from_env("9.9.9", None),
3050 Duration::from_secs(30),
3051 );
3052 let reader = report.reader_contention;
3053 assert_eq!(reader.reader_admission_capacity, 1);
3054 assert_eq!(reader.available_reader_admission_slots, 1);
3055 assert_eq!(reader.reader_acquisitions, 1);
3056 assert_eq!(reader.pooled_reader_checkouts, 1);
3057 assert_eq!(reader.standalone_reader_opens, 0);
3058 assert_eq!(reader.infrastructure_standalone_reader_opens, 0);
3059 assert_eq!(reader.reader_checkout_timeouts, 1);
3060 assert_eq!(reader.active_pooled_reader_checkouts, 0);
3061 assert_eq!(reader.peak_active_pooled_reader_checkouts, 1);
3062 assert_eq!(reader.completed_pooled_reader_checkouts, 1);
3063 assert!(reader.max_completed_reader_hold_micros > 0);
3064
3065 let json = serde_json::to_value(&report).expect("report serializes");
3066 assert_eq!(
3067 json.pointer("/reader_contention/reader_admission_capacity"),
3068 Some(&serde_json::json!(1)),
3069 "the operator wire payload must expose the reader admission budget"
3070 );
3071 assert_eq!(
3072 json.pointer("/reader_contention/reader_checkout_timeouts"),
3073 Some(&serde_json::json!(1)),
3074 "the operator wire payload must expose the reader timeout phase"
3075 );
3076 assert!(
3077 json.pointer("/reader_contention/max_completed_reader_hold_micros")
3078 .is_some(),
3079 "the operator wire payload must expose completed hold-time evidence"
3080 );
3081 }
3082
3083 #[test]
3089 fn diagnostics_counts_reader_busy_handler_timeouts_apart_from_checkout_timeouts() {
3090 let dir = tempfile::tempdir().expect("tempdir");
3091 let path = dir.path().join("reader_busy_timeouts.db");
3092 let pool = ConnectionPool::new(PoolConfig {
3093 path: Some(path.clone()),
3094 wal_mode: false,
3095 write_queue_enabled: Some(false),
3096 busy_timeout: Duration::from_millis(50),
3097 ..PoolConfig::default()
3098 })
3099 .expect("rollback-journal file-backed pool");
3100 pool.writer()
3101 .expect("writer")
3102 .conn()
3103 .execute_batch("CREATE TABLE busy_fixture (id INTEGER PRIMARY KEY)")
3104 .expect("fixture table");
3105
3106 let reader = pool.reader().expect("reader checkout");
3108 let rows = reader
3109 .query_row("SELECT count(*) FROM busy_fixture", [], |row| {
3110 row.get::<_, i64>(0)
3111 })
3112 .expect("unlocked read");
3113 assert_eq!(rows, 0);
3114 drop(reader);
3115 assert_eq!(
3116 ReaderContentionDiagnostics::snapshot(&pool).reader_busy_timeouts,
3117 0
3118 );
3119
3120 let holder = Connection::open(&path).expect("second connection");
3121 holder
3122 .execute_batch("BEGIN EXCLUSIVE")
3123 .expect("exclusive lock");
3124 let reader = pool.reader().expect("reader checkout");
3125 let refused = reader
3126 .query_row("SELECT count(*) FROM busy_fixture", [], |row| {
3127 row.get::<_, i64>(0)
3128 })
3129 .expect_err("a read behind an exclusive lock must be refused");
3130 assert!(
3131 matches!(
3132 &refused,
3133 crate::SqliteError::Rusqlite(rusqlite::Error::SqliteFailure(code, _))
3134 if code.code == rusqlite::ErrorCode::DatabaseBusy
3135 ),
3136 "the refusal must be SQLITE_BUSY: {refused}"
3137 );
3138 holder.execute_batch("ROLLBACK").expect("release lock");
3139 drop(reader);
3140
3141 let snapshot = ReaderContentionDiagnostics::snapshot(&pool);
3142 assert_eq!(snapshot.reader_busy_timeouts, 1);
3143 assert_eq!(
3144 snapshot.reader_checkout_timeouts, 0,
3145 "a busy-handler refusal after checkout is not a checkout timeout"
3146 );
3147 let json = serde_json::to_value(snapshot).expect("snapshot serializes");
3148 assert_eq!(
3149 json.pointer("/reader_busy_timeouts"),
3150 Some(&serde_json::json!(1)),
3151 "the operator wire payload must expose the busy-handler count"
3152 );
3153 }
3154
3155 #[test]
3156 fn diagnostics_reports_configured_reader_budget_and_both_deadlines() {
3157 let pool = ConnectionPool::new(PoolConfig {
3158 max_readers: 6,
3159 checkout_timeout: Duration::from_millis(17),
3160 busy_timeout: Duration::from_millis(31),
3161 ..PoolConfig::default()
3162 })
3163 .expect("in-memory pool");
3164 let report = collect(
3165 &pool,
3166 BuildIdentity::from_env("9.9.9", None),
3167 Duration::from_secs(30),
3168 );
3169 let reader = report.reader_contention;
3170 assert_eq!(reader.reader_admission_capacity, 1);
3171
3172 let json = serde_json::to_value(&report).expect("report serializes");
3173 assert_eq!(
3174 json.pointer("/reader_contention/configured_reader_cap"),
3175 Some(&serde_json::json!(6))
3176 );
3177 assert_eq!(
3178 json.pointer("/reader_contention/configured_checkout_timeout_ms"),
3179 Some(&serde_json::json!(17))
3180 );
3181 assert_eq!(
3182 json.pointer("/reader_contention/configured_busy_timeout_ms"),
3183 Some(&serde_json::json!(31))
3184 );
3185 }
3186
3187 #[test]
3188 fn diagnostics_composes_file_backed_standalone_acquisitions_without_counting_its_probe() {
3189 let dir = tempfile::tempdir().expect("tempdir");
3190 let (pool, _path) = seeded_pool(&dir);
3191
3192 drop(
3193 pool.open_standalone_writer()
3194 .expect("write-traffic standalone connection"),
3195 );
3196
3197 let report = collect_with_audit_append_failures(
3198 &pool,
3199 BuildIdentity::from_env("9.9.9", None),
3200 Duration::from_secs(30),
3201 0,
3202 );
3203 assert_eq!(report.writer_contention.writer_acquisitions, 2);
3204 assert_eq!(report.writer_contention.pooled_writer_acquisitions, 1);
3205 assert_eq!(report.writer_contention.standalone_writer_acquisitions, 1);
3206 assert_eq!(report.writer_contention.writer_task_acquisitions, 0);
3207 assert_eq!(report.writer_contention.writer_acquisition_timeouts, 0);
3208
3209 let second = collect_with_audit_append_failures(
3210 &pool,
3211 BuildIdentity::from_env("9.9.9", None),
3212 Duration::from_secs(30),
3213 0,
3214 );
3215 assert_eq!(
3216 second.writer_contention, report.writer_contention,
3217 "the diagnostics PASSIVE probe must not inflate write-traffic counters"
3218 );
3219 }
3220
3221 #[test]
3222 fn runtime_aware_collect_exposes_the_supplied_audit_failure_counter() {
3223 let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool");
3224
3225 let report = collect_with_audit_append_failures(
3226 &pool,
3227 BuildIdentity::from_env("9.9.9", None),
3228 Duration::from_secs(30),
3229 17,
3230 );
3231
3232 assert_eq!(report.writer_contention.audit_append_failures, Some(17));
3233 assert!(
3234 report
3235 .writer_contention
3236 .audit_append_failures_unavailable_reason
3237 .is_none(),
3238 "a supplied runtime counter must not carry an unavailable reason"
3239 );
3240 }
3241
3242 #[test]
3243 fn diagnostics_exposes_an_induced_writer_checkout_timeout() {
3244 let pool = ConnectionPool::new(PoolConfig {
3245 checkout_timeout: Duration::from_millis(1),
3246 ..PoolConfig::default()
3247 })
3248 .expect("in-memory pool");
3249
3250 let held = pool.writer().expect("first writer checkout succeeds");
3251 assert!(
3252 matches!(
3253 pool.writer(),
3254 Err(crate::SqliteError::WriterPoolCheckoutTimeout { .. })
3255 ),
3256 "holding the sole writer must exercise the typed timeout path"
3257 );
3258 drop(held);
3259
3260 let report = collect_with_audit_append_failures(
3261 &pool,
3262 BuildIdentity::from_env("9.9.9", None),
3263 Duration::from_secs(30),
3264 0,
3265 );
3266 assert_eq!(report.writer_contention.writer_acquisitions, 1);
3267 assert_eq!(report.writer_contention.pooled_writer_acquisitions, 1);
3268 assert_eq!(report.writer_contention.standalone_writer_acquisitions, 0);
3269 assert_eq!(report.writer_contention.writer_task_acquisitions, 0);
3270 assert_eq!(report.writer_contention.writer_acquisition_timeouts, 1);
3271 }
3272
3273 #[test]
3276 fn never_observed_sentinels_serialize_as_null() {
3277 let counters = CheckpointCounters {
3278 last_observed_wal_pages: None,
3279 truncate_attempts: 0,
3280 truncate_consecutive_failures: 0,
3281 checkpoint_skipped_ticks: 0,
3282 checkpoint_consecutive_skips: 0,
3283 checkpoint_last_skip_wal_pages: None,
3284 checkpoint_pressure_elevated_ticks: 0,
3285 checkpoint_pressure_episodes_started: 0,
3286 checkpoint_pressure_episodes_recovered: 0,
3287 checkpoint_lifecycle_append_attempts: 0,
3288 checkpoint_lifecycle_append_failures: 0,
3289 checkpoint_lifecycle_enqueue_drops: 0,
3290 read_tx_max_age_evictions: 0,
3291 };
3292 let json = serde_json::to_value(counters).expect("serializes");
3293 assert!(json["last_observed_wal_pages"].is_null());
3294 assert!(json["checkpoint_last_skip_wal_pages"].is_null());
3295 }
3296
3297 #[test]
3298 fn graph_edge_integrity_detects_the_pre_v14_duplicate_state() {
3299 let conn = Connection::open_in_memory().expect("in-memory sqlite");
3300 conn.execute_batch(
3301 "CREATE TABLE graph_edges (
3302 namespace TEXT NOT NULL,
3303 id TEXT NOT NULL,
3304 PRIMARY KEY (namespace, id)
3305 );
3306 CREATE TABLE graph_edges_seq (
3307 seq INTEGER PRIMARY KEY AUTOINCREMENT,
3308 edge_id TEXT NOT NULL UNIQUE
3309 );
3310 CREATE TABLE entities (id TEXT PRIMARY KEY, deleted_at INTEGER, merged_into TEXT);
3311 INSERT INTO graph_edges(namespace, id)
3312 VALUES ('alpha', 'shared-edge'), ('beta', 'shared-edge');
3313 INSERT INTO graph_edges_seq(edge_id) VALUES ('shared-edge');",
3314 )
3315 .expect("seed the state possible before the V14 uniqueness guard");
3316
3317 let integrity = graph_edge_integrity(&conn).expect("integrity query succeeds");
3318
3319 assert_eq!(integrity.duplicate_edge_id_groups, 1);
3320 assert_eq!(integrity.graph_edges_rows, 2);
3321 assert_eq!(integrity.graph_edges_seq_rows, 1);
3322 assert!(integrity.pre_v14_duplicate_edge_state_detected);
3323 }
3324
3325 #[test]
3326 fn graph_edge_integrity_counts_only_live_rows_that_still_carry_merge_provenance() {
3327 let conn = Connection::open_in_memory().expect("in-memory sqlite");
3328 conn.execute_batch(
3329 "CREATE TABLE graph_edges (
3330 namespace TEXT NOT NULL,
3331 id TEXT NOT NULL,
3332 PRIMARY KEY (namespace, id)
3333 );
3334 CREATE TABLE graph_edges_seq (
3335 seq INTEGER PRIMARY KEY AUTOINCREMENT,
3336 edge_id TEXT NOT NULL UNIQUE
3337 );
3338 CREATE TABLE entities (id TEXT PRIMARY KEY, deleted_at INTEGER, merged_into TEXT);
3339 INSERT INTO entities(id, deleted_at, merged_into) VALUES
3340 ('kept', NULL, NULL),
3341 ('tombstoned-source', 1, 'kept'),
3342 ('left-live-by-an-old-restore', NULL, 'kept'),
3343 ('plain-soft-delete', 1, NULL);",
3344 )
3345 .expect("seed one row of each shape");
3346
3347 let integrity = graph_edge_integrity(&conn).expect("integrity query succeeds");
3348
3349 assert_eq!(
3350 integrity.live_entities_carrying_merged_into, 1,
3351 "a merge tombstone and a plain live row are both in order; only the live row \
3352 carrying merged_into is the invariant violation"
3353 );
3354 assert_eq!(integrity.duplicate_edge_id_groups, 0);
3355 }
3356
3357 #[test]
3358 fn graph_edge_integrity_does_not_mislabel_retained_delete_history_as_a_duplicate() {
3359 let conn = Connection::open_in_memory().expect("in-memory sqlite");
3360 conn.execute_batch(
3361 "CREATE TABLE graph_edges (
3362 namespace TEXT NOT NULL,
3363 id TEXT NOT NULL,
3364 PRIMARY KEY (namespace, id)
3365 );
3366 CREATE TABLE graph_edges_seq (
3367 seq INTEGER PRIMARY KEY AUTOINCREMENT,
3368 edge_id TEXT NOT NULL UNIQUE
3369 );
3370 CREATE TABLE entities (id TEXT PRIMARY KEY, deleted_at INTEGER, merged_into TEXT);
3371 INSERT INTO graph_edges(namespace, id) VALUES ('local', 'live-edge');
3372 INSERT INTO graph_edges_seq(edge_id)
3373 VALUES ('deleted-edge'), ('live-edge');",
3374 )
3375 .expect("seed a retained sequence row for a hard-deleted edge");
3376
3377 let integrity = graph_edge_integrity(&conn).expect("integrity query succeeds");
3378
3379 assert_eq!(integrity.duplicate_edge_id_groups, 0);
3380 assert_eq!(integrity.graph_edges_rows, 1);
3381 assert_eq!(integrity.graph_edges_seq_rows, 2);
3382 assert!(
3383 !integrity.pre_v14_duplicate_edge_state_detected,
3384 "ledger rows intentionally survive hard deletion; count mismatch alone is not the \
3385 pre-V14 duplicate state"
3386 );
3387 }
3388
3389 #[cfg(unix)]
3390 #[test]
3391 fn unmeasured_sidecar_cleanup_fields_are_absent_from_the_wire_payload() {
3392 let pin = wal_pin_attribution_from_census(crate::walpin::CensusResult {
3393 holders: std::collections::HashSet::new(),
3394 uninspectable_pids: Vec::new(),
3395 truncated: false,
3396 budget_exhausted: false,
3397 });
3398
3399 assert_eq!(pin.sidecar_listing_truncated, None);
3400 assert_eq!(pin.sidecar_entries_cleanup_would_reap, None);
3401 let json = serde_json::to_value(pin).expect("attribution serializes");
3402 assert!(
3403 json.get("sidecar_listing_truncated").is_none(),
3404 "a skipped enumeration must omit sidecar_listing_truncated, not fabricate false"
3405 );
3406 assert!(
3407 json.get("sidecar_entries_cleanup_would_reap").is_none(),
3408 "a skipped enumeration must omit sidecar_entries_cleanup_would_reap, not fabricate 0"
3409 );
3410 }
3411
3412 #[cfg(unix)]
3413 #[test]
3414 fn wal_pin_census_serializes_only_under_the_nested_carrier() {
3415 let pin = wal_pin_attribution_from_census(crate::walpin::CensusResult {
3416 holders: std::collections::HashSet::from([41, 7]),
3417 uninspectable_pids: vec![99],
3418 truncated: true,
3419 budget_exhausted: false,
3420 });
3421
3422 let json = serde_json::to_value(pin).expect("attribution serializes");
3423 assert_eq!(json["census"]["holder_pids"], serde_json::json!([7, 41]));
3424 assert_eq!(
3425 json["census"]["uninspectable_pids"],
3426 serde_json::json!([99])
3427 );
3428 assert_eq!(json["census"]["truncated"], true);
3429
3430 for duplicate in [
3431 "census_holder_pids",
3432 "census_uninspectable_pids",
3433 "census_truncated",
3434 "census_is_complete",
3435 ] {
3436 assert!(
3437 json.get(duplicate).is_none(),
3438 "wal_pin.{duplicate} must not duplicate wal_pin.census: {json}"
3439 );
3440 }
3441 }
3442
3443 #[test]
3447 fn probe_refuses_a_missing_configured_path_without_creating_it() {
3448 let dir = tempfile::tempdir().expect("tempdir");
3449 let (pool, path) = seeded_pool(&dir);
3450
3451 for suffix in ["", "-wal", "-shm"] {
3452 let mut p = path.as_os_str().to_os_string();
3453 p.push(suffix);
3454 let _ = std::fs::remove_file(PathBuf::from(p));
3455 }
3456 assert!(!path.exists(), "precondition: the database file is gone");
3457
3458 let report = collect(
3459 &pool,
3460 BuildIdentity::from_env("0.0.0", None),
3461 Duration::from_secs(30),
3462 );
3463
3464 assert!(
3465 report.checkpoint_probe.is_none(),
3466 "a missing database must not yield a probe result: {report:?}"
3467 );
3468 assert!(
3469 report.checkpoint_probe_error.is_some(),
3470 "a missing database must say why there is no probe: {report:?}"
3471 );
3472 assert!(
3473 !path.exists(),
3474 "a diagnostics request must never create the database it was asked about"
3475 );
3476 }
3477
3478 #[test]
3481 fn collect_degrades_gracefully_for_an_in_memory_backend() {
3482 let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool");
3483 let report = collect(
3484 &pool,
3485 BuildIdentity::from_env("0.0.0", None),
3486 Duration::from_secs(30),
3487 );
3488
3489 assert!(report.db_path.is_none());
3490 assert!(report.wal_file.is_none());
3491 assert!(report.checkpoint_probe.is_none());
3492 assert!(report.fts_segments.is_none());
3493 assert!(report.fts_segments_error.is_some());
3494 assert!(
3495 report.checkpoint_probe_error.is_some(),
3496 "an in-memory report must say WHY there is no probe"
3497 );
3498 assert!(!report.wal_pin.available);
3499 assert!(report.wal_pin.unavailable_reason.is_some());
3500 assert_eq!(report.wal_pin.status, WalPinAttributionStatus::Unavailable);
3501 assert!(matches!(
3502 report.wal_pin.census,
3503 WalPinCensus::Unavailable { .. }
3504 ));
3505 assert!(report.size_composition.is_none());
3506 assert!(report
3507 .size_composition_error
3508 .as_deref()
3509 .is_some_and(|reason| reason.contains("no file-backed page composition")));
3510 }
3511
3512 #[cfg(unix)]
3513 #[test]
3514 fn incomplete_holder_census_is_a_tagged_degraded_result() {
3515 let census = crate::walpin::CensusResult {
3516 holders: std::collections::HashSet::from([41, 7]),
3517 uninspectable_pids: vec![99, 99],
3518 truncated: true,
3519 budget_exhausted: false,
3520 };
3521
3522 let pin = wal_pin_attribution_from_census(census);
3523
3524 assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
3525 assert!(!pin.available);
3526 assert!(!pin.census_is_complete);
3527 assert_eq!(pin.census_holder_pids, vec![7, 41]);
3528 assert_eq!(pin.census_uninspectable_pids, vec![99]);
3529 assert!(
3530 pin.unavailable_reason.as_deref().is_some_and(
3531 |reason| reason.contains("additional database holders cannot be ruled out")
3532 ),
3533 "the legacy reason must also fail loud for old consumers: {pin:?}"
3534 );
3535 match &pin.census {
3536 WalPinCensus::Incomplete {
3537 holder_pids,
3538 uninspectable_pids,
3539 truncated,
3540 reason,
3541 } => {
3542 assert_eq!(holder_pids, &vec![7, 41]);
3543 assert_eq!(uninspectable_pids, &vec![99]);
3544 assert!(*truncated);
3545 assert!(reason.contains("additional database holders cannot be ruled out"));
3546 }
3547 other => panic!("incomplete scan must serialize as incomplete, got {other:?}"),
3548 }
3549
3550 let json = serde_json::to_value(&pin).expect("serializes");
3551 assert_eq!(json["status"], "degraded");
3552 assert_eq!(json["census"]["status"], "incomplete");
3553 }
3554
3555 #[cfg(unix)]
3556 #[test]
3557 fn complete_holder_census_stays_explicit_while_attribution_is_degraded() {
3558 let census = crate::walpin::CensusResult {
3559 holders: std::collections::HashSet::from([7]),
3560 uninspectable_pids: Vec::new(),
3561 truncated: false,
3562 budget_exhausted: false,
3563 };
3564
3565 let pin = wal_pin_attribution_from_census(census);
3566
3567 assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
3568 assert!(pin.census_is_complete);
3569 assert!(matches!(
3570 pin.census,
3571 WalPinCensus::Complete { ref holder_pids } if holder_pids == &vec![7]
3572 ));
3573 assert_eq!(
3574 pin.status_reasons.len(),
3575 1,
3576 "only missing sidecar reconciliation degrades a complete OS census"
3577 );
3578 }
3579
3580 #[cfg(unix)]
3581 #[test]
3582 fn complete_holder_and_read_only_sidecar_evidence_reconcile_to_complete() {
3583 let census = crate::walpin::CensusResult {
3584 holders: std::collections::HashSet::from([7]),
3585 uninspectable_pids: Vec::new(),
3586 truncated: false,
3587 budget_exhausted: false,
3588 };
3589 let sidecar = crate::walpin::WalpinReport {
3590 entries: vec![crate::walpin::WalpinPidHealth::RegisteredSilent { pid: 7 }],
3591 sidecar_listing_truncated: false,
3592 cleanup_would_reap: 0,
3593 orphan_temps_reaped: 0,
3594 };
3595
3596 let pin = wal_pin_attribution_from_evidence(census, sidecar);
3597
3598 assert_eq!(pin.status, WalPinAttributionStatus::Complete);
3599 assert!(pin.available);
3600 assert!(pin.fully_attributed);
3601 assert!(pin.status_reasons.is_empty());
3602 assert_eq!(pin.registered_silent_pids, vec![7]);
3603 assert!(pin.census_pids_without_attribution.is_empty());
3604 assert_eq!(pin.sidecar_listing_truncated, Some(false));
3605 assert_eq!(pin.sidecar_entries_cleanup_would_reap, Some(0));
3606 }
3607
3608 #[cfg(unix)]
3609 #[test]
3610 fn complete_census_with_an_unregistered_holder_is_degraded_not_exonerated() {
3611 let census = crate::walpin::CensusResult {
3612 holders: std::collections::HashSet::from([7, 41]),
3613 uninspectable_pids: Vec::new(),
3614 truncated: false,
3615 budget_exhausted: false,
3616 };
3617 let sidecar = crate::walpin::WalpinReport {
3618 entries: vec![crate::walpin::WalpinPidHealth::RegisteredSilent { pid: 7 }],
3619 sidecar_listing_truncated: false,
3620 cleanup_would_reap: 0,
3621 orphan_temps_reaped: 0,
3622 };
3623
3624 let pin = wal_pin_attribution_from_evidence(census, sidecar);
3625
3626 assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
3627 assert!(!pin.available);
3628 assert!(!pin.fully_attributed);
3629 assert_eq!(pin.census_pids_without_attribution, vec![41]);
3630 assert!(pin
3631 .status_reasons
3632 .iter()
3633 .any(|reason| reason.contains("holder(s) have no sidecar attribution")));
3634 }
3635
3636 #[test]
3637 fn database_size_composition_reports_tables_indexes_fts_and_vectors_separately() {
3638 let conn = Connection::open_in_memory().expect("in-memory sqlite");
3639 conn.execute_batch(
3640 "CREATE TABLE docs(id INTEGER PRIMARY KEY, body TEXT NOT NULL); \
3641 CREATE INDEX idx_docs_body ON docs(body); \
3642 CREATE TABLE fts_demo_data(id INTEGER PRIMARY KEY, block BLOB); \
3643 CREATE TABLE vec_demo_chunks(id INTEGER PRIMARY KEY, vectors BLOB); \
3644 CREATE TABLE knowledge_sections(id INTEGER PRIMARY KEY, embedding BLOB); \
3645 INSERT INTO docs(body) VALUES (zeroblob(8192)); \
3646 INSERT INTO fts_demo_data(block) VALUES (zeroblob(8192)); \
3647 INSERT INTO vec_demo_chunks(vectors) VALUES (zeroblob(8192)); \
3648 INSERT INTO knowledge_sections(embedding) VALUES (zeroblob(8192));",
3649 )
3650 .expect("seed size classes");
3651
3652 let composition = database_size_composition(&conn).expect("dbstat composition");
3653 let class_for = |name: &str| {
3654 composition
3655 .objects
3656 .iter()
3657 .find(|object| object.name == name)
3658 .map(|object| object.storage_class)
3659 };
3660
3661 assert_eq!(class_for("docs"), Some(DatabaseStorageClass::RowTable));
3662 assert_eq!(
3663 class_for("idx_docs_body"),
3664 Some(DatabaseStorageClass::Index)
3665 );
3666 assert_eq!(
3667 class_for("fts_demo_data"),
3668 Some(DatabaseStorageClass::FullText)
3669 );
3670 assert_eq!(
3671 class_for("vec_demo_chunks"),
3672 Some(DatabaseStorageClass::Vector)
3673 );
3674 assert_eq!(
3675 class_for("knowledge_sections"),
3676 Some(DatabaseStorageClass::MixedRowAndEmbedding)
3677 );
3678 assert!(composition.vector_bytes > 0);
3679 assert!(composition.full_text_bytes > 0);
3680 assert!(composition.mixed_embedding_bytes > 0);
3681 assert_eq!(
3682 composition
3683 .accounted_bytes
3684 .saturating_add(composition.freelist_bytes)
3685 .saturating_add(composition.unaccounted_bytes),
3686 composition.database_bytes
3687 );
3688 }
3689
3690 #[cfg(unix)]
3693 #[test]
3694 #[serial(khive_walpin_sidecar_env)]
3695 fn wal_pin_attribution_degrades_when_a_holder_has_no_sidecar_registration() {
3696 let dir = tempfile::tempdir().expect("tempdir");
3697 let (pool, path) = seeded_pool(&dir);
3698 let _ = &pool;
3699
3700 let pin = wal_pin_attribution(&path, Duration::from_secs(30));
3701
3702 assert!(
3703 !pin.fully_attributed,
3704 "the pool's OS holder has no test sidecar registration"
3705 );
3706 assert!(
3707 pin.unavailable_reason.is_some(),
3708 "the missing holder attribution must be explained: {pin:?}"
3709 );
3710 assert!(pin.sidecar_entries.is_empty());
3711 assert!(pin.reporting.is_empty());
3712 assert_eq!(pin.sidecar_listing_truncated, Some(false));
3713 assert_eq!(pin.sidecar_entries_cleanup_would_reap, Some(0));
3714 assert_eq!(pin.status, WalPinAttributionStatus::Degraded);
3715 assert!(matches!(
3716 pin.census,
3717 WalPinCensus::Complete { .. } | WalPinCensus::Incomplete { .. }
3718 ));
3719 }
3720
3721 #[cfg(unix)]
3728 #[test]
3729 #[serial(khive_walpin_sidecar_env)]
3730 fn diagnostics_finds_sidecar_evidence_through_an_aliased_database_path() {
3731 let dir = tempfile::tempdir().expect("tempdir");
3732 let real_dir = dir.path().join("real");
3733 std::fs::create_dir(&real_dir).expect("mkdir real dir");
3734 let real_path = real_dir.join("diag.db");
3735 std::fs::write(&real_path, b"").expect("create real file");
3736 let alias_path = dir.path().join("alias.db");
3737 std::os::unix::fs::symlink(&real_path, &alias_path).expect("symlink alias");
3738
3739 let pool = ConnectionPool::new(PoolConfig {
3740 path: Some(alias_path.clone()),
3741 ..PoolConfig::for_test()
3742 })
3743 .expect("pool open through symlinked path");
3744 {
3745 let writer = pool.try_writer().expect("writer");
3746 writer
3747 .conn()
3748 .execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
3749 .expect("seed a write so the WAL file exists");
3750 }
3751
3752 let canonical = pool
3753 .canonical_path()
3754 .expect("file-backed pool")
3755 .to_path_buf();
3756 assert_ne!(
3757 canonical, alias_path,
3758 "the alias must actually differ from the canonical path for this test to mean \
3759 anything"
3760 );
3761
3762 let pid = std::process::id();
3763 let sidecar_dir = crate::walpin::sidecar_dir_for(&canonical);
3764 let beacon = crate::walpin::WalpinBeacon {
3765 pid,
3766 process_role: "session".to_string(),
3767 started_at: crate::walpin::process_start_time_secs(pid).unwrap_or(0),
3768 sweep_interval_ms: 5_000,
3769 };
3770 crate::walpin::write_beacon(&sidecar_dir, &beacon).expect("seed this process's beacon");
3771
3772 let report = collect(
3773 &pool,
3774 BuildIdentity::from_env("test", None),
3775 Duration::from_secs(30),
3776 );
3777
3778 assert_eq!(
3787 report.wal_pin.sidecar_listing_truncated,
3788 Some(false),
3789 "the sidecar enumeration must run to completion: {:?}",
3790 report.wal_pin
3791 );
3792 assert!(
3793 report.wal_pin.registered_silent_pids.contains(&pid),
3794 "the beacon written beside the canonical path must be found: {:?}",
3795 report.wal_pin
3796 );
3797 assert!(
3798 report.wal_pin.census_holder_pids.contains(&pid),
3799 "the OS census must find this process holding its own database open: {:?}",
3800 report.wal_pin
3801 );
3802 assert!(
3803 !report
3804 .wal_pin
3805 .census_pids_without_attribution
3806 .contains(&pid),
3807 "this process's own holder entry must be attributed by its own sidecar evidence, \
3808 not left unexplained: {:?}",
3809 report.wal_pin
3810 );
3811 }
3812
3813 #[cfg(unix)]
3821 #[tokio::test]
3822 #[serial(khive_walpin_sidecar_env)]
3823 async fn diagnostics_finds_sidecar_evidence_through_an_aliased_database_path_async() {
3824 let dir = tempfile::tempdir().expect("tempdir");
3825 let real_dir = dir.path().join("real");
3826 std::fs::create_dir(&real_dir).expect("mkdir real dir");
3827 let real_path = real_dir.join("diag.db");
3828 std::fs::write(&real_path, b"").expect("create real file");
3829 let alias_path = dir.path().join("alias.db");
3830 std::os::unix::fs::symlink(&real_path, &alias_path).expect("symlink alias");
3831
3832 let pool = ConnectionPool::new(PoolConfig {
3833 path: Some(alias_path.clone()),
3834 ..PoolConfig::for_test()
3835 })
3836 .expect("pool open through symlinked path");
3837 {
3838 let writer = pool.try_writer().expect("writer");
3839 writer
3840 .conn()
3841 .execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
3842 .expect("seed a write so the WAL file exists");
3843 }
3844
3845 let canonical = pool
3846 .canonical_path()
3847 .expect("file-backed pool")
3848 .to_path_buf();
3849 assert_ne!(
3850 canonical, alias_path,
3851 "the alias must actually differ from the canonical path for this test to mean \
3852 anything"
3853 );
3854
3855 let pid = std::process::id();
3856 let sidecar_dir = crate::walpin::sidecar_dir_for(&canonical);
3857 let beacon = crate::walpin::WalpinBeacon {
3858 pid,
3859 process_role: "session".to_string(),
3860 started_at: crate::walpin::process_start_time_secs(pid).unwrap_or(0),
3861 sweep_interval_ms: 5_000,
3862 };
3863 crate::walpin::write_beacon(&sidecar_dir, &beacon).expect("seed this process's beacon");
3864
3865 let pool = Arc::new(pool);
3866 let report = collect_with_audit_append_failures_interruptibly(
3867 Arc::clone(&pool),
3868 BuildIdentity::from_env("test", None),
3869 Duration::from_secs(30),
3870 0,
3871 )
3872 .await
3873 .expect("diagnostics succeed");
3874
3875 assert_eq!(
3880 report.wal_pin.sidecar_listing_truncated,
3881 Some(false),
3882 "the sidecar enumeration must run to completion: {:?}",
3883 report.wal_pin
3884 );
3885 assert!(
3886 report.wal_pin.registered_silent_pids.contains(&pid),
3887 "the beacon written beside the canonical path must be found: {:?}",
3888 report.wal_pin
3889 );
3890 assert!(
3891 report.wal_pin.census_holder_pids.contains(&pid),
3892 "the OS census must find this process holding its own database open: {:?}",
3893 report.wal_pin
3894 );
3895 assert!(
3896 !report
3897 .wal_pin
3898 .census_pids_without_attribution
3899 .contains(&pid),
3900 "this process's own holder entry must be attributed by its own sidecar evidence, \
3901 not left unexplained: {:?}",
3902 report.wal_pin
3903 );
3904 }
3905}