1use std::path::{Path, PathBuf};
27use std::sync::atomic::{AtomicU64, Ordering};
28use std::sync::Arc;
29use std::time::{Duration, Instant};
30
31use crate::pool::{ConnectionPool, WriterGuard};
32
33static LAST_WAL_PAGES: AtomicU64 = AtomicU64::new(u64::MAX);
42
43static TRUNCATE_ATTEMPTS: AtomicU64 = AtomicU64::new(0);
46
47static TRUNCATE_CONSECUTIVE_FAILURES: AtomicU64 = AtomicU64::new(0);
51
52static CHECKPOINT_SKIPPED_TICKS: AtomicU64 = AtomicU64::new(0);
56
57static CHECKPOINT_CONSECUTIVE_SKIPS: AtomicU64 = AtomicU64::new(0);
61
62static CHECKPOINT_LAST_SKIP_WAL_PAGES: AtomicU64 = AtomicU64::new(u64::MAX);
66
67pub fn last_observed_wal_pages() -> Option<u64> {
70 match LAST_WAL_PAGES.load(Ordering::Relaxed) {
71 u64::MAX => None,
72 pages => Some(pages),
73 }
74}
75
76pub fn truncate_attempts() -> u64 {
78 TRUNCATE_ATTEMPTS.load(Ordering::Relaxed)
79}
80
81pub fn truncate_consecutive_failures() -> u64 {
83 TRUNCATE_CONSECUTIVE_FAILURES.load(Ordering::Relaxed)
84}
85
86pub fn checkpoint_skipped_ticks() -> u64 {
88 CHECKPOINT_SKIPPED_TICKS.load(Ordering::Relaxed)
89}
90
91pub fn checkpoint_consecutive_skips() -> u64 {
93 CHECKPOINT_CONSECUTIVE_SKIPS.load(Ordering::Relaxed)
94}
95
96pub fn checkpoint_last_skip_wal_pages() -> Option<u64> {
99 match CHECKPOINT_LAST_SKIP_WAL_PAGES.load(Ordering::Relaxed) {
100 u64::MAX => None,
101 pages => Some(pages),
102 }
103}
104
105fn note_checkpoint_skipped() {
109 CHECKPOINT_SKIPPED_TICKS.fetch_add(1, Ordering::Relaxed);
110 CHECKPOINT_CONSECUTIVE_SKIPS.fetch_add(1, Ordering::Relaxed);
111 if let Some(pages) = last_observed_wal_pages() {
112 CHECKPOINT_LAST_SKIP_WAL_PAGES.store(pages, Ordering::Relaxed);
113 }
114}
115
116fn note_checkpoint_observed(_wal_pages: u64) {
121 CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
122}
123
124#[cfg(test)]
128pub(crate) fn reset_checkpoint_metrics_for_tests() {
129 CHECKPOINT_SKIPPED_TICKS.store(0, Ordering::Relaxed);
130 CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
131 CHECKPOINT_LAST_SKIP_WAL_PAGES.store(u64::MAX, Ordering::Relaxed);
132}
133
134#[derive(Debug, Clone, Copy, PartialEq, Eq)]
142pub enum CheckpointTick {
143 Skipped,
145 Observed(u64),
147}
148
149pub const DEFAULT_WARN_SUSTAINED_CYCLES: u8 = 3;
152
153#[derive(Clone, Debug)]
158pub struct CheckpointConfig {
159 pub interval: Duration,
164
165 pub warn_pages: u64,
170
171 pub warn_sustained_cycles: u8,
179
180 pub high_water_pages: u64,
190
191 pub truncate_high_water_pages: u64,
202
203 pub truncate_min_interval: Duration,
213
214 pub truncate_busy_timeout: Duration,
221
222 pub tx_warn_secs: Duration,
230
231 pub tx_max_age_secs: Duration,
239}
240
241impl Default for CheckpointConfig {
242 fn default() -> Self {
243 Self {
244 interval: Duration::from_millis(500),
245 warn_pages: 2000,
246 warn_sustained_cycles: DEFAULT_WARN_SUSTAINED_CYCLES,
247 high_water_pages: 6000,
248 truncate_high_water_pages: 20_000,
249 truncate_min_interval: Duration::from_secs(300),
250 truncate_busy_timeout: Duration::from_millis(2000),
251 tx_warn_secs: Duration::from_secs(30),
252 tx_max_age_secs: Duration::from_secs(120),
253 }
254 }
255}
256
257impl CheckpointConfig {
258 pub fn from_env() -> Self {
262 let mut cfg = Self::default();
263
264 if let Ok(ms) = std::env::var("KHIVE_CHECKPOINT_INTERVAL_MS") {
265 if let Ok(v) = ms.parse::<u64>() {
266 if v > 0 {
267 cfg.interval = Duration::from_millis(v);
268 }
269 }
270 }
271
272 if let Ok(v) = std::env::var("KHIVE_WAL_WARN_PAGES") {
273 if let Ok(n) = v.parse::<u64>() {
274 if n > 0 {
275 cfg.warn_pages = n;
276 }
277 }
278 }
279
280 if let Ok(v) = std::env::var("KHIVE_WAL_WARN_SUSTAINED_CYCLES") {
281 if let Ok(n) = v.parse::<u8>() {
282 if n > 0 {
283 cfg.warn_sustained_cycles = n;
284 }
285 }
286 }
287
288 if let Ok(v) = std::env::var("KHIVE_WAL_HIGH_WATER_PAGES") {
289 if let Ok(n) = v.parse::<u64>() {
290 if n > 0 {
291 cfg.high_water_pages = n;
292 }
293 }
294 }
295
296 if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES") {
297 if let Ok(n) = v.parse::<u64>() {
298 if n > 0 {
299 cfg.truncate_high_water_pages = n;
300 }
301 }
302 }
303
304 if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS") {
305 if let Ok(n) = v.parse::<u64>() {
306 if n > 0 {
307 cfg.truncate_min_interval = Duration::from_secs(n);
308 }
309 }
310 }
311
312 if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_BUSY_MS") {
313 if let Ok(n) = v.parse::<u64>() {
314 if n > 0 {
315 cfg.truncate_busy_timeout = Duration::from_millis(n);
316 }
317 }
318 }
319
320 (cfg.tx_warn_secs, cfg.tx_max_age_secs) =
321 tx_age_thresholds_from_env(cfg.tx_warn_secs, cfg.tx_max_age_secs);
322
323 cfg
324 }
325}
326
327fn tx_age_thresholds_from_env(
341 default_warn: Duration,
342 default_max: Duration,
343) -> (Duration, Duration) {
344 let mut warn_secs = default_warn;
345 let mut max_age_secs = default_max;
346
347 if let Ok(v) = std::env::var("KHIVE_TX_WARN_SECS") {
348 if let Ok(n) = v.parse::<u64>() {
349 if n > 0 {
350 warn_secs = Duration::from_secs(n);
351 }
352 }
353 }
354
355 if let Ok(v) = std::env::var("KHIVE_TX_MAX_AGE_SECS") {
356 if let Ok(n) = v.parse::<u64>() {
357 if n > 0 {
358 max_age_secs = Duration::from_secs(n);
359 }
360 }
361 }
362
363 if warn_secs >= max_age_secs {
364 tracing::warn!(
365 configured_tx_warn_secs = warn_secs.as_secs_f64(),
366 configured_tx_max_age_secs = max_age_secs.as_secs_f64(),
367 fallback_tx_warn_secs = default_warn.as_secs_f64(),
368 fallback_tx_max_age_secs = default_max.as_secs_f64(),
369 "KHIVE_TX_WARN_SECS must be strictly less than KHIVE_TX_MAX_AGE_SECS; \
370 both transaction-age thresholds were rejected and reset to their defaults"
371 );
372 return (default_warn, default_max);
373 }
374
375 (warn_secs, max_age_secs)
376}
377
378#[derive(Debug, Default)]
385pub struct TruncateState {
386 last_attempt: Option<Instant>,
390 consecutive_failures: u32,
395}
396
397#[derive(Debug, Clone, Copy, PartialEq, Eq)]
404pub enum CheckpointSeverityRung {
405 Info,
407 Warn,
410 Alarm,
413}
414
415#[derive(Debug, Default, Clone)]
419pub struct CheckpointSeverityState {
420 was_above_warn: bool,
423 consecutive_above_warn: u8,
426 warn_emitted_for_episode: bool,
430}
431
432#[derive(Debug, Clone, Copy, PartialEq, Eq)]
435pub struct CheckpointSeverityEmission {
436 pub rung: CheckpointSeverityRung,
440 pub wal_pages: u64,
442 pub threshold_pages: u64,
444 pub consecutive_cycles: u8,
447}
448
449impl CheckpointSeverityState {
450 pub fn observe_wal_pages(
461 &mut self,
462 wal_pages: u64,
463 config: &CheckpointConfig,
464 ) -> Vec<CheckpointSeverityEmission> {
465 let mut emissions = Vec::new();
466 let above_warn = wal_pages >= config.warn_pages;
467
468 if above_warn {
469 self.consecutive_above_warn = self.consecutive_above_warn.saturating_add(1);
470
471 if !self.was_above_warn {
472 emissions.push(CheckpointSeverityEmission {
473 rung: CheckpointSeverityRung::Info,
474 wal_pages,
475 threshold_pages: config.warn_pages,
476 consecutive_cycles: self.consecutive_above_warn,
477 });
478 }
479
480 if !self.warn_emitted_for_episode
481 && self.consecutive_above_warn >= config.warn_sustained_cycles
482 {
483 emissions.push(CheckpointSeverityEmission {
484 rung: CheckpointSeverityRung::Warn,
485 wal_pages,
486 threshold_pages: config.warn_pages,
487 consecutive_cycles: self.consecutive_above_warn,
488 });
489 self.warn_emitted_for_episode = true;
490 }
491 } else {
492 self.consecutive_above_warn = 0;
493 self.warn_emitted_for_episode = false;
494 }
495
496 self.was_above_warn = above_warn;
497 emissions
498 }
499}
500
501#[derive(Debug, Clone, Copy, PartialEq, Eq)]
505pub enum TxAgeRung {
506 Warn,
508 Stale,
513}
514
515#[derive(Debug, Clone, PartialEq, Eq)]
517pub struct TxAgeEmission {
518 pub rung: TxAgeRung,
519 pub age: Duration,
520 pub label: Option<String>,
521}
522
523#[derive(Debug, Default, Clone)]
534pub struct TxAgeSweepState {
535 was_above_warn: bool,
538 was_above_max_age: bool,
541 tracked_id: Option<khive_storage::tx_registry::TxId>,
546}
547
548impl TxAgeSweepState {
549 pub fn observe(
561 &mut self,
562 oldest: Option<(khive_storage::tx_registry::TxId, Duration, Option<String>)>,
563 tx_warn_secs: Duration,
564 tx_max_age_secs: Duration,
565 ) -> Vec<TxAgeEmission> {
566 let mut emissions = Vec::new();
567
568 let Some((id, age, label)) = oldest else {
569 self.was_above_warn = false;
570 self.was_above_max_age = false;
571 self.tracked_id = None;
572 return emissions;
573 };
574
575 if self.tracked_id != Some(id) {
576 self.was_above_warn = false;
577 self.was_above_max_age = false;
578 }
579 self.tracked_id = Some(id);
580
581 let above_warn = age >= tx_warn_secs;
582 let above_max_age = age >= tx_max_age_secs;
583
584 if above_warn && !self.was_above_warn {
585 emissions.push(TxAgeEmission {
586 rung: TxAgeRung::Warn,
587 age,
588 label: label.clone(),
589 });
590 }
591 if above_max_age && !self.was_above_max_age {
592 emissions.push(TxAgeEmission {
593 rung: TxAgeRung::Stale,
594 age,
595 label,
596 });
597 }
598
599 self.was_above_warn = above_warn;
600 self.was_above_max_age = above_max_age;
601 emissions
602 }
603}
604
605fn log_tx_age_emission(emission: &TxAgeEmission) {
610 let label = emission.label.as_deref().unwrap_or("<unlabeled>");
611 match emission.rung {
612 TxAgeRung::Warn => {
613 tracing::warn!(
614 tx_age_secs = emission.age.as_secs_f64(),
615 tx_label = label,
616 "ADR-091 Plank 1: open transaction registry entry exceeded soft-cap age"
617 );
618 }
619 TxAgeRung::Stale => {
620 tracing::error!(
621 tx_age_secs = emission.age.as_secs_f64(),
622 tx_label = label,
623 "ADR-091 Plank 1: open transaction registry entry exceeded the cooperative \
624 stale-op cap; no in-process mechanism can force-close it — investigate the \
625 labeled caller directly"
626 );
627 }
628 }
629}
630
631struct WalpinSidecarState {
639 dir: PathBuf,
640 pid: u32,
641 role: &'static str,
642 started_at: i64,
643 sweep_interval_ms: u64,
648 wrote: bool,
649 beacon_registered: bool,
655 last_heartbeat: Option<LastHeartbeatState>,
662}
663
664struct LastHeartbeatState {
670 span_id: khive_storage::tx_registry::TxId,
671 label: Option<String>,
672 attribution_basis: &'static str,
673 sweep_interval_ms: u64,
674 oldest_tx_started_at: i64,
675}
676
677impl LastHeartbeatState {
678 fn content_matches(
684 &self,
685 span_id: khive_storage::tx_registry::TxId,
686 label: &Option<String>,
687 attribution_basis: &str,
688 sweep_interval_ms: u64,
689 ) -> bool {
690 self.span_id == span_id
691 && self.label == *label
692 && self.attribution_basis == attribution_basis
693 && self.sweep_interval_ms == sweep_interval_ms
694 }
695}
696
697impl WalpinSidecarState {
698 fn new(
701 db_path: Option<&Path>,
702 is_file_backed: bool,
703 role: &'static str,
704 interval: Duration,
705 ) -> Option<Self> {
706 let path = db_path?;
707 if !crate::walpin::sidecar_enabled(is_file_backed) {
708 return None;
709 }
710 let pid = std::process::id();
711 Some(Self {
712 dir: crate::walpin::sidecar_dir_for(path),
713 pid,
714 role,
715 started_at: crate::walpin::process_start_time_secs(pid).unwrap_or(0),
716 sweep_interval_ms: interval.as_millis().min(u64::MAX as u128) as u64,
717 wrote: false,
718 last_heartbeat: None,
719 beacon_registered: false,
720 })
721 }
722
723 async fn register_beacon(&mut self) {
732 let dir = self.dir.clone();
733 let beacon = crate::walpin::WalpinBeacon {
734 pid: self.pid,
735 process_role: self.role.to_string(),
736 started_at: self.started_at,
737 sweep_interval_ms: self.sweep_interval_ms,
738 };
739 let result =
740 tokio::task::spawn_blocking(move || crate::walpin::write_beacon(&dir, &beacon)).await;
741 match result {
742 Ok(Ok(())) => {
743 self.beacon_registered = true;
744 }
745 Ok(Err(e)) => {
746 tracing::warn!(
747 error = %e,
748 "ADR-091 Amendment 2: failed to write walpin registration beacon; \
749 this process's sidecar health will read as unknown, not registered-silent"
750 );
751 }
752 Err(join_err) => {
753 tracing::warn!(
754 error = %join_err,
755 "ADR-091 Amendment 2: walpin beacon write task panicked"
756 );
757 }
758 }
759 }
760
761 async fn refresh_beacon(&mut self) {
771 if !self.beacon_registered {
772 self.register_beacon().await;
773 return;
774 }
775 let dir = self.dir.clone();
776 let pid = self.pid;
777 let result =
778 tokio::task::spawn_blocking(move || crate::walpin::touch_beacon(&dir, pid)).await;
779 match result {
780 Ok(Ok(())) => {}
781 Ok(Err(e)) => {
782 self.beacon_registered = false;
783 tracing::warn!(
784 error = %e,
785 "ADR-091 Amendment 2: failed to refresh walpin registration beacon; \
786 this process's sidecar health will read as unknown, not registered-silent"
787 );
788 }
789 Err(join_err) => {
790 self.beacon_registered = false;
791 tracing::warn!(
792 error = %join_err,
793 "ADR-091 Amendment 2: walpin beacon refresh task panicked"
794 );
795 }
796 }
797 }
798
799 async fn drop_beacon_fail_closed(&mut self) {
810 let dir = self.dir.clone();
811 let pid = self.pid;
812 self.beacon_registered = false;
813 let result =
814 tokio::task::spawn_blocking(move || crate::walpin::remove_beacon(&dir, pid)).await;
815 match result {
816 Ok(Ok(())) => {}
817 Ok(Err(e)) => {
818 tracing::warn!(
819 error = %e,
820 "ADR-091 Amendment 2: failed to remove walpin beacon after a failed \
821 heartbeat write; beacon will age out of the freshness window instead"
822 );
823 }
824 Err(join_err) => {
825 tracing::warn!(
826 error = %join_err,
827 "ADR-091 Amendment 2: walpin beacon removal task panicked"
828 );
829 }
830 }
831 }
832
833 async fn observe(
837 &mut self,
838 oldest: Option<khive_storage::tx_registry::OldestSpan>,
839 tx_warn_secs: Duration,
840 ) {
841 match oldest {
842 Some(span) if span.age >= tx_warn_secs => {
843 let attribution_basis = match span.origin {
850 khive_storage::tx_registry::TxOrigin::Database(_) => "origin",
851 khive_storage::tx_registry::TxOrigin::Unscoped
852 | khive_storage::tx_registry::TxOrigin::Memory => "fallback",
853 };
854
855 let content_unchanged = self.wrote
861 && self.last_heartbeat.as_ref().is_some_and(|last| {
862 last.content_matches(
863 span.id,
864 &span.label,
865 attribution_basis,
866 self.sweep_interval_ms,
867 )
868 });
869
870 if content_unchanged {
871 let dir = self.dir.clone();
872 let pid = self.pid;
873 let touch_result = tokio::task::spawn_blocking(move || {
874 crate::walpin::touch_heartbeat(&dir, pid)
875 })
876 .await;
877 match touch_result {
878 Ok(Ok(())) => {
879 self.refresh_beacon().await;
880 return;
881 }
882 Ok(Err(e)) => {
883 tracing::warn!(
884 error = %e,
885 "ADR-091 Amendment 3 Plank F1: walpin heartbeat touch failed; \
886 recreating with a full body write"
887 );
888 }
889 Err(join_err) => {
890 tracing::warn!(
891 error = %join_err,
892 "ADR-091 Amendment 3 Plank F1: walpin heartbeat touch task \
893 panicked; recreating with a full body write"
894 );
895 }
896 }
897 }
902
903 let oldest_tx_started_at = self
910 .last_heartbeat
911 .as_ref()
912 .filter(|last| last.span_id == span.id)
913 .map(|last| last.oldest_tx_started_at)
914 .unwrap_or_else(|| now_epoch_secs().saturating_sub(span.age.as_secs() as i64));
915
916 let heartbeat = crate::walpin::WalpinHeartbeat {
917 pid: self.pid,
918 process_role: self.role.to_string(),
919 started_at: self.started_at,
920 oldest_tx_age_secs: span.age.as_secs_f64(),
921 oldest_tx_label: span.label.clone(),
922 oldest_tx_started_at: Some(oldest_tx_started_at),
923 updated_at: now_epoch_secs(),
924 sweep_interval_ms: self.sweep_interval_ms,
925 attribution_basis: Some(attribution_basis.to_string()),
926 };
927 let dir = self.dir.clone();
928 let result = tokio::task::spawn_blocking(move || {
929 crate::walpin::write_heartbeat(&dir, &heartbeat)
930 })
931 .await;
932 match result {
942 Ok(Ok(())) => {
943 self.wrote = true;
944 self.last_heartbeat = Some(LastHeartbeatState {
945 span_id: span.id,
946 label: span.label,
947 attribution_basis,
948 sweep_interval_ms: self.sweep_interval_ms,
949 oldest_tx_started_at,
950 });
951 self.refresh_beacon().await;
952 }
953 Ok(Err(e)) => {
954 tracing::warn!(
955 error = %e,
956 "ADR-091 Amendment 2 Plank B: failed to write walpin heartbeat; \
957 removing beacon so this process cannot read as \
958 registered-silent while over threshold"
959 );
960 self.last_heartbeat = None;
964 self.drop_beacon_fail_closed().await;
965 }
966 Err(join_err) => {
967 tracing::warn!(
968 error = %join_err,
969 "ADR-091 Amendment 2 Plank B: walpin heartbeat write task panicked"
970 );
971 self.last_heartbeat = None;
972 self.drop_beacon_fail_closed().await;
973 }
974 }
975 }
976 _ => {
977 self.refresh_beacon().await;
978 if self.wrote {
979 let dir = self.dir.clone();
980 let pid = self.pid;
981 let result = tokio::task::spawn_blocking(move || {
982 crate::walpin::remove_heartbeat(&dir, pid)
983 })
984 .await;
985 match result {
986 Ok(Ok(())) => {}
987 Ok(Err(e)) => tracing::warn!(
988 error = %e,
989 "ADR-091 Amendment 2 Plank B: failed to remove walpin heartbeat"
990 ),
991 Err(join_err) => tracing::warn!(
992 error = %join_err,
993 "ADR-091 Amendment 2 Plank B: walpin heartbeat removal task panicked"
994 ),
995 }
996 self.wrote = false;
997 self.last_heartbeat = None;
998 }
999 }
1000 }
1001 }
1002
1003 async fn shutdown(&mut self) {
1004 if self.wrote {
1005 let dir = self.dir.clone();
1006 let pid = self.pid;
1007 let _ = tokio::task::spawn_blocking(move || crate::walpin::remove_heartbeat(&dir, pid))
1008 .await;
1009 self.wrote = false;
1010 }
1011 }
1012}
1013
1014fn now_epoch_secs() -> i64 {
1015 std::time::SystemTime::now()
1016 .duration_since(std::time::UNIX_EPOCH)
1017 .map(|d| d.as_secs() as i64)
1018 .unwrap_or(0)
1019}
1020
1021#[derive(Clone, Debug)]
1026pub struct SessionSweepConfig {
1027 pub interval: Duration,
1032 pub tx_warn_secs: Duration,
1034 pub tx_max_age_secs: Duration,
1036}
1037
1038impl Default for SessionSweepConfig {
1039 fn default() -> Self {
1040 Self {
1041 interval: Duration::from_secs(5),
1042 tx_warn_secs: Duration::from_secs(30),
1043 tx_max_age_secs: Duration::from_secs(120),
1044 }
1045 }
1046}
1047
1048impl SessionSweepConfig {
1049 pub fn from_env() -> Self {
1053 let mut cfg = Self::default();
1054
1055 if let Ok(ms) = std::env::var("KHIVE_SESSION_SWEEP_INTERVAL_MS") {
1056 if let Ok(v) = ms.parse::<u64>() {
1057 if v > 0 {
1058 cfg.interval = Duration::from_millis(v);
1059 }
1060 }
1061 }
1062 (cfg.tx_warn_secs, cfg.tx_max_age_secs) =
1067 tx_age_thresholds_from_env(cfg.tx_warn_secs, cfg.tx_max_age_secs);
1068
1069 cfg
1070 }
1071}
1072
1073pub struct SweepBackend {
1083 pub pool: Arc<ConnectionPool>,
1084 pub is_main: bool,
1085}
1086
1087struct BackendSweep {
1093 filter: khive_storage::tx_registry::TxOriginFilter,
1094 tx_age_state: TxAgeSweepState,
1095 sidecar: Option<WalpinSidecarState>,
1096}
1097
1098pub async fn run_session_sweep_task(
1112 backends: Vec<SweepBackend>,
1113 config: SessionSweepConfig,
1114 mut shutdown_rx: tokio::sync::watch::Receiver<()>,
1115) {
1116 let mut interval = tokio::time::interval(config.interval);
1117 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
1118
1119 let mut sweeps: Vec<BackendSweep> = Vec::with_capacity(backends.len());
1120 for backend in backends {
1121 let identity = match backend.pool.origin() {
1122 khive_storage::tx_registry::TxOrigin::Database(id) => id,
1123 khive_storage::tx_registry::TxOrigin::Memory
1126 | khive_storage::tx_registry::TxOrigin::Unscoped => continue,
1127 };
1128 let filter = if backend.is_main {
1129 khive_storage::tx_registry::TxOriginFilter::Main(identity)
1130 } else {
1131 khive_storage::tx_registry::TxOriginFilter::Secondary(identity)
1132 };
1133 let sidecar = WalpinSidecarState::new(
1134 backend.pool.canonical_path(),
1135 true,
1136 "session",
1137 config.interval,
1138 );
1139 sweeps.push(BackendSweep {
1140 filter,
1141 tx_age_state: TxAgeSweepState::default(),
1142 sidecar,
1143 });
1144 }
1145 for sweep in sweeps.iter_mut() {
1146 if let Some(sidecar) = sweep.sidecar.as_mut() {
1147 sidecar.register_beacon().await;
1148 }
1149 }
1150
1151 loop {
1152 tokio::select! {
1153 _ = interval.tick() => {}
1154 _ = shutdown_rx.changed() => break,
1155 }
1156
1157 for sweep in sweeps.iter_mut() {
1158 let oldest = khive_storage::tx_registry::oldest_for(&sweep.filter);
1159 for emission in sweep.tx_age_state.observe(
1160 oldest.as_ref().map(|s| (s.id, s.age, s.label.clone())),
1161 config.tx_warn_secs,
1162 config.tx_max_age_secs,
1163 ) {
1164 log_tx_age_emission(&emission);
1165 }
1166 if let Some(sidecar) = sweep.sidecar.as_mut() {
1167 sidecar.observe(oldest, config.tx_warn_secs).await;
1168 }
1169 }
1170 }
1171
1172 for sweep in sweeps.iter_mut() {
1173 if let Some(sidecar) = sweep.sidecar.as_mut() {
1174 sidecar.shutdown().await;
1175 }
1176 }
1177}
1178
1179#[derive(Clone)]
1184pub struct CheckpointLifecycleOwner {
1185 event_store: Arc<dyn khive_storage::EventStore>,
1186 namespace: String,
1187}
1188
1189impl CheckpointLifecycleOwner {
1190 pub fn new(
1192 event_store: Arc<dyn khive_storage::EventStore>,
1193 namespace: impl Into<String>,
1194 ) -> Self {
1195 Self {
1196 event_store,
1197 namespace: namespace.into(),
1198 }
1199 }
1200}
1201
1202const CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY: usize = 1;
1206
1207struct CheckpointLifecycleEmitter {
1216 namespace: Option<String>,
1217 sender: Option<tokio::sync::mpsc::Sender<khive_storage::Event>>,
1218 worker: Option<tokio::task::JoinHandle<()>>,
1219 busy_warning_emitted: bool,
1220}
1221
1222impl CheckpointLifecycleEmitter {
1223 fn new(owner: Option<CheckpointLifecycleOwner>) -> Self {
1224 let Some(owner) = owner else {
1225 return Self {
1226 namespace: None,
1227 sender: None,
1228 worker: None,
1229 busy_warning_emitted: false,
1230 };
1231 };
1232
1233 let namespace = owner.namespace.clone();
1234 let (sender, mut receiver) =
1235 tokio::sync::mpsc::channel::<khive_storage::Event>(CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY);
1236 let worker = tokio::spawn(async move {
1237 while let Some(event) = receiver.recv().await {
1238 let kind = event.kind;
1239 if let Err(err) = owner.event_store.append_event(event).await {
1240 tracing::warn!(
1241 error = %err,
1242 event_kind = %kind.name(),
1243 "checkpoint lifecycle event append failed"
1244 );
1245 }
1246 }
1247 });
1248
1249 Self {
1250 namespace: Some(namespace),
1251 sender: Some(sender),
1252 worker: Some(worker),
1253 busy_warning_emitted: false,
1254 }
1255 }
1256
1257 fn try_emit<P: serde::Serialize>(&mut self, kind: khive_types::EventKind, payload: P) -> bool {
1260 let (Some(namespace), Some(sender)) = (&self.namespace, &self.sender) else {
1261 return true;
1262 };
1263 let payload_value = match serde_json::to_value(&payload) {
1264 Ok(value) => value,
1265 Err(err) => {
1266 tracing::warn!(
1267 error = %err,
1268 event_kind = %kind.name(),
1269 "failed to serialize checkpoint lifecycle event payload"
1270 );
1271 return false;
1272 }
1273 };
1274 let event = khive_storage::Event::new(
1275 namespace,
1276 "checkpoint.lifecycle",
1277 kind,
1278 khive_types::SubstrateKind::Event,
1279 "daemon:checkpoint_task",
1280 )
1281 .with_payload(payload_value);
1282
1283 match sender.try_send(event) {
1284 Ok(()) => {
1285 self.busy_warning_emitted = false;
1286 true
1287 }
1288 Err(tokio::sync::mpsc::error::TrySendError::Full(event)) => {
1289 if !self.busy_warning_emitted {
1290 tracing::warn!(
1291 event_kind = %event.kind.name(),
1292 queue_capacity = CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY,
1293 "checkpoint lifecycle event dropped because the append worker is busy"
1294 );
1295 self.busy_warning_emitted = true;
1296 }
1297 false
1298 }
1299 Err(tokio::sync::mpsc::error::TrySendError::Closed(event)) => {
1300 tracing::warn!(
1301 event_kind = %event.kind.name(),
1302 "checkpoint lifecycle event dropped because the append worker stopped"
1303 );
1304 false
1305 }
1306 }
1307 }
1308
1309 async fn shutdown(mut self) {
1317 drop(self.sender.take());
1318 let Some(worker) = self.worker.take() else {
1319 return;
1320 };
1321 worker.abort();
1322 match worker.await {
1323 Ok(()) => {}
1324 Err(err) if err.is_cancelled() => {}
1325 Err(err) => tracing::warn!(
1326 error = %err,
1327 "checkpoint lifecycle event append worker terminated unexpectedly"
1328 ),
1329 }
1330 }
1331}
1332
1333impl Drop for CheckpointLifecycleEmitter {
1334 fn drop(&mut self) {
1335 if let Some(worker) = &self.worker {
1341 worker.abort();
1342 }
1343 }
1344}
1345
1346pub async fn run_checkpoint_task(
1374 pool: Arc<ConnectionPool>,
1375 config: CheckpointConfig,
1376 lifecycle_owner: Option<CheckpointLifecycleOwner>,
1377 mut shutdown_rx: tokio::sync::watch::Receiver<()>,
1378 is_main: bool,
1379) {
1380 let mut interval = tokio::time::interval(config.interval);
1381 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
1382 let mut severity_state = CheckpointSeverityState::default();
1383 let mut tx_age_state = TxAgeSweepState::default();
1384 let mut was_above_high_water = false;
1385 let mut truncate_state = TruncateState::default();
1386 let mut lifecycle_emitter = CheckpointLifecycleEmitter::new(lifecycle_owner);
1387 let mut event_elevation_open = false;
1393 let tx_filter = match pool.origin() {
1406 khive_storage::tx_registry::TxOrigin::Database(id) => Some(if is_main {
1407 khive_storage::tx_registry::TxOriginFilter::Main(id)
1408 } else {
1409 khive_storage::tx_registry::TxOriginFilter::Secondary(id)
1410 }),
1411 khive_storage::tx_registry::TxOrigin::Memory
1412 | khive_storage::tx_registry::TxOrigin::Unscoped => None,
1413 };
1414 #[cfg(unix)]
1420 let mut walpin_state =
1421 WalpinSidecarState::new(pool.canonical_path(), true, "daemon", config.interval);
1422 #[cfg(unix)]
1423 if let Some(sidecar) = walpin_state.as_mut() {
1424 sidecar.register_beacon().await;
1425 }
1426
1427 loop {
1428 tokio::select! {
1433 _ = interval.tick() => {}
1434 _ = shutdown_rx.changed() => break,
1435 }
1436
1437 let tick = checkpoint_once(&pool, &config, &mut truncate_state);
1438
1439 let oldest_tx = tx_filter
1455 .as_ref()
1456 .and_then(khive_storage::tx_registry::oldest_for);
1457 for emission in tx_age_state.observe(
1458 oldest_tx.as_ref().map(|s| (s.id, s.age, s.label.clone())),
1459 config.tx_warn_secs,
1460 config.tx_max_age_secs,
1461 ) {
1462 log_tx_age_emission(&emission);
1463 }
1464 #[cfg(unix)]
1468 if let Some(sidecar) = walpin_state.as_mut() {
1469 sidecar
1470 .observe(oldest_tx.clone(), config.tx_warn_secs)
1471 .await;
1472 }
1473
1474 let wal_pages = match tick {
1477 CheckpointTick::Skipped => continue,
1478 CheckpointTick::Observed(n) => n,
1479 };
1480
1481 let above_warn = wal_pages >= config.warn_pages;
1482 let above_high_water = wal_pages >= config.high_water_pages;
1483 let above_truncate_high_water = wal_pages >= config.truncate_high_water_pages;
1484
1485 log_tx_registry_oldest_debug(wal_pages, oldest_tx.as_ref());
1491
1492 for emission in severity_state.observe_wal_pages(wal_pages, &config) {
1497 match emission.rung {
1498 CheckpointSeverityRung::Info => {
1499 log_tx_registry_oldest_warn(wal_pages, oldest_tx.as_ref());
1500 tracing::info!(
1501 wal_pages = emission.wal_pages,
1502 warn_threshold = emission.threshold_pages,
1503 "WAL page count crossed warn threshold"
1504 );
1505 }
1506 CheckpointSeverityRung::Warn => {
1507 tracing::warn!(
1508 wal_pages = emission.wal_pages,
1509 warn_threshold = emission.threshold_pages,
1510 consecutive_cycles = emission.consecutive_cycles,
1511 "WAL page count failed to drain below warn threshold"
1512 );
1513 }
1514 CheckpointSeverityRung::Alarm => {
1515 }
1517 }
1518 }
1519
1520 let high_water_crossed = crossing_warn(above_high_water, &mut was_above_high_water);
1521 if high_water_crossed {
1522 log_tx_registry_snapshot_warn(wal_pages);
1523 tracing::warn!(
1524 wal_pages,
1525 high_water = config.high_water_pages,
1526 "WAL high-water mark exceeded; sustained WAL pressure — \
1527 a long-lived reader may be pinning an old snapshot that PASSIVE cannot reclaim"
1528 );
1529 }
1530
1531 if checkpoint_outcome_should_emit(above_warn, event_elevation_open) {
1535 let payload = khive_storage::CheckpointOutcomeRecordedPayload {
1536 wal_pages,
1537 warn_pages: config.warn_pages,
1538 high_water_pages: config.high_water_pages,
1539 truncate_high_water_pages: config.truncate_high_water_pages,
1540 above_warn,
1541 above_high_water,
1542 above_truncate_high_water,
1543 };
1544 if lifecycle_emitter
1545 .try_emit(khive_types::EventKind::CheckpointOutcomeRecorded, payload)
1546 {
1547 event_elevation_open = above_warn;
1548 }
1549 }
1550 }
1551
1552 lifecycle_emitter.shutdown().await;
1553
1554 #[cfg(unix)]
1555 if let Some(sidecar) = walpin_state.as_mut() {
1556 sidecar.shutdown().await;
1557 }
1558}
1559
1560fn checkpoint_outcome_should_emit(above_warn: bool, was_elevated: bool) -> bool {
1566 above_warn || was_elevated
1567}
1568
1569fn log_tx_registry_oldest_debug(
1578 wal_pages: u64,
1579 oldest: Option<&khive_storage::tx_registry::OldestSpan>,
1580) {
1581 if let Some(span) = oldest {
1582 tracing::debug!(
1583 wal_pages,
1584 oldest_tx_age_secs = span.age.as_secs_f64(),
1585 oldest_tx_label = span.label.as_deref().unwrap_or("<unlabeled>"),
1586 "WAL checkpoint tick: oldest open transaction registry entry"
1587 );
1588 }
1589}
1590
1591fn log_tx_registry_oldest_warn(
1595 wal_pages: u64,
1596 oldest: Option<&khive_storage::tx_registry::OldestSpan>,
1597) {
1598 if let Some(span) = oldest {
1599 tracing::warn!(
1600 wal_pages,
1601 oldest_tx_age_secs = span.age.as_secs_f64(),
1602 oldest_tx_label = span.label.as_deref().unwrap_or("<unlabeled>"),
1603 "WAL checkpoint tick: oldest open transaction registry entry"
1604 );
1605 }
1606}
1607
1608fn log_tx_registry_snapshot_warn(wal_pages: u64) {
1612 for (age, label) in khive_storage::tx_registry::snapshot() {
1613 tracing::warn!(
1614 wal_pages,
1615 tx_age_secs = age.as_secs_f64(),
1616 tx_label = label.as_deref().unwrap_or("<unlabeled>"),
1617 "WAL high-water: open transaction registry entry"
1618 );
1619 }
1620}
1621
1622pub fn checkpoint_once(
1639 pool: &ConnectionPool,
1640 config: &CheckpointConfig,
1641 truncate_state: &mut TruncateState,
1642) -> CheckpointTick {
1643 let writer = match pool.try_writer_nowait() {
1644 Ok(w) => w,
1645 Err(_) => {
1646 note_checkpoint_skipped();
1647 return CheckpointTick::Skipped;
1648 }
1649 };
1650
1651 let wal_pages = query_wal_pages(writer.conn());
1652
1653 if let Err(e) = writer
1654 .conn()
1655 .execute_batch("PRAGMA wal_checkpoint(PASSIVE)")
1656 {
1657 tracing::warn!(error = %e, "WAL checkpoint failed");
1658 } else {
1659 tracing::debug!(wal_pages, "WAL checkpoint issued");
1660 }
1661
1662 maybe_truncate(pool, &writer, config, wal_pages, truncate_state);
1663
1664 CheckpointTick::Observed(wal_pages)
1665}
1666
1667fn maybe_truncate(
1672 pool: &ConnectionPool,
1673 writer: &WriterGuard<'_>,
1674 config: &CheckpointConfig,
1675 wal_pages_before: u64,
1676 truncate_state: &mut TruncateState,
1677) {
1678 if wal_pages_before < config.truncate_high_water_pages {
1679 return;
1680 }
1681
1682 if let Some(last) = truncate_state.last_attempt {
1683 if last.elapsed() < config.truncate_min_interval {
1684 return;
1685 }
1686 }
1687
1688 log_tx_registry_snapshot_warn(wal_pages_before);
1691
1692 let conn = writer.conn();
1693 let original_busy_timeout = pool.config().busy_timeout;
1694
1695 if let Err(e) = conn.busy_timeout(config.truncate_busy_timeout) {
1696 tracing::warn!(error = %e, "failed to lower busy_timeout for TRUNCATE attempt; skipping");
1702 return;
1703 }
1704
1705 #[cfg(unix)]
1706 let holder_census = capture_wal_holder_census(pool);
1707
1708 truncate_state.last_attempt = Some(Instant::now());
1712
1713 let start = Instant::now();
1714 let outcome = conn.execute_batch("PRAGMA wal_checkpoint(TRUNCATE)");
1715 let elapsed = start.elapsed();
1716
1717 if let Err(e) = conn.busy_timeout(original_busy_timeout) {
1720 tracing::warn!(error = %e, "failed to restore busy_timeout after TRUNCATE attempt");
1721 }
1722
1723 match outcome {
1724 Ok(()) => {
1725 let wal_pages_after = query_wal_pages(conn);
1726 tracing::info!(
1727 wal_pages_before,
1728 wal_pages_after,
1729 elapsed_ms = elapsed.as_millis() as u64,
1730 "WAL TRUNCATE checkpoint attempted"
1731 );
1732
1733 let made_progress = wal_pages_after < wal_pages_before;
1734 if !made_progress {
1735 tracing::warn!(
1736 wal_pages_before,
1737 wal_pages_after,
1738 "WAL TRUNCATE attempt made no progress; \
1739 a long-lived reader may still be pinning the WAL snapshot"
1740 );
1741 log_tx_registry_snapshot_warn(wal_pages_after);
1742 #[cfg(test)]
1743 if let Some(path) = pool.canonical_path() {
1744 truncate_report_test_sync::after_no_progress_before_report(path);
1745 }
1746 #[cfg(unix)]
1747 log_walpin_sidecar_report(pool, holder_census);
1748 log_wal_pin_depth(conn);
1749 }
1750
1751 note_truncate_outcome(config, wal_pages_after, truncate_state);
1752 }
1753 Err(e) => {
1754 tracing::warn!(error = %e, wal_pages_before, "WAL TRUNCATE attempt failed");
1755 log_tx_registry_snapshot_warn(wal_pages_before);
1756 note_truncate_outcome(config, wal_pages_before, truncate_state);
1757 }
1758 }
1759}
1760
1761#[cfg(test)]
1762mod truncate_report_test_sync {
1763 use std::path::{Path, PathBuf};
1764 use std::sync::mpsc::{sync_channel, Receiver, SyncSender};
1765 use std::sync::Mutex;
1766
1767 struct Hook {
1768 db_path: PathBuf,
1769 reached_tx: SyncSender<()>,
1770 proceed_rx: Receiver<()>,
1771 }
1772
1773 static HOOK: Mutex<Option<Hook>> = Mutex::new(None);
1774
1775 pub(crate) fn install(db_path: PathBuf) -> (Receiver<()>, SyncSender<()>) {
1776 let (reached_tx, reached_rx) = sync_channel(0);
1777 let (proceed_tx, proceed_rx) = sync_channel(0);
1778 let replaced = HOOK
1779 .lock()
1780 .unwrap_or_else(|poisoned| poisoned.into_inner())
1781 .replace(Hook {
1782 db_path,
1783 reached_tx,
1784 proceed_rx,
1785 });
1786 assert!(replaced.is_none(), "truncate report hook already installed");
1787 (reached_rx, proceed_tx)
1788 }
1789
1790 pub(crate) fn uninstall() {
1791 *HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
1792 }
1793
1794 pub(crate) fn after_no_progress_before_report(db_path: &Path) {
1795 let hook = {
1796 let mut guard = HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
1797 match guard.as_ref() {
1798 Some(hook) if hook.db_path == db_path => guard.take(),
1799 _ => None,
1800 }
1801 };
1802 let Some(hook) = hook else {
1803 return;
1804 };
1805 let _ = hook.reached_tx.send(());
1806 let _ = hook.proceed_rx.recv();
1807 }
1808}
1809
1810fn note_truncate_outcome(
1816 config: &CheckpointConfig,
1817 wal_pages_after: u64,
1818 state: &mut TruncateState,
1819) {
1820 TRUNCATE_ATTEMPTS.fetch_add(1, Ordering::Relaxed);
1825
1826 if wal_pages_after >= config.warn_pages {
1827 state.consecutive_failures = state.consecutive_failures.saturating_add(1);
1828 if state.consecutive_failures == 3 {
1829 tracing::warn!(
1830 wal_pages_after,
1831 warn_threshold = config.warn_pages,
1832 "WAL TRUNCATE has failed to clear WAL pressure for 3 consecutive attempts"
1833 );
1834 }
1835 } else {
1836 state.consecutive_failures = 0;
1837 }
1838
1839 TRUNCATE_CONSECUTIVE_FAILURES.store(state.consecutive_failures as u64, Ordering::Relaxed);
1840}
1841
1842#[cfg(unix)]
1847fn capture_wal_holder_census(
1848 pool: &ConnectionPool,
1849) -> Option<Result<crate::walpin::CensusResult, String>> {
1850 let path = pool.canonical_path()?;
1851 if !crate::walpin::sidecar_enabled(true) {
1852 return None;
1853 }
1854 Some(crate::walpin::census_holders(path).map_err(|e| e.to_string()))
1855}
1856
1857#[cfg(unix)]
1871fn log_walpin_sidecar_report(
1872 pool: &ConnectionPool,
1873 census: Option<Result<crate::walpin::CensusResult, String>>,
1874) {
1875 let Some(census) = census else {
1876 return;
1877 };
1878 let Some(path) = pool.canonical_path() else {
1879 return;
1880 };
1881 let dir = crate::walpin::sidecar_dir_for(path);
1882 let sweep_interval = SessionSweepConfig::from_env().interval;
1887 let report = match crate::walpin::enumerate_live(&dir, sweep_interval) {
1888 Ok(report) => report,
1889 Err(e) => {
1890 tracing::warn!(
1891 error = %e,
1892 "ADR-091 Amendment 2 Plank B: sidecar directory failed the trust-boundary \
1893 check; cross-process WAL-pin attribution is unestablished for this tick"
1894 );
1895 return;
1896 }
1897 };
1898 let now = now_epoch_secs();
1899 for hb in report.reporting() {
1900 tracing::warn!(
1906 walpin_pid = hb.pid,
1907 walpin_role = %hb.process_role,
1908 walpin_oldest_tx_age_secs = hb.current_oldest_tx_age_secs(now),
1909 walpin_oldest_tx_label = hb.oldest_tx_label.as_deref().unwrap_or("<unlabeled>"),
1910 walpin_attribution_basis = hb.attribution_basis.as_deref().unwrap_or("<unspecified>"),
1911 walpin_attribution_evidence_backed = hb.attribution_is_evidence_backed(),
1912 walpin_health = "reporting",
1913 "ADR-091 Amendment 2 Plank B: live cross-process WAL-pin attribution report"
1914 );
1915 }
1916 for pid in report.registered_silent_pids() {
1917 tracing::debug!(
1918 walpin_pid = pid,
1919 walpin_health = "registered_silent",
1920 "ADR-091 Amendment 2 Plank B: process affirmatively reports no over-threshold span"
1921 );
1922 }
1923 let mut unknown_pids: Vec<u32> = report.unknown_pids().collect();
1924
1925 match census {
1930 Ok(census) => {
1931 let sidecar_known: std::collections::HashSet<u32> = report
1932 .reporting()
1933 .map(|hb| hb.pid)
1934 .chain(report.registered_silent_pids())
1935 .chain(unknown_pids.iter().copied())
1936 .collect();
1937 let mut census_only: Vec<u32> =
1938 census.holders.difference(&sidecar_known).copied().collect();
1939 if !census_only.is_empty() {
1940 census_only.sort_unstable();
1941 tracing::warn!(
1942 ?census_only,
1943 "ADR-091 Amendment 2: these PIDs hold the database file open \
1944 at the OS level but have no sidecar data at all (pre-feature binary, \
1945 sidecar disabled, or wedged before its first write)"
1946 );
1947 unknown_pids.extend(census_only);
1948 }
1949 if !census.is_complete() {
1950 let mut uninspectable = census.uninspectable_pids.clone();
1951 uninspectable.sort_unstable();
1952 tracing::warn!(
1953 ?uninspectable,
1954 truncated = census.truncated,
1955 "ADR-091 Amendment 2: the OS-derived holder census is \
1956 INCOMPLETE — either specific PIDs' open file descriptors could not be \
1957 inspected (permission denied, or a listing race), or the enumeration walk \
1958 itself has positive evidence it did not see the full live-process universe \
1959 (namespace/visibility check, directory-iterator error, self-canary, or a \
1960 libproc buffer that stayed at capacity after bounded retries) — cannot \
1961 rule out an unregistered holder"
1962 );
1963 if uninspectable.is_empty() {
1964 unknown_pids.push(0);
1971 } else {
1972 unknown_pids.extend(uninspectable);
1973 }
1974 }
1975 }
1976 Err(e) => {
1977 tracing::warn!(
1978 error = %e,
1979 "ADR-091 Amendment 2: OS-derived holder census failed; \
1980 attribution cannot rule out an unregistered database holder this tick"
1981 );
1982 unknown_pids.push(0);
1986 }
1987 }
1988
1989 if !unknown_pids.is_empty() {
1990 tracing::warn!(
1991 ?unknown_pids,
1992 "ADR-091 Amendment 2 Plank B: sidecar health unestablished for these PIDs; \
1993 attribution is inconclusive and the native/unregistered-mechanism conclusion \
1994 is NOT licensed this tick"
1995 );
1996 } else if report.reporting().next().is_none() {
1997 tracing::info!(
1998 "ADR-091 Amendment 2 Plank B: every live PID is reporting or registered-silent \
1999 with none pinning; the WAL pin is not attributable to any in-process registry \
2000 span this sidecar covers"
2001 );
2002 }
2003}
2004
2005fn log_wal_pin_depth(conn: &rusqlite::Connection) {
2011 match query_wal_pin_depth(conn) {
2012 Ok((log, checkpointed)) => {
2013 tracing::warn!(
2014 wal_log_frames = log,
2015 wal_checkpointed_frames = checkpointed,
2016 wal_pin_depth = (log - checkpointed).max(0),
2017 "ADR-091 Amendment 2 Plank C: WAL pin depth after TRUNCATE no-progress"
2018 );
2019 }
2020 Err(e) => {
2021 tracing::warn!(
2022 error = %e,
2023 "ADR-091 Amendment 2 Plank C: failed to query WAL pin depth"
2024 );
2025 }
2026 }
2027}
2028
2029fn query_wal_pin_depth(conn: &rusqlite::Connection) -> rusqlite::Result<(i64, i64)> {
2036 conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
2037 Ok((row.get::<_, i64>(1)?, row.get::<_, i64>(2)?))
2038 })
2039}
2040
2041fn crossing_warn(now_above: bool, was_above: &mut bool) -> bool {
2050 let fire = now_above && !*was_above;
2051 *was_above = now_above;
2052 fire
2053}
2054
2055fn query_wal_pages(conn: &rusqlite::Connection) -> u64 {
2068 let pages = conn
2069 .query_row("PRAGMA wal_checkpoint", [], |row| row.get::<_, i64>(1))
2070 .unwrap_or(0)
2071 .max(0) as u64;
2072 LAST_WAL_PAGES.store(pages, Ordering::Relaxed);
2076 note_checkpoint_observed(pages);
2077 pages
2078}
2079
2080#[cfg(test)]
2081mod tests {
2082 use super::*;
2083 use crate::pool::PoolConfig;
2084 use serial_test::serial;
2085 use tracing::field::{Field, Visit};
2086
2087 #[derive(Clone, Debug, Default)]
2088 struct CapturedEvent {
2089 message: Option<String>,
2090 oldest_tx_label: Option<String>,
2091 tx_label: Option<String>,
2092 census_only: Option<String>,
2093 }
2094
2095 #[derive(Default)]
2096 struct CapturedEventVisitor(CapturedEvent);
2097
2098 impl Visit for CapturedEventVisitor {
2099 fn record_str(&mut self, field: &Field, value: &str) {
2100 match field.name() {
2101 "message" => self.0.message = Some(value.to_string()),
2102 "oldest_tx_label" => self.0.oldest_tx_label = Some(value.to_string()),
2103 "tx_label" => self.0.tx_label = Some(value.to_string()),
2104 _ => {}
2105 }
2106 }
2107
2108 fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
2109 let formatted = format!("{value:?}");
2110 let cleaned = formatted
2111 .trim_start_matches('"')
2112 .trim_end_matches('"')
2113 .to_string();
2114 match field.name() {
2115 "message" => self.0.message = Some(cleaned),
2116 "oldest_tx_label" => self.0.oldest_tx_label = Some(cleaned),
2117 "tx_label" => self.0.tx_label = Some(cleaned),
2118 "census_only" => self.0.census_only = Some(cleaned),
2119 _ => {}
2120 }
2121 }
2122 }
2123
2124 struct CaptureSubscriber {
2129 events: std::sync::Arc<std::sync::Mutex<Vec<CapturedEvent>>>,
2130 }
2131
2132 impl tracing::Subscriber for CaptureSubscriber {
2133 fn enabled(&self, _: &tracing::Metadata<'_>) -> bool {
2134 true
2135 }
2136 fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
2137 tracing::span::Id::from_u64(1)
2138 }
2139 fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
2140 fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
2141 fn event(&self, event: &tracing::Event<'_>) {
2142 let mut visitor = CapturedEventVisitor::default();
2143 event.record(&mut visitor);
2144 self.events.lock().unwrap().push(visitor.0);
2145 }
2146 fn enter(&self, _: &tracing::span::Id) {}
2147 fn exit(&self, _: &tracing::span::Id) {}
2148 }
2149
2150 #[test]
2153 #[serial(tx_registry)]
2154 fn log_tx_registry_oldest_debug_reports_oldest_open_entry() {
2155 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
2156 let subscriber = CaptureSubscriber {
2157 events: std::sync::Arc::clone(&buffer),
2158 };
2159
2160 let _handle =
2161 khive_storage::tx_registry::register(Some("checkpoint_tick_test".to_string()));
2162
2163 let oldest = khive_storage::tx_registry::oldest().map(|(id, age, label)| {
2164 khive_storage::tx_registry::OldestSpan {
2165 id,
2166 age,
2167 label,
2168 origin: khive_storage::tx_registry::TxOrigin::Unscoped,
2169 }
2170 });
2171 let expected_label = oldest
2172 .as_ref()
2173 .and_then(|s| s.label.clone())
2174 .unwrap_or_else(|| "<unlabeled>".to_string());
2175
2176 tracing::subscriber::with_default(subscriber, || {
2177 log_tx_registry_oldest_debug(100, oldest.as_ref());
2178 });
2179
2180 let events = buffer.lock().unwrap();
2181 assert!(
2182 events.iter().any(|e| {
2183 e.message.as_deref()
2184 == Some("WAL checkpoint tick: oldest open transaction registry entry")
2185 && e.oldest_tx_label.as_deref() == Some(expected_label.as_str())
2186 }),
2187 "expected a log line naming the open registry entry's label, got: {events:?}"
2188 );
2189 }
2190
2191 #[test]
2197 #[serial(tx_registry)]
2198 fn registry_warns_fire_on_crossing_and_do_not_repeat() {
2199 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
2200 let subscriber = CaptureSubscriber {
2201 events: std::sync::Arc::clone(&buffer),
2202 };
2203
2204 let _handle =
2205 khive_storage::tx_registry::register(Some("registry_warn_crossing_test".to_string()));
2206 let oldest = khive_storage::tx_registry::oldest().map(|(id, age, label)| {
2207 khive_storage::tx_registry::OldestSpan {
2208 id,
2209 age,
2210 label,
2211 origin: khive_storage::tx_registry::TxOrigin::Unscoped,
2212 }
2213 });
2214
2215 let mut was_above_warn = false;
2216 let mut was_above_high_water = false;
2217
2218 tracing::subscriber::with_default(subscriber, || {
2219 if crossing_warn(true, &mut was_above_warn) {
2221 log_tx_registry_oldest_warn(6000, oldest.as_ref());
2222 }
2223 if crossing_warn(true, &mut was_above_high_water) {
2224 log_tx_registry_snapshot_warn(6000);
2225 }
2226
2227 if crossing_warn(true, &mut was_above_warn) {
2229 log_tx_registry_oldest_warn(6000, oldest.as_ref());
2230 }
2231 if crossing_warn(true, &mut was_above_high_water) {
2232 log_tx_registry_snapshot_warn(6000);
2233 }
2234 });
2235
2236 let events = buffer.lock().unwrap();
2237
2238 let oldest_warn_count = events
2248 .iter()
2249 .filter(|e| {
2250 e.message.as_deref()
2251 == Some("WAL checkpoint tick: oldest open transaction registry entry")
2252 })
2253 .count();
2254 assert_eq!(
2255 oldest_warn_count, 1,
2256 "oldest-entry WARN must fire exactly once across two above-threshold ticks, got: {events:?}"
2257 );
2258
2259 let snapshot_warn_count = events
2260 .iter()
2261 .filter(|e| {
2262 e.message.as_deref() == Some("WAL high-water: open transaction registry entry")
2263 && e.tx_label.as_deref() == Some("registry_warn_crossing_test")
2264 })
2265 .count();
2266 assert_eq!(
2267 snapshot_warn_count, 1,
2268 "high-water snapshot WARN must fire exactly once across two above-threshold ticks, got: {events:?}"
2269 );
2270 }
2271
2272 #[test]
2275 fn log_tx_age_emission_carries_label_for_both_rungs() {
2276 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
2277 let subscriber = CaptureSubscriber {
2278 events: std::sync::Arc::clone(&buffer),
2279 };
2280
2281 tracing::subscriber::with_default(subscriber, || {
2282 log_tx_age_emission(&TxAgeEmission {
2283 rung: TxAgeRung::Warn,
2284 age: Duration::from_secs(45),
2285 label: Some("plank1_warn_test".to_string()),
2286 });
2287 log_tx_age_emission(&TxAgeEmission {
2288 rung: TxAgeRung::Stale,
2289 age: Duration::from_secs(150),
2290 label: Some("plank1_stale_test".to_string()),
2291 });
2292 });
2293
2294 let events = buffer.lock().unwrap();
2295 assert!(
2296 events.iter().any(|e| {
2297 e.message.as_deref()
2298 == Some(
2299 "ADR-091 Plank 1: open transaction registry entry exceeded soft-cap age",
2300 )
2301 && e.tx_label.as_deref() == Some("plank1_warn_test")
2302 }),
2303 "expected a Warn-rung log line naming the entry, got: {events:?}"
2304 );
2305 assert!(
2306 events.iter().any(|e| {
2307 e.message.as_deref().is_some_and(|m| {
2308 m.starts_with(
2309 "ADR-091 Plank 1: open transaction registry entry exceeded the cooperative",
2310 )
2311 }) && e.tx_label.as_deref() == Some("plank1_stale_test")
2312 }),
2313 "expected a Stale-rung log line naming the entry, got: {events:?}"
2314 );
2315 }
2316
2317 fn file_pool(path: &std::path::Path) -> Arc<ConnectionPool> {
2318 let cfg = PoolConfig {
2319 path: Some(path.to_path_buf()),
2320 ..PoolConfig::default()
2321 };
2322 Arc::new(ConnectionPool::new(cfg).expect("pool open"))
2323 }
2324
2325 struct TruncateReportHookGuard;
2326
2327 impl Drop for TruncateReportHookGuard {
2328 fn drop(&mut self) {
2329 truncate_report_test_sync::uninstall();
2330 }
2331 }
2332
2333 struct ReaderProcess {
2334 child: std::process::Child,
2335 _stdout: std::io::BufReader<std::process::ChildStdout>,
2336 }
2337
2338 impl ReaderProcess {
2339 fn spawn(db_path: &std::path::Path) -> Self {
2340 use std::io::BufRead;
2341 use std::process::Stdio;
2342
2343 let mut child = std::process::Command::new(
2344 std::env::current_exe().expect("resolve current test executable"),
2345 )
2346 .args([
2347 "--exact",
2348 "checkpoint::tests::walpin_transient_reader_process_helper",
2349 "--nocapture",
2350 ])
2351 .env("KHIVE_CHECKPOINT_READER_HELPER_PATH", db_path)
2352 .stdin(Stdio::piped())
2353 .stdout(Stdio::piped())
2354 .spawn()
2355 .expect("spawn transient WAL reader helper");
2356
2357 let stdout = child.stdout.take().expect("capture helper stdout");
2358 let mut reader = std::io::BufReader::new(stdout);
2359 let mut line = String::new();
2360 loop {
2361 line.clear();
2362 let bytes = reader
2363 .read_line(&mut line)
2364 .expect("read transient reader readiness signal");
2365 assert!(bytes > 0, "reader helper exited before readiness signal");
2366 if line.contains("KHIVE_CHECKPOINT_READER_READY") {
2367 break;
2368 }
2369 }
2370 Self {
2371 child,
2372 _stdout: reader,
2373 }
2374 }
2375
2376 fn pid(&self) -> u32 {
2377 self.child.id()
2378 }
2379
2380 fn release(&mut self) {
2381 use std::io::Write;
2382
2383 let mut stdin = self.child.stdin.take().expect("helper stdin is available");
2384 stdin
2385 .write_all(b"release\n")
2386 .expect("release transient reader");
2387 drop(stdin);
2388 let status = self.child.wait().expect("wait for transient reader helper");
2389 assert!(status.success(), "transient reader helper failed: {status}");
2390 }
2391 }
2392
2393 impl Drop for ReaderProcess {
2394 fn drop(&mut self) {
2395 if self.child.try_wait().ok().flatten().is_none() {
2396 let _ = self.child.kill();
2397 let _ = self.child.wait();
2398 }
2399 }
2400 }
2401
2402 #[test]
2403 fn walpin_transient_reader_process_helper() {
2404 use std::io::Write;
2405
2406 let Some(path) = std::env::var_os("KHIVE_CHECKPOINT_READER_HELPER_PATH") else {
2407 return;
2408 };
2409 let conn = rusqlite::Connection::open(path).expect("helper opens database");
2410 conn.execute_batch("BEGIN DEFERRED; SELECT * FROM t;")
2411 .expect("helper pins a read snapshot");
2412 println!("KHIVE_CHECKPOINT_READER_READY");
2413 std::io::stdout().flush().expect("flush readiness signal");
2414 let mut release = String::new();
2415 std::io::stdin()
2416 .read_line(&mut release)
2417 .expect("wait for release signal");
2418 conn.execute_batch("COMMIT")
2419 .expect("helper releases read snapshot");
2420 }
2421
2422 #[test]
2423 #[cfg(unix)]
2424 #[serial(checkpoint_skip_metrics, walpin_report_seam)]
2425 fn no_progress_report_keeps_holder_released_after_truncate_timeout() {
2426 let dir = tempfile::tempdir().expect("tempdir");
2427 let path = dir.path().join("transient-reader.db");
2428 let pool = file_pool(&path);
2429 {
2430 let writer = pool.try_writer().expect("writer");
2431 writer
2432 .conn()
2433 .execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
2434 .expect("seed WAL before reader snapshot");
2435 }
2436
2437 let mut reader = ReaderProcess::spawn(&path);
2438 let reader_pid = reader.pid();
2439 {
2440 let writer = pool.try_writer().expect("writer");
2441 writer
2442 .conn()
2443 .execute_batch("INSERT INTO t VALUES (2);")
2444 .expect("append WAL behind reader snapshot");
2445 }
2446
2447 let canonical_path = pool
2448 .canonical_path()
2449 .expect("file-backed pool has canonical path")
2450 .to_path_buf();
2451 let (reached_rx, proceed_tx) = truncate_report_test_sync::install(canonical_path.clone());
2452 let _hook_guard = TruncateReportHookGuard;
2453 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
2454 let thread_buffer = std::sync::Arc::clone(&buffer);
2455 let checkpoint_pool = Arc::clone(&pool);
2456 let checkpoint = std::thread::spawn(move || {
2457 let subscriber = CaptureSubscriber {
2458 events: thread_buffer,
2459 };
2460 tracing::subscriber::with_default(subscriber, || {
2461 checkpoint_once(
2462 &checkpoint_pool,
2463 &CheckpointConfig {
2464 truncate_high_water_pages: 0,
2465 truncate_min_interval: Duration::ZERO,
2466 truncate_busy_timeout: Duration::from_millis(50),
2467 ..CheckpointConfig::default()
2468 },
2469 &mut TruncateState::default(),
2470 )
2471 })
2472 });
2473
2474 reached_rx
2475 .recv_timeout(Duration::from_secs(5))
2476 .expect("TRUNCATE must report no progress while the reader is pinned");
2477 reader.release();
2478 let post_attempt_census =
2479 crate::walpin::census_holders(&canonical_path).expect("post-attempt holder census");
2480 assert!(
2481 !post_attempt_census.holders.contains(&reader_pid),
2482 "released reader PID must be absent from a post-attempt census"
2483 );
2484 proceed_tx
2485 .send(())
2486 .expect("allow no-progress reporting to continue");
2487 checkpoint.join().expect("checkpoint thread");
2488
2489 let events = buffer.lock().expect("captured events");
2490 assert!(
2491 events.iter().any(|event| {
2492 event
2493 .census_only
2494 .as_deref()
2495 .is_some_and(|pids| pids.contains(&reader_pid.to_string()))
2496 }),
2497 "the no-progress report must retain PID {reader_pid} from the pre-attempt census: {events:?}"
2498 );
2499 }
2500
2501 #[test]
2507 #[serial(checkpoint_skip_metrics)]
2508 fn checkpoint_once_succeeds_on_file_backed_pool() {
2509 let dir = tempfile::tempdir().unwrap();
2510 let path = dir.path().join("wal_test.db");
2511 let pool = file_pool(&path);
2512
2513 {
2515 let writer = pool.try_writer().unwrap();
2516 writer
2517 .conn()
2518 .execute_batch("CREATE TABLE IF NOT EXISTS t (x INTEGER);")
2519 .unwrap();
2520 writer
2521 .conn()
2522 .execute_batch("INSERT INTO t VALUES (1);")
2523 .unwrap();
2524 }
2525
2526 checkpoint_once(
2527 &pool,
2528 &CheckpointConfig::default(),
2529 &mut TruncateState::default(),
2530 );
2531 }
2532
2533 #[test]
2534 #[serial(checkpoint_skip_metrics)]
2535 fn checkpoint_once_is_noop_on_in_memory_pool() {
2536 let cfg = PoolConfig {
2538 path: None,
2539 ..PoolConfig::default()
2540 };
2541 let pool = Arc::new(ConnectionPool::new(cfg).expect("in-memory pool"));
2542 checkpoint_once(
2543 &pool,
2544 &CheckpointConfig::default(),
2545 &mut TruncateState::default(),
2546 );
2547 }
2548
2549 #[tokio::test]
2550 #[serial(checkpoint_skip_metrics)]
2551 async fn checkpoint_task_exits_on_shutdown_signal() {
2552 let dir = tempfile::tempdir().unwrap();
2553 let path = dir.path().join("wal_task_shutdown.db");
2554 let pool = file_pool(&path);
2555
2556 let cfg = CheckpointConfig {
2558 interval: Duration::from_millis(10),
2559 ..Default::default()
2560 };
2561
2562 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
2563 let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
2564
2565 shutdown_tx.send(()).expect("send shutdown signal");
2566
2567 tokio::time::timeout(Duration::from_secs(1), handle)
2568 .await
2569 .expect("checkpoint task should exit within 1s")
2570 .expect("checkpoint task panicked");
2571 }
2572
2573 #[tokio::test]
2577 #[serial(checkpoint_skip_metrics)]
2578 async fn checkpoint_task_exits_via_shutdown_signal_with_live_event_store_pool_clone() {
2579 let dir = tempfile::tempdir().unwrap();
2580 let path = dir.path().join("wal_task_event_store.db");
2581 let pool = file_pool(&path);
2582
2583 let cfg = CheckpointConfig {
2584 interval: Duration::from_millis(10),
2585 ..Default::default()
2586 };
2587
2588 let event_store: Arc<dyn khive_storage::EventStore> =
2589 Arc::new(crate::stores::event::SqlEventStore::new_scoped(
2590 Arc::clone(&pool),
2591 true,
2592 "local".to_string(),
2593 ));
2594 let sibling_pool_clone = Arc::clone(&pool);
2599
2600 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
2601 let handle = tokio::spawn(run_checkpoint_task(
2602 pool,
2603 cfg,
2604 Some(CheckpointLifecycleOwner::new(event_store, "local")),
2605 shutdown_rx,
2606 true,
2607 ));
2608
2609 assert!(
2613 Arc::strong_count(&sibling_pool_clone) > 1,
2614 "test setup must reproduce the multi-owner shape the bug depends on"
2615 );
2616
2617 shutdown_tx.send(()).expect("send shutdown signal");
2618
2619 tokio::time::timeout(Duration::from_secs(1), handle)
2620 .await
2621 .expect(
2622 "checkpoint task should exit within 1s via the watch signal, \
2623 even with a live sibling Arc<ConnectionPool> clone held by \
2624 the event store",
2625 )
2626 .expect("checkpoint task panicked");
2627 }
2628
2629 #[test]
2630 #[serial]
2631 fn checkpoint_config_env_override() {
2632 std::env::set_var("KHIVE_CHECKPOINT_INTERVAL_MS", "250");
2633 std::env::set_var("KHIVE_WAL_WARN_PAGES", "1500");
2634 std::env::set_var("KHIVE_WAL_HIGH_WATER_PAGES", "8000");
2635 std::env::set_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES", "12000");
2636 std::env::set_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS", "60");
2637 std::env::set_var("KHIVE_WAL_TRUNCATE_BUSY_MS", "500");
2638 std::env::set_var("KHIVE_TX_WARN_SECS", "15");
2639 std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "90");
2640
2641 let cfg = CheckpointConfig::from_env();
2642
2643 std::env::remove_var("KHIVE_CHECKPOINT_INTERVAL_MS");
2644 std::env::remove_var("KHIVE_WAL_WARN_PAGES");
2645 std::env::remove_var("KHIVE_WAL_HIGH_WATER_PAGES");
2646 std::env::remove_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES");
2647 std::env::remove_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS");
2648 std::env::remove_var("KHIVE_WAL_TRUNCATE_BUSY_MS");
2649 std::env::remove_var("KHIVE_TX_WARN_SECS");
2650 std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
2651
2652 assert_eq!(cfg.interval, Duration::from_millis(250));
2653 assert_eq!(cfg.warn_pages, 1500);
2654 assert_eq!(cfg.high_water_pages, 8000);
2655 assert_eq!(cfg.truncate_high_water_pages, 12000);
2656 assert_eq!(cfg.truncate_min_interval, Duration::from_secs(60));
2657 assert_eq!(cfg.truncate_busy_timeout, Duration::from_millis(500));
2658 assert_eq!(cfg.tx_warn_secs, Duration::from_secs(15));
2659 assert_eq!(cfg.tx_max_age_secs, Duration::from_secs(90));
2660 }
2661
2662 #[test]
2663 #[serial]
2664 fn checkpoint_config_defaults_on_invalid_env() {
2665 let default = CheckpointConfig::default();
2666
2667 std::env::set_var("KHIVE_CHECKPOINT_INTERVAL_MS", "not_a_number");
2668 std::env::set_var("KHIVE_WAL_WARN_PAGES", "");
2669 std::env::set_var("KHIVE_WAL_HIGH_WATER_PAGES", "0");
2670 std::env::set_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES", "not_a_number");
2671 std::env::set_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS", "");
2672 std::env::set_var("KHIVE_WAL_TRUNCATE_BUSY_MS", "0");
2673 std::env::set_var("KHIVE_TX_WARN_SECS", "not_a_number");
2674 std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "0");
2675
2676 let cfg = CheckpointConfig::from_env();
2677
2678 std::env::remove_var("KHIVE_CHECKPOINT_INTERVAL_MS");
2679 std::env::remove_var("KHIVE_WAL_WARN_PAGES");
2680 std::env::remove_var("KHIVE_WAL_HIGH_WATER_PAGES");
2681 std::env::remove_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES");
2682 std::env::remove_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS");
2683 std::env::remove_var("KHIVE_WAL_TRUNCATE_BUSY_MS");
2684 std::env::remove_var("KHIVE_TX_WARN_SECS");
2685 std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
2686
2687 assert_eq!(cfg.interval, default.interval);
2688 assert_eq!(cfg.warn_pages, default.warn_pages);
2689 assert_eq!(cfg.high_water_pages, default.high_water_pages);
2690 assert_eq!(
2691 cfg.truncate_high_water_pages,
2692 default.truncate_high_water_pages
2693 );
2694 assert_eq!(cfg.truncate_min_interval, default.truncate_min_interval);
2695 assert_eq!(cfg.truncate_busy_timeout, default.truncate_busy_timeout);
2696 assert_eq!(cfg.tx_warn_secs, default.tx_warn_secs);
2697 assert_eq!(cfg.tx_max_age_secs, default.tx_max_age_secs);
2698 }
2699
2700 #[test]
2705 #[serial(checkpoint_skip_metrics)]
2706 fn checkpoint_high_water_does_not_block_behind_reader() {
2707 let dir = tempfile::tempdir().unwrap();
2708 let path = dir.path().join("high_water_test.db");
2709
2710 let pool = Arc::new(
2714 ConnectionPool::new(PoolConfig {
2715 path: Some(path.clone()),
2716 busy_timeout: Duration::from_millis(2000),
2717 ..PoolConfig::default()
2718 })
2719 .expect("pool open"),
2720 );
2721
2722 {
2724 let writer = pool.try_writer().unwrap();
2725 writer
2726 .conn()
2727 .execute_batch(
2728 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
2729 )
2730 .unwrap();
2731 }
2732
2733 let reader = pool.reader().expect("reader");
2737 reader
2738 .execute_batch("BEGIN DEFERRED; SELECT * FROM t;")
2739 .expect("begin read tx");
2740
2741 {
2745 let writer = pool.try_writer().unwrap();
2746 writer
2747 .conn()
2748 .execute_batch("INSERT INTO t VALUES (2);")
2749 .unwrap();
2750 }
2751
2752 let start = std::time::Instant::now();
2753 checkpoint_once(
2754 &pool,
2755 &CheckpointConfig::default(),
2756 &mut TruncateState::default(),
2757 );
2758 let elapsed = start.elapsed();
2759
2760 reader.execute_batch("COMMIT;").ok();
2762 drop(reader);
2763
2764 assert!(
2768 elapsed < std::time::Duration::from_millis(500),
2769 "checkpoint_once with active reader snapshot took {:?}; \
2770 expected <500ms (PASSIVE must not block on readers; \
2771 a TRUNCATE regression would block ~2000ms)",
2772 elapsed
2773 );
2774 }
2775
2776 #[test]
2777 #[serial]
2778 fn checkpoint_config_rejects_zero_for_all_fields() {
2779 let default = CheckpointConfig::default();
2780 std::env::set_var("KHIVE_CHECKPOINT_INTERVAL_MS", "0");
2781 std::env::set_var("KHIVE_WAL_WARN_PAGES", "0");
2782 std::env::set_var("KHIVE_WAL_HIGH_WATER_PAGES", "0");
2783 std::env::set_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES", "0");
2784 std::env::set_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS", "0");
2785 std::env::set_var("KHIVE_WAL_TRUNCATE_BUSY_MS", "0");
2786 std::env::set_var("KHIVE_TX_WARN_SECS", "0");
2787 std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "0");
2788
2789 let cfg = CheckpointConfig::from_env();
2790
2791 std::env::remove_var("KHIVE_CHECKPOINT_INTERVAL_MS");
2792 std::env::remove_var("KHIVE_WAL_WARN_PAGES");
2793 std::env::remove_var("KHIVE_WAL_HIGH_WATER_PAGES");
2794 std::env::remove_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES");
2795 std::env::remove_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS");
2796 std::env::remove_var("KHIVE_WAL_TRUNCATE_BUSY_MS");
2797 std::env::remove_var("KHIVE_TX_WARN_SECS");
2798 std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
2799
2800 assert_eq!(
2801 cfg.interval, default.interval,
2802 "zero interval must fall back to default"
2803 );
2804 assert_eq!(
2805 cfg.warn_pages, default.warn_pages,
2806 "zero warn_pages must fall back to default"
2807 );
2808 assert_eq!(
2809 cfg.high_water_pages, default.high_water_pages,
2810 "zero high_water_pages must fall back to default"
2811 );
2812 assert_eq!(
2813 cfg.truncate_high_water_pages, default.truncate_high_water_pages,
2814 "zero truncate_high_water_pages must fall back to default"
2815 );
2816 assert_eq!(
2817 cfg.truncate_min_interval, default.truncate_min_interval,
2818 "zero truncate_min_interval must fall back to default"
2819 );
2820 assert_eq!(
2821 cfg.truncate_busy_timeout, default.truncate_busy_timeout,
2822 "zero truncate_busy_timeout must fall back to default"
2823 );
2824 assert_eq!(
2825 cfg.tx_warn_secs, default.tx_warn_secs,
2826 "zero tx_warn_secs must fall back to default"
2827 );
2828 assert_eq!(
2829 cfg.tx_max_age_secs, default.tx_max_age_secs,
2830 "zero tx_max_age_secs must fall back to default"
2831 );
2832 }
2833
2834 #[test]
2837 #[serial]
2838 fn checkpoint_config_rejects_reversed_tx_thresholds() {
2839 let default = CheckpointConfig::default();
2840 std::env::set_var("KHIVE_TX_WARN_SECS", "120");
2841 std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "30");
2842
2843 let cfg = CheckpointConfig::from_env();
2844
2845 std::env::remove_var("KHIVE_TX_WARN_SECS");
2846 std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
2847
2848 assert_eq!(
2849 cfg.tx_warn_secs, default.tx_warn_secs,
2850 "a reversed pair must fall back tx_warn_secs to its default, got: {:?}",
2851 cfg.tx_warn_secs
2852 );
2853 assert_eq!(
2854 cfg.tx_max_age_secs, default.tx_max_age_secs,
2855 "a reversed pair must fall back tx_max_age_secs to its default, got: {:?}",
2856 cfg.tx_max_age_secs
2857 );
2858 }
2859
2860 #[test]
2863 #[serial]
2864 fn checkpoint_config_rejects_equal_tx_thresholds() {
2865 let default = CheckpointConfig::default();
2866 std::env::set_var("KHIVE_TX_WARN_SECS", "60");
2867 std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "60");
2868
2869 let cfg = CheckpointConfig::from_env();
2870
2871 std::env::remove_var("KHIVE_TX_WARN_SECS");
2872 std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
2873
2874 assert_eq!(
2875 cfg.tx_warn_secs, default.tx_warn_secs,
2876 "an equal pair must fall back tx_warn_secs to its default, got: {:?}",
2877 cfg.tx_warn_secs
2878 );
2879 assert_eq!(
2880 cfg.tx_max_age_secs, default.tx_max_age_secs,
2881 "an equal pair must fall back tx_max_age_secs to its default, got: {:?}",
2882 cfg.tx_max_age_secs
2883 );
2884 }
2885
2886 #[test]
2889 fn skipped_tick_does_not_reset_high_water_crossing_state() {
2890 let mut was_above = false;
2891
2892 assert!(
2894 crossing_warn(true, &mut was_above),
2895 "should fire on first crossing"
2896 );
2897 assert!(was_above);
2898
2899 assert!(was_above, "was_above must stay true across skipped ticks");
2906
2907 let fired = crossing_warn(true, &mut was_above);
2909 assert!(!fired, "WARN must not re-fire while still above threshold");
2910
2911 let fired = crossing_warn(false, &mut was_above);
2913 assert!(!fired);
2914 assert!(!was_above);
2915
2916 let fired = crossing_warn(true, &mut was_above);
2918 assert!(fired, "WARN must fire again on a new below→above crossing");
2919 }
2920
2921 #[test]
2928 fn warn_pages_fires_once_on_crossing_not_every_tick() {
2929 let mut was_above_warn = false;
2930
2931 let fired_1 = crossing_warn(true, &mut was_above_warn);
2933 let fired_2 = crossing_warn(true, &mut was_above_warn);
2934 let fired_3 = crossing_warn(true, &mut was_above_warn);
2935
2936 assert!(fired_1, "WARN must fire on the first in-band tick");
2937 assert!(
2938 !fired_2,
2939 "WARN must not fire on the second consecutive in-band tick"
2940 );
2941 assert!(
2942 !fired_3,
2943 "WARN must not fire on the third consecutive in-band tick"
2944 );
2945
2946 crossing_warn(false, &mut was_above_warn);
2948 assert!(!was_above_warn);
2949
2950 let fired_reentry = crossing_warn(true, &mut was_above_warn);
2952 assert!(
2953 fired_reentry,
2954 "WARN must fire again on re-entry into warn band"
2955 );
2956 }
2957
2958 #[test]
2964 #[serial(tx_registry, checkpoint_skip_metrics)]
2965 fn truncate_attempts_when_high_water_crossed_with_no_prior_attempt() {
2966 let dir = tempfile::tempdir().unwrap();
2967 let path = dir.path().join("truncate_trigger.db");
2968 let pool = file_pool(&path);
2969
2970 {
2971 let writer = pool.try_writer().unwrap();
2972 writer
2973 .conn()
2974 .execute_batch(
2975 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
2976 )
2977 .unwrap();
2978 }
2979
2980 let config = CheckpointConfig {
2981 truncate_high_water_pages: 0,
2985 truncate_min_interval: Duration::from_secs(300),
2986 ..CheckpointConfig::default()
2987 };
2988 let mut state = TruncateState::default();
2989
2990 assert!(
2991 state.last_attempt.is_none(),
2992 "precondition: no attempt has run yet"
2993 );
2994
2995 let tick = checkpoint_once(&pool, &config, &mut state);
2996 assert!(matches!(tick, CheckpointTick::Observed(_)));
2997 assert!(
2998 state.last_attempt.is_some(),
2999 "an attempt must be stamped once the high-water threshold is crossed"
3000 );
3001 }
3002
3003 #[test]
3006 #[serial(tx_registry, checkpoint_skip_metrics)]
3007 fn truncate_does_not_attempt_below_high_water() {
3008 let dir = tempfile::tempdir().unwrap();
3009 let path = dir.path().join("truncate_below_threshold.db");
3010 let pool = file_pool(&path);
3011
3012 {
3013 let writer = pool.try_writer().unwrap();
3014 writer
3015 .conn()
3016 .execute_batch(
3017 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
3018 )
3019 .unwrap();
3020 }
3021
3022 let config = CheckpointConfig {
3024 truncate_high_water_pages: u64::MAX,
3025 ..CheckpointConfig::default()
3026 };
3027 let mut state = TruncateState::default();
3028
3029 checkpoint_once(&pool, &config, &mut state);
3030
3031 assert!(
3032 state.last_attempt.is_none(),
3033 "a below-threshold tick must never stamp last_attempt"
3034 );
3035 }
3036
3037 #[test]
3041 #[serial(tx_registry, checkpoint_skip_metrics)]
3042 fn truncate_min_interval_skip_does_not_restamp_last_attempt() {
3043 let dir = tempfile::tempdir().unwrap();
3044 let path = dir.path().join("truncate_min_interval.db");
3045 let pool = file_pool(&path);
3046
3047 {
3048 let writer = pool.try_writer().unwrap();
3049 writer
3050 .conn()
3051 .execute_batch(
3052 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
3053 )
3054 .unwrap();
3055 }
3056
3057 let config = CheckpointConfig {
3058 truncate_high_water_pages: 0,
3059 truncate_min_interval: Duration::from_secs(300),
3060 ..CheckpointConfig::default()
3061 };
3062 let mut state = TruncateState::default();
3063
3064 checkpoint_once(&pool, &config, &mut state);
3065 let first_attempt = state.last_attempt.expect("first tick must attempt");
3066
3067 checkpoint_once(&pool, &config, &mut state);
3071 let second_attempt = state.last_attempt.expect("attempt timestamp must persist");
3072
3073 assert_eq!(
3074 first_attempt, second_attempt,
3075 "a tick within truncate_min_interval must not re-stamp last_attempt"
3076 );
3077 }
3078
3079 #[test]
3086 #[serial(tx_registry, checkpoint_skip_metrics)]
3087 fn busy_writer_skips_both_passive_and_truncate() {
3088 reset_checkpoint_metrics_for_tests();
3089
3090 let dir = tempfile::tempdir().unwrap();
3091 let path = dir.path().join("truncate_busy_skip.db");
3092 let pool = file_pool(&path);
3093
3094 {
3095 let writer = pool.try_writer().unwrap();
3096 writer
3097 .conn()
3098 .execute_batch(
3099 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
3100 )
3101 .unwrap();
3102 }
3103
3104 let mut warmup_state = TruncateState::default();
3107 let warmup_tick = checkpoint_once(&pool, &CheckpointConfig::default(), &mut warmup_state);
3108 let observed_pages = match warmup_tick {
3109 CheckpointTick::Observed(n) => n,
3110 CheckpointTick::Skipped => panic!("warmup tick must observe, not skip"),
3111 };
3112 assert_eq!(
3113 checkpoint_consecutive_skips(),
3114 0,
3115 "an observed tick must not itself count as a skip"
3116 );
3117
3118 let _held = pool.try_writer().unwrap();
3121
3122 let config = CheckpointConfig {
3123 truncate_high_water_pages: 0,
3124 ..CheckpointConfig::default()
3125 };
3126 let mut state = TruncateState::default();
3127
3128 let tick = checkpoint_once(&pool, &config, &mut state);
3129
3130 assert_eq!(
3131 tick,
3132 CheckpointTick::Skipped,
3133 "a busy writer must skip the tick entirely"
3134 );
3135 assert!(
3136 state.last_attempt.is_none(),
3137 "a skipped tick (writer busy) must never stamp last_attempt, \
3138 even with a threshold that would otherwise arm immediately"
3139 );
3140
3141 assert_eq!(
3142 checkpoint_skipped_ticks(),
3143 1,
3144 "one skipped tick must bump the lifetime skipped-tick counter"
3145 );
3146 assert_eq!(
3147 checkpoint_consecutive_skips(),
3148 1,
3149 "one skipped tick must bump the consecutive-skip run length"
3150 );
3151 assert_eq!(
3152 checkpoint_last_skip_wal_pages(),
3153 Some(observed_pages),
3154 "the skip must snapshot the last-observed WAL pressure"
3155 );
3156 }
3157
3158 #[test]
3173 fn all_checkpoint_metrics_callers_are_serial_tagged() {
3174 const SELF_SRC: &str = include_str!("checkpoint.rs");
3175 let lines: Vec<&str> = SELF_SRC.lines().collect();
3176
3177 let attr_starts: Vec<usize> = lines
3178 .iter()
3179 .enumerate()
3180 .filter(|(_, l)| {
3181 let t = l.trim();
3182 t == "#[test]" || t.starts_with("#[tokio::test")
3183 })
3184 .map(|(i, _)| i)
3185 .collect();
3186
3187 let mut offenders = Vec::new();
3188
3189 for (idx, &start) in attr_starts.iter().enumerate() {
3190 let end = attr_starts.get(idx + 1).copied().unwrap_or(lines.len());
3191 let span = &lines[start..end];
3192
3193 let touches_shared_metrics = span
3194 .iter()
3195 .any(|l| l.contains("checkpoint_once(") || l.contains("run_checkpoint_task("));
3196 if !touches_shared_metrics {
3197 continue;
3198 }
3199
3200 let has_group_tag = span
3201 .iter()
3202 .any(|l| l.contains("#[serial") && l.contains("checkpoint_skip_metrics"));
3203
3204 if !has_group_tag {
3205 let name = span
3206 .iter()
3207 .find_map(|l| {
3208 let t = l.trim_start();
3209 let t = t.strip_prefix("pub(crate) ").unwrap_or(t);
3210 let t = t.strip_prefix("pub ").unwrap_or(t);
3211 let t = t.strip_prefix("async ").unwrap_or(t);
3212 t.strip_prefix("fn ")
3213 .map(|rest| rest.split(['(', '<']).next().unwrap_or("").trim())
3214 })
3215 .unwrap_or("<unknown test>");
3216 offenders.push(name.to_string());
3217 }
3218 }
3219
3220 assert!(
3221 offenders.is_empty(),
3222 "these tests call checkpoint_once/run_checkpoint_task (which write the \
3223 process-wide LAST_WAL_PAGES/CHECKPOINT_* atomics via query_wal_pages) but \
3224 are not tagged #[serial(checkpoint_skip_metrics)] (or a group including it); \
3225 an untagged caller running concurrently on cargo's default test thread pool \
3226 can clobber those atomics mid-assertion in another test (the #828/#845 race): \
3227 {offenders:?}"
3228 );
3229 }
3230
3231 #[test]
3235 #[serial(tx_registry, checkpoint_skip_metrics)]
3236 fn observed_tick_resets_consecutive_skips_but_not_lifetime_total() {
3237 reset_checkpoint_metrics_for_tests();
3238
3239 let dir = tempfile::tempdir().unwrap();
3240 let path = dir.path().join("skip_then_observe.db");
3241 let pool = file_pool(&path);
3242
3243 {
3244 let writer = pool.try_writer().unwrap();
3245 writer
3246 .conn()
3247 .execute_batch(
3248 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
3249 )
3250 .unwrap();
3251 }
3252
3253 {
3255 let _held = pool.try_writer().unwrap();
3256 let mut state = TruncateState::default();
3257 for _ in 0..2 {
3258 let tick = checkpoint_once(&pool, &CheckpointConfig::default(), &mut state);
3259 assert_eq!(tick, CheckpointTick::Skipped);
3260 }
3261 }
3262 assert_eq!(checkpoint_skipped_ticks(), 2);
3263 assert_eq!(checkpoint_consecutive_skips(), 2);
3264
3265 let mut state = TruncateState::default();
3267 let tick = checkpoint_once(&pool, &CheckpointConfig::default(), &mut state);
3268 assert!(matches!(tick, CheckpointTick::Observed(_)));
3269
3270 assert_eq!(
3271 checkpoint_skipped_ticks(),
3272 2,
3273 "an observed tick must not change the lifetime skipped-tick total"
3274 );
3275 assert_eq!(
3276 checkpoint_consecutive_skips(),
3277 0,
3278 "an observed tick must reset the consecutive-skip run length"
3279 );
3280 }
3281
3282 #[test]
3287 fn note_truncate_outcome_warns_once_at_third_consecutive_failure() {
3288 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
3289 let subscriber = CaptureSubscriber {
3290 events: std::sync::Arc::clone(&buffer),
3291 };
3292
3293 let config = CheckpointConfig {
3294 warn_pages: 2000,
3295 ..CheckpointConfig::default()
3296 };
3297 let mut state = TruncateState::default();
3298
3299 tracing::subscriber::with_default(subscriber, || {
3300 note_truncate_outcome(&config, 5000, &mut state);
3302 note_truncate_outcome(&config, 5000, &mut state);
3303 note_truncate_outcome(&config, 5000, &mut state);
3304 note_truncate_outcome(&config, 5000, &mut state);
3306 });
3307
3308 assert_eq!(state.consecutive_failures, 4);
3309
3310 let events = buffer.lock().unwrap();
3311 let escalation_count = events
3312 .iter()
3313 .filter(|e| {
3314 e.message.as_deref()
3315 == Some(
3316 "WAL TRUNCATE has failed to clear WAL pressure for 3 consecutive attempts",
3317 )
3318 })
3319 .count();
3320 assert_eq!(
3321 escalation_count, 1,
3322 "escalation WARN must fire exactly once at the 3rd consecutive failure, got: {events:?}"
3323 );
3324
3325 note_truncate_outcome(&config, 100, &mut state);
3327 assert_eq!(
3328 state.consecutive_failures, 0,
3329 "an attempt that clears warn_pages must reset the consecutive-failure counter"
3330 );
3331 }
3332
3333 fn severity_test_config() -> CheckpointConfig {
3336 CheckpointConfig {
3337 warn_pages: 100,
3338 warn_sustained_cycles: 3,
3339 ..CheckpointConfig::default()
3340 }
3341 }
3342
3343 #[test]
3346 fn severity_ladder_info_on_first_crossing_no_warn() {
3347 let config = severity_test_config();
3348 let mut state = CheckpointSeverityState::default();
3349
3350 let below = state.observe_wal_pages(10, &config);
3351 assert!(below.is_empty(), "below-warn tick must emit nothing");
3352
3353 let above = state.observe_wal_pages(150, &config);
3354 assert_eq!(
3355 above,
3356 vec![CheckpointSeverityEmission {
3357 rung: CheckpointSeverityRung::Info,
3358 wal_pages: 150,
3359 threshold_pages: 100,
3360 consecutive_cycles: 1,
3361 }],
3362 "first below->above crossing must emit exactly one INFO and no WARN"
3363 );
3364 }
3365
3366 #[test]
3369 fn severity_ladder_warn_on_third_consecutive_cycle() {
3370 let config = severity_test_config();
3371 let mut state = CheckpointSeverityState::default();
3372
3373 let tick1 = state.observe_wal_pages(150, &config);
3374 assert_eq!(tick1.len(), 1);
3375 assert_eq!(tick1[0].rung, CheckpointSeverityRung::Info);
3376
3377 let tick2 = state.observe_wal_pages(150, &config);
3378 assert!(
3379 tick2.is_empty(),
3380 "second consecutive above-warn tick must emit nothing yet"
3381 );
3382
3383 let tick3 = state.observe_wal_pages(150, &config);
3384 assert_eq!(
3385 tick3,
3386 vec![CheckpointSeverityEmission {
3387 rung: CheckpointSeverityRung::Warn,
3388 wal_pages: 150,
3389 threshold_pages: 100,
3390 consecutive_cycles: 3,
3391 }],
3392 "WARN must fire exactly on the third consecutive above-warn tick"
3393 );
3394
3395 let tick4 = state.observe_wal_pages(150, &config);
3396 assert!(
3397 tick4.is_empty(),
3398 "WARN must not repeat on a fourth consecutive above-warn tick"
3399 );
3400 }
3401
3402 #[test]
3405 fn severity_ladder_rearms_warn_after_drain() {
3406 let config = severity_test_config();
3407 let mut state = CheckpointSeverityState::default();
3408
3409 for _ in 0..3 {
3411 state.observe_wal_pages(150, &config);
3412 }
3413 assert!(state.warn_emitted_for_episode);
3414
3415 let drain = state.observe_wal_pages(10, &config);
3417 assert!(drain.is_empty(), "a draining tick must emit nothing");
3418
3419 let reentry = state.observe_wal_pages(150, &config);
3421 assert_eq!(reentry.len(), 1);
3422 assert_eq!(reentry[0].rung, CheckpointSeverityRung::Info);
3423
3424 let mid = state.observe_wal_pages(150, &config);
3425 assert!(mid.is_empty());
3426
3427 let second_warn = state.observe_wal_pages(150, &config);
3428 assert_eq!(
3429 second_warn,
3430 vec![CheckpointSeverityEmission {
3431 rung: CheckpointSeverityRung::Warn,
3432 wal_pages: 150,
3433 threshold_pages: 100,
3434 consecutive_cycles: 3,
3435 }],
3436 "a fresh elevation episode after a drain must WARN again"
3437 );
3438 }
3439
3440 #[test]
3443 fn severity_ladder_isolated_crossings_never_warn() {
3444 let config = severity_test_config();
3445 let mut state = CheckpointSeverityState::default();
3446
3447 for _ in 0..3 {
3448 let crossing = state.observe_wal_pages(150, &config);
3449 assert_eq!(
3450 crossing.len(),
3451 1,
3452 "each isolated crossing must emit exactly one INFO"
3453 );
3454 assert_eq!(crossing[0].rung, CheckpointSeverityRung::Info);
3455
3456 let drain = state.observe_wal_pages(10, &config);
3457 assert!(drain.is_empty(), "the drain tick must emit nothing");
3458 }
3459
3460 assert!(
3461 !state.warn_emitted_for_episode,
3462 "isolated single-tick crossings must never accumulate into a WARN"
3463 );
3464 }
3465
3466 #[test]
3471 fn severity_ladder_never_emits_alarm() {
3472 let config = CheckpointConfig {
3473 warn_pages: 100,
3474 warn_sustained_cycles: 1,
3475 ..CheckpointConfig::default()
3476 };
3477 let mut state = CheckpointSeverityState::default();
3478
3479 for wal_pages in [150, 200, 250, u64::MAX] {
3480 let emissions = state.observe_wal_pages(wal_pages, &config);
3481 assert!(
3482 emissions
3483 .iter()
3484 .all(|e| e.rung != CheckpointSeverityRung::Alarm),
3485 "observe_wal_pages must never emit the ALARM rung, got: {emissions:?}"
3486 );
3487 }
3488 }
3489
3490 fn tx_age_test_config() -> CheckpointConfig {
3494 CheckpointConfig {
3495 tx_warn_secs: Duration::from_secs(30),
3496 tx_max_age_secs: Duration::from_secs(120),
3497 ..CheckpointConfig::default()
3498 }
3499 }
3500
3501 fn tx_id(n: u64) -> khive_storage::tx_registry::TxId {
3506 khive_storage::tx_registry::TxId(n)
3507 }
3508
3509 #[test]
3511 fn tx_age_sweep_empty_registry_emits_nothing() {
3512 let config = tx_age_test_config();
3513 let mut state = TxAgeSweepState::default();
3514
3515 let emissions = state.observe(None, config.tx_warn_secs, config.tx_max_age_secs);
3516 assert!(emissions.is_empty(), "no open entry must emit nothing");
3517 }
3518
3519 #[test]
3521 fn tx_age_sweep_fresh_entry_emits_nothing() {
3522 let config = tx_age_test_config();
3523 let mut state = TxAgeSweepState::default();
3524
3525 let emissions = state.observe(
3526 Some((
3527 tx_id(1),
3528 Duration::from_secs(5),
3529 Some("fresh_span".to_string()),
3530 )),
3531 config.tx_warn_secs,
3532 config.tx_max_age_secs,
3533 );
3534 assert!(emissions.is_empty(), "a fresh entry must emit nothing");
3535 }
3536
3537 #[test]
3541 fn tx_age_sweep_warn_fires_once_on_crossing() {
3542 let config = tx_age_test_config();
3543 let mut state = TxAgeSweepState::default();
3544
3545 let tick1 = state.observe(
3546 Some((
3547 tx_id(1),
3548 Duration::from_secs(45),
3549 Some("stale_span".to_string()),
3550 )),
3551 config.tx_warn_secs,
3552 config.tx_max_age_secs,
3553 );
3554 assert_eq!(
3555 tick1,
3556 vec![TxAgeEmission {
3557 rung: TxAgeRung::Warn,
3558 age: Duration::from_secs(45),
3559 label: Some("stale_span".to_string()),
3560 }],
3561 "crossing tx_warn_secs must emit exactly one Warn"
3562 );
3563
3564 let tick2 = state.observe(
3565 Some((
3566 tx_id(1),
3567 Duration::from_secs(50),
3568 Some("stale_span".to_string()),
3569 )),
3570 config.tx_warn_secs,
3571 config.tx_max_age_secs,
3572 );
3573 assert!(
3574 tick2.is_empty(),
3575 "Warn must not repeat while the entry stays in the warn band"
3576 );
3577 }
3578
3579 #[test]
3582 fn tx_age_sweep_stale_fires_once_on_crossing() {
3583 let config = tx_age_test_config();
3584 let mut state = TxAgeSweepState::default();
3585
3586 state.observe(
3589 Some((
3590 tx_id(1),
3591 Duration::from_secs(45),
3592 Some("stuck_writer_task_tx".to_string()),
3593 )),
3594 config.tx_warn_secs,
3595 config.tx_max_age_secs,
3596 );
3597
3598 let tick = state.observe(
3599 Some((
3600 tx_id(1),
3601 Duration::from_secs(130),
3602 Some("stuck_writer_task_tx".to_string()),
3603 )),
3604 config.tx_warn_secs,
3605 config.tx_max_age_secs,
3606 );
3607 assert_eq!(
3608 tick,
3609 vec![TxAgeEmission {
3610 rung: TxAgeRung::Stale,
3611 age: Duration::from_secs(130),
3612 label: Some("stuck_writer_task_tx".to_string()),
3613 }],
3614 "crossing tx_max_age_secs must emit exactly one Stale"
3615 );
3616
3617 let tick_repeat = state.observe(
3618 Some((
3619 tx_id(1),
3620 Duration::from_secs(200),
3621 Some("stuck_writer_task_tx".to_string()),
3622 )),
3623 config.tx_warn_secs,
3624 config.tx_max_age_secs,
3625 );
3626 assert!(
3627 tick_repeat.is_empty(),
3628 "Stale must not repeat while the entry stays above tx_max_age_secs"
3629 );
3630 }
3631
3632 #[test]
3636 fn tx_age_sweep_already_stale_entry_emits_both_rungs_same_tick() {
3637 let config = tx_age_test_config();
3638 let mut state = TxAgeSweepState::default();
3639
3640 let tick = state.observe(
3641 Some((
3642 tx_id(1),
3643 Duration::from_secs(300),
3644 Some("ancient_tx".to_string()),
3645 )),
3646 config.tx_warn_secs,
3647 config.tx_max_age_secs,
3648 );
3649 assert_eq!(
3650 tick,
3651 vec![
3652 TxAgeEmission {
3653 rung: TxAgeRung::Warn,
3654 age: Duration::from_secs(300),
3655 label: Some("ancient_tx".to_string()),
3656 },
3657 TxAgeEmission {
3658 rung: TxAgeRung::Stale,
3659 age: Duration::from_secs(300),
3660 label: Some("ancient_tx".to_string()),
3661 },
3662 ],
3663 "an already-stale entry must cross both rungs on its first observed tick"
3664 );
3665 }
3666
3667 #[test]
3670 fn tx_age_sweep_rearms_after_entry_clears() {
3671 let config = tx_age_test_config();
3672 let mut state = TxAgeSweepState::default();
3673
3674 state.observe(
3675 Some((
3676 tx_id(1),
3677 Duration::from_secs(150),
3678 Some("first_span".to_string()),
3679 )),
3680 config.tx_warn_secs,
3681 config.tx_max_age_secs,
3682 );
3683
3684 let cleared = state.observe(None, config.tx_warn_secs, config.tx_max_age_secs);
3686 assert!(cleared.is_empty(), "a clearing tick must emit nothing");
3687
3688 let fresh = state.observe(
3690 Some((
3691 tx_id(2),
3692 Duration::from_secs(2),
3693 Some("second_span".to_string()),
3694 )),
3695 config.tx_warn_secs,
3696 config.tx_max_age_secs,
3697 );
3698 assert!(fresh.is_empty(), "a fresh oldest entry must emit nothing");
3699
3700 let rewarn = state.observe(
3702 Some((
3703 tx_id(2),
3704 Duration::from_secs(35),
3705 Some("second_span".to_string()),
3706 )),
3707 config.tx_warn_secs,
3708 config.tx_max_age_secs,
3709 );
3710 assert_eq!(
3711 rewarn,
3712 vec![TxAgeEmission {
3713 rung: TxAgeRung::Warn,
3714 age: Duration::from_secs(35),
3715 label: Some("second_span".to_string()),
3716 }],
3717 "a fresh stale episode after a clear must Warn again"
3718 );
3719 }
3720
3721 #[test]
3725 fn tx_age_sweep_stale_replacement_without_intervening_clear_still_names_new_entry() {
3726 let config = tx_age_test_config();
3727 let mut state = TxAgeSweepState::default();
3728
3729 let tick_a = state.observe(
3730 Some((
3731 tx_id(1),
3732 Duration::from_secs(300),
3733 Some("stale_entry_a".to_string()),
3734 )),
3735 config.tx_warn_secs,
3736 config.tx_max_age_secs,
3737 );
3738 assert_eq!(
3739 tick_a.len(),
3740 2,
3741 "entry A must cross both rungs on its first observed tick, got: {tick_a:?}"
3742 );
3743
3744 let tick_b = state.observe(
3747 Some((
3748 tx_id(2),
3749 Duration::from_secs(400),
3750 Some("stale_entry_b".to_string()),
3751 )),
3752 config.tx_warn_secs,
3753 config.tx_max_age_secs,
3754 );
3755 assert_eq!(
3756 tick_b,
3757 vec![
3758 TxAgeEmission {
3759 rung: TxAgeRung::Warn,
3760 age: Duration::from_secs(400),
3761 label: Some("stale_entry_b".to_string()),
3762 },
3763 TxAgeEmission {
3764 rung: TxAgeRung::Stale,
3765 age: Duration::from_secs(400),
3766 label: Some("stale_entry_b".to_string()),
3767 },
3768 ],
3769 "a same-tick identity change to an already-stale successor must re-emit both \
3770 rungs naming the NEW entry, got: {tick_b:?}"
3771 );
3772 }
3773
3774 #[test]
3777 fn tx_age_sweep_uses_configured_thresholds_not_hardcoded_defaults() {
3778 let config = CheckpointConfig {
3779 tx_warn_secs: Duration::from_millis(1),
3780 tx_max_age_secs: Duration::from_millis(2),
3781 ..CheckpointConfig::default()
3782 };
3783 let mut state = TxAgeSweepState::default();
3784
3785 let tick = state.observe(
3786 Some((
3787 tx_id(1),
3788 Duration::from_millis(5),
3789 Some("fast_cap_span".to_string()),
3790 )),
3791 config.tx_warn_secs,
3792 config.tx_max_age_secs,
3793 );
3794 assert_eq!(
3795 tick.len(),
3796 2,
3797 "a millisecond-scale cap must cross both rungs immediately, got: {tick:?}"
3798 );
3799 }
3800
3801 #[test]
3804 #[serial(tx_registry, checkpoint_skip_metrics)]
3805 fn tx_age_sweep_names_long_lived_reader_pinning_wal_past_high_water() {
3806 let dir = tempfile::tempdir().unwrap();
3807 let path = dir.path().join("tx_age_sweep_reader_pin.db");
3808 let pool = file_pool(&path);
3809
3810 {
3811 let writer = pool.try_writer().unwrap();
3812 writer
3813 .conn()
3814 .execute_batch(
3815 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
3816 )
3817 .unwrap();
3818 }
3819
3820 let reader = pool.reader().expect("reader");
3825 reader
3826 .execute_batch("BEGIN DEFERRED; SELECT * FROM t;")
3827 .expect("begin read tx");
3828 let _tx_handle =
3829 khive_storage::tx_registry::register(Some("tx_age_sweep_reader_pin_test".to_string()));
3830
3831 let config = CheckpointConfig {
3834 high_water_pages: 1,
3835 tx_warn_secs: Duration::from_millis(1),
3836 tx_max_age_secs: Duration::from_millis(1),
3837 ..CheckpointConfig::default()
3838 };
3839 {
3840 let writer = pool.try_writer().unwrap();
3841 for i in 0..50 {
3842 writer
3843 .conn()
3844 .execute_batch(&format!("INSERT INTO t VALUES ({i});"))
3845 .unwrap();
3846 }
3847 }
3848
3849 let tick = checkpoint_once(&pool, &config, &mut TruncateState::default());
3850 let wal_pages = match tick {
3851 CheckpointTick::Observed(n) => n,
3852 CheckpointTick::Skipped => panic!("writer must not be busy in this test"),
3853 };
3854 assert!(
3855 wal_pages >= config.high_water_pages,
3856 "test setup must actually drive wal_pages ({wal_pages}) past high_water_pages \
3857 ({}) for this regression to mean anything",
3858 config.high_water_pages
3859 );
3860
3861 std::thread::sleep(Duration::from_millis(5));
3868 let our_entry = khive_storage::tx_registry::snapshot()
3879 .into_iter()
3880 .find(|(_, label)| label.as_deref() == Some("tx_age_sweep_reader_pin_test"))
3881 .expect("this test's own tx_registry entry must still be open");
3882 let mut tx_age_state = TxAgeSweepState::default();
3883 let emissions = tx_age_state.observe(
3884 Some((tx_id(1), our_entry.0, our_entry.1)),
3885 config.tx_warn_secs,
3886 config.tx_max_age_secs,
3887 );
3888 assert!(
3889 emissions.iter().any(|e| e.rung == TxAgeRung::Stale
3890 && e.label.as_deref() == Some("tx_age_sweep_reader_pin_test")),
3891 "expected a Stale emission naming the pinning reader, got: {emissions:?}"
3892 );
3893
3894 reader.execute_batch("COMMIT;").ok();
3895 drop(reader);
3896 drop(_tx_handle);
3897 }
3898
3899 #[test]
3902 #[serial(tx_registry, checkpoint_skip_metrics)]
3903 fn tx_age_sweep_own_entry_survives_concurrent_older_registration() {
3904 let _decoy = khive_storage::tx_registry::register(Some("decoy_unrelated_span".to_string()));
3905 std::thread::sleep(Duration::from_millis(2));
3906 let _own = khive_storage::tx_registry::register(Some("this_test_own_span".to_string()));
3907 std::thread::sleep(Duration::from_millis(5));
3908
3909 let global_oldest = khive_storage::tx_registry::oldest().expect("registry not empty");
3915 assert_ne!(
3916 global_oldest.2.as_deref(),
3917 Some("this_test_own_span"),
3918 "test setup must reproduce the race: an older, unrelated entry must be \
3919 the current global oldest, got: {global_oldest:?}"
3920 );
3921
3922 let our_entry = khive_storage::tx_registry::snapshot()
3923 .into_iter()
3924 .find(|(_, label)| label.as_deref() == Some("this_test_own_span"))
3925 .expect("this test's own tx_registry entry must still be open");
3926
3927 let config = CheckpointConfig {
3928 tx_warn_secs: Duration::from_millis(1),
3929 tx_max_age_secs: Duration::from_millis(1),
3930 ..CheckpointConfig::default()
3931 };
3932 let mut state = TxAgeSweepState::default();
3933 let emissions = state.observe(
3934 Some((tx_id(2), our_entry.0, our_entry.1)),
3935 config.tx_warn_secs,
3936 config.tx_max_age_secs,
3937 );
3938 assert!(
3939 emissions
3940 .iter()
3941 .any(|e| e.rung == TxAgeRung::Stale
3942 && e.label.as_deref() == Some("this_test_own_span")),
3943 "expected a Stale emission naming this test's own span despite an older, \
3944 unrelated concurrent registration, got: {emissions:?}"
3945 );
3946 }
3947
3948 #[test]
3950 #[serial]
3951 fn checkpoint_config_warn_sustained_cycles_env_override() {
3952 let default = CheckpointConfig::default();
3953 assert_eq!(default.warn_sustained_cycles, DEFAULT_WARN_SUSTAINED_CYCLES);
3954
3955 std::env::set_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES", "5");
3956 let cfg = CheckpointConfig::from_env();
3957 std::env::remove_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES");
3958 assert_eq!(cfg.warn_sustained_cycles, 5);
3959
3960 std::env::set_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES", "0");
3961 let cfg_zero = CheckpointConfig::from_env();
3962 std::env::remove_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES");
3963 assert_eq!(
3964 cfg_zero.warn_sustained_cycles, DEFAULT_WARN_SUSTAINED_CYCLES,
3965 "zero must fall back to the default"
3966 );
3967
3968 std::env::set_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES", "not_a_number");
3969 let cfg_invalid = CheckpointConfig::from_env();
3970 std::env::remove_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES");
3971 assert_eq!(
3972 cfg_invalid.warn_sustained_cycles,
3973 DEFAULT_WARN_SUSTAINED_CYCLES
3974 );
3975 }
3976
3977 #[derive(Clone, Copy)]
3980 enum FakeAppendBehavior {
3981 Record,
3982 Fail,
3983 }
3984
3985 struct FakeEventStore {
3986 events: std::sync::Mutex<Vec<khive_storage::Event>>,
3987 append_attempts: std::sync::atomic::AtomicUsize,
3988 append_behavior: FakeAppendBehavior,
3989 }
3990
3991 impl Default for FakeEventStore {
3992 fn default() -> Self {
3993 Self {
3994 events: std::sync::Mutex::new(Vec::new()),
3995 append_attempts: std::sync::atomic::AtomicUsize::new(0),
3996 append_behavior: FakeAppendBehavior::Record,
3997 }
3998 }
3999 }
4000
4001 impl FakeEventStore {
4002 fn failing() -> Self {
4003 Self {
4004 append_behavior: FakeAppendBehavior::Fail,
4005 ..Self::default()
4006 }
4007 }
4008 }
4009
4010 #[async_trait::async_trait]
4011 impl khive_storage::EventStore for FakeEventStore {
4012 async fn append_event(
4013 &self,
4014 event: khive_storage::Event,
4015 ) -> khive_storage::StorageResult<()> {
4016 self.append_attempts
4017 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
4018 match self.append_behavior {
4019 FakeAppendBehavior::Record => {
4020 self.events.lock().unwrap().push(event);
4021 Ok(())
4022 }
4023 FakeAppendBehavior::Fail => Err(khive_storage::StorageError::Internal(
4024 "synthetic checkpoint lifecycle append failure".to_string(),
4025 )),
4026 }
4027 }
4028
4029 async fn append_events(
4030 &self,
4031 events: Vec<khive_storage::Event>,
4032 ) -> khive_storage::StorageResult<khive_storage::BatchWriteSummary> {
4033 let count = events.len() as u64;
4034 self.events.lock().unwrap().extend(events);
4035 Ok(khive_storage::BatchWriteSummary {
4036 attempted: count,
4037 affected: count,
4038 failed: 0,
4039 first_error: String::new(),
4040 })
4041 }
4042
4043 async fn get_event(
4044 &self,
4045 id: uuid::Uuid,
4046 ) -> khive_storage::StorageResult<Option<khive_storage::Event>> {
4047 Ok(self
4048 .events
4049 .lock()
4050 .unwrap()
4051 .iter()
4052 .find(|e| e.id == id)
4053 .cloned())
4054 }
4055
4056 async fn query_events(
4057 &self,
4058 _filter: khive_storage::EventFilter,
4059 _page: khive_storage::PageRequest,
4060 ) -> khive_storage::StorageResult<khive_storage::Page<khive_storage::Event>> {
4061 unimplemented!("not exercised by the checkpoint lifecycle-event tests")
4062 }
4063
4064 async fn count_events(
4065 &self,
4066 _filter: khive_storage::EventFilter,
4067 ) -> khive_storage::StorageResult<u64> {
4068 Ok(self.events.lock().unwrap().len() as u64)
4069 }
4070 }
4071
4072 #[test]
4077 fn checkpoint_outcome_should_emit_covers_all_transitions() {
4078 assert!(
4079 checkpoint_outcome_should_emit(true, false),
4080 "first elevated tick must emit"
4081 );
4082 assert!(
4083 checkpoint_outcome_should_emit(true, true),
4084 "sustained elevated tick must emit"
4085 );
4086 assert!(
4087 checkpoint_outcome_should_emit(false, true),
4088 "the single drain row (elevated -> healthy) must emit"
4089 );
4090 assert!(
4091 !checkpoint_outcome_should_emit(false, false),
4092 "an ordinary below-warn tick must not emit"
4093 );
4094 }
4095
4096 #[tokio::test]
4097 #[serial(checkpoint_skip_metrics)]
4098 async fn checkpoint_task_emits_outcome_events_while_elevated_and_stops_after_drain() {
4099 let dir = tempfile::tempdir().unwrap();
4100 let path = dir.path().join("outcome_emit.db");
4101 let pool = file_pool(&path);
4102
4103 let cfg = CheckpointConfig {
4106 interval: Duration::from_millis(10),
4107 warn_pages: 0,
4108 ..CheckpointConfig::default()
4109 };
4110 let store = Arc::new(FakeEventStore::default());
4111 let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
4112
4113 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4114 let handle = tokio::spawn(run_checkpoint_task(
4115 pool,
4116 cfg,
4117 Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
4118 shutdown_rx,
4119 true,
4120 ));
4121
4122 let emitted = wait_for(Duration::from_secs(10), || {
4125 !store.events.lock().unwrap().is_empty()
4126 })
4127 .await;
4128 shutdown_tx.send(()).expect("send shutdown signal");
4129 tokio::time::timeout(Duration::from_secs(1), handle)
4130 .await
4131 .expect("checkpoint task should exit within 1s")
4132 .expect("checkpoint task panicked");
4133
4134 let events = store.events.lock().unwrap();
4135 assert!(
4136 emitted,
4137 "an always-elevated config must append at least one CheckpointOutcomeRecorded event \
4138 within the poll deadline"
4139 );
4140 assert!(
4141 events
4142 .iter()
4143 .all(|e| e.kind == khive_types::EventKind::CheckpointOutcomeRecorded),
4144 "every appended event must be CheckpointOutcomeRecorded, got: {events:?}"
4145 );
4146 assert!(
4147 events.iter().all(|e| e.namespace == "local"),
4148 "events must be stamped with the namespace passed to run_checkpoint_task"
4149 );
4150 }
4151
4152 #[tokio::test]
4159 #[serial(checkpoint_skip_metrics)]
4160 async fn checkpoint_cycles_and_task_shutdown_do_not_wait_for_a_contended_lifecycle_writer() {
4161 let dir = tempfile::tempdir().unwrap();
4162 let path = dir.path().join("outcome_contended_sink.db");
4163 let checkpoint_pool = file_pool(&path);
4164
4165 let event_pool = Arc::new(
4166 ConnectionPool::new(PoolConfig {
4167 path: None,
4168 checkout_timeout: Duration::from_secs(5),
4169 write_queue_enabled: false,
4170 ..PoolConfig::default()
4171 })
4172 .expect("event pool"),
4173 );
4174 {
4175 let writer = event_pool.try_writer().expect("initialize event schema");
4176 crate::stores::event::ensure_events_schema(writer.conn())
4177 .expect("initialize event schema");
4178 }
4179 let event_store: Arc<dyn khive_storage::EventStore> =
4180 Arc::new(crate::stores::event::SqlEventStore::new_scoped(
4181 Arc::clone(&event_pool),
4182 false,
4183 "local",
4184 ));
4185 let held_event_writer = event_pool
4186 .try_writer()
4187 .expect("hold the event-store writer");
4188
4189 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
4190 let subscriber = CaptureSubscriber {
4191 events: std::sync::Arc::clone(&buffer),
4192 };
4193 let _tracing_guard = tracing::subscriber::set_default(subscriber);
4194
4195 let cfg = CheckpointConfig {
4196 interval: Duration::from_millis(10),
4197 warn_pages: 0,
4198 ..CheckpointConfig::default()
4199 };
4200 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4201 let handle = tokio::spawn(run_checkpoint_task(
4202 checkpoint_pool,
4203 cfg,
4204 Some(CheckpointLifecycleOwner::new(event_store, "local")),
4205 shutdown_rx,
4206 true,
4207 ));
4208
4209 let dropped = wait_for(Duration::from_secs(2), || {
4210 buffer.lock().unwrap().iter().any(|event| {
4211 event.message.as_deref()
4212 == Some("checkpoint lifecycle event dropped because the append worker is busy")
4213 })
4214 })
4215 .await;
4216 assert!(
4217 dropped,
4218 "later elevated ticks must reach the non-blocking enqueue while the first append is \
4219 still waiting for the held event writer; got: {:?}",
4220 buffer.lock().unwrap()
4221 );
4222
4223 shutdown_tx.send(()).expect("send shutdown signal");
4224 tokio::time::timeout(Duration::from_secs(1), handle)
4225 .await
4226 .expect(
4227 "the run_checkpoint_task handle must not wait for the event store's \
4228 five-second writer checkout",
4229 )
4230 .expect("checkpoint task panicked");
4231
4232 drop(held_event_writer);
4237 }
4238
4239 #[tokio::test]
4242 #[serial(checkpoint_skip_metrics)]
4243 async fn checkpoint_task_continues_after_lifecycle_append_failure() {
4244 let dir = tempfile::tempdir().unwrap();
4245 let path = dir.path().join("outcome_failing_sink.db");
4246 let pool = file_pool(&path);
4247 let store = Arc::new(FakeEventStore::failing());
4248 let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
4249
4250 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
4251 let subscriber = CaptureSubscriber {
4252 events: std::sync::Arc::clone(&buffer),
4253 };
4254 let _tracing_guard = tracing::subscriber::set_default(subscriber);
4255
4256 let cfg = CheckpointConfig {
4257 interval: Duration::from_millis(10),
4258 warn_pages: 0,
4259 ..CheckpointConfig::default()
4260 };
4261 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4262 let handle = tokio::spawn(run_checkpoint_task(
4263 pool,
4264 cfg,
4265 Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
4266 shutdown_rx,
4267 true,
4268 ));
4269
4270 let retried = wait_for(Duration::from_secs(2), || {
4271 store
4272 .append_attempts
4273 .load(std::sync::atomic::Ordering::Relaxed)
4274 >= 2
4275 })
4276 .await;
4277 shutdown_tx.send(()).expect("send shutdown signal");
4278 tokio::time::timeout(Duration::from_secs(1), handle)
4279 .await
4280 .expect("checkpoint task should remain responsive after sink failure")
4281 .expect("checkpoint task panicked");
4282
4283 assert!(
4284 retried,
4285 "a failed append must not terminate the worker or checkpoint task"
4286 );
4287 let captured = buffer.lock().unwrap().clone();
4288 assert!(
4289 captured.iter().any(|event| event.message.as_deref()
4290 == Some("checkpoint lifecycle event append failed")),
4291 "lifecycle append failures must remain observable; got: {:?}",
4292 captured
4293 );
4294 }
4295
4296 #[tokio::test]
4297 #[serial(checkpoint_skip_metrics)]
4298 async fn secondary_checkpoint_task_with_lifecycle_ownership_emits_outcome_events() {
4299 let dir = tempfile::tempdir().unwrap();
4300 let path = dir.path().join("secondary_outcome.db");
4301 let pool = file_pool(&path);
4302 let cfg = CheckpointConfig {
4303 interval: Duration::from_millis(10),
4304 warn_pages: 0,
4305 ..CheckpointConfig::default()
4306 };
4307 let store = Arc::new(FakeEventStore::default());
4308 let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
4309
4310 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4311 let handle = tokio::spawn(run_checkpoint_task(
4312 pool,
4313 cfg,
4314 Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
4315 shutdown_rx,
4316 false,
4317 ));
4318
4319 let emitted = wait_for(Duration::from_secs(10), || {
4322 !store.events.lock().unwrap().is_empty()
4323 })
4324 .await;
4325 shutdown_tx.send(()).expect("send shutdown signal");
4326 tokio::time::timeout(Duration::from_secs(1), handle)
4327 .await
4328 .expect("checkpoint task should exit within 1s")
4329 .expect("checkpoint task panicked");
4330
4331 assert!(
4332 emitted,
4333 "a designated secondary lifecycle owner must append outcome events within the poll \
4334 deadline"
4335 );
4336 }
4337
4338 #[tokio::test]
4339 #[serial(checkpoint_skip_metrics)]
4340 async fn checkpoint_task_emits_nothing_while_healthy() {
4341 let dir = tempfile::tempdir().unwrap();
4342 let path = dir.path().join("outcome_no_emit.db");
4343 let pool = file_pool(&path);
4344
4345 let cfg = CheckpointConfig {
4348 interval: Duration::from_millis(10),
4349 warn_pages: u64::MAX,
4350 ..CheckpointConfig::default()
4351 };
4352 let store = Arc::new(FakeEventStore::default());
4353 let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
4354
4355 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4356 let handle = tokio::spawn(run_checkpoint_task(
4357 pool,
4358 cfg,
4359 Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
4360 shutdown_rx,
4361 true,
4362 ));
4363
4364 tokio::time::sleep(Duration::from_millis(60)).await;
4365 shutdown_tx.send(()).expect("send shutdown signal");
4366 tokio::time::timeout(Duration::from_secs(1), handle)
4367 .await
4368 .expect("checkpoint task should exit within 1s")
4369 .expect("checkpoint task panicked");
4370
4371 assert!(
4372 store.events.lock().unwrap().is_empty(),
4373 "a config that never crosses warn_pages must never append a lifecycle event"
4374 );
4375 }
4376
4377 #[tokio::test]
4378 #[serial(checkpoint_skip_metrics)]
4379 async fn checkpoint_task_with_no_event_store_does_not_panic() {
4380 let dir = tempfile::tempdir().unwrap();
4381 let path = dir.path().join("outcome_none_store.db");
4382 let pool = file_pool(&path);
4383
4384 let cfg = CheckpointConfig {
4385 interval: Duration::from_millis(10),
4386 warn_pages: 0,
4387 ..CheckpointConfig::default()
4388 };
4389
4390 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4391 let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
4392
4393 tokio::time::sleep(Duration::from_millis(40)).await;
4394 shutdown_tx.send(()).expect("send shutdown signal");
4395 tokio::time::timeout(Duration::from_secs(1), handle)
4396 .await
4397 .expect("checkpoint task should exit within 1s")
4398 .expect("checkpoint task panicked");
4399 }
4400
4401 #[tokio::test]
4418 #[serial(tx_registry, checkpoint_skip_metrics)]
4419 async fn checkpoint_task_sweeps_stale_registry_entry_while_wal_is_healthy() {
4420 let dir = tempfile::tempdir().unwrap();
4421 let path = dir.path().join("tx_age_sweep_task_healthy_wal.db");
4422 let pool = file_pool(&path);
4423
4424 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
4425 let subscriber = CaptureSubscriber {
4426 events: std::sync::Arc::clone(&buffer),
4427 };
4428 let _tracing_guard = tracing::subscriber::set_default(subscriber);
4429
4430 let _tx_handle = khive_storage::tx_registry::register(Some(
4431 "checkpoint_task_healthy_wal_sweep_test".to_string(),
4432 ));
4433
4434 let cfg = CheckpointConfig {
4435 interval: Duration::from_millis(10),
4436 warn_pages: u64::MAX,
4437 high_water_pages: u64::MAX,
4438 truncate_high_water_pages: u64::MAX,
4439 tx_warn_secs: Duration::from_millis(1),
4440 tx_max_age_secs: Duration::from_millis(1),
4441 ..CheckpointConfig::default()
4442 };
4443
4444 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4445 let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
4446
4447 let swept = wait_for(Duration::from_secs(10), || {
4454 buffer.lock().unwrap().iter().any(|e| {
4455 e.tx_label.as_deref() == Some("checkpoint_task_healthy_wal_sweep_test")
4456 && e.message
4457 .as_deref()
4458 .is_some_and(|m| m.contains("stale-op cap"))
4459 })
4460 })
4461 .await;
4462
4463 shutdown_tx.send(()).expect("send shutdown signal");
4464 tokio::time::timeout(Duration::from_secs(1), handle)
4465 .await
4466 .expect("checkpoint task should exit within 1s")
4467 .expect("checkpoint task panicked");
4468
4469 drop(_tx_handle);
4470
4471 let events = buffer.lock().unwrap();
4472 assert!(
4473 swept,
4474 "expected the spawned task to sweep and escalate the stale registry entry \
4475 to Stale on its own within the poll deadline, got: {events:?}"
4476 );
4477 }
4478
4479 #[tokio::test]
4483 #[serial(tx_registry, checkpoint_skip_metrics)]
4484 async fn checkpoint_task_emits_no_age_alert_for_an_empty_registry() {
4485 let dir = tempfile::tempdir().unwrap();
4486 let path = dir.path().join("tx_age_sweep_task_empty_registry.db");
4487 let pool = file_pool(&path);
4488
4489 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
4490 let subscriber = CaptureSubscriber {
4491 events: std::sync::Arc::clone(&buffer),
4492 };
4493 let _tracing_guard = tracing::subscriber::set_default(subscriber);
4494
4495 let cfg = CheckpointConfig {
4496 interval: Duration::from_millis(10),
4497 tx_warn_secs: Duration::from_millis(1),
4498 tx_max_age_secs: Duration::from_millis(1),
4499 ..CheckpointConfig::default()
4500 };
4501
4502 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4503 let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
4504
4505 tokio::time::sleep(Duration::from_millis(40)).await;
4506 shutdown_tx.send(()).expect("send shutdown signal");
4507 tokio::time::timeout(Duration::from_secs(1), handle)
4508 .await
4509 .expect("checkpoint task should exit within 1s")
4510 .expect("checkpoint task panicked");
4511
4512 let events = buffer.lock().unwrap();
4513 assert!(
4514 events.iter().all(|e| e
4515 .message
4516 .as_deref()
4517 .is_none_or(|m| !m.contains("ADR-091 Plank 1"))),
4518 "an empty registry must never produce a Plank 1 age emission, got: {events:?}"
4519 );
4520 }
4521
4522 #[tokio::test]
4530 #[serial(tx_registry, checkpoint_skip_metrics)]
4531 async fn checkpoint_task_sweeps_stale_entry_even_when_writer_is_busy_every_tick() {
4532 reset_checkpoint_metrics_for_tests();
4533
4534 let dir = tempfile::tempdir().unwrap();
4535 let path = dir.path().join("tx_age_sweep_task_writer_busy.db");
4536 let pool = file_pool(&path);
4537 {
4538 let writer = pool.try_writer().unwrap();
4539 writer
4540 .conn()
4541 .execute_batch("CREATE TABLE IF NOT EXISTS t (x INTEGER);")
4542 .unwrap();
4543 }
4544
4545 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
4546 let subscriber = CaptureSubscriber {
4547 events: std::sync::Arc::clone(&buffer),
4548 };
4549 let _tracing_guard = tracing::subscriber::set_default(subscriber);
4550
4551 let _tx_handle = khive_storage::tx_registry::register(Some(
4552 "checkpoint_task_writer_busy_sweep_test".to_string(),
4553 ));
4554
4555 let _writer_guard = pool.try_writer().expect("acquire writer for busy hold");
4559
4560 let cfg = CheckpointConfig {
4561 interval: Duration::from_millis(10),
4562 tx_warn_secs: Duration::from_millis(1),
4563 tx_max_age_secs: Duration::from_millis(1),
4564 ..CheckpointConfig::default()
4565 };
4566
4567 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4568 let handle = tokio::spawn(run_checkpoint_task(
4569 Arc::clone(&pool),
4570 cfg,
4571 None,
4572 shutdown_rx,
4573 true,
4574 ));
4575
4576 let deadline = tokio::time::Instant::now() + Duration::from_secs(10);
4582 while checkpoint_skipped_ticks() == 0 {
4583 assert!(
4584 tokio::time::Instant::now() < deadline,
4585 "test setup must actually drive at least one Skipped tick for this \
4586 regression to mean anything (none within 10s)"
4587 );
4588 tokio::time::sleep(Duration::from_millis(10)).await;
4589 }
4590
4591 loop {
4592 let events = buffer.lock().unwrap().clone();
4593 if events.iter().any(|e| {
4594 e.tx_label.as_deref() == Some("checkpoint_task_writer_busy_sweep_test")
4595 && e.message
4596 .as_deref()
4597 .is_some_and(|m| m.contains("stale-op cap"))
4598 }) {
4599 break;
4600 }
4601 assert!(
4602 tokio::time::Instant::now() < deadline,
4603 "expected the age sweep to fire even though every tick's writer checkout \
4604 was skipped within 10s, got: {events:?}"
4605 );
4606 tokio::time::sleep(Duration::from_millis(10)).await;
4607 }
4608
4609 shutdown_tx.send(()).expect("send shutdown signal");
4610 tokio::time::timeout(Duration::from_secs(1), handle)
4611 .await
4612 .expect("checkpoint task should exit within 1s")
4613 .expect("checkpoint task panicked");
4614
4615 drop(_writer_guard);
4616 drop(_tx_handle);
4617 }
4618
4619 #[tokio::test]
4623 async fn session_sweep_task_exits_on_shutdown_signal() {
4624 let cfg = SessionSweepConfig {
4625 interval: Duration::from_millis(10),
4626 ..SessionSweepConfig::default()
4627 };
4628 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4629 let handle = tokio::spawn(run_session_sweep_task(Vec::new(), cfg, shutdown_rx));
4630
4631 shutdown_tx.send(()).expect("send shutdown signal");
4632
4633 tokio::time::timeout(Duration::from_secs(1), handle)
4634 .await
4635 .expect("session sweep task should exit within 1s")
4636 .expect("session sweep task panicked");
4637 }
4638
4639 async fn wait_for(deadline: Duration, mut cond: impl FnMut() -> bool) -> bool {
4643 let start = std::time::Instant::now();
4644 while start.elapsed() < deadline {
4645 if cond() {
4646 return true;
4647 }
4648 tokio::time::sleep(Duration::from_millis(5)).await;
4649 }
4650 cond()
4651 }
4652
4653 #[tokio::test]
4654 #[serial(khive_walpin_sidecar_env)]
4655 async fn walpin_observe_drops_beacon_when_heartbeat_write_fails() {
4656 let dir = tempfile::tempdir().unwrap();
4657 let db_path = dir.path().join("observe_gate.db");
4658 let sidecar_dir = crate::walpin::sidecar_dir_for(&db_path);
4659 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
4660 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
4661
4662 let mut state = WalpinSidecarState::new(
4663 Some(db_path.as_path()),
4664 true,
4665 "session",
4666 Duration::from_millis(500),
4667 )
4668 .expect("sidecar enabled for a file-backed path");
4669 state.register_beacon().await;
4670 let pid = std::process::id();
4671 let beacon_path = sidecar_dir.join(format!("{pid}.beacon"));
4672 let before = std::fs::metadata(&beacon_path)
4673 .expect("beacon registered")
4674 .modified()
4675 .unwrap();
4676
4677 let obstruction = sidecar_dir.join(format!(".{pid}.json.tmp"));
4682 std::fs::create_dir(&obstruction).unwrap();
4683
4684 tokio::time::sleep(Duration::from_millis(20)).await;
4685 let over_threshold = Some(khive_storage::tx_registry::OldestSpan {
4686 id: khive_storage::tx_registry::TxId(1),
4687 age: Duration::from_secs(60),
4688 label: None,
4689 origin: khive_storage::tx_registry::TxOrigin::Unscoped,
4690 });
4691 state
4692 .observe(over_threshold.clone(), Duration::from_secs(30))
4693 .await;
4694
4695 assert!(
4696 !sidecar_dir.join(format!("{pid}.json")).exists(),
4697 "heartbeat write must have failed"
4698 );
4699 assert!(
4702 !beacon_path.exists(),
4703 "a failed heartbeat write must remove the beacon — a still-fresh \
4704 beacon with no heartbeat would classify registered-silent \
4705 (before-mtime {before:?})"
4706 );
4707
4708 std::fs::remove_dir(&obstruction).unwrap();
4711 state.observe(over_threshold, Duration::from_secs(30)).await;
4712 assert!(
4713 sidecar_dir.join(format!("{pid}.json")).exists(),
4714 "heartbeat must land once the write path recovers"
4715 );
4716 assert!(
4717 beacon_path.exists(),
4718 "beacon must re-register on the first healthy tick after removal"
4719 );
4720 }
4721
4722 #[tokio::test]
4723 #[serial(khive_walpin_sidecar_env)]
4724 async fn walpin_observe_touches_mtime_without_rewriting_body_when_content_unchanged() {
4725 let dir = tempfile::tempdir().unwrap();
4726 let db_path = dir.path().join("observe_touch.db");
4727 let sidecar_dir = crate::walpin::sidecar_dir_for(&db_path);
4728 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
4729 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
4730
4731 let mut state = WalpinSidecarState::new(
4732 Some(db_path.as_path()),
4733 true,
4734 "session",
4735 Duration::from_millis(500),
4736 )
4737 .expect("sidecar enabled for a file-backed path");
4738 let pid = std::process::id();
4739 let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
4740 let span = khive_storage::tx_registry::OldestSpan {
4741 id: khive_storage::tx_registry::TxId(1),
4742 age: Duration::from_secs(60),
4743 label: None,
4744 origin: khive_storage::tx_registry::TxOrigin::Unscoped,
4745 };
4746
4747 state
4748 .observe(Some(span.clone()), Duration::from_secs(30))
4749 .await;
4750 let body_after_create = std::fs::read(&heartbeat_path).expect("heartbeat written");
4751
4752 let backdated = std::time::SystemTime::now() - Duration::from_secs(120);
4758 std::fs::File::open(&heartbeat_path)
4759 .unwrap()
4760 .set_modified(backdated)
4761 .unwrap();
4762
4763 state.observe(Some(span), Duration::from_secs(30)).await;
4764
4765 let body_after_second_observe =
4766 std::fs::read(&heartbeat_path).expect("heartbeat still present");
4767 assert_eq!(
4768 body_after_create, body_after_second_observe,
4769 "unchanged oldest-span identity/label/attribution/cadence must touch mtime, \
4770 not rewrite the body"
4771 );
4772 let mtime_after = std::fs::metadata(&heartbeat_path)
4773 .unwrap()
4774 .modified()
4775 .unwrap();
4776 assert!(
4777 mtime_after > backdated,
4778 "the touch must advance mtime past the backdated value"
4779 );
4780 }
4781
4782 #[tokio::test]
4783 #[serial(khive_walpin_sidecar_env)]
4784 async fn walpin_observe_recreates_heartbeat_after_it_is_deleted_while_span_still_live() {
4785 let dir = tempfile::tempdir().unwrap();
4786 let db_path = dir.path().join("observe_recreate.db");
4787 let sidecar_dir = crate::walpin::sidecar_dir_for(&db_path);
4788 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
4789 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
4790
4791 let mut state = WalpinSidecarState::new(
4792 Some(db_path.as_path()),
4793 true,
4794 "session",
4795 Duration::from_millis(500),
4796 )
4797 .expect("sidecar enabled for a file-backed path");
4798 let pid = std::process::id();
4799 let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
4800 let span = khive_storage::tx_registry::OldestSpan {
4801 id: khive_storage::tx_registry::TxId(1),
4802 age: Duration::from_secs(60),
4803 label: None,
4804 origin: khive_storage::tx_registry::TxOrigin::Unscoped,
4805 };
4806
4807 state
4808 .observe(Some(span.clone()), Duration::from_secs(30))
4809 .await;
4810 assert!(heartbeat_path.exists(), "heartbeat written on first tick");
4811
4812 std::fs::remove_file(&heartbeat_path).unwrap();
4818 assert!(!heartbeat_path.exists());
4819
4820 state.observe(Some(span), Duration::from_secs(30)).await;
4821
4822 assert!(
4823 heartbeat_path.exists(),
4824 "a touch failure against a deleted heartbeat must recreate it via a full write"
4825 );
4826 let recreated: crate::walpin::WalpinHeartbeat =
4827 serde_json::from_slice(&std::fs::read(&heartbeat_path).unwrap()).unwrap();
4828 assert_eq!(recreated.pid, pid);
4829 assert_eq!(recreated.oldest_tx_age_secs, 60.0);
4830 }
4831
4832 #[tokio::test]
4833 #[serial(tx_registry, khive_walpin_sidecar_env)]
4834 async fn session_sweep_task_writes_and_clears_walpin_heartbeat() {
4835 let dir = tempfile::tempdir().unwrap();
4836 let db_path = dir.path().join("session_sweep.db");
4837 let pool = file_pool(&db_path);
4838 let sidecar_dir =
4839 crate::walpin::sidecar_dir_for(pool.canonical_path().expect("file-backed pool"));
4840 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
4841 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
4842
4843 let cfg = SessionSweepConfig {
4844 interval: Duration::from_millis(10),
4845 tx_warn_secs: Duration::from_millis(20),
4846 tx_max_age_secs: Duration::from_millis(500),
4847 };
4848 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4849 let handle = tokio::spawn(run_session_sweep_task(
4850 vec![SweepBackend {
4851 pool: Arc::clone(&pool),
4852 is_main: true,
4853 }],
4854 cfg,
4855 shutdown_rx,
4856 ));
4857
4858 let pid = std::process::id();
4865 let beacon = crate::walpin::beacon_path(&sidecar_dir, pid);
4866 assert!(
4867 wait_for(Duration::from_secs(2), || beacon.exists()).await,
4868 "a quiet process must still register its one-time beacon"
4869 );
4870 assert!(
4871 !sidecar_dir.join(format!("{pid}.json")).exists(),
4872 "a quiet process must not write a walpin heartbeat"
4873 );
4874
4875 let tx_handle =
4876 khive_storage::tx_registry::register(Some("session_sweep_walpin_test".to_string()));
4877 let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
4878 assert!(
4879 wait_for(Duration::from_secs(2), || heartbeat_path.exists()).await,
4880 "expected a walpin heartbeat once the span crossed tx_warn_secs"
4881 );
4882 let body = std::fs::read_to_string(&heartbeat_path).unwrap();
4883 let hb: crate::walpin::WalpinHeartbeat = serde_json::from_str(&body).unwrap();
4884 assert_eq!(hb.pid, pid);
4885 assert_eq!(hb.process_role, "session");
4886 assert_eq!(
4887 hb.oldest_tx_label.as_deref(),
4888 Some("session_sweep_walpin_test")
4889 );
4890 assert_eq!(
4891 hb.attribution_basis.as_deref(),
4892 Some("fallback"),
4893 "an Unscoped span observed only through the main view's fallback \
4894 must carry attribution_basis=\"fallback\", never \"origin\""
4895 );
4896
4897 drop(tx_handle);
4898 assert!(
4899 wait_for(Duration::from_secs(2), || !heartbeat_path.exists()).await,
4900 "heartbeat must be removed once the stale span clears"
4901 );
4902
4903 shutdown_tx.send(()).expect("send shutdown signal");
4904 tokio::time::timeout(Duration::from_secs(1), handle)
4905 .await
4906 .expect("session sweep task should exit within 1s")
4907 .expect("session sweep task panicked");
4908 }
4909
4910 #[tokio::test]
4921 #[serial(tx_registry, khive_walpin_sidecar_env)]
4922 async fn session_sweep_fan_out_scopes_secondary_span_to_secondary_sidecar_only() {
4923 let main_dir = tempfile::tempdir().unwrap();
4924 let secondary_dir = tempfile::tempdir().unwrap();
4925 let main_pool = file_pool(&main_dir.path().join("main.db"));
4926 let secondary_pool = file_pool(&secondary_dir.path().join("secondary.db"));
4927 let main_sidecar =
4928 crate::walpin::sidecar_dir_for(main_pool.canonical_path().expect("file-backed"));
4929 let secondary_sidecar =
4930 crate::walpin::sidecar_dir_for(secondary_pool.canonical_path().expect("file-backed"));
4931 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
4932 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
4933
4934 let cfg = SessionSweepConfig {
4935 interval: Duration::from_millis(10),
4936 tx_warn_secs: Duration::from_millis(20),
4937 tx_max_age_secs: Duration::from_millis(500),
4938 };
4939 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4940 let handle = tokio::spawn(run_session_sweep_task(
4941 vec![
4942 SweepBackend {
4943 pool: Arc::clone(&main_pool),
4944 is_main: true,
4945 },
4946 SweepBackend {
4947 pool: Arc::clone(&secondary_pool),
4948 is_main: false,
4949 },
4950 ],
4951 cfg,
4952 shutdown_rx,
4953 ));
4954
4955 let pid = std::process::id();
4956 let secondary_heartbeat = secondary_sidecar.join(format!("{pid}.json"));
4957 let main_heartbeat = main_sidecar.join(format!("{pid}.json"));
4958
4959 let tx_handle = khive_storage::tx_registry::register_scoped(
4960 Some("graph_traverse_read".to_string()),
4961 secondary_pool.origin(),
4962 );
4963 assert!(
4964 wait_for(Duration::from_secs(2), || secondary_heartbeat.exists()).await,
4965 "expected a walpin heartbeat in the secondary backend's own sidecar"
4966 );
4967 assert!(
4968 !main_heartbeat.exists(),
4969 "a span scoped to the secondary backend's origin must never produce \
4970 a heartbeat in the main backend's sidecar"
4971 );
4972
4973 let body = std::fs::read_to_string(&secondary_heartbeat).unwrap();
4974 let hb: crate::walpin::WalpinHeartbeat = serde_json::from_str(&body).unwrap();
4975 assert_eq!(hb.oldest_tx_label.as_deref(), Some("graph_traverse_read"));
4976 assert_eq!(
4977 hb.attribution_basis.as_deref(),
4978 Some("origin"),
4979 "a Secondary-view winner is always Database-origin-backed — never fallback"
4980 );
4981
4982 drop(tx_handle);
4983 assert!(
4984 wait_for(Duration::from_secs(2), || !secondary_heartbeat.exists()).await,
4985 "secondary heartbeat must be removed once its span clears"
4986 );
4987 assert!(
4988 !main_heartbeat.exists(),
4989 "the main sidecar must have stayed untouched for the whole tick sequence"
4990 );
4991
4992 shutdown_tx.send(()).expect("send shutdown signal");
4993 tokio::time::timeout(Duration::from_secs(1), handle)
4994 .await
4995 .expect("session sweep task should exit within 1s")
4996 .expect("session sweep task panicked");
4997 }
4998
4999 #[tokio::test]
5008 #[serial(tx_registry, checkpoint_skip_metrics, khive_walpin_sidecar_env)]
5009 async fn checkpoint_task_ignores_span_registered_against_other_backend_origin_and_unscoped() {
5010 let dir_a = tempfile::tempdir().unwrap();
5011 let dir_b = tempfile::tempdir().unwrap();
5012 let pool_a = file_pool(&dir_a.path().join("backend_a.db"));
5013 let pool_b = file_pool(&dir_b.path().join("backend_b.db"));
5016 let sidecar_a =
5017 crate::walpin::sidecar_dir_for(pool_a.canonical_path().expect("file-backed"));
5018 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
5019 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
5020
5021 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5022 let subscriber = CaptureSubscriber {
5023 events: std::sync::Arc::clone(&buffer),
5024 };
5025 let _tracing_guard = tracing::subscriber::set_default(subscriber);
5026
5027 let _b_origin_handle = khive_storage::tx_registry::register_scoped(
5028 Some("b_origin_span_ignored_by_a".to_string()),
5029 pool_b.origin(),
5030 );
5031 let _unscoped_handle = khive_storage::tx_registry::register(Some(
5032 "unscoped_span_ignored_by_secondary".to_string(),
5033 ));
5034
5035 let cfg = CheckpointConfig {
5036 interval: Duration::from_millis(10),
5037 tx_warn_secs: Duration::from_millis(1),
5038 tx_max_age_secs: Duration::from_millis(1),
5039 ..CheckpointConfig::default()
5040 };
5041 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
5042 let handle = tokio::spawn(run_checkpoint_task(
5043 pool_a,
5044 cfg,
5045 None,
5046 shutdown_rx,
5047 false, ));
5049
5050 tokio::time::sleep(Duration::from_millis(60)).await;
5055 shutdown_tx.send(()).expect("send shutdown signal");
5056 tokio::time::timeout(Duration::from_secs(1), handle)
5057 .await
5058 .expect("checkpoint task should exit within 1s")
5059 .expect("checkpoint task panicked");
5060
5061 let events = buffer.lock().unwrap();
5062 assert!(
5063 events.iter().all(|e| {
5064 e.tx_label.as_deref() != Some("b_origin_span_ignored_by_a")
5065 && e.tx_label.as_deref() != Some("unscoped_span_ignored_by_secondary")
5066 }),
5067 "backend A's Secondary filter must never emit an age alert naming a span \
5068 registered against a different backend's origin or an Unscoped span, got: \
5069 {events:?}"
5070 );
5071 assert!(
5072 !sidecar_a
5073 .join(format!("{}.json", std::process::id()))
5074 .exists(),
5075 "backend A's own sidecar must never gain a heartbeat from a span it does not own"
5076 );
5077 }
5078
5079 #[tokio::test]
5086 #[serial(tx_registry, checkpoint_skip_metrics, khive_walpin_sidecar_env)]
5087 async fn checkpoint_task_detects_and_enumerates_secondary_backend_stall() {
5088 let dir = tempfile::tempdir().unwrap();
5089 let pool = file_pool(&dir.path().join("secondary_stall.db"));
5090 let sidecar_dir =
5091 crate::walpin::sidecar_dir_for(pool.canonical_path().expect("file-backed"));
5092 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
5093 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
5094
5095 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5096 let subscriber = CaptureSubscriber {
5097 events: std::sync::Arc::clone(&buffer),
5098 };
5099 let _tracing_guard = tracing::subscriber::set_default(subscriber);
5100
5101 let tx_handle = khive_storage::tx_registry::register_scoped(
5102 Some("secondary_stall_test".to_string()),
5103 pool.origin(),
5104 );
5105
5106 let cfg = CheckpointConfig {
5107 interval: Duration::from_millis(10),
5108 tx_warn_secs: Duration::from_millis(5),
5109 tx_max_age_secs: Duration::from_millis(500),
5110 ..CheckpointConfig::default()
5111 };
5112 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
5113 let pid = std::process::id();
5114 let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
5115 let handle = tokio::spawn(run_checkpoint_task(
5116 pool,
5117 cfg,
5118 None,
5119 shutdown_rx,
5120 false, ));
5122
5123 assert!(
5124 wait_for(Duration::from_secs(2), || heartbeat_path.exists()).await,
5125 "expected a walpin heartbeat once the secondary backend's own span crossed \
5126 tx_warn_secs"
5127 );
5128 let body = std::fs::read_to_string(&heartbeat_path).unwrap();
5129 let hb: crate::walpin::WalpinHeartbeat = serde_json::from_str(&body).unwrap();
5130 assert_eq!(hb.oldest_tx_label.as_deref(), Some("secondary_stall_test"));
5131 assert_eq!(
5132 hb.attribution_basis.as_deref(),
5133 Some("origin"),
5134 "a Secondary-view winner is always Database-origin-backed — never fallback"
5135 );
5136 assert!(
5137 hb.oldest_tx_age_secs > 0.0,
5138 "the heartbeat must reflect a nonzero stale age for the secondary backend's own \
5139 span, got {hb:?}"
5140 );
5141
5142 shutdown_tx.send(()).expect("send shutdown signal");
5143 tokio::time::timeout(Duration::from_secs(1), handle)
5144 .await
5145 .expect("checkpoint task should exit within 1s")
5146 .expect("checkpoint task panicked");
5147
5148 drop(tx_handle);
5149
5150 let events = buffer.lock().unwrap();
5151 assert!(
5152 events.iter().any(|e| {
5153 e.tx_label.as_deref() == Some("secondary_stall_test")
5154 && e.message
5155 .as_deref()
5156 .is_some_and(|m| m.contains("ADR-091 Plank 1"))
5157 }),
5158 "expected the secondary backend's own checkpoint task to emit a Plank 1 age alert \
5159 for its own stalled span, got: {events:?}"
5160 );
5161 }
5162
5163 #[test]
5164 fn wal_pin_depth_arithmetic_against_real_connection() {
5165 let dir = tempfile::tempdir().unwrap();
5166 let path = dir.path().join("pin_depth.db");
5167 let pool = file_pool(&path);
5168 let writer = pool.try_writer().expect("acquire writer");
5169 let conn = writer.conn();
5170
5171 conn.execute_batch("CREATE TABLE t (v INTEGER)").unwrap();
5172 conn.execute_batch("INSERT INTO t (v) VALUES (1)").unwrap();
5173
5174 let (log, checkpointed) =
5175 query_wal_pin_depth(conn).expect("PRAGMA wal_checkpoint(PASSIVE) must succeed");
5176 assert!(
5180 log >= checkpointed,
5181 "checkpointed frames cannot exceed log frames"
5182 );
5183 assert_eq!(
5184 log - checkpointed,
5185 0,
5186 "an unpinned WAL must fully checkpoint under PASSIVE"
5187 );
5188 }
5189
5190 #[test]
5191 fn wal_pin_depth_arithmetic_on_in_memory_pool_errors_cleanly() {
5192 let cfg = PoolConfig {
5196 path: None,
5197 ..PoolConfig::default()
5198 };
5199 let pool = ConnectionPool::new(cfg).expect("in-memory pool");
5200 let writer = pool.try_writer().expect("acquire writer");
5201 let _ = query_wal_pin_depth(writer.conn());
5204 }
5205}