1use std::collections::HashMap;
40use std::path::{Path, PathBuf};
41use std::sync::atomic::{AtomicU64, Ordering};
42use std::sync::{Arc, Mutex, OnceLock};
43use std::time::{Duration, Instant};
44
45use crate::pool::ConnectionPool;
46
47static LAST_WAL_PAGES: AtomicU64 = AtomicU64::new(u64::MAX);
56
57static TRUNCATE_ATTEMPTS: AtomicU64 = AtomicU64::new(0);
60
61static TRUNCATE_CONSECUTIVE_FAILURES: AtomicU64 = AtomicU64::new(0);
65
66static CHECKPOINT_SKIPPED_TICKS: AtomicU64 = AtomicU64::new(0);
72
73static CHECKPOINT_CONSECUTIVE_SKIPS: AtomicU64 = AtomicU64::new(0);
78
79static CHECKPOINT_LAST_SKIP_WAL_PAGES: AtomicU64 = AtomicU64::new(u64::MAX);
83
84static CHECKPOINT_PRESSURE_ELEVATED_TICKS: AtomicU64 = AtomicU64::new(0);
87
88static CHECKPOINT_PRESSURE_EPISODES_STARTED: AtomicU64 = AtomicU64::new(0);
90
91static CHECKPOINT_PRESSURE_EPISODES_RECOVERED: AtomicU64 = AtomicU64::new(0);
93
94static CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS: AtomicU64 = AtomicU64::new(0);
96
97static CHECKPOINT_LIFECYCLE_APPEND_FAILURES: AtomicU64 = AtomicU64::new(0);
99
100static CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS: AtomicU64 = AtomicU64::new(0);
103
104static READ_TX_MAX_AGE_EVICTIONS: AtomicU64 = AtomicU64::new(0);
110
111#[derive(Debug, Clone, PartialEq, Eq)]
116pub struct RoutineWalObservation {
117 pub busy: i64,
118 pub log_frames: u64,
119 pub checkpointed_frames: u64,
120 pub pending_frames: u64,
121 pub physical_wal_bytes: Option<u64>,
122 pub observed_at_unix_ms: u64,
123}
124
125static ROUTINE_WAL_OBSERVATIONS: OnceLock<Mutex<HashMap<Option<PathBuf>, RoutineWalObservation>>> =
129 OnceLock::new();
130
131fn routine_wal_observations() -> &'static Mutex<HashMap<Option<PathBuf>, RoutineWalObservation>> {
132 ROUTINE_WAL_OBSERVATIONS.get_or_init(|| Mutex::new(HashMap::new()))
133}
134
135fn checkpoint_db_key_from_path(path: Option<&Path>) -> Option<PathBuf> {
136 path.map(Path::to_path_buf)
137}
138
139fn checkpoint_db_key(pool: &ConnectionPool) -> Option<PathBuf> {
140 checkpoint_db_key_from_path(pool.canonical_path())
141}
142
143fn observed_at_unix_ms() -> u64 {
144 std::time::SystemTime::now()
145 .duration_since(std::time::UNIX_EPOCH)
146 .map(|duration| duration.as_millis() as u64)
147 .unwrap_or(0)
148}
149
150fn physical_wal_bytes(pool: &ConnectionPool) -> Option<u64> {
151 let path = pool.canonical_path()?;
152 let mut sidecar = path.as_os_str().to_os_string();
153 sidecar.push("-wal");
154 std::fs::metadata(PathBuf::from(sidecar))
155 .ok()
156 .map(|metadata| metadata.len())
157}
158
159fn record_routine_wal_observation(
160 pool: &ConnectionPool,
161 raw: RawCheckpointObservation,
162) -> RoutineWalObservation {
163 let log_frames = raw.log_frames.max(0) as u64;
164 let checkpointed_frames = raw.checkpointed_frames.max(0) as u64;
165 let observation = RoutineWalObservation {
166 busy: raw.busy,
167 log_frames,
168 checkpointed_frames,
169 pending_frames: log_frames.saturating_sub(checkpointed_frames),
170 physical_wal_bytes: physical_wal_bytes(pool),
171 observed_at_unix_ms: observed_at_unix_ms(),
172 };
173 routine_wal_observations()
174 .lock()
175 .unwrap_or_else(std::sync::PoisonError::into_inner)
176 .insert(checkpoint_db_key(pool), observation.clone());
177 observation
178}
179
180pub fn routine_wal_observation(pool: &ConnectionPool) -> Option<RoutineWalObservation> {
183 routine_wal_observations()
184 .lock()
185 .unwrap_or_else(std::sync::PoisonError::into_inner)
186 .get(&checkpoint_db_key(pool))
187 .cloned()
188}
189
190pub fn last_observed_wal_pages() -> Option<u64> {
193 match LAST_WAL_PAGES.load(Ordering::Relaxed) {
194 u64::MAX => None,
195 pages => Some(pages),
196 }
197}
198
199pub fn truncate_attempts() -> u64 {
201 TRUNCATE_ATTEMPTS.load(Ordering::Relaxed)
202}
203
204pub fn truncate_consecutive_failures() -> u64 {
206 TRUNCATE_CONSECUTIVE_FAILURES.load(Ordering::Relaxed)
207}
208
209pub fn checkpoint_skipped_ticks() -> u64 {
212 CHECKPOINT_SKIPPED_TICKS.load(Ordering::Relaxed)
213}
214
215pub fn checkpoint_consecutive_skips() -> u64 {
217 CHECKPOINT_CONSECUTIVE_SKIPS.load(Ordering::Relaxed)
218}
219
220pub fn checkpoint_last_skip_wal_pages() -> Option<u64> {
223 match CHECKPOINT_LAST_SKIP_WAL_PAGES.load(Ordering::Relaxed) {
224 u64::MAX => None,
225 pages => Some(pages),
226 }
227}
228
229pub fn checkpoint_pressure_elevated_ticks() -> u64 {
231 CHECKPOINT_PRESSURE_ELEVATED_TICKS.load(Ordering::Relaxed)
232}
233
234pub fn checkpoint_pressure_episodes_started() -> u64 {
236 CHECKPOINT_PRESSURE_EPISODES_STARTED.load(Ordering::Relaxed)
237}
238
239pub fn checkpoint_pressure_episodes_recovered() -> u64 {
241 CHECKPOINT_PRESSURE_EPISODES_RECOVERED.load(Ordering::Relaxed)
242}
243
244pub fn checkpoint_lifecycle_append_attempts() -> u64 {
246 CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.load(Ordering::Relaxed)
247}
248
249pub fn checkpoint_lifecycle_append_failures() -> u64 {
251 CHECKPOINT_LIFECYCLE_APPEND_FAILURES.load(Ordering::Relaxed)
252}
253
254pub fn read_tx_max_age_evictions() -> u64 {
257 READ_TX_MAX_AGE_EVICTIONS.load(Ordering::Relaxed)
258}
259
260pub(crate) fn note_read_tx_max_age_eviction() {
266 READ_TX_MAX_AGE_EVICTIONS.fetch_add(1, Ordering::Relaxed);
267}
268
269pub fn checkpoint_lifecycle_enqueue_drops() -> u64 {
271 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.load(Ordering::Relaxed)
272}
273
274fn note_checkpoint_skipped() {
279 CHECKPOINT_SKIPPED_TICKS.fetch_add(1, Ordering::Relaxed);
280 CHECKPOINT_CONSECUTIVE_SKIPS.fetch_add(1, Ordering::Relaxed);
281 if let Some(pages) = last_observed_wal_pages() {
282 CHECKPOINT_LAST_SKIP_WAL_PAGES.store(pages, Ordering::Relaxed);
283 }
284}
285
286fn note_checkpoint_observed(_wal_pages: u64) {
291 CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
292}
293
294fn note_checkpoint_pressure_observation(above_warn: bool, was_above_warn: bool) {
295 if above_warn {
296 CHECKPOINT_PRESSURE_ELEVATED_TICKS.fetch_add(1, Ordering::Relaxed);
297 if !was_above_warn {
298 CHECKPOINT_PRESSURE_EPISODES_STARTED.fetch_add(1, Ordering::Relaxed);
299 }
300 } else if was_above_warn {
301 CHECKPOINT_PRESSURE_EPISODES_RECOVERED.fetch_add(1, Ordering::Relaxed);
302 }
303}
304
305#[cfg(test)]
309pub(crate) fn reset_checkpoint_metrics_for_tests() {
310 CHECKPOINT_SKIPPED_TICKS.store(0, Ordering::Relaxed);
311 CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
312 CHECKPOINT_LAST_SKIP_WAL_PAGES.store(u64::MAX, Ordering::Relaxed);
313 CHECKPOINT_PRESSURE_ELEVATED_TICKS.store(0, Ordering::Relaxed);
314 CHECKPOINT_PRESSURE_EPISODES_STARTED.store(0, Ordering::Relaxed);
315 CHECKPOINT_PRESSURE_EPISODES_RECOVERED.store(0, Ordering::Relaxed);
316 CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.store(0, Ordering::Relaxed);
317 CHECKPOINT_LIFECYCLE_APPEND_FAILURES.store(0, Ordering::Relaxed);
318 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.store(0, Ordering::Relaxed);
319 READ_TX_MAX_AGE_EVICTIONS.store(0, Ordering::Relaxed);
320}
321
322#[derive(Debug, Clone, Copy, PartialEq, Eq)]
333pub enum CheckpointTick {
334 Skipped,
337 Observed(u64),
339}
340
341pub const DEFAULT_WARN_SUSTAINED_CYCLES: u8 = 3;
344
345#[derive(Clone, Debug)]
350pub struct CheckpointConfig {
351 pub interval: Duration,
356
357 pub warn_pages: u64,
362
363 pub warn_sustained_cycles: u8,
371
372 pub high_water_pages: u64,
382
383 pub truncate_high_water_pages: u64,
394
395 pub truncate_min_interval: Duration,
406
407 pub truncate_busy_timeout: Duration,
414
415 pub tx_warn_secs: Duration,
423
424 pub tx_max_age_secs: Duration,
439}
440
441impl Default for CheckpointConfig {
442 fn default() -> Self {
443 Self {
444 interval: Duration::from_millis(500),
445 warn_pages: 2000,
446 warn_sustained_cycles: DEFAULT_WARN_SUSTAINED_CYCLES,
447 high_water_pages: 6000,
448 truncate_high_water_pages: 20_000,
449 truncate_min_interval: Duration::from_secs(300),
450 truncate_busy_timeout: Duration::from_millis(2000),
451 tx_warn_secs: Duration::from_secs(30),
452 tx_max_age_secs: Duration::from_secs(120),
453 }
454 }
455}
456
457impl CheckpointConfig {
458 pub fn from_env() -> Self {
462 let mut cfg = Self::default();
463
464 if let Ok(ms) = std::env::var("KHIVE_CHECKPOINT_INTERVAL_MS") {
465 if let Ok(v) = ms.parse::<u64>() {
466 if v > 0 {
467 cfg.interval = Duration::from_millis(v);
468 }
469 }
470 }
471
472 if let Ok(v) = std::env::var("KHIVE_WAL_WARN_PAGES") {
473 if let Ok(n) = v.parse::<u64>() {
474 if n > 0 {
475 cfg.warn_pages = n;
476 }
477 }
478 }
479
480 if let Ok(v) = std::env::var("KHIVE_WAL_WARN_SUSTAINED_CYCLES") {
481 if let Ok(n) = v.parse::<u8>() {
482 if n > 0 {
483 cfg.warn_sustained_cycles = n;
484 }
485 }
486 }
487
488 if let Ok(v) = std::env::var("KHIVE_WAL_HIGH_WATER_PAGES") {
489 if let Ok(n) = v.parse::<u64>() {
490 if n > 0 {
491 cfg.high_water_pages = n;
492 }
493 }
494 }
495
496 if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES") {
497 if let Ok(n) = v.parse::<u64>() {
498 if n > 0 {
499 cfg.truncate_high_water_pages = n;
500 }
501 }
502 }
503
504 if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS") {
505 if let Ok(n) = v.parse::<u64>() {
506 if n > 0 {
507 cfg.truncate_min_interval = Duration::from_secs(n);
508 }
509 }
510 }
511
512 if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_BUSY_MS") {
513 if let Ok(n) = v.parse::<u64>() {
514 if n > 0 {
515 cfg.truncate_busy_timeout = Duration::from_millis(n);
516 }
517 }
518 }
519
520 (cfg.tx_warn_secs, cfg.tx_max_age_secs) =
521 tx_age_thresholds_from_env(cfg.tx_warn_secs, cfg.tx_max_age_secs);
522
523 cfg
524 }
525}
526
527pub(crate) fn tx_age_thresholds_from_env(
541 default_warn: Duration,
542 default_max: Duration,
543) -> (Duration, Duration) {
544 let mut warn_secs = default_warn;
545 let mut max_age_secs = default_max;
546
547 if let Ok(v) = std::env::var("KHIVE_TX_WARN_SECS") {
548 if let Ok(n) = v.parse::<u64>() {
549 if n > 0 {
550 warn_secs = Duration::from_secs(n);
551 }
552 }
553 }
554
555 if let Ok(v) = std::env::var("KHIVE_TX_MAX_AGE_SECS") {
556 if let Ok(n) = v.parse::<u64>() {
557 if n > 0 {
558 max_age_secs = Duration::from_secs(n);
559 }
560 }
561 }
562
563 if warn_secs >= max_age_secs {
564 tracing::warn!(
565 configured_tx_warn_secs = warn_secs.as_secs_f64(),
566 configured_tx_max_age_secs = max_age_secs.as_secs_f64(),
567 fallback_tx_warn_secs = default_warn.as_secs_f64(),
568 fallback_tx_max_age_secs = default_max.as_secs_f64(),
569 "KHIVE_TX_WARN_SECS must be strictly less than KHIVE_TX_MAX_AGE_SECS; \
570 both transaction-age thresholds were rejected and reset to their defaults"
571 );
572 return (default_warn, default_max);
573 }
574
575 (warn_secs, max_age_secs)
576}
577
578#[cfg(unix)]
579const DEFAULT_WALPIN_FULL_SCAN_INTERVAL: Duration = Duration::from_secs(30);
580
581#[cfg(unix)]
582#[derive(Debug, Clone)]
583struct CachedWalpinAttribution {
584 report: crate::walpin::WalpinReport,
585 census: Result<crate::walpin::CensusResult, String>,
586 captured_at: Instant,
587}
588
589#[cfg(unix)]
590#[derive(Debug)]
591enum WalpinFullScanPlan {
592 Refresh {
593 previous_last_attempt: Option<Instant>,
594 },
595 Cached(CachedWalpinAttribution),
596 Suppressed,
597}
598
599#[derive(Debug)]
606pub struct TruncateState {
607 last_attempt: Option<Instant>,
611 consecutive_failures: u32,
616 #[cfg(unix)]
621 legacy_walpin_fallback_interval: Duration,
622 #[cfg(unix)]
627 walpin_full_scan_interval: Duration,
628 #[cfg(unix)]
629 walpin_full_scan_last_attempt: Option<Instant>,
630 #[cfg(unix)]
631 walpin_cached_attribution: Option<CachedWalpinAttribution>,
632 #[cfg(unix)]
635 sidecar_attribution_attempted_this_tick: bool,
636}
637
638impl Default for TruncateState {
639 fn default() -> Self {
640 Self {
641 last_attempt: None,
642 consecutive_failures: 0,
643 #[cfg(unix)]
644 legacy_walpin_fallback_interval: DEFAULT_SESSION_SWEEP_INTERVAL,
645 #[cfg(unix)]
646 walpin_full_scan_interval: DEFAULT_WALPIN_FULL_SCAN_INTERVAL,
647 #[cfg(unix)]
648 walpin_full_scan_last_attempt: None,
649 #[cfg(unix)]
650 walpin_cached_attribution: None,
651 #[cfg(unix)]
652 sidecar_attribution_attempted_this_tick: false,
653 }
654 }
655}
656
657impl TruncateState {
658 #[cfg(unix)]
659 fn with_legacy_walpin_fallback(interval: Duration) -> Self {
660 Self {
661 legacy_walpin_fallback_interval: interval,
662 ..Self::default()
663 }
664 }
665
666 #[cfg(all(test, unix))]
667 fn with_walpin_full_scan_cadence(interval: Duration) -> Self {
668 Self {
669 walpin_full_scan_interval: interval,
670 ..Self::default()
671 }
672 }
673
674 #[cfg(unix)]
675 fn begin_tick(&mut self) {
676 self.sidecar_attribution_attempted_this_tick = false;
677 }
678
679 #[cfg(unix)]
680 fn housekeeping_due(&self) -> bool {
681 !self.sidecar_attribution_attempted_this_tick
682 && self.walpin_full_scan_due_at(Instant::now())
683 }
684
685 #[cfg(unix)]
686 fn walpin_full_scan_due_at(&self, now: Instant) -> bool {
687 self.walpin_full_scan_last_attempt.is_none_or(|last| {
688 now.saturating_duration_since(last) >= self.walpin_full_scan_interval
689 })
690 }
691
692 #[cfg(unix)]
693 fn claim_walpin_full_scan_at(&mut self, now: Instant) -> bool {
694 if !self.walpin_full_scan_due_at(now) {
695 return false;
696 }
697 self.walpin_full_scan_last_attempt = Some(now);
698 true
699 }
700
701 #[cfg(unix)]
702 fn plan_walpin_attribution_at(&mut self, now: Instant) -> WalpinFullScanPlan {
703 if self.walpin_full_scan_due_at(now) {
704 let previous_last_attempt = self.walpin_full_scan_last_attempt.replace(now);
705 WalpinFullScanPlan::Refresh {
706 previous_last_attempt,
707 }
708 } else if let Some(cached) = self.walpin_cached_attribution.clone() {
709 WalpinFullScanPlan::Cached(cached)
710 } else {
711 WalpinFullScanPlan::Suppressed
712 }
713 }
714
715 #[cfg(unix)]
716 fn restore_walpin_full_scan_reservation(&mut self, previous_last_attempt: Option<Instant>) {
717 self.walpin_full_scan_last_attempt = previous_last_attempt;
718 }
719
720 #[cfg(unix)]
721 fn cache_walpin_attribution(
722 &mut self,
723 report: crate::walpin::WalpinReport,
724 census: Result<crate::walpin::CensusResult, String>,
725 captured_at: Instant,
726 ) {
727 self.walpin_cached_attribution = Some(CachedWalpinAttribution {
728 report,
729 census,
730 captured_at,
731 });
732 }
733}
734
735#[derive(Debug, Clone, Copy, PartialEq, Eq)]
742pub enum CheckpointSeverityRung {
743 Info,
745 Warn,
748 Alarm,
751}
752
753#[derive(Debug, Default, Clone)]
757pub struct CheckpointSeverityState {
758 was_above_warn: bool,
761 consecutive_above_warn: u8,
764 warn_emitted_for_episode: bool,
768}
769
770#[derive(Debug, Clone, Copy, PartialEq, Eq)]
773pub struct CheckpointSeverityEmission {
774 pub rung: CheckpointSeverityRung,
778 pub wal_pages: u64,
780 pub threshold_pages: u64,
782 pub consecutive_cycles: u8,
785}
786
787impl CheckpointSeverityState {
788 pub fn observe_wal_pages(
799 &mut self,
800 wal_pages: u64,
801 config: &CheckpointConfig,
802 ) -> Vec<CheckpointSeverityEmission> {
803 let mut emissions = Vec::new();
804 let above_warn = wal_pages >= config.warn_pages;
805
806 if above_warn {
807 self.consecutive_above_warn = self.consecutive_above_warn.saturating_add(1);
808
809 if !self.was_above_warn {
810 emissions.push(CheckpointSeverityEmission {
811 rung: CheckpointSeverityRung::Info,
812 wal_pages,
813 threshold_pages: config.warn_pages,
814 consecutive_cycles: self.consecutive_above_warn,
815 });
816 }
817
818 if !self.warn_emitted_for_episode
819 && self.consecutive_above_warn >= config.warn_sustained_cycles
820 {
821 emissions.push(CheckpointSeverityEmission {
822 rung: CheckpointSeverityRung::Warn,
823 wal_pages,
824 threshold_pages: config.warn_pages,
825 consecutive_cycles: self.consecutive_above_warn,
826 });
827 self.warn_emitted_for_episode = true;
828 }
829 } else {
830 self.consecutive_above_warn = 0;
831 self.warn_emitted_for_episode = false;
832 }
833
834 self.was_above_warn = above_warn;
835 emissions
836 }
837}
838
839#[derive(Debug, Clone, Copy, PartialEq, Eq)]
843pub enum TxAgeRung {
844 Warn,
846 Stale,
851}
852
853#[derive(Debug, Clone, PartialEq, Eq)]
855pub struct TxAgeEmission {
856 pub rung: TxAgeRung,
857 pub age: Duration,
858 pub label: Option<String>,
859}
860
861#[derive(Debug, Default, Clone)]
872pub struct TxAgeSweepState {
873 was_above_warn: bool,
876 was_above_max_age: bool,
879 tracked_id: Option<khive_storage::tx_registry::TxId>,
884}
885
886impl TxAgeSweepState {
887 pub fn observe(
899 &mut self,
900 oldest: Option<(khive_storage::tx_registry::TxId, Duration, Option<String>)>,
901 tx_warn_secs: Duration,
902 tx_max_age_secs: Duration,
903 ) -> Vec<TxAgeEmission> {
904 let mut emissions = Vec::new();
905
906 let Some((id, age, label)) = oldest else {
907 self.was_above_warn = false;
908 self.was_above_max_age = false;
909 self.tracked_id = None;
910 return emissions;
911 };
912
913 if self.tracked_id != Some(id) {
914 self.was_above_warn = false;
915 self.was_above_max_age = false;
916 }
917 self.tracked_id = Some(id);
918
919 let above_warn = age >= tx_warn_secs;
920 let above_max_age = age >= tx_max_age_secs;
921
922 if above_warn && !self.was_above_warn {
923 emissions.push(TxAgeEmission {
924 rung: TxAgeRung::Warn,
925 age,
926 label: label.clone(),
927 });
928 }
929 if above_max_age && !self.was_above_max_age {
930 emissions.push(TxAgeEmission {
931 rung: TxAgeRung::Stale,
932 age,
933 label,
934 });
935 }
936
937 self.was_above_warn = above_warn;
938 self.was_above_max_age = above_max_age;
939 emissions
940 }
941}
942
943fn log_tx_age_emission(emission: &TxAgeEmission) {
948 let label = emission.label.as_deref().unwrap_or("<unlabeled>");
949 match emission.rung {
950 TxAgeRung::Warn => {
951 tracing::warn!(
952 tx_age_secs = emission.age.as_secs_f64(),
953 tx_label = label,
954 "ADR-091 Plank 1: open transaction registry entry exceeded soft-cap age"
955 );
956 }
957 TxAgeRung::Stale => {
958 tracing::error!(
959 tx_age_secs = emission.age.as_secs_f64(),
960 tx_label = label,
961 "ADR-091 Plank 1: open transaction registry entry exceeded the cooperative \
962 stale-op cap; no in-process mechanism can force-close it — investigate the \
963 labeled caller directly"
964 );
965 }
966 }
967}
968
969struct WalpinSidecarState {
977 dir: PathBuf,
978 pid: u32,
979 role: &'static str,
980 started_at: i64,
981 sweep_interval_ms: u64,
986 wrote: bool,
987 beacon_registered: bool,
993 last_heartbeat: Option<LastHeartbeatState>,
1000}
1001
1002struct LastHeartbeatState {
1008 span_id: khive_storage::tx_registry::TxId,
1009 label: Option<String>,
1010 attribution_basis: &'static str,
1011 sweep_interval_ms: u64,
1012 oldest_tx_started_at: i64,
1013}
1014
1015impl LastHeartbeatState {
1016 fn content_matches(
1022 &self,
1023 span_id: khive_storage::tx_registry::TxId,
1024 label: &Option<String>,
1025 attribution_basis: &str,
1026 sweep_interval_ms: u64,
1027 ) -> bool {
1028 self.span_id == span_id
1029 && self.label == *label
1030 && self.attribution_basis == attribution_basis
1031 && self.sweep_interval_ms == sweep_interval_ms
1032 }
1033}
1034
1035impl WalpinSidecarState {
1036 fn new(
1039 db_path: Option<&Path>,
1040 is_file_backed: bool,
1041 role: &'static str,
1042 interval: Duration,
1043 ) -> Option<Self> {
1044 let path = db_path?;
1045 if !crate::walpin::sidecar_enabled(is_file_backed) {
1046 return None;
1047 }
1048 let pid = std::process::id();
1049 Some(Self {
1050 dir: crate::walpin::sidecar_dir_for(path),
1051 pid,
1052 role,
1053 started_at: crate::walpin::process_start_time_secs(pid).unwrap_or(0),
1054 sweep_interval_ms: interval.as_millis().min(u64::MAX as u128) as u64,
1055 wrote: false,
1056 last_heartbeat: None,
1057 beacon_registered: false,
1058 })
1059 }
1060
1061 async fn register_beacon(&mut self) {
1070 let dir = self.dir.clone();
1071 let beacon = crate::walpin::WalpinBeacon {
1072 pid: self.pid,
1073 process_role: self.role.to_string(),
1074 started_at: self.started_at,
1075 sweep_interval_ms: self.sweep_interval_ms,
1076 };
1077 let result =
1078 tokio::task::spawn_blocking(move || crate::walpin::write_beacon(&dir, &beacon)).await;
1079 match result {
1080 Ok(Ok(())) => {
1081 self.beacon_registered = true;
1082 }
1083 Ok(Err(e)) => {
1084 tracing::warn!(
1085 error = %e,
1086 "ADR-091 Amendment 2: failed to write walpin registration beacon; \
1087 this process's sidecar health will read as unknown, not registered-silent"
1088 );
1089 }
1090 Err(join_err) => {
1091 tracing::warn!(
1092 error = %join_err,
1093 "ADR-091 Amendment 2: walpin beacon write task panicked"
1094 );
1095 }
1096 }
1097 }
1098
1099 #[cfg(unix)]
1105 async fn reap_dead_entries_bounded(
1106 &self,
1107 legacy_fallback_interval: Duration,
1108 ) -> Option<crate::walpin::WalpinReport> {
1109 let dir = self.dir.clone();
1110 let result = tokio::task::spawn_blocking(move || {
1111 crate::walpin::housekeep_live(&dir, legacy_fallback_interval)
1112 })
1113 .await;
1114 match result {
1115 Ok(Ok(report)) => Some(report),
1116 Ok(Err(e)) => {
1117 tracing::warn!(
1118 error = %e,
1119 "ADR-091 Amendment 6: bounded walpin sidecar cleanup failed"
1120 );
1121 None
1122 }
1123 Err(join_err) => {
1124 tracing::warn!(
1125 error = %join_err,
1126 "ADR-091 Amendment 6: walpin sidecar cleanup task panicked"
1127 );
1128 None
1129 }
1130 }
1131 }
1132
1133 async fn refresh_beacon(&mut self) {
1143 if !self.beacon_registered {
1144 self.register_beacon().await;
1145 return;
1146 }
1147 let dir = self.dir.clone();
1148 let pid = self.pid;
1149 let result =
1150 tokio::task::spawn_blocking(move || crate::walpin::touch_beacon(&dir, pid)).await;
1151 match result {
1152 Ok(Ok(())) => {}
1153 Ok(Err(e)) => {
1154 self.beacon_registered = false;
1155 tracing::warn!(
1156 error = %e,
1157 "ADR-091 Amendment 2: failed to refresh walpin registration beacon; \
1158 this process's sidecar health will read as unknown, not registered-silent"
1159 );
1160 }
1161 Err(join_err) => {
1162 self.beacon_registered = false;
1163 tracing::warn!(
1164 error = %join_err,
1165 "ADR-091 Amendment 2: walpin beacon refresh task panicked"
1166 );
1167 }
1168 }
1169 }
1170
1171 async fn drop_beacon_fail_closed(&mut self) {
1182 let dir = self.dir.clone();
1183 let pid = self.pid;
1184 self.beacon_registered = false;
1185 let result =
1186 tokio::task::spawn_blocking(move || crate::walpin::remove_beacon(&dir, pid)).await;
1187 match result {
1188 Ok(Ok(())) => {}
1189 Ok(Err(e)) => {
1190 tracing::warn!(
1191 error = %e,
1192 "ADR-091 Amendment 2: failed to remove walpin beacon after a failed \
1193 heartbeat write; beacon will age out of the freshness window instead"
1194 );
1195 }
1196 Err(join_err) => {
1197 tracing::warn!(
1198 error = %join_err,
1199 "ADR-091 Amendment 2: walpin beacon removal task panicked"
1200 );
1201 }
1202 }
1203 }
1204
1205 async fn observe(
1209 &mut self,
1210 oldest: Option<khive_storage::tx_registry::OldestSpan>,
1211 tx_warn_secs: Duration,
1212 ) {
1213 match oldest {
1214 Some(span) if span.age >= tx_warn_secs => {
1215 let attribution_basis = match span.origin {
1222 khive_storage::tx_registry::TxOrigin::Database(_) => "origin",
1223 khive_storage::tx_registry::TxOrigin::Unscoped
1224 | khive_storage::tx_registry::TxOrigin::Memory => "fallback",
1225 };
1226
1227 let content_unchanged = self.wrote
1233 && self.last_heartbeat.as_ref().is_some_and(|last| {
1234 last.content_matches(
1235 span.id,
1236 &span.label,
1237 attribution_basis,
1238 self.sweep_interval_ms,
1239 )
1240 });
1241
1242 if content_unchanged {
1243 let dir = self.dir.clone();
1244 let pid = self.pid;
1245 let touch_result = tokio::task::spawn_blocking(move || {
1246 crate::walpin::touch_heartbeat(&dir, pid)
1247 })
1248 .await;
1249 match touch_result {
1250 Ok(Ok(())) => {
1251 self.refresh_beacon().await;
1252 return;
1253 }
1254 Ok(Err(e)) => {
1255 tracing::warn!(
1256 error = %e,
1257 "ADR-091 Amendment 3 Plank F1: walpin heartbeat touch failed; \
1258 recreating with a full body write"
1259 );
1260 }
1261 Err(join_err) => {
1262 tracing::warn!(
1263 error = %join_err,
1264 "ADR-091 Amendment 3 Plank F1: walpin heartbeat touch task \
1265 panicked; recreating with a full body write"
1266 );
1267 }
1268 }
1269 }
1274
1275 let oldest_tx_started_at = self
1282 .last_heartbeat
1283 .as_ref()
1284 .filter(|last| last.span_id == span.id)
1285 .map(|last| last.oldest_tx_started_at)
1286 .unwrap_or_else(|| now_epoch_secs().saturating_sub(span.age.as_secs() as i64));
1287
1288 let heartbeat = crate::walpin::WalpinHeartbeat {
1289 pid: self.pid,
1290 process_role: self.role.to_string(),
1291 started_at: self.started_at,
1292 oldest_tx_age_secs: span.age.as_secs_f64(),
1293 oldest_tx_label: span.label.clone(),
1294 oldest_tx_started_at: Some(oldest_tx_started_at),
1295 updated_at: now_epoch_secs(),
1296 sweep_interval_ms: self.sweep_interval_ms,
1297 attribution_basis: Some(attribution_basis.to_string()),
1298 };
1299 let dir = self.dir.clone();
1300 let result = tokio::task::spawn_blocking(move || {
1301 crate::walpin::write_heartbeat(&dir, &heartbeat)
1302 })
1303 .await;
1304 match result {
1314 Ok(Ok(())) => {
1315 self.wrote = true;
1316 self.last_heartbeat = Some(LastHeartbeatState {
1317 span_id: span.id,
1318 label: span.label,
1319 attribution_basis,
1320 sweep_interval_ms: self.sweep_interval_ms,
1321 oldest_tx_started_at,
1322 });
1323 self.refresh_beacon().await;
1324 }
1325 Ok(Err(e)) => {
1326 tracing::warn!(
1327 error = %e,
1328 "ADR-091 Amendment 2 Plank B: failed to write walpin heartbeat; \
1329 removing beacon so this process cannot read as \
1330 registered-silent while over threshold"
1331 );
1332 self.last_heartbeat = None;
1336 self.drop_beacon_fail_closed().await;
1337 }
1338 Err(join_err) => {
1339 tracing::warn!(
1340 error = %join_err,
1341 "ADR-091 Amendment 2 Plank B: walpin heartbeat write task panicked"
1342 );
1343 self.last_heartbeat = None;
1344 self.drop_beacon_fail_closed().await;
1345 }
1346 }
1347 }
1348 _ => {
1349 self.refresh_beacon().await;
1350 if self.wrote {
1351 let dir = self.dir.clone();
1352 let pid = self.pid;
1353 let result = tokio::task::spawn_blocking(move || {
1354 crate::walpin::remove_heartbeat(&dir, pid)
1355 })
1356 .await;
1357 match result {
1358 Ok(Ok(())) => {}
1359 Ok(Err(e)) => tracing::warn!(
1360 error = %e,
1361 "ADR-091 Amendment 2 Plank B: failed to remove walpin heartbeat"
1362 ),
1363 Err(join_err) => tracing::warn!(
1364 error = %join_err,
1365 "ADR-091 Amendment 2 Plank B: walpin heartbeat removal task panicked"
1366 ),
1367 }
1368 self.wrote = false;
1369 self.last_heartbeat = None;
1370 }
1371 }
1372 }
1373 }
1374
1375 async fn shutdown(&mut self) {
1376 if self.wrote {
1377 let dir = self.dir.clone();
1378 let pid = self.pid;
1379 let _ = tokio::task::spawn_blocking(move || crate::walpin::remove_heartbeat(&dir, pid))
1380 .await;
1381 self.wrote = false;
1382 }
1383 }
1384}
1385
1386#[cfg(unix)]
1387async fn run_walpin_housekeeping_if_due(
1388 sidecar: &WalpinSidecarState,
1389 state: &mut TruncateState,
1390 legacy_fallback_interval: Duration,
1391) -> bool {
1392 if !state.housekeeping_due() || !state.claim_walpin_full_scan_at(Instant::now()) {
1393 return false;
1394 }
1395 if let Some(report) = sidecar
1396 .reap_dead_entries_bounded(legacy_fallback_interval)
1397 .await
1398 {
1399 state.cache_walpin_attribution(
1400 report,
1401 Err("OS holder census is unavailable for a housekeeping-only scan".to_string()),
1402 Instant::now(),
1403 );
1404 }
1405 true
1406}
1407
1408fn now_epoch_secs() -> i64 {
1409 std::time::SystemTime::now()
1410 .duration_since(std::time::UNIX_EPOCH)
1411 .map(|d| d.as_secs() as i64)
1412 .unwrap_or(0)
1413}
1414
1415const DEFAULT_SESSION_SWEEP_INTERVAL: Duration = Duration::from_secs(5);
1420
1421#[derive(Clone, Debug)]
1422pub struct SessionSweepConfig {
1423 pub interval: Duration,
1428 pub tx_warn_secs: Duration,
1430 pub tx_max_age_secs: Duration,
1432}
1433
1434impl Default for SessionSweepConfig {
1435 fn default() -> Self {
1436 Self {
1437 interval: DEFAULT_SESSION_SWEEP_INTERVAL,
1438 tx_warn_secs: Duration::from_secs(30),
1439 tx_max_age_secs: Duration::from_secs(120),
1440 }
1441 }
1442}
1443
1444impl SessionSweepConfig {
1445 pub fn from_env() -> Self {
1449 let mut cfg = Self {
1450 interval: session_sweep_interval_from_env(),
1451 ..Self::default()
1452 };
1453 (cfg.tx_warn_secs, cfg.tx_max_age_secs) =
1458 tx_age_thresholds_from_env(cfg.tx_warn_secs, cfg.tx_max_age_secs);
1459
1460 cfg
1461 }
1462}
1463
1464fn session_sweep_interval_from_env() -> Duration {
1465 std::env::var("KHIVE_SESSION_SWEEP_INTERVAL_MS")
1466 .ok()
1467 .and_then(|ms| ms.parse::<u64>().ok())
1468 .filter(|ms| *ms > 0)
1469 .map(Duration::from_millis)
1470 .unwrap_or(DEFAULT_SESSION_SWEEP_INTERVAL)
1471}
1472
1473pub struct SweepBackend {
1483 pub pool: Arc<ConnectionPool>,
1484 pub is_main: bool,
1485}
1486
1487struct BackendSweep {
1493 filter: khive_storage::tx_registry::TxOriginFilter,
1494 tx_age_state: TxAgeSweepState,
1495 sidecar: Option<WalpinSidecarState>,
1496}
1497
1498pub async fn run_session_sweep_task(
1512 backends: Vec<SweepBackend>,
1513 config: SessionSweepConfig,
1514 mut shutdown_rx: tokio::sync::watch::Receiver<()>,
1515) {
1516 let mut interval = tokio::time::interval(config.interval);
1517 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
1518
1519 let mut sweeps: Vec<BackendSweep> = Vec::with_capacity(backends.len());
1520 for backend in backends {
1521 let identity = match backend.pool.origin() {
1522 khive_storage::tx_registry::TxOrigin::Database(id) => id,
1523 khive_storage::tx_registry::TxOrigin::Memory
1526 | khive_storage::tx_registry::TxOrigin::Unscoped => continue,
1527 };
1528 let filter = if backend.is_main {
1529 khive_storage::tx_registry::TxOriginFilter::Main(identity)
1530 } else {
1531 khive_storage::tx_registry::TxOriginFilter::Secondary(identity)
1532 };
1533 let sidecar = WalpinSidecarState::new(
1534 backend.pool.canonical_path(),
1535 true,
1536 "session",
1537 config.interval,
1538 );
1539 sweeps.push(BackendSweep {
1540 filter,
1541 tx_age_state: TxAgeSweepState::default(),
1542 sidecar,
1543 });
1544 }
1545 for sweep in sweeps.iter_mut() {
1546 if let Some(sidecar) = sweep.sidecar.as_mut() {
1547 sidecar.register_beacon().await;
1548 }
1549 }
1550
1551 loop {
1552 tokio::select! {
1553 _ = interval.tick() => {}
1554 _ = shutdown_rx.changed() => break,
1555 }
1556
1557 for sweep in sweeps.iter_mut() {
1558 let oldest = khive_storage::tx_registry::oldest_for(&sweep.filter);
1559 for emission in sweep.tx_age_state.observe(
1560 oldest.as_ref().map(|s| (s.id, s.age, s.label.clone())),
1561 config.tx_warn_secs,
1562 config.tx_max_age_secs,
1563 ) {
1564 log_tx_age_emission(&emission);
1565 }
1566 if let Some(sidecar) = sweep.sidecar.as_mut() {
1567 sidecar.observe(oldest, config.tx_warn_secs).await;
1568 }
1569 }
1570 }
1571
1572 for sweep in sweeps.iter_mut() {
1573 if let Some(sidecar) = sweep.sidecar.as_mut() {
1574 sidecar.shutdown().await;
1575 }
1576 }
1577}
1578
1579#[derive(Clone)]
1584pub struct CheckpointLifecycleOwner {
1585 event_store: Arc<dyn khive_storage::EventStore>,
1586 namespace: String,
1587}
1588
1589impl CheckpointLifecycleOwner {
1590 pub fn new(
1592 event_store: Arc<dyn khive_storage::EventStore>,
1593 namespace: impl Into<String>,
1594 ) -> Self {
1595 Self {
1596 event_store,
1597 namespace: namespace.into(),
1598 }
1599 }
1600}
1601
1602const CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY: usize = 1;
1606
1607#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1608struct CheckpointPressureEpisode {
1609 elevated_ticks: u64,
1610 peak_wal_pages: u64,
1611}
1612
1613impl CheckpointPressureEpisode {
1614 fn start(wal_pages: u64) -> Self {
1615 Self {
1616 elevated_ticks: 1,
1617 peak_wal_pages: wal_pages,
1618 }
1619 }
1620
1621 fn observe(&mut self, wal_pages: u64) {
1622 self.elevated_ticks = self.elevated_ticks.saturating_add(1);
1623 self.peak_wal_pages = self.peak_wal_pages.max(wal_pages);
1624 }
1625}
1626
1627struct CheckpointLifecycleEmitter {
1636 namespace: Option<String>,
1637 sender: Option<tokio::sync::mpsc::Sender<khive_storage::Event>>,
1638 worker: Option<tokio::task::JoinHandle<()>>,
1639 busy_warning_emitted: bool,
1640}
1641
1642impl CheckpointLifecycleEmitter {
1643 fn new(owner: Option<CheckpointLifecycleOwner>) -> Self {
1644 let Some(owner) = owner else {
1645 return Self {
1646 namespace: None,
1647 sender: None,
1648 worker: None,
1649 busy_warning_emitted: false,
1650 };
1651 };
1652
1653 let namespace = owner.namespace.clone();
1654 let (sender, mut receiver) =
1655 tokio::sync::mpsc::channel::<khive_storage::Event>(CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY);
1656 let worker = tokio::spawn(async move {
1657 while let Some(event) = receiver.recv().await {
1658 let kind = event.kind;
1659 CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.fetch_add(1, Ordering::Relaxed);
1660 if let Err(err) = owner.event_store.append_event(event).await {
1661 CHECKPOINT_LIFECYCLE_APPEND_FAILURES.fetch_add(1, Ordering::Relaxed);
1662 tracing::warn!(
1663 error = %err,
1664 event_kind = %kind.name(),
1665 "checkpoint lifecycle event append failed"
1666 );
1667 }
1668 }
1669 });
1670
1671 Self {
1672 namespace: Some(namespace),
1673 sender: Some(sender),
1674 worker: Some(worker),
1675 busy_warning_emitted: false,
1676 }
1677 }
1678
1679 fn try_emit<P: serde::Serialize>(&mut self, kind: khive_types::EventKind, payload: P) -> bool {
1682 let (Some(namespace), Some(sender)) = (&self.namespace, &self.sender) else {
1683 return true;
1684 };
1685 let payload_value = match serde_json::to_value(&payload) {
1686 Ok(value) => value,
1687 Err(err) => {
1688 tracing::warn!(
1689 error = %err,
1690 event_kind = %kind.name(),
1691 "failed to serialize checkpoint lifecycle event payload"
1692 );
1693 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.fetch_add(1, Ordering::Relaxed);
1694 return false;
1695 }
1696 };
1697 let payload_schema_version = match kind {
1698 khive_types::EventKind::CheckpointOutcomeRecorded => 2,
1699 _ => 1,
1700 };
1701 let event = khive_storage::Event::new(
1702 namespace,
1703 "checkpoint.lifecycle",
1704 kind,
1705 khive_types::SubstrateKind::Event,
1706 "daemon:checkpoint_task",
1707 )
1708 .with_payload(payload_value)
1709 .with_payload_schema_version(payload_schema_version);
1710
1711 match sender.try_send(event) {
1712 Ok(()) => {
1713 self.busy_warning_emitted = false;
1714 true
1715 }
1716 Err(tokio::sync::mpsc::error::TrySendError::Full(event)) => {
1717 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.fetch_add(1, Ordering::Relaxed);
1718 if !self.busy_warning_emitted {
1719 tracing::warn!(
1720 event_kind = %event.kind.name(),
1721 queue_capacity = CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY,
1722 "checkpoint lifecycle event dropped because the append worker is busy"
1723 );
1724 self.busy_warning_emitted = true;
1725 }
1726 false
1727 }
1728 Err(tokio::sync::mpsc::error::TrySendError::Closed(event)) => {
1729 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.fetch_add(1, Ordering::Relaxed);
1730 tracing::warn!(
1731 event_kind = %event.kind.name(),
1732 "checkpoint lifecycle event dropped because the append worker stopped"
1733 );
1734 false
1735 }
1736 }
1737 }
1738
1739 async fn shutdown(mut self) {
1747 drop(self.sender.take());
1748 let Some(worker) = self.worker.take() else {
1749 return;
1750 };
1751 worker.abort();
1752 match worker.await {
1753 Ok(()) => {}
1754 Err(err) if err.is_cancelled() => {}
1755 Err(err) => tracing::warn!(
1756 error = %err,
1757 "checkpoint lifecycle event append worker terminated unexpectedly"
1758 ),
1759 }
1760 }
1761}
1762
1763impl Drop for CheckpointLifecycleEmitter {
1764 fn drop(&mut self) {
1765 if let Some(worker) = &self.worker {
1771 worker.abort();
1772 }
1773 }
1774}
1775
1776struct CheckpointConnection {
1806 conn: Option<rusqlite::Connection>,
1807 consecutive_open_failures: u32,
1815}
1816
1817impl CheckpointConnection {
1818 fn new() -> Self {
1819 Self {
1820 conn: None,
1821 consecutive_open_failures: 0,
1822 }
1823 }
1824
1825 fn ensure_open(&mut self, pool: &ConnectionPool) -> Option<&rusqlite::Connection> {
1841 if self.conn.is_none() {
1842 match pool.open_standalone_writer_untracked() {
1843 Ok(conn) => {
1844 if let Err(e) = conn.pragma_update(None, "wal_autocheckpoint", 0) {
1851 tracing::warn!(
1852 error = %e,
1853 "could not disable autocheckpoint on the dedicated checkpoint \
1854 connection"
1855 );
1856 }
1857 if self.consecutive_open_failures > 0 {
1858 tracing::info!(
1859 prior_consecutive_failures = self.consecutive_open_failures,
1860 "dedicated checkpoint connection opened successfully, ending a \
1861 failure streak"
1862 );
1863 }
1864 self.consecutive_open_failures = 0;
1865 self.conn = Some(conn);
1866 }
1867 Err(e) => {
1868 if self.consecutive_open_failures == 0 {
1869 tracing::warn!(
1870 error = %e,
1871 "failed to open the dedicated checkpoint connection; \
1872 this tick is skipped and the open retried next tick"
1873 );
1874 } else {
1875 tracing::debug!(
1876 error = %e,
1877 consecutive_failures = self.consecutive_open_failures,
1878 "dedicated checkpoint connection still unavailable; \
1879 this tick is skipped and the open retried next tick"
1880 );
1881 }
1882 self.consecutive_open_failures =
1883 self.consecutive_open_failures.saturating_add(1);
1884 return None;
1885 }
1886 }
1887 }
1888 self.conn.as_ref()
1889 }
1890
1891 fn drop_connection(&mut self) {
1895 self.conn = None;
1896 }
1897}
1898
1899pub async fn run_checkpoint_task(
1933 pool: Arc<ConnectionPool>,
1934 config: CheckpointConfig,
1935 lifecycle_owner: Option<CheckpointLifecycleOwner>,
1936 mut shutdown_rx: tokio::sync::watch::Receiver<()>,
1937 is_main: bool,
1938) {
1939 match pool.claim_checkpoint_ownership() {
1946 Ok(()) => {
1947 if let Err(e) = pool.propagate_checkpoint_claim_to_writer_task().await {
1948 tracing::warn!(
1949 error = %e,
1950 "checkpoint task could not reach the writer task's connection; it keeps the \
1951 bounded autocheckpoint fallback"
1952 );
1953 }
1954 }
1955 Err(e) => {
1956 tracing::warn!(
1957 error = %e,
1958 "checkpoint task could not re-apply the ownership pragma on the pooled writer; \
1959 writer connections keep the bounded autocheckpoint fallback unless ownership is \
1960 claimed later"
1961 );
1962 }
1963 }
1964 let mut interval = tokio::time::interval(config.interval);
1965 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
1966 let mut severity_state = CheckpointSeverityState::default();
1967 let mut tx_age_state = TxAgeSweepState::default();
1968 let mut was_above_high_water = false;
1969 #[cfg(unix)]
1970 let legacy_walpin_fallback_interval = DEFAULT_SESSION_SWEEP_INTERVAL;
1971 #[cfg(unix)]
1972 let mut truncate_state =
1973 TruncateState::with_legacy_walpin_fallback(legacy_walpin_fallback_interval);
1974 #[cfg(not(unix))]
1975 let mut truncate_state = TruncateState::default();
1976 let mut lifecycle_emitter = CheckpointLifecycleEmitter::new(lifecycle_owner);
1977 let mut event_elevation_open = false;
1982 let mut pressure_episode: Option<CheckpointPressureEpisode> = None;
1983 let mut pending_recovery: Option<khive_storage::CheckpointOutcomeRecordedPayload> = None;
1988 let mut was_observed_above_warn = false;
1989 let tx_filter = match pool.origin() {
2002 khive_storage::tx_registry::TxOrigin::Database(id) => Some(if is_main {
2003 khive_storage::tx_registry::TxOriginFilter::Main(id)
2004 } else {
2005 khive_storage::tx_registry::TxOriginFilter::Secondary(id)
2006 }),
2007 khive_storage::tx_registry::TxOrigin::Memory
2008 | khive_storage::tx_registry::TxOrigin::Unscoped => None,
2009 };
2010 #[cfg(unix)]
2016 let mut walpin_state =
2017 WalpinSidecarState::new(pool.canonical_path(), true, "daemon", config.interval);
2018 #[cfg(unix)]
2019 if let Some(sidecar) = walpin_state.as_mut() {
2020 sidecar.register_beacon().await;
2021 }
2022
2023 let mut checkpoint_conn = CheckpointConnection::new();
2027 checkpoint_conn.ensure_open(&pool);
2028
2029 loop {
2030 tokio::select! {
2035 _ = interval.tick() => {}
2036 _ = shutdown_rx.changed() => break,
2037 }
2038
2039 #[cfg(unix)]
2040 truncate_state.begin_tick();
2041
2042 #[cfg(unix)]
2043 let mut pending_sidecar_attribution = None;
2044
2045 let tick = match checkpoint_conn.ensure_open(&pool) {
2046 None => {
2047 note_checkpoint_skipped();
2048 CheckpointTick::Skipped
2049 }
2050 Some(conn) => match checkpoint_once_core(&pool, conn, &config, &mut truncate_state) {
2051 Ok(outcome) => {
2052 #[cfg(unix)]
2053 {
2054 pending_sidecar_attribution = outcome.sidecar_attribution;
2055 }
2056 #[cfg(not(unix))]
2057 let _ = outcome.sidecar_attribution;
2058 CheckpointTick::Observed(outcome.wal_pages)
2059 }
2060 Err(e) => {
2061 tracing::warn!(
2062 error = %e,
2063 "dedicated checkpoint connection failed a pragma; \
2064 dropping it for a fresh reopen next tick"
2065 );
2066 checkpoint_conn.drop_connection();
2067 note_checkpoint_skipped();
2068 CheckpointTick::Skipped
2069 }
2070 },
2071 };
2072
2073 #[cfg(unix)]
2082 if let Err(error) =
2083 complete_walpin_attribution(pending_sidecar_attribution, &mut truncate_state).await
2084 {
2085 tracing::warn!(
2086 error = %error,
2087 failure_kind = error.kind(),
2088 "ADR-091 Amendment 2 Plank B: no-progress sidecar attribution failed"
2089 );
2090 }
2091
2092 let oldest_tx = tx_filter
2107 .as_ref()
2108 .and_then(khive_storage::tx_registry::oldest_for);
2109 for emission in tx_age_state.observe(
2110 oldest_tx.as_ref().map(|s| (s.id, s.age, s.label.clone())),
2111 config.tx_warn_secs,
2112 config.tx_max_age_secs,
2113 ) {
2114 log_tx_age_emission(&emission);
2115 }
2116 #[cfg(unix)]
2120 if let Some(sidecar) = walpin_state.as_mut() {
2121 sidecar
2122 .observe(oldest_tx.clone(), config.tx_warn_secs)
2123 .await;
2124 let _ = run_walpin_housekeeping_if_due(
2125 sidecar,
2126 &mut truncate_state,
2127 legacy_walpin_fallback_interval,
2128 )
2129 .await;
2130 }
2131
2132 let wal_pages = match tick {
2135 CheckpointTick::Skipped => continue,
2136 CheckpointTick::Observed(n) => n,
2137 };
2138
2139 let above_warn = wal_pages >= config.warn_pages;
2140 let above_high_water = wal_pages >= config.high_water_pages;
2141 let above_truncate_high_water = wal_pages >= config.truncate_high_water_pages;
2142 note_checkpoint_pressure_observation(above_warn, was_observed_above_warn);
2143 was_observed_above_warn = above_warn;
2144
2145 log_tx_registry_oldest_debug(wal_pages, oldest_tx.as_ref());
2151
2152 for emission in severity_state.observe_wal_pages(wal_pages, &config) {
2157 match emission.rung {
2158 CheckpointSeverityRung::Info => {
2159 log_tx_registry_oldest_warn(wal_pages, oldest_tx.as_ref());
2160 tracing::info!(
2161 wal_pages = emission.wal_pages,
2162 warn_threshold = emission.threshold_pages,
2163 "WAL page count crossed warn threshold"
2164 );
2165 }
2166 CheckpointSeverityRung::Warn => {
2167 tracing::warn!(
2168 wal_pages = emission.wal_pages,
2169 warn_threshold = emission.threshold_pages,
2170 consecutive_cycles = emission.consecutive_cycles,
2171 "WAL page count failed to drain below warn threshold"
2172 );
2173 }
2174 CheckpointSeverityRung::Alarm => {
2175 }
2177 }
2178 }
2179
2180 let high_water_crossed = crossing_warn(above_high_water, &mut was_above_high_water);
2181 if high_water_crossed {
2182 log_tx_registry_snapshot_warn(wal_pages);
2183 tracing::warn!(
2184 wal_pages,
2185 high_water = config.high_water_pages,
2186 "WAL high-water mark exceeded; sustained WAL pressure — \
2187 a long-lived reader may be pinning an old snapshot that PASSIVE cannot reclaim"
2188 );
2189 }
2190
2191 observe_checkpoint_pressure_tick(
2197 above_warn,
2198 wal_pages,
2199 above_high_water,
2200 above_truncate_high_water,
2201 &config,
2202 &mut event_elevation_open,
2203 &mut pressure_episode,
2204 &mut pending_recovery,
2205 |payload| {
2206 lifecycle_emitter
2207 .try_emit(khive_types::EventKind::CheckpointOutcomeRecorded, payload)
2208 },
2209 );
2210 }
2211
2212 lifecycle_emitter.shutdown().await;
2213
2214 #[cfg(unix)]
2215 if let Some(sidecar) = walpin_state.as_mut() {
2216 sidecar.shutdown().await;
2217 }
2218}
2219
2220fn checkpoint_outcome_should_emit(above_warn: bool, was_elevated: bool) -> bool {
2224 above_warn != was_elevated
2225}
2226
2227#[allow(clippy::too_many_arguments)]
2240fn observe_checkpoint_pressure_tick(
2241 above_warn: bool,
2242 wal_pages: u64,
2243 above_high_water: bool,
2244 above_truncate_high_water: bool,
2245 config: &CheckpointConfig,
2246 event_elevation_open: &mut bool,
2247 pressure_episode: &mut Option<CheckpointPressureEpisode>,
2248 pending_recovery: &mut Option<khive_storage::CheckpointOutcomeRecordedPayload>,
2249 mut try_emit: impl FnMut(khive_storage::CheckpointOutcomeRecordedPayload) -> bool,
2250) {
2251 let pending_blocks_emission = if let Some(payload) = pending_recovery.clone() {
2259 if try_emit(payload) {
2260 *pending_recovery = None;
2261 false
2262 } else {
2263 true
2264 }
2265 } else {
2266 false
2267 };
2268
2269 if above_warn {
2270 match pressure_episode.as_mut() {
2271 Some(episode) => episode.observe(wal_pages),
2272 None => *pressure_episode = Some(CheckpointPressureEpisode::start(wal_pages)),
2273 }
2274 } else if !*event_elevation_open {
2275 if pending_blocks_emission && pressure_episode.is_some() {
2282 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.fetch_add(1, Ordering::Relaxed);
2283 tracing::warn!(
2284 wal_pages,
2285 "checkpoint pressure episode elapsed unreported behind an undelivered recovery summary"
2286 );
2287 }
2288 *pressure_episode = None;
2289 }
2290
2291 if pending_blocks_emission || !checkpoint_outcome_should_emit(above_warn, *event_elevation_open)
2292 {
2293 return;
2294 }
2295 let Some(episode) = *pressure_episode else {
2296 tracing::warn!(
2297 above_warn,
2298 event_elevation_open = *event_elevation_open,
2299 "checkpoint pressure transition has no episode aggregate"
2300 );
2301 return;
2302 };
2303 let payload = khive_storage::CheckpointOutcomeRecordedPayload {
2304 wal_pages,
2305 warn_pages: config.warn_pages,
2306 high_water_pages: config.high_water_pages,
2307 truncate_high_water_pages: config.truncate_high_water_pages,
2308 above_warn,
2309 above_high_water,
2310 above_truncate_high_water,
2311 episode_elevated_ticks: Some(episode.elevated_ticks),
2312 episode_peak_wal_pages: Some(episode.peak_wal_pages),
2313 };
2314 if try_emit(payload.clone()) {
2315 *event_elevation_open = above_warn;
2316 if !above_warn {
2317 *pressure_episode = None;
2318 }
2319 } else if !above_warn {
2320 debug_assert!(
2331 pending_recovery.is_none(),
2332 "recovery emission attempted while an earlier summary was still pending"
2333 );
2334 *event_elevation_open = false;
2335 *pressure_episode = None;
2336 *pending_recovery = Some(payload);
2337 }
2338}
2339
2340fn log_tx_registry_oldest_debug(
2349 wal_pages: u64,
2350 oldest: Option<&khive_storage::tx_registry::OldestSpan>,
2351) {
2352 if let Some(span) = oldest {
2353 tracing::debug!(
2354 wal_pages,
2355 oldest_tx_age_secs = span.age.as_secs_f64(),
2356 oldest_tx_label = span.label.as_deref().unwrap_or("<unlabeled>"),
2357 "WAL checkpoint tick: oldest open transaction registry entry"
2358 );
2359 }
2360}
2361
2362fn log_tx_registry_oldest_warn(
2366 wal_pages: u64,
2367 oldest: Option<&khive_storage::tx_registry::OldestSpan>,
2368) {
2369 if let Some(span) = oldest {
2370 tracing::warn!(
2371 wal_pages,
2372 oldest_tx_age_secs = span.age.as_secs_f64(),
2373 oldest_tx_label = span.label.as_deref().unwrap_or("<unlabeled>"),
2374 "WAL checkpoint tick: oldest open transaction registry entry"
2375 );
2376 }
2377}
2378
2379fn log_tx_registry_snapshot_warn(wal_pages: u64) {
2383 for (age, label) in khive_storage::tx_registry::snapshot() {
2384 tracing::warn!(
2385 wal_pages,
2386 tx_age_secs = age.as_secs_f64(),
2387 tx_label = label.as_deref().unwrap_or("<unlabeled>"),
2388 "WAL high-water: open transaction registry entry"
2389 );
2390 }
2391}
2392
2393#[derive(Debug)]
2398#[must_use]
2399struct CheckpointCoreOutcome {
2400 wal_pages: u64,
2401 sidecar_attribution: Option<WalpinAttributionRequest>,
2402}
2403
2404pub fn checkpoint_once(
2428 pool: &ConnectionPool,
2429 conn: &rusqlite::Connection,
2430 config: &CheckpointConfig,
2431 truncate_state: &mut TruncateState,
2432) -> Result<u64, rusqlite::Error> {
2433 checkpoint_once_core(pool, conn, config, truncate_state).map(|outcome| outcome.wal_pages)
2434}
2435
2436fn checkpoint_once_core(
2440 pool: &ConnectionPool,
2441 conn: &rusqlite::Connection,
2442 config: &CheckpointConfig,
2443 truncate_state: &mut TruncateState,
2444) -> Result<CheckpointCoreOutcome, rusqlite::Error> {
2445 #[cfg(unix)]
2446 truncate_state.begin_tick();
2447 let raw_observation = match query_checkpoint_observation(conn) {
2448 Ok(observation) => observation,
2449 Err(e) => {
2450 tracing::warn!(error = %e, "WAL checkpoint failed");
2451 return Err(e);
2452 }
2453 };
2454 let observation = record_routine_wal_observation(pool, raw_observation);
2455 let wal_pages = observation.log_frames;
2456 LAST_WAL_PAGES.store(wal_pages, Ordering::Relaxed);
2457 note_checkpoint_observed(wal_pages);
2458
2459 if raw_observation.busy != 0 {
2460 tracing::debug!(
2461 busy = raw_observation.busy,
2462 wal_log_frames = raw_observation.log_frames,
2463 wal_checkpointed_frames = raw_observation.checkpointed_frames,
2464 wal_pending_frames = observation.pending_frames,
2465 wal_physical_bytes = ?observation.physical_wal_bytes,
2466 "WAL PASSIVE checkpoint reported incomplete progress"
2467 );
2468 }
2469 tracing::debug!(
2470 wal_pages,
2471 wal_checkpointed_frames = observation.checkpointed_frames,
2472 wal_pending_frames = observation.pending_frames,
2473 wal_physical_bytes = ?observation.physical_wal_bytes,
2474 "WAL checkpoint issued"
2475 );
2476
2477 let sidecar_attribution = maybe_truncate(pool, conn, config, wal_pages, truncate_state);
2478
2479 Ok(CheckpointCoreOutcome {
2480 wal_pages,
2481 sidecar_attribution,
2482 })
2483}
2484
2485fn maybe_truncate(
2491 pool: &ConnectionPool,
2492 conn: &rusqlite::Connection,
2493 config: &CheckpointConfig,
2494 wal_pages_before: u64,
2495 truncate_state: &mut TruncateState,
2496) -> Option<WalpinAttributionRequest> {
2497 if wal_pages_before < config.truncate_high_water_pages {
2498 return None;
2499 }
2500
2501 if let Some(last) = truncate_state.last_attempt {
2502 if last.elapsed() < config.truncate_min_interval {
2503 return None;
2504 }
2505 }
2506
2507 log_tx_registry_snapshot_warn(wal_pages_before);
2510
2511 let original_busy_timeout = pool.config().busy_timeout;
2512
2513 if let Err(e) = conn.busy_timeout(config.truncate_busy_timeout) {
2514 tracing::warn!(error = %e, "failed to lower busy_timeout for TRUNCATE attempt; skipping");
2520 return None;
2521 }
2522
2523 #[cfg(unix)]
2524 let mut holder_attribution = capture_walpin_attribution_request(pool, truncate_state);
2525 #[cfg(unix)]
2526 let mut sidecar_attribution = None;
2527 #[cfg(not(unix))]
2528 let sidecar_attribution = None;
2529
2530 truncate_state.last_attempt = Some(Instant::now());
2534
2535 let start = Instant::now();
2536 let outcome = conn.execute_batch("PRAGMA wal_checkpoint(TRUNCATE)");
2537 let elapsed = start.elapsed();
2538
2539 if let Err(e) = conn.busy_timeout(original_busy_timeout) {
2542 tracing::warn!(error = %e, "failed to restore busy_timeout after TRUNCATE attempt");
2543 }
2544
2545 match outcome {
2546 Ok(()) => {
2547 let wal_pages_after = query_wal_pages(conn);
2548 tracing::info!(
2549 wal_pages_before,
2550 wal_pages_after,
2551 elapsed_ms = elapsed.as_millis() as u64,
2552 "WAL TRUNCATE checkpoint attempted"
2553 );
2554
2555 let made_progress = wal_pages_after < wal_pages_before;
2556 if !made_progress {
2557 tracing::warn!(
2558 wal_pages_before,
2559 wal_pages_after,
2560 "WAL TRUNCATE attempt made no progress; \
2561 a long-lived reader may still be pinning the WAL snapshot"
2562 );
2563 log_tx_registry_snapshot_warn(wal_pages_after);
2564 #[cfg(test)]
2565 if let Some(path) = pool.canonical_path() {
2566 truncate_report_test_sync::after_no_progress_before_report(path);
2567 }
2568 #[cfg(unix)]
2569 {
2570 sidecar_attribution = holder_attribution.take();
2577 }
2578 log_wal_pin_depth(conn);
2579 }
2580
2581 note_truncate_outcome(config, wal_pages_after, truncate_state);
2582 }
2583 Err(e) => {
2584 tracing::warn!(error = %e, wal_pages_before, "WAL TRUNCATE attempt failed");
2585 log_tx_registry_snapshot_warn(wal_pages_before);
2586 note_truncate_outcome(config, wal_pages_before, truncate_state);
2587 }
2588 }
2589 #[cfg(unix)]
2590 if let Some(WalpinAttributionRequest::Fresh {
2591 previous_last_attempt,
2592 ..
2593 }) = holder_attribution.as_ref()
2594 {
2595 truncate_state.restore_walpin_full_scan_reservation(*previous_last_attempt);
2596 }
2597 sidecar_attribution
2598}
2599
2600#[cfg(test)]
2601mod truncate_report_test_sync {
2602 use std::path::{Path, PathBuf};
2603 use std::sync::mpsc::{sync_channel, Receiver, SyncSender};
2604 use std::sync::Mutex;
2605
2606 struct Hook {
2607 db_path: PathBuf,
2608 reached_tx: SyncSender<()>,
2609 proceed_rx: Receiver<()>,
2610 }
2611
2612 static HOOK: Mutex<Option<Hook>> = Mutex::new(None);
2613
2614 pub(crate) fn install(db_path: PathBuf) -> (Receiver<()>, SyncSender<()>) {
2615 let (reached_tx, reached_rx) = sync_channel(0);
2616 let (proceed_tx, proceed_rx) = sync_channel(0);
2617 let replaced = HOOK
2618 .lock()
2619 .unwrap_or_else(|poisoned| poisoned.into_inner())
2620 .replace(Hook {
2621 db_path,
2622 reached_tx,
2623 proceed_rx,
2624 });
2625 assert!(replaced.is_none(), "truncate report hook already installed");
2626 (reached_rx, proceed_tx)
2627 }
2628
2629 pub(crate) fn uninstall() {
2630 *HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
2631 }
2632
2633 pub(crate) fn after_no_progress_before_report(db_path: &Path) {
2634 let hook = {
2635 let mut guard = HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
2636 match guard.as_ref() {
2637 Some(hook) if hook.db_path == db_path => guard.take(),
2638 _ => None,
2639 }
2640 };
2641 let Some(hook) = hook else {
2642 return;
2643 };
2644 let _ = hook.reached_tx.send(());
2645 let _ = hook.proceed_rx.recv();
2646 }
2647}
2648
2649#[cfg(all(test, unix))]
2654mod walpin_attribution_test_sync {
2655 use std::path::{Path, PathBuf};
2656 use std::sync::atomic::{AtomicUsize, Ordering};
2657 use std::sync::mpsc::{sync_channel, Receiver, SyncSender};
2658 use std::sync::{Arc, Mutex};
2659
2660 enum Behavior {
2661 Pause {
2662 reached_tx: tokio::sync::oneshot::Sender<std::thread::ThreadId>,
2663 proceed_rx: Receiver<()>,
2664 },
2665 Panic,
2666 }
2667
2668 struct Hook {
2669 dir: PathBuf,
2670 behavior: Behavior,
2671 }
2672
2673 static HOOK: Mutex<Option<Hook>> = Mutex::new(None);
2674 static REPORT_COUNTER: Mutex<Option<Arc<AtomicUsize>>> = Mutex::new(None);
2675
2676 pub(crate) fn install_pause(
2677 dir: PathBuf,
2678 ) -> (
2679 tokio::sync::oneshot::Receiver<std::thread::ThreadId>,
2680 SyncSender<()>,
2681 Arc<AtomicUsize>,
2682 ) {
2683 let (reached_tx, reached_rx) = tokio::sync::oneshot::channel();
2684 let (proceed_tx, proceed_rx) = sync_channel(0);
2685 let report_counter = Arc::new(AtomicUsize::new(0));
2686 let replaced = HOOK
2687 .lock()
2688 .unwrap_or_else(|poisoned| poisoned.into_inner())
2689 .replace(Hook {
2690 dir,
2691 behavior: Behavior::Pause {
2692 reached_tx,
2693 proceed_rx,
2694 },
2695 });
2696 assert!(
2697 replaced.is_none(),
2698 "walpin attribution hook already installed"
2699 );
2700 *REPORT_COUNTER
2701 .lock()
2702 .unwrap_or_else(|poisoned| poisoned.into_inner()) = Some(Arc::clone(&report_counter));
2703 (reached_rx, proceed_tx, report_counter)
2704 }
2705
2706 pub(crate) fn install_panic(dir: PathBuf) {
2707 let replaced = HOOK
2708 .lock()
2709 .unwrap_or_else(|poisoned| poisoned.into_inner())
2710 .replace(Hook {
2711 dir,
2712 behavior: Behavior::Panic,
2713 });
2714 assert!(
2715 replaced.is_none(),
2716 "walpin attribution hook already installed"
2717 );
2718 }
2719
2720 pub(crate) fn uninstall() {
2721 *HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
2722 *REPORT_COUNTER
2723 .lock()
2724 .unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
2725 }
2726
2727 pub(crate) fn before_enumeration(dir: &Path) {
2728 let hook = {
2729 let mut guard = HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
2730 match guard.as_ref() {
2731 Some(hook) if hook.dir == dir => guard.take(),
2732 _ => None,
2733 }
2734 };
2735 let Some(hook) = hook else {
2736 return;
2737 };
2738 match hook.behavior {
2739 Behavior::Pause {
2740 reached_tx,
2741 proceed_rx,
2742 } => {
2743 if reached_tx.send(std::thread::current().id()).is_ok() {
2744 let _ = proceed_rx.recv();
2745 }
2746 }
2747 Behavior::Panic => panic!("injected walpin attribution worker panic"),
2748 }
2749 }
2750
2751 pub(crate) fn report_used() {
2752 if let Some(counter) = REPORT_COUNTER
2753 .lock()
2754 .unwrap_or_else(|poisoned| poisoned.into_inner())
2755 .as_ref()
2756 {
2757 counter.fetch_add(1, Ordering::SeqCst);
2758 }
2759 }
2760}
2761
2762fn note_truncate_outcome(
2768 config: &CheckpointConfig,
2769 wal_pages_after: u64,
2770 state: &mut TruncateState,
2771) {
2772 TRUNCATE_ATTEMPTS.fetch_add(1, Ordering::Relaxed);
2777
2778 if wal_pages_after >= config.warn_pages {
2779 state.consecutive_failures = state.consecutive_failures.saturating_add(1);
2780 if state.consecutive_failures == 3 {
2781 tracing::warn!(
2782 wal_pages_after,
2783 warn_threshold = config.warn_pages,
2784 "WAL TRUNCATE has failed to clear WAL pressure for 3 consecutive attempts"
2785 );
2786 }
2787 } else {
2788 state.consecutive_failures = 0;
2789 }
2790
2791 TRUNCATE_CONSECUTIVE_FAILURES.store(state.consecutive_failures as u64, Ordering::Relaxed);
2792}
2793
2794#[cfg(unix)]
2802#[derive(Debug)]
2803enum WalpinAttributionRequest {
2804 Fresh {
2805 dir: PathBuf,
2806 census: Result<crate::walpin::CensusResult, String>,
2807 legacy_fallback_interval: Duration,
2808 previous_last_attempt: Option<Instant>,
2809 },
2810 Cached(CachedWalpinAttribution),
2811 Suppressed,
2812}
2813
2814#[cfg(unix)]
2815#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2816enum WalpinReportFreshness {
2817 Fresh,
2818 Cached { age: Duration },
2819}
2820
2821#[cfg(unix)]
2822impl WalpinReportFreshness {
2823 fn is_fresh(self) -> bool {
2824 self == Self::Fresh
2825 }
2826}
2827
2828#[cfg(not(unix))]
2831type WalpinAttributionRequest = ();
2832
2833#[cfg(unix)]
2838#[derive(Debug, Clone, PartialEq, Eq)]
2839enum WalpinAttributionFailure {
2840 Enumeration(String),
2841 Worker(String),
2842}
2843
2844#[cfg(unix)]
2845impl WalpinAttributionFailure {
2846 fn kind(&self) -> &'static str {
2847 match self {
2848 Self::Enumeration(_) => "enumeration",
2849 Self::Worker(_) => "blocking_worker",
2850 }
2851 }
2852}
2853
2854#[cfg(unix)]
2855impl std::fmt::Display for WalpinAttributionFailure {
2856 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2857 match self {
2858 Self::Enumeration(error) => write!(
2859 formatter,
2860 "sidecar directory failed the trust-boundary enumeration; cross-process \
2861 WAL-pin attribution is unestablished for this tick: {error}"
2862 ),
2863 Self::Worker(error) => write!(
2864 formatter,
2865 "sidecar attribution blocking worker failed; cross-process WAL-pin \
2866 attribution is unestablished for this tick: {error}"
2867 ),
2868 }
2869 }
2870}
2871
2872#[cfg(unix)]
2875fn capture_walpin_attribution_request(
2876 pool: &ConnectionPool,
2877 state: &mut TruncateState,
2878) -> Option<WalpinAttributionRequest> {
2879 let path = pool.canonical_path()?;
2880 if !crate::walpin::sidecar_enabled(true) {
2881 return None;
2882 }
2883 let legacy_fallback_interval = state.legacy_walpin_fallback_interval;
2884 Some(match state.plan_walpin_attribution_at(Instant::now()) {
2885 WalpinFullScanPlan::Refresh {
2886 previous_last_attempt,
2887 } => WalpinAttributionRequest::Fresh {
2888 dir: crate::walpin::sidecar_dir_for(path),
2889 census: crate::walpin::census_holders(path).map_err(|error| error.to_string()),
2890 legacy_fallback_interval,
2891 previous_last_attempt,
2892 },
2893 WalpinFullScanPlan::Cached(cached) => WalpinAttributionRequest::Cached(cached),
2894 WalpinFullScanPlan::Suppressed => WalpinAttributionRequest::Suppressed,
2895 })
2896}
2897
2898#[cfg(unix)]
2904async fn complete_walpin_attribution(
2905 request: Option<WalpinAttributionRequest>,
2906 state: &mut TruncateState,
2907) -> Result<bool, WalpinAttributionFailure> {
2908 let Some(request) = request else {
2909 return Ok(false);
2910 };
2911 match request {
2912 WalpinAttributionRequest::Suppressed => Ok(false),
2913 WalpinAttributionRequest::Cached(cached) => {
2914 log_walpin_sidecar_report(
2915 &cached.report,
2916 cached.census,
2917 WalpinReportFreshness::Cached {
2918 age: Instant::now().saturating_duration_since(cached.captured_at),
2919 },
2920 );
2921 Ok(true)
2922 }
2923 WalpinAttributionRequest::Fresh {
2924 dir,
2925 census,
2926 legacy_fallback_interval,
2927 previous_last_attempt: _,
2928 } => {
2929 state.sidecar_attribution_attempted_this_tick = true;
2930 if state.walpin_full_scan_last_attempt.is_none() {
2931 state.walpin_full_scan_last_attempt = Some(Instant::now());
2932 }
2933 let fallback = state.walpin_cached_attribution.clone();
2934 let result = tokio::task::spawn_blocking(move || {
2935 #[cfg(test)]
2936 walpin_attribution_test_sync::before_enumeration(&dir);
2937 crate::walpin::enumerate_live(&dir, legacy_fallback_interval)
2938 })
2939 .await
2940 .map_err(|error| WalpinAttributionFailure::Worker(error.to_string()))
2941 .and_then(|result| {
2942 result.map_err(|error| WalpinAttributionFailure::Enumeration(error.to_string()))
2943 });
2944
2945 match result {
2946 Ok(report) => {
2947 let captured_at = Instant::now();
2948 log_walpin_sidecar_report(
2949 &report,
2950 census.clone(),
2951 WalpinReportFreshness::Fresh,
2952 );
2953 state.cache_walpin_attribution(report, census, captured_at);
2954 Ok(true)
2955 }
2956 Err(error) => {
2957 if let Some(cached) = fallback {
2958 log_walpin_sidecar_report(
2959 &cached.report,
2960 cached.census,
2961 WalpinReportFreshness::Cached {
2962 age: Instant::now().saturating_duration_since(cached.captured_at),
2963 },
2964 );
2965 }
2966 Err(error)
2967 }
2968 }
2969 }
2970 }
2971}
2972
2973#[cfg(unix)]
2988fn log_walpin_sidecar_report(
2989 report: &crate::walpin::WalpinReport,
2990 census: Result<crate::walpin::CensusResult, String>,
2991 freshness: WalpinReportFreshness,
2992) {
2993 #[cfg(test)]
2994 walpin_attribution_test_sync::report_used();
2995 let now = now_epoch_secs();
2996 for hb in report.reporting() {
2997 tracing::warn!(
3003 walpin_pid = hb.pid,
3004 walpin_role = %hb.process_role,
3005 walpin_oldest_tx_age_secs = hb.current_oldest_tx_age_secs(now),
3006 walpin_oldest_tx_label = hb.oldest_tx_label.as_deref().unwrap_or("<unlabeled>"),
3007 walpin_attribution_basis = hb.attribution_basis.as_deref().unwrap_or("<unspecified>"),
3008 walpin_attribution_evidence_backed = hb.attribution_is_evidence_backed(),
3009 walpin_attribution_fresh = freshness.is_fresh(),
3010 walpin_health = "reporting",
3011 "ADR-091 Amendment 2 Plank B: live cross-process WAL-pin attribution report"
3012 );
3013 }
3014 for pid in report.registered_silent_pids() {
3015 tracing::debug!(
3016 walpin_pid = pid,
3017 walpin_health = "registered_silent",
3018 walpin_attribution_fresh = freshness.is_fresh(),
3019 "ADR-091 Amendment 2 Plank B: process affirmatively reports no over-threshold span"
3020 );
3021 }
3022 let mut unknown_pids: Vec<u32> = report.unknown_pids().collect();
3023 if let WalpinReportFreshness::Cached { age } = freshness {
3024 tracing::warn!(
3025 walpin_cache_age_ms = age.as_millis() as u64,
3026 "cached WAL-pin attribution is diagnostic-only; fully-attributed \
3027 conclusion is not licensed"
3028 );
3029 unknown_pids.push(0);
3030 }
3031
3032 match census {
3037 Ok(census) => {
3038 let sidecar_known: std::collections::HashSet<u32> = report
3039 .reporting()
3040 .map(|hb| hb.pid)
3041 .chain(report.registered_silent_pids())
3042 .chain(unknown_pids.iter().copied())
3043 .collect();
3044 let mut census_only: Vec<u32> =
3045 census.holders.difference(&sidecar_known).copied().collect();
3046 if !census_only.is_empty() {
3047 census_only.sort_unstable();
3048 tracing::warn!(
3049 ?census_only,
3050 "ADR-091 Amendment 2: these PIDs hold the database file open \
3051 at the OS level but have no sidecar data at all (pre-feature binary, \
3052 sidecar disabled, or wedged before its first write)"
3053 );
3054 unknown_pids.extend(census_only);
3055 }
3056 if !census.is_complete() {
3057 let mut uninspectable = census.uninspectable_pids.clone();
3058 uninspectable.sort_unstable();
3059 tracing::warn!(
3060 ?uninspectable,
3061 truncated = census.truncated,
3062 "ADR-091 Amendment 2: the OS-derived holder census is \
3063 INCOMPLETE — either specific PIDs' open file descriptors could not be \
3064 inspected (permission denied, or a listing race), or the enumeration walk \
3065 itself has positive evidence it did not see the full live-process universe \
3066 (namespace/visibility check, directory-iterator error, self-canary, or a \
3067 libproc buffer that stayed at capacity after bounded retries) — cannot \
3068 rule out an unregistered holder"
3069 );
3070 if uninspectable.is_empty() {
3071 unknown_pids.push(0);
3078 } else {
3079 unknown_pids.extend(uninspectable);
3080 }
3081 }
3082 }
3083 Err(e) => {
3084 tracing::warn!(
3085 error = %e,
3086 "ADR-091 Amendment 2: OS-derived holder census failed; \
3087 attribution cannot rule out an unregistered database holder this tick"
3088 );
3089 unknown_pids.push(0);
3093 }
3094 }
3095
3096 unknown_pids.sort_unstable();
3097 unknown_pids.dedup();
3098 if !unknown_pids.is_empty() {
3099 tracing::warn!(
3100 ?unknown_pids,
3101 "ADR-091 Amendment 2 Plank B: sidecar health unestablished for these PIDs; \
3102 attribution is inconclusive and the native/unregistered-mechanism conclusion \
3103 is NOT licensed this tick"
3104 );
3105 } else if report.reporting().next().is_none() {
3106 tracing::info!(
3107 "ADR-091 Amendment 2 Plank B: every live PID is reporting or registered-silent \
3108 with none pinning; the WAL pin is not attributable to any in-process registry \
3109 span this sidecar covers"
3110 );
3111 }
3112}
3113
3114fn log_wal_pin_depth(conn: &rusqlite::Connection) {
3120 match query_wal_pin_depth(conn) {
3121 Ok((log, checkpointed)) => {
3122 tracing::warn!(
3123 wal_log_frames = log,
3124 wal_checkpointed_frames = checkpointed,
3125 wal_pin_depth = (log - checkpointed).max(0),
3126 "ADR-091 Amendment 2 Plank C: WAL pin depth after TRUNCATE no-progress"
3127 );
3128 }
3129 Err(e) => {
3130 tracing::warn!(
3131 error = %e,
3132 "ADR-091 Amendment 2 Plank C: failed to query WAL pin depth"
3133 );
3134 }
3135 }
3136}
3137
3138fn query_wal_pin_depth(conn: &rusqlite::Connection) -> rusqlite::Result<(i64, i64)> {
3145 conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
3146 Ok((row.get::<_, i64>(1)?, row.get::<_, i64>(2)?))
3147 })
3148}
3149
3150fn crossing_warn(now_above: bool, was_above: &mut bool) -> bool {
3159 let fire = now_above && !*was_above;
3160 *was_above = now_above;
3161 fire
3162}
3163
3164#[derive(Debug, Clone, Copy)]
3165struct RawCheckpointObservation {
3166 busy: i64,
3167 log_frames: i64,
3168 checkpointed_frames: i64,
3169}
3170
3171fn query_checkpoint_observation(
3175 conn: &rusqlite::Connection,
3176) -> rusqlite::Result<RawCheckpointObservation> {
3177 conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
3178 Ok(RawCheckpointObservation {
3179 busy: row.get(0)?,
3180 log_frames: row.get(1)?,
3181 checkpointed_frames: row.get(2)?,
3182 })
3183 })
3184}
3185
3186fn query_wal_pages(conn: &rusqlite::Connection) -> u64 {
3192 let pages = query_checkpoint_observation(conn)
3193 .map(|observation| observation.log_frames)
3194 .unwrap_or(0)
3195 .max(0) as u64;
3196 LAST_WAL_PAGES.store(pages, Ordering::Relaxed);
3200 note_checkpoint_observed(pages);
3201 pages
3202}
3203
3204#[cfg(test)]
3205mod tests {
3206 use super::*;
3207 use crate::pool::PoolConfig;
3208 use crate::writer_task::WriterTaskHandle;
3209 use rusqlite::hooks::{AuthAction, Authorization};
3210 use serial_test::serial;
3211 use tracing::field::{Field, Visit};
3212
3213 #[derive(Clone, Debug, Default)]
3214 struct CapturedEvent {
3215 message: Option<String>,
3216 oldest_tx_label: Option<String>,
3217 tx_label: Option<String>,
3218 census_only: Option<String>,
3219 }
3220
3221 #[derive(Default)]
3222 struct CapturedEventVisitor(CapturedEvent);
3223
3224 impl Visit for CapturedEventVisitor {
3225 fn record_str(&mut self, field: &Field, value: &str) {
3226 match field.name() {
3227 "message" => self.0.message = Some(value.to_string()),
3228 "oldest_tx_label" => self.0.oldest_tx_label = Some(value.to_string()),
3229 "tx_label" => self.0.tx_label = Some(value.to_string()),
3230 _ => {}
3231 }
3232 }
3233
3234 fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
3235 let formatted = format!("{value:?}");
3236 let cleaned = formatted
3237 .trim_start_matches('"')
3238 .trim_end_matches('"')
3239 .to_string();
3240 match field.name() {
3241 "message" => self.0.message = Some(cleaned),
3242 "oldest_tx_label" => self.0.oldest_tx_label = Some(cleaned),
3243 "tx_label" => self.0.tx_label = Some(cleaned),
3244 "census_only" => self.0.census_only = Some(cleaned),
3245 _ => {}
3246 }
3247 }
3248 }
3249
3250 struct CaptureSubscriber {
3255 events: std::sync::Arc<std::sync::Mutex<Vec<CapturedEvent>>>,
3256 }
3257
3258 impl tracing::Subscriber for CaptureSubscriber {
3259 fn enabled(&self, _: &tracing::Metadata<'_>) -> bool {
3260 true
3261 }
3262 fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
3263 tracing::span::Id::from_u64(1)
3264 }
3265 fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
3266 fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
3267 fn event(&self, event: &tracing::Event<'_>) {
3268 let mut visitor = CapturedEventVisitor::default();
3269 event.record(&mut visitor);
3270 self.events.lock().unwrap().push(visitor.0);
3271 }
3272 fn enter(&self, _: &tracing::span::Id) {}
3273 fn exit(&self, _: &tracing::span::Id) {}
3274 }
3275
3276 #[test]
3279 #[serial(tx_registry)]
3280 fn log_tx_registry_oldest_debug_reports_oldest_open_entry() {
3281 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
3282 let subscriber = CaptureSubscriber {
3283 events: std::sync::Arc::clone(&buffer),
3284 };
3285
3286 let _handle =
3287 khive_storage::tx_registry::register(Some("checkpoint_tick_test".to_string()));
3288
3289 let oldest = khive_storage::tx_registry::oldest().map(|(id, age, label)| {
3290 khive_storage::tx_registry::OldestSpan {
3291 id,
3292 age,
3293 label,
3294 origin: khive_storage::tx_registry::TxOrigin::Unscoped,
3295 }
3296 });
3297 let expected_label = oldest
3298 .as_ref()
3299 .and_then(|s| s.label.clone())
3300 .unwrap_or_else(|| "<unlabeled>".to_string());
3301
3302 tracing::subscriber::with_default(subscriber, || {
3303 log_tx_registry_oldest_debug(100, oldest.as_ref());
3304 });
3305
3306 let events = buffer.lock().unwrap();
3307 assert!(
3308 events.iter().any(|e| {
3309 e.message.as_deref()
3310 == Some("WAL checkpoint tick: oldest open transaction registry entry")
3311 && e.oldest_tx_label.as_deref() == Some(expected_label.as_str())
3312 }),
3313 "expected a log line naming the open registry entry's label, got: {events:?}"
3314 );
3315 }
3316
3317 #[test]
3323 #[serial(tx_registry)]
3324 fn registry_warns_fire_on_crossing_and_do_not_repeat() {
3325 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
3326 let subscriber = CaptureSubscriber {
3327 events: std::sync::Arc::clone(&buffer),
3328 };
3329
3330 let _handle =
3331 khive_storage::tx_registry::register(Some("registry_warn_crossing_test".to_string()));
3332 let oldest = khive_storage::tx_registry::oldest().map(|(id, age, label)| {
3333 khive_storage::tx_registry::OldestSpan {
3334 id,
3335 age,
3336 label,
3337 origin: khive_storage::tx_registry::TxOrigin::Unscoped,
3338 }
3339 });
3340
3341 let mut was_above_warn = false;
3342 let mut was_above_high_water = false;
3343
3344 tracing::subscriber::with_default(subscriber, || {
3345 if crossing_warn(true, &mut was_above_warn) {
3347 log_tx_registry_oldest_warn(6000, oldest.as_ref());
3348 }
3349 if crossing_warn(true, &mut was_above_high_water) {
3350 log_tx_registry_snapshot_warn(6000);
3351 }
3352
3353 if crossing_warn(true, &mut was_above_warn) {
3355 log_tx_registry_oldest_warn(6000, oldest.as_ref());
3356 }
3357 if crossing_warn(true, &mut was_above_high_water) {
3358 log_tx_registry_snapshot_warn(6000);
3359 }
3360 });
3361
3362 let events = buffer.lock().unwrap();
3363
3364 let oldest_warn_count = events
3374 .iter()
3375 .filter(|e| {
3376 e.message.as_deref()
3377 == Some("WAL checkpoint tick: oldest open transaction registry entry")
3378 })
3379 .count();
3380 assert_eq!(
3381 oldest_warn_count, 1,
3382 "oldest-entry WARN must fire exactly once across two above-threshold ticks, got: {events:?}"
3383 );
3384
3385 let snapshot_warn_count = events
3386 .iter()
3387 .filter(|e| {
3388 e.message.as_deref() == Some("WAL high-water: open transaction registry entry")
3389 && e.tx_label.as_deref() == Some("registry_warn_crossing_test")
3390 })
3391 .count();
3392 assert_eq!(
3393 snapshot_warn_count, 1,
3394 "high-water snapshot WARN must fire exactly once across two above-threshold ticks, got: {events:?}"
3395 );
3396 }
3397
3398 #[test]
3401 fn log_tx_age_emission_carries_label_for_both_rungs() {
3402 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
3403 let subscriber = CaptureSubscriber {
3404 events: std::sync::Arc::clone(&buffer),
3405 };
3406
3407 tracing::subscriber::with_default(subscriber, || {
3408 log_tx_age_emission(&TxAgeEmission {
3409 rung: TxAgeRung::Warn,
3410 age: Duration::from_secs(45),
3411 label: Some("plank1_warn_test".to_string()),
3412 });
3413 log_tx_age_emission(&TxAgeEmission {
3414 rung: TxAgeRung::Stale,
3415 age: Duration::from_secs(150),
3416 label: Some("plank1_stale_test".to_string()),
3417 });
3418 });
3419
3420 let events = buffer.lock().unwrap();
3421 assert!(
3422 events.iter().any(|e| {
3423 e.message.as_deref()
3424 == Some(
3425 "ADR-091 Plank 1: open transaction registry entry exceeded soft-cap age",
3426 )
3427 && e.tx_label.as_deref() == Some("plank1_warn_test")
3428 }),
3429 "expected a Warn-rung log line naming the entry, got: {events:?}"
3430 );
3431 assert!(
3432 events.iter().any(|e| {
3433 e.message.as_deref().is_some_and(|m| {
3434 m.starts_with(
3435 "ADR-091 Plank 1: open transaction registry entry exceeded the cooperative",
3436 )
3437 }) && e.tx_label.as_deref() == Some("plank1_stale_test")
3438 }),
3439 "expected a Stale-rung log line naming the entry, got: {events:?}"
3440 );
3441 }
3442
3443 fn file_pool(path: &std::path::Path) -> Arc<ConnectionPool> {
3444 let cfg = PoolConfig {
3445 path: Some(path.to_path_buf()),
3446 ..PoolConfig::default()
3447 };
3448 Arc::new(ConnectionPool::new(cfg).expect("pool open"))
3449 }
3450
3451 async fn writer_task_wal_autocheckpoint_pages(handle: &WriterTaskHandle) -> u32 {
3452 handle
3453 .send_top_level(|conn| {
3454 conn.pragma_query_value(None, "wal_autocheckpoint", |row| row.get::<_, u32>(0))
3455 .map_err(|error| khive_storage::error::StorageError::Pool {
3456 operation: "test_wal_autocheckpoint".into(),
3457 message: error.to_string(),
3458 })
3459 })
3460 .await
3461 .expect("query writer-task connection pragma")
3462 }
3463
3464 fn checkpoint_conn(pool: &ConnectionPool) -> rusqlite::Connection {
3468 pool.open_standalone_writer()
3469 .expect("open dedicated checkpoint connection")
3470 }
3471
3472 struct TruncateReportHookGuard;
3473
3474 impl Drop for TruncateReportHookGuard {
3475 fn drop(&mut self) {
3476 truncate_report_test_sync::uninstall();
3477 }
3478 }
3479
3480 #[cfg(unix)]
3481 struct WalpinAttributionHookGuard;
3482
3483 #[cfg(unix)]
3484 impl Drop for WalpinAttributionHookGuard {
3485 fn drop(&mut self) {
3486 walpin_attribution_test_sync::uninstall();
3487 }
3488 }
3489
3490 #[test]
3491 #[cfg(unix)]
3492 fn walpin_full_scan_cadence_refreshes_first_then_reuses_until_boundary() {
3493 let cadence = Duration::from_secs(30);
3494 let started_at = Instant::now();
3495 let mut state = TruncateState::with_walpin_full_scan_cadence(cadence);
3496
3497 assert!(matches!(
3498 state.plan_walpin_attribution_at(started_at),
3499 WalpinFullScanPlan::Refresh { .. }
3500 ));
3501 state.cache_walpin_attribution(
3502 crate::walpin::WalpinReport::default(),
3503 Ok(crate::walpin::CensusResult::default()),
3504 started_at,
3505 );
3506
3507 assert!(matches!(
3508 state.plan_walpin_attribution_at(started_at + cadence - Duration::from_nanos(1)),
3509 WalpinFullScanPlan::Cached(_)
3510 ));
3511 assert!(matches!(
3512 state.plan_walpin_attribution_at(started_at + cadence),
3513 WalpinFullScanPlan::Refresh { .. }
3514 ));
3515 }
3516
3517 #[test]
3518 #[cfg(unix)]
3519 fn walpin_full_scan_failure_retries_only_after_cadence() {
3520 let cadence = Duration::from_secs(30);
3521 let started_at = Instant::now();
3522 let mut state = TruncateState::with_walpin_full_scan_cadence(cadence);
3523
3524 assert!(matches!(
3525 state.plan_walpin_attribution_at(started_at),
3526 WalpinFullScanPlan::Refresh { .. }
3527 ));
3528 assert!(matches!(
3531 state.plan_walpin_attribution_at(started_at + cadence - Duration::from_nanos(1)),
3532 WalpinFullScanPlan::Suppressed
3533 ));
3534 assert!(matches!(
3535 state.plan_walpin_attribution_at(started_at + cadence),
3536 WalpinFullScanPlan::Refresh { .. }
3537 ));
3538 }
3539
3540 #[test]
3541 #[cfg(unix)]
3542 fn cached_walpin_report_is_diagnostic_only_even_when_fully_attributed() {
3543 let report = crate::walpin::WalpinReport::default();
3544 assert!(
3545 report.fully_attributed(),
3546 "the fixture must otherwise license the sharp conclusion"
3547 );
3548 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
3549 let subscriber = CaptureSubscriber {
3550 events: std::sync::Arc::clone(&buffer),
3551 };
3552
3553 tracing::subscriber::with_default(subscriber, || {
3554 log_walpin_sidecar_report(
3555 &report,
3556 Ok(crate::walpin::CensusResult::default()),
3557 WalpinReportFreshness::Cached {
3558 age: Duration::from_secs(1),
3559 },
3560 );
3561 });
3562
3563 let events = buffer.lock().unwrap();
3564 assert!(
3565 events.iter().any(|event| {
3566 event.message.as_deref()
3567 == Some(
3568 "cached WAL-pin attribution is diagnostic-only; fully-attributed \
3569 conclusion is not licensed",
3570 )
3571 }),
3572 "cached attribution must declare its fail-closed status: {events:?}"
3573 );
3574 assert!(
3575 !events.iter().any(|event| {
3576 event.message.as_deref().is_some_and(|message| {
3577 message.starts_with("ADR-091 Amendment 2 Plank B: every live PID is reporting")
3578 })
3579 }),
3580 "cached attribution must never authorize the fully-attributed conclusion: {events:?}"
3581 );
3582 }
3583
3584 #[tokio::test(flavor = "current_thread")]
3585 #[cfg(unix)]
3586 #[serial(checkpoint_skip_metrics, khive_walpin_sidecar_env)]
3587 async fn progressing_truncate_releases_full_scan_reservation_to_housekeeping() {
3588 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
3589 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
3590 let dir = tempfile::tempdir().unwrap();
3591 let path = dir.path().join("walpin-progress-reservation.db");
3592 let pool = file_pool(&path);
3593 {
3594 let writer = pool.try_writer().unwrap();
3595 writer
3596 .conn()
3597 .execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
3598 .unwrap();
3599 }
3600 let conn = checkpoint_conn(&pool);
3601 let mut state = TruncateState::default();
3602 let config = CheckpointConfig {
3603 truncate_high_water_pages: 0,
3604 truncate_min_interval: Duration::ZERO,
3605 ..CheckpointConfig::default()
3606 };
3607
3608 assert!(
3609 maybe_truncate(&pool, &conn, &config, u64::MAX, &mut state).is_none(),
3610 "a progressing TRUNCATE must not schedule no-progress attribution"
3611 );
3612 assert!(
3613 state.housekeeping_due(),
3614 "unused pre-TRUNCATE reservation must be restored before housekeeping"
3615 );
3616 let sidecar =
3617 WalpinSidecarState::new(pool.canonical_path(), true, "daemon", config.interval)
3618 .expect("file-backed test sidecar");
3619 assert!(
3620 run_walpin_housekeeping_if_due(&sidecar, &mut state, DEFAULT_SESSION_SWEEP_INTERVAL,)
3621 .await,
3622 "the production housekeeping arm must consume one full scan"
3623 );
3624 assert!(state.walpin_cached_attribution.is_some());
3625 }
3626
3627 #[tokio::test(flavor = "current_thread")]
3628 #[cfg(unix)]
3629 #[serial(checkpoint_skip_metrics, khive_walpin_sidecar_env)]
3630 async fn erroring_truncate_releases_full_scan_reservation_to_housekeeping() {
3631 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
3632 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
3633 let dir = tempfile::tempdir().unwrap();
3634 let path = dir.path().join("walpin-error-reservation.db");
3635 let pool = file_pool(&path);
3636 {
3637 let writer = pool.try_writer().unwrap();
3638 writer
3639 .conn()
3640 .execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
3641 .unwrap();
3642 }
3643 let conn = checkpoint_conn(&pool);
3644 conn.authorizer(Some(
3645 |context: rusqlite::hooks::AuthContext<'_>| match context.action {
3646 AuthAction::Pragma { pragma_name, .. }
3647 if pragma_name.eq_ignore_ascii_case("wal_checkpoint") =>
3648 {
3649 Authorization::Deny
3650 }
3651 _ => Authorization::Allow,
3652 },
3653 ))
3654 .unwrap();
3655 let mut state = TruncateState::default();
3656 let config = CheckpointConfig {
3657 truncate_high_water_pages: 0,
3658 truncate_min_interval: Duration::ZERO,
3659 ..CheckpointConfig::default()
3660 };
3661
3662 assert!(
3663 maybe_truncate(&pool, &conn, &config, u64::MAX, &mut state).is_none(),
3664 "an erroring TRUNCATE must not schedule no-progress attribution"
3665 );
3666 conn.authorizer(None::<fn(rusqlite::hooks::AuthContext<'_>) -> Authorization>)
3667 .unwrap();
3668 assert!(
3669 state.housekeeping_due(),
3670 "failed TRUNCATE must restore its unused full-scan reservation"
3671 );
3672 let sidecar =
3673 WalpinSidecarState::new(pool.canonical_path(), true, "daemon", config.interval)
3674 .expect("file-backed test sidecar");
3675 assert!(
3676 run_walpin_housekeeping_if_due(&sidecar, &mut state, DEFAULT_SESSION_SWEEP_INTERVAL,)
3677 .await,
3678 "the production housekeeping arm must consume one full scan"
3679 );
3680 assert!(state.walpin_cached_attribution.is_some());
3681 }
3682
3683 struct ReaderProcess {
3684 child: std::process::Child,
3685 _stdout: std::io::BufReader<std::process::ChildStdout>,
3686 }
3687
3688 impl ReaderProcess {
3689 fn spawn(db_path: &std::path::Path) -> Self {
3690 use std::io::BufRead;
3691 use std::process::Stdio;
3692
3693 let mut child = std::process::Command::new(
3694 std::env::current_exe().expect("resolve current test executable"),
3695 )
3696 .args([
3697 "--exact",
3698 "checkpoint::tests::walpin_transient_reader_process_helper",
3699 "--nocapture",
3700 ])
3701 .env("KHIVE_CHECKPOINT_READER_HELPER_PATH", db_path)
3702 .stdin(Stdio::piped())
3703 .stdout(Stdio::piped())
3704 .spawn()
3705 .expect("spawn transient WAL reader helper");
3706
3707 let stdout = child.stdout.take().expect("capture helper stdout");
3708 let mut reader = std::io::BufReader::new(stdout);
3709 let mut line = String::new();
3710 loop {
3711 line.clear();
3712 let bytes = reader
3713 .read_line(&mut line)
3714 .expect("read transient reader readiness signal");
3715 assert!(bytes > 0, "reader helper exited before readiness signal");
3716 if line.contains("KHIVE_CHECKPOINT_READER_READY") {
3717 break;
3718 }
3719 }
3720 Self {
3721 child,
3722 _stdout: reader,
3723 }
3724 }
3725
3726 fn pid(&self) -> u32 {
3727 self.child.id()
3728 }
3729
3730 fn release(&mut self) {
3731 use std::io::Write;
3732
3733 let mut stdin = self.child.stdin.take().expect("helper stdin is available");
3734 stdin
3735 .write_all(b"release\n")
3736 .expect("release transient reader");
3737 drop(stdin);
3738 let status = self.child.wait().expect("wait for transient reader helper");
3739 assert!(status.success(), "transient reader helper failed: {status}");
3740 }
3741 }
3742
3743 impl Drop for ReaderProcess {
3744 fn drop(&mut self) {
3745 if self.child.try_wait().ok().flatten().is_none() {
3746 let _ = self.child.kill();
3747 let _ = self.child.wait();
3748 }
3749 }
3750 }
3751
3752 #[test]
3753 fn walpin_transient_reader_process_helper() {
3754 use std::io::Write;
3755
3756 let Some(path) = std::env::var_os("KHIVE_CHECKPOINT_READER_HELPER_PATH") else {
3757 return;
3758 };
3759 let conn = rusqlite::Connection::open(path).expect("helper opens database");
3760 conn.execute_batch("BEGIN DEFERRED; SELECT * FROM t;")
3761 .expect("helper pins a read snapshot");
3762 println!("KHIVE_CHECKPOINT_READER_READY");
3763 std::io::stdout().flush().expect("flush readiness signal");
3764 let mut release = String::new();
3765 std::io::stdin()
3766 .read_line(&mut release)
3767 .expect("wait for release signal");
3768 conn.execute_batch("COMMIT")
3769 .expect("helper releases read snapshot");
3770 }
3771
3772 #[tokio::test(flavor = "current_thread")]
3773 #[cfg(unix)]
3774 #[serial(
3775 checkpoint_skip_metrics,
3776 khive_walpin_sidecar_env,
3777 walpin_attribution_async,
3778 walpin_report_seam
3779 )]
3780 async fn no_progress_report_keeps_holder_released_after_truncate_timeout() {
3781 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
3782 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
3783 let dir = tempfile::tempdir().expect("tempdir");
3784 let path = dir.path().join("transient-reader.db");
3785 let pool = file_pool(&path);
3786 {
3787 let writer = pool.try_writer().expect("writer");
3788 writer
3789 .conn()
3790 .execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
3791 .expect("seed WAL before reader snapshot");
3792 }
3793
3794 let mut reader = ReaderProcess::spawn(&path);
3795 let reader_pid = reader.pid();
3796 {
3797 let writer = pool.try_writer().expect("writer");
3798 writer
3799 .conn()
3800 .execute_batch("INSERT INTO t VALUES (2);")
3801 .expect("append WAL behind reader snapshot");
3802 }
3803
3804 let canonical_path = pool
3805 .canonical_path()
3806 .expect("file-backed pool has canonical path")
3807 .to_path_buf();
3808 let (reached_rx, proceed_tx) = truncate_report_test_sync::install(canonical_path.clone());
3809 let _hook_guard = TruncateReportHookGuard;
3810 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
3811 let checkpoint_pool = Arc::clone(&pool);
3812 let dedicated_conn = checkpoint_conn(&checkpoint_pool);
3813 let checkpoint = std::thread::spawn(move || {
3814 let mut state = TruncateState::default();
3815 let result = checkpoint_once_core(
3816 &checkpoint_pool,
3817 &dedicated_conn,
3818 &CheckpointConfig {
3819 truncate_high_water_pages: 0,
3820 truncate_min_interval: Duration::ZERO,
3821 truncate_busy_timeout: Duration::from_millis(50),
3822 ..CheckpointConfig::default()
3823 },
3824 &mut state,
3825 );
3826 (result, state)
3827 });
3828
3829 reached_rx
3830 .recv_timeout(Duration::from_secs(5))
3831 .expect("TRUNCATE must report no progress while the reader is pinned");
3832 reader.release();
3833 let post_attempt_census =
3834 crate::walpin::census_holders(&canonical_path).expect("post-attempt holder census");
3835 assert!(
3836 !post_attempt_census.holders.contains(&reader_pid),
3837 "released reader PID must be absent from a post-attempt census"
3838 );
3839 proceed_tx
3840 .send(())
3841 .expect("allow no-progress reporting to continue");
3842 let (checkpoint_result, mut state) = checkpoint.join().expect("checkpoint thread");
3843 let outcome = checkpoint_result.expect("checkpoint succeeds");
3844 assert!(
3845 outcome.sidecar_attribution.is_some(),
3846 "the synchronous checkpoint result must carry a separate attribution request"
3847 );
3848 assert!(
3849 !state.sidecar_attribution_attempted_this_tick,
3850 "capturing a request is not the same as attempting its directory enumeration"
3851 );
3852
3853 let subscriber = CaptureSubscriber {
3854 events: std::sync::Arc::clone(&buffer),
3855 };
3856 let _subscriber_guard = tracing::subscriber::set_default(subscriber);
3857 let attribution_attempted =
3858 complete_walpin_attribution(outcome.sidecar_attribution, &mut state)
3859 .await
3860 .expect("deferred attribution succeeds");
3861 assert!(
3862 attribution_attempted,
3863 "a no-progress attribution pass must suppress the redundant healthy-housekeeping \
3864 pass for the same tick"
3865 );
3866
3867 let events = buffer.lock().expect("captured events");
3868 assert!(
3869 events.iter().any(|event| {
3870 event
3871 .census_only
3872 .as_deref()
3873 .is_some_and(|pids| pids.contains(&reader_pid.to_string()))
3874 }),
3875 "the no-progress report must retain PID {reader_pid} from the pre-attempt census: {events:?}"
3876 );
3877 }
3878
3879 #[tokio::test(flavor = "current_thread")]
3885 #[cfg(unix)]
3886 #[serial(walpin_attribution_async)]
3887 async fn no_progress_attribution_is_off_runtime_and_awaited_before_report_use() {
3888 use std::sync::atomic::Ordering;
3889
3890 let dir = tempfile::tempdir().expect("tempdir");
3891 let sidecar_dir = dir.path().join("checkpoint.db.walpin");
3892 let (reached_rx, proceed_tx, report_counter) =
3893 walpin_attribution_test_sync::install_pause(sidecar_dir.clone());
3894 let _hook_guard = WalpinAttributionHookGuard;
3895
3896 let runtime_thread = std::thread::current().id();
3897 let mut state = TruncateState::default();
3898 let request = Some(WalpinAttributionRequest::Fresh {
3899 dir: sidecar_dir,
3900 census: Ok(crate::walpin::CensusResult::default()),
3901 legacy_fallback_interval: DEFAULT_SESSION_SWEEP_INTERVAL,
3902 previous_last_attempt: None,
3903 });
3904 let completion = tokio::spawn(async move {
3905 let result = complete_walpin_attribution(request, &mut state).await;
3906 (result, state)
3907 });
3908
3909 let blocking_thread = reached_rx
3910 .await
3911 .expect("spawn_blocking attribution reached test seam");
3912 assert_ne!(
3913 blocking_thread, runtime_thread,
3914 "sidecar enumeration must not execute on the current-thread Tokio runtime worker"
3915 );
3916 assert!(
3917 !completion.is_finished(),
3918 "the async attribution owner must await the still-paused blocking enumeration"
3919 );
3920 assert_eq!(
3921 report_counter.load(Ordering::SeqCst),
3922 0,
3923 "the attribution report must not be consumed before enumeration completes"
3924 );
3925
3926 proceed_tx
3927 .send(())
3928 .expect("release blocking attribution enumeration");
3929 let (result, state) = completion.await.expect("attribution task joins");
3930 assert_eq!(result, Ok(true));
3931 assert_eq!(
3932 report_counter.load(Ordering::SeqCst),
3933 1,
3934 "the completed enumeration must feed exactly one report use"
3935 );
3936 assert!(state.sidecar_attribution_attempted_this_tick);
3937 assert!(
3938 !state.housekeeping_due(),
3939 "completed attribution must suppress same-tick housekeeping"
3940 );
3941 }
3942
3943 #[tokio::test(flavor = "current_thread")]
3948 #[cfg(unix)]
3949 #[serial(walpin_attribution_async)]
3950 async fn no_progress_attribution_join_failure_is_honest_and_suppresses_retry() {
3951 let dir = tempfile::tempdir().expect("tempdir");
3952 let sidecar_dir = dir.path().join("checkpoint.db.walpin");
3953 walpin_attribution_test_sync::install_panic(sidecar_dir.clone());
3954 let _hook_guard = WalpinAttributionHookGuard;
3955
3956 let mut state = TruncateState::default();
3957 let request = Some(WalpinAttributionRequest::Fresh {
3958 dir: sidecar_dir,
3959 census: Ok(crate::walpin::CensusResult::default()),
3960 legacy_fallback_interval: DEFAULT_SESSION_SWEEP_INTERVAL,
3961 previous_last_attempt: None,
3962 });
3963
3964 let error = complete_walpin_attribution(request, &mut state)
3965 .await
3966 .expect_err("injected worker panic must surface as failure");
3967 assert!(
3968 matches!(error, WalpinAttributionFailure::Worker(_)),
3969 "join failure must retain its worker classification: {error:?}"
3970 );
3971 assert!(state.sidecar_attribution_attempted_this_tick);
3972 assert!(
3973 !state.housekeeping_due(),
3974 "an indeterminate partial pass must not authorize a second scan"
3975 );
3976 }
3977
3978 #[test]
3983 #[cfg(unix)]
3984 #[serial(checkpoint_skip_metrics)]
3985 fn async_checkpoint_source_keeps_enumeration_behind_awaited_spawn_blocking() {
3986 fn section<'a>(source: &'a str, start: &str, end: &str) -> &'a str {
3987 source
3988 .split_once(start)
3989 .unwrap_or_else(|| panic!("missing source marker {start:?}"))
3990 .1
3991 .split_once(end)
3992 .unwrap_or_else(|| panic!("missing source marker {end:?}"))
3993 .0
3994 }
3995
3996 let source = include_str!("checkpoint.rs");
3997 let checkpoint_core = section(
3998 source,
3999 "fn checkpoint_once_core(",
4000 "/// Evaluate and, if due, attempt a TRUNCATE escalation",
4001 );
4002 let truncate_core = section(
4003 source,
4004 "fn maybe_truncate(",
4005 "#[cfg(test)]\nmod truncate_report_test_sync",
4006 );
4007 let report_logger = section(
4008 source,
4009 "fn log_walpin_sidecar_report(",
4010 "/// ADR-091 Amendment 2 Plank C",
4011 );
4012 for (name, body) in [
4013 ("checkpoint_once_core", checkpoint_core),
4014 ("maybe_truncate", truncate_core),
4015 ("log_walpin_sidecar_report", report_logger),
4016 ] {
4017 assert!(
4018 !body.contains("enumerate_live("),
4019 "{name} must not perform direct sidecar enumeration"
4020 );
4021 }
4022
4023 let async_completion = section(
4024 source,
4025 "async fn complete_walpin_attribution(",
4026 "/// When a TRUNCATE attempt makes no progress",
4027 );
4028 let spawn = async_completion
4029 .find("tokio::task::spawn_blocking")
4030 .expect("completion must spawn blocking work");
4031 let enumerate = async_completion
4032 .find("crate::walpin::enumerate_live")
4033 .expect("blocking closure must perform the attribution enumeration");
4034 let awaited = async_completion[enumerate..]
4035 .find(".await")
4036 .map(|offset| enumerate + offset)
4037 .expect("blocking worker must be awaited");
4038 assert!(spawn < enumerate && enumerate < awaited);
4039
4040 let task = section(
4041 source,
4042 "pub async fn run_checkpoint_task(",
4043 "/// Whether a `CheckpointOutcomeRecorded` transition should be enqueued",
4044 );
4045 let checkpoint = task
4046 .find("checkpoint_once_core(")
4047 .expect("checkpoint core call");
4048 let completion = task
4049 .find("complete_walpin_attribution(")
4050 .expect("awaited attribution completion");
4051 let housekeeping = task
4052 .find("run_walpin_housekeeping_if_due(")
4053 .expect("fallback housekeeping");
4054 let outcome = task
4055 .find("observe_checkpoint_pressure_tick(")
4056 .expect("lifecycle outcome use");
4057 assert!(
4058 checkpoint < completion && completion < housekeeping && housekeeping < outcome,
4059 "tick ordering must be checkpoint -> awaited attribution -> housekeeping decision -> outcome"
4060 );
4061 let housekeeping_helper = section(
4062 source,
4063 "async fn run_walpin_housekeeping_if_due(",
4064 "fn now_epoch_secs()",
4065 );
4066 assert!(
4067 housekeeping_helper.contains("reap_dead_entries_bounded(legacy_fallback_interval)"),
4068 "the ordered housekeeping arm must retain the bounded full scan"
4069 );
4070
4071 let pressure_tick = section(
4076 source,
4077 "fn observe_checkpoint_pressure_tick(",
4078 "/// ADR-091 Plank 0",
4079 );
4080 assert!(
4081 pressure_tick.contains("checkpoint_outcome_should_emit"),
4082 "extracted pressure tick helper must gate on the lifecycle emit decision"
4083 );
4084 }
4085
4086 #[test]
4092 #[serial(checkpoint_skip_metrics)]
4093 fn checkpoint_once_succeeds_on_file_backed_pool() {
4094 let dir = tempfile::tempdir().unwrap();
4095 let path = dir.path().join("wal_test.db");
4096 let pool = file_pool(&path);
4097
4098 {
4100 let writer = pool.try_writer().unwrap();
4101 writer
4102 .conn()
4103 .execute_batch("CREATE TABLE IF NOT EXISTS t (x INTEGER);")
4104 .unwrap();
4105 writer
4106 .conn()
4107 .execute_batch("INSERT INTO t VALUES (1);")
4108 .unwrap();
4109 }
4110
4111 let conn = checkpoint_conn(&pool);
4112 checkpoint_once(
4113 &pool,
4114 &conn,
4115 &CheckpointConfig::default(),
4116 &mut TruncateState::default(),
4117 )
4118 .expect("checkpoint_once must succeed against a healthy dedicated connection");
4119 }
4120
4121 #[test]
4127 fn open_standalone_writer_fails_on_in_memory_pool() {
4128 let cfg = PoolConfig {
4129 path: None,
4130 ..PoolConfig::default()
4131 };
4132 let pool = Arc::new(ConnectionPool::new(cfg).expect("in-memory pool"));
4133 assert!(
4134 pool.open_standalone_writer().is_err(),
4135 "an in-memory pool must not be able to open a dedicated checkpoint connection"
4136 );
4137 }
4138
4139 #[test]
4145 fn ensure_open_does_not_move_writer_acquisition_counters() {
4146 let dir = tempfile::tempdir().unwrap();
4147 let path = dir.path().join("checkpoint_ensure_open.db");
4148 let pool = file_pool(&path);
4149
4150 let before_first_open = pool.writer_acquisition_snapshot();
4151 let mut checkpoint_conn = CheckpointConnection::new();
4152 checkpoint_conn
4153 .ensure_open(&pool)
4154 .expect("dedicated checkpoint connection must open against a file-backed pool");
4155 assert_eq!(
4156 pool.writer_acquisition_snapshot(),
4157 before_first_open,
4158 "the checkpoint connection's initial open must not count as a writer acquisition"
4159 );
4160
4161 checkpoint_conn.conn = None;
4162 let before_reopen = pool.writer_acquisition_snapshot();
4163 checkpoint_conn
4164 .ensure_open(&pool)
4165 .expect("dedicated checkpoint connection must reopen after invalidation");
4166 assert_eq!(
4167 pool.writer_acquisition_snapshot(),
4168 before_reopen,
4169 "reopening the checkpoint connection must not count as a writer acquisition either"
4170 );
4171 }
4172
4173 #[test]
4174 fn checkpoint_connection_disables_wal_autocheckpoint_on_open_and_reopen() {
4175 let dir = tempfile::tempdir().unwrap();
4176 let path = dir.path().join("checkpoint_autocheckpoint.db");
4177 let pool = file_pool(&path);
4178 let mut checkpoint_conn = CheckpointConnection::new();
4179
4180 let initial: u32 = checkpoint_conn
4181 .ensure_open(&pool)
4182 .expect("dedicated checkpoint connection must open")
4183 .pragma_query_value(None, "wal_autocheckpoint", |row| row.get(0))
4184 .expect("read initial autocheckpoint setting");
4185 assert_eq!(initial, 0);
4186
4187 checkpoint_conn.conn = None;
4188 let reopened: u32 = checkpoint_conn
4189 .ensure_open(&pool)
4190 .expect("dedicated checkpoint connection must reopen")
4191 .pragma_query_value(None, "wal_autocheckpoint", |row| row.get(0))
4192 .expect("read reopened autocheckpoint setting");
4193 assert_eq!(reopened, 0);
4194 }
4195
4196 #[tokio::test(flavor = "current_thread")]
4197 #[serial(checkpoint_skip_metrics)]
4198 async fn failed_checkpoint_claim_keeps_existing_writer_task_on_fallback() {
4199 let dir = tempfile::tempdir().unwrap();
4200 let path = dir.path().join("failed_claim_writer_task.db");
4201 let pool = Arc::new(
4202 ConnectionPool::new(PoolConfig {
4203 path: Some(path),
4204 checkout_timeout: Duration::from_millis(1),
4205 write_queue_enabled: Some(true),
4206 ..PoolConfig::default()
4207 })
4208 .expect("pool open"),
4209 );
4210 let writer_task = pool
4211 .writer_task_handle()
4212 .expect("writer-task resolution")
4213 .expect("writer task enabled");
4214 assert_eq!(
4215 writer_task_wal_autocheckpoint_pages(&writer_task).await,
4216 crate::pool::FALLBACK_WAL_AUTOCHECKPOINT_PAGES
4217 );
4218
4219 let legacy_conn = pool.legacy_conn();
4220 let (held_tx, held_rx) = tokio::sync::oneshot::channel();
4221 let (release_tx, release_rx) = std::sync::mpsc::channel();
4222 let holder = tokio::task::spawn_blocking(move || {
4223 let _held_writer = legacy_conn.lock();
4224 held_tx.send(()).expect("signal held pooled writer");
4225 release_rx.recv().expect("release held pooled writer");
4226 });
4227 held_rx.await.expect("pooled writer holder started");
4228 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4229 drop(shutdown_tx);
4230 run_checkpoint_task(
4231 Arc::clone(&pool),
4232 CheckpointConfig {
4233 interval: Duration::from_secs(60),
4234 ..CheckpointConfig::default()
4235 },
4236 None,
4237 shutdown_rx,
4238 true,
4239 )
4240 .await;
4241
4242 assert_eq!(pool.writer_acquisition_snapshot().timeouts, 1);
4243 assert_eq!(
4244 writer_task_wal_autocheckpoint_pages(&writer_task).await,
4245 crate::pool::FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
4246 "failed pooled-writer claim must not partially propagate ownership"
4247 );
4248 release_tx.send(()).expect("release pooled writer");
4249 holder.await.expect("pooled writer holder joined");
4250 }
4251
4252 #[tokio::test]
4253 #[serial(checkpoint_skip_metrics)]
4254 async fn checkpoint_task_exits_on_shutdown_signal() {
4255 let dir = tempfile::tempdir().unwrap();
4256 let path = dir.path().join("wal_task_shutdown.db");
4257 let pool = file_pool(&path);
4258
4259 let cfg = CheckpointConfig {
4261 interval: Duration::from_millis(10),
4262 ..Default::default()
4263 };
4264
4265 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4266 let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
4267
4268 shutdown_tx.send(()).expect("send shutdown signal");
4269
4270 tokio::time::timeout(Duration::from_secs(1), handle)
4271 .await
4272 .expect("checkpoint task should exit within 1s")
4273 .expect("checkpoint task panicked");
4274 }
4275
4276 #[cfg(unix)]
4277 #[tokio::test]
4278 #[serial(checkpoint_skip_metrics, khive_walpin_sidecar_env)]
4279 async fn healthy_checkpoint_tick_reaps_a_dead_walpin_beacon_without_truncate() {
4280 let dir = tempfile::tempdir().unwrap();
4281 let path = dir.path().join("healthy_sidecar_reap.db");
4282 let pool = file_pool(&path);
4283 let sidecar_dir =
4284 crate::walpin::sidecar_dir_for(pool.canonical_path().expect("file-backed pool"));
4285 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
4286 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
4287
4288 let dead_pid = 2_000_000_000;
4289 let dead_beacon = crate::walpin::WalpinBeacon {
4290 pid: dead_pid,
4291 process_role: "session".to_string(),
4292 started_at: 1,
4293 sweep_interval_ms: 5_000,
4294 };
4295 crate::walpin::write_beacon(&sidecar_dir, &dead_beacon)
4296 .expect("seed a crashed process's orphan beacon");
4297 let dead_beacon_path = crate::walpin::beacon_path(&sidecar_dir, dead_pid);
4298 assert!(dead_beacon_path.exists(), "orphan fixture must exist");
4299
4300 let cfg = CheckpointConfig {
4301 interval: Duration::from_millis(10),
4302 warn_pages: u64::MAX,
4303 high_water_pages: u64::MAX,
4304 truncate_high_water_pages: u64::MAX,
4305 ..CheckpointConfig::default()
4306 };
4307 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4308 let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
4309
4310 let reaped = wait_for(Duration::from_secs(2), || !dead_beacon_path.exists()).await;
4311 shutdown_tx.send(()).expect("send shutdown signal");
4312 tokio::time::timeout(Duration::from_secs(1), handle)
4313 .await
4314 .expect("checkpoint task should exit within 1s")
4315 .expect("checkpoint task panicked");
4316
4317 assert!(
4318 reaped,
4319 "the ordinary healthy tick must reap positively dead sidecar residue independently of \
4320 TRUNCATE diagnostics"
4321 );
4322 }
4323
4324 #[tokio::test]
4328 #[serial(checkpoint_skip_metrics)]
4329 async fn checkpoint_task_exits_via_shutdown_signal_with_live_event_store_pool_clone() {
4330 let dir = tempfile::tempdir().unwrap();
4331 let path = dir.path().join("wal_task_event_store.db");
4332 let pool = file_pool(&path);
4333
4334 let cfg = CheckpointConfig {
4335 interval: Duration::from_millis(10),
4336 ..Default::default()
4337 };
4338
4339 let event_store: Arc<dyn khive_storage::EventStore> =
4340 Arc::new(crate::stores::event::SqlEventStore::new_scoped(
4341 Arc::clone(&pool),
4342 true,
4343 "local".to_string(),
4344 ));
4345 let sibling_pool_clone = Arc::clone(&pool);
4350
4351 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4352 let handle = tokio::spawn(run_checkpoint_task(
4353 pool,
4354 cfg,
4355 Some(CheckpointLifecycleOwner::new(event_store, "local")),
4356 shutdown_rx,
4357 true,
4358 ));
4359
4360 assert!(
4364 Arc::strong_count(&sibling_pool_clone) > 1,
4365 "test setup must reproduce the multi-owner shape the bug depends on"
4366 );
4367
4368 shutdown_tx.send(()).expect("send shutdown signal");
4369
4370 tokio::time::timeout(Duration::from_secs(1), handle)
4371 .await
4372 .expect(
4373 "checkpoint task should exit within 1s via the watch signal, \
4374 even with a live sibling Arc<ConnectionPool> clone held by \
4375 the event store",
4376 )
4377 .expect("checkpoint task panicked");
4378 }
4379
4380 #[test]
4381 #[serial]
4382 fn checkpoint_config_env_override() {
4383 std::env::set_var("KHIVE_CHECKPOINT_INTERVAL_MS", "250");
4384 std::env::set_var("KHIVE_WAL_WARN_PAGES", "1500");
4385 std::env::set_var("KHIVE_WAL_HIGH_WATER_PAGES", "8000");
4386 std::env::set_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES", "12000");
4387 std::env::set_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS", "60");
4388 std::env::set_var("KHIVE_WAL_TRUNCATE_BUSY_MS", "500");
4389 std::env::set_var("KHIVE_TX_WARN_SECS", "15");
4390 std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "90");
4391
4392 let cfg = CheckpointConfig::from_env();
4393
4394 std::env::remove_var("KHIVE_CHECKPOINT_INTERVAL_MS");
4395 std::env::remove_var("KHIVE_WAL_WARN_PAGES");
4396 std::env::remove_var("KHIVE_WAL_HIGH_WATER_PAGES");
4397 std::env::remove_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES");
4398 std::env::remove_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS");
4399 std::env::remove_var("KHIVE_WAL_TRUNCATE_BUSY_MS");
4400 std::env::remove_var("KHIVE_TX_WARN_SECS");
4401 std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
4402
4403 assert_eq!(cfg.interval, Duration::from_millis(250));
4404 assert_eq!(cfg.warn_pages, 1500);
4405 assert_eq!(cfg.high_water_pages, 8000);
4406 assert_eq!(cfg.truncate_high_water_pages, 12000);
4407 assert_eq!(cfg.truncate_min_interval, Duration::from_secs(60));
4408 assert_eq!(cfg.truncate_busy_timeout, Duration::from_millis(500));
4409 assert_eq!(cfg.tx_warn_secs, Duration::from_secs(15));
4410 assert_eq!(cfg.tx_max_age_secs, Duration::from_secs(90));
4411 }
4412
4413 #[test]
4414 #[serial]
4415 fn checkpoint_config_defaults_on_invalid_env() {
4416 let default = CheckpointConfig::default();
4417
4418 std::env::set_var("KHIVE_CHECKPOINT_INTERVAL_MS", "not_a_number");
4419 std::env::set_var("KHIVE_WAL_WARN_PAGES", "");
4420 std::env::set_var("KHIVE_WAL_HIGH_WATER_PAGES", "0");
4421 std::env::set_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES", "not_a_number");
4422 std::env::set_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS", "");
4423 std::env::set_var("KHIVE_WAL_TRUNCATE_BUSY_MS", "0");
4424 std::env::set_var("KHIVE_TX_WARN_SECS", "not_a_number");
4425 std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "0");
4426
4427 let cfg = CheckpointConfig::from_env();
4428
4429 std::env::remove_var("KHIVE_CHECKPOINT_INTERVAL_MS");
4430 std::env::remove_var("KHIVE_WAL_WARN_PAGES");
4431 std::env::remove_var("KHIVE_WAL_HIGH_WATER_PAGES");
4432 std::env::remove_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES");
4433 std::env::remove_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS");
4434 std::env::remove_var("KHIVE_WAL_TRUNCATE_BUSY_MS");
4435 std::env::remove_var("KHIVE_TX_WARN_SECS");
4436 std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
4437
4438 assert_eq!(cfg.interval, default.interval);
4439 assert_eq!(cfg.warn_pages, default.warn_pages);
4440 assert_eq!(cfg.high_water_pages, default.high_water_pages);
4441 assert_eq!(
4442 cfg.truncate_high_water_pages,
4443 default.truncate_high_water_pages
4444 );
4445 assert_eq!(cfg.truncate_min_interval, default.truncate_min_interval);
4446 assert_eq!(cfg.truncate_busy_timeout, default.truncate_busy_timeout);
4447 assert_eq!(cfg.tx_warn_secs, default.tx_warn_secs);
4448 assert_eq!(cfg.tx_max_age_secs, default.tx_max_age_secs);
4449 }
4450
4451 #[test]
4456 #[serial(checkpoint_skip_metrics)]
4457 fn checkpoint_high_water_does_not_block_behind_reader() {
4458 let dir = tempfile::tempdir().unwrap();
4459 let path = dir.path().join("high_water_test.db");
4460
4461 let pool = Arc::new(
4465 ConnectionPool::new(PoolConfig {
4466 path: Some(path.clone()),
4467 busy_timeout: Duration::from_millis(2000),
4468 ..PoolConfig::default()
4469 })
4470 .expect("pool open"),
4471 );
4472
4473 {
4475 let writer = pool.try_writer().unwrap();
4476 writer
4477 .conn()
4478 .execute_batch(
4479 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
4480 )
4481 .unwrap();
4482 }
4483
4484 let reader = pool.reader().expect("reader");
4488 reader
4489 .execute_batch("BEGIN DEFERRED; SELECT * FROM t;")
4490 .expect("begin read tx");
4491
4492 {
4496 let writer = pool.try_writer().unwrap();
4497 writer
4498 .conn()
4499 .execute_batch("INSERT INTO t VALUES (2);")
4500 .unwrap();
4501 }
4502
4503 let conn = checkpoint_conn(&pool);
4504 let start = std::time::Instant::now();
4505 checkpoint_once(
4506 &pool,
4507 &conn,
4508 &CheckpointConfig::default(),
4509 &mut TruncateState::default(),
4510 )
4511 .expect("checkpoint_once must succeed against a healthy dedicated connection");
4512 let elapsed = start.elapsed();
4513
4514 reader.execute_batch("COMMIT;").ok();
4516 drop(reader);
4517
4518 assert!(
4522 elapsed < std::time::Duration::from_millis(500),
4523 "checkpoint_once with active reader snapshot took {:?}; \
4524 expected <500ms (PASSIVE must not block on readers; \
4525 a TRUNCATE regression would block ~2000ms)",
4526 elapsed
4527 );
4528 }
4529
4530 #[test]
4531 #[serial]
4532 fn checkpoint_config_rejects_zero_for_all_fields() {
4533 let default = CheckpointConfig::default();
4534 std::env::set_var("KHIVE_CHECKPOINT_INTERVAL_MS", "0");
4535 std::env::set_var("KHIVE_WAL_WARN_PAGES", "0");
4536 std::env::set_var("KHIVE_WAL_HIGH_WATER_PAGES", "0");
4537 std::env::set_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES", "0");
4538 std::env::set_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS", "0");
4539 std::env::set_var("KHIVE_WAL_TRUNCATE_BUSY_MS", "0");
4540 std::env::set_var("KHIVE_TX_WARN_SECS", "0");
4541 std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "0");
4542
4543 let cfg = CheckpointConfig::from_env();
4544
4545 std::env::remove_var("KHIVE_CHECKPOINT_INTERVAL_MS");
4546 std::env::remove_var("KHIVE_WAL_WARN_PAGES");
4547 std::env::remove_var("KHIVE_WAL_HIGH_WATER_PAGES");
4548 std::env::remove_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES");
4549 std::env::remove_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS");
4550 std::env::remove_var("KHIVE_WAL_TRUNCATE_BUSY_MS");
4551 std::env::remove_var("KHIVE_TX_WARN_SECS");
4552 std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
4553
4554 assert_eq!(
4555 cfg.interval, default.interval,
4556 "zero interval must fall back to default"
4557 );
4558 assert_eq!(
4559 cfg.warn_pages, default.warn_pages,
4560 "zero warn_pages must fall back to default"
4561 );
4562 assert_eq!(
4563 cfg.high_water_pages, default.high_water_pages,
4564 "zero high_water_pages must fall back to default"
4565 );
4566 assert_eq!(
4567 cfg.truncate_high_water_pages, default.truncate_high_water_pages,
4568 "zero truncate_high_water_pages must fall back to default"
4569 );
4570 assert_eq!(
4571 cfg.truncate_min_interval, default.truncate_min_interval,
4572 "zero truncate_min_interval must fall back to default"
4573 );
4574 assert_eq!(
4575 cfg.truncate_busy_timeout, default.truncate_busy_timeout,
4576 "zero truncate_busy_timeout must fall back to default"
4577 );
4578 assert_eq!(
4579 cfg.tx_warn_secs, default.tx_warn_secs,
4580 "zero tx_warn_secs must fall back to default"
4581 );
4582 assert_eq!(
4583 cfg.tx_max_age_secs, default.tx_max_age_secs,
4584 "zero tx_max_age_secs must fall back to default"
4585 );
4586 }
4587
4588 #[test]
4591 #[serial]
4592 fn checkpoint_config_rejects_reversed_tx_thresholds() {
4593 let default = CheckpointConfig::default();
4594 std::env::set_var("KHIVE_TX_WARN_SECS", "120");
4595 std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "30");
4596
4597 let cfg = CheckpointConfig::from_env();
4598
4599 std::env::remove_var("KHIVE_TX_WARN_SECS");
4600 std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
4601
4602 assert_eq!(
4603 cfg.tx_warn_secs, default.tx_warn_secs,
4604 "a reversed pair must fall back tx_warn_secs to its default, got: {:?}",
4605 cfg.tx_warn_secs
4606 );
4607 assert_eq!(
4608 cfg.tx_max_age_secs, default.tx_max_age_secs,
4609 "a reversed pair must fall back tx_max_age_secs to its default, got: {:?}",
4610 cfg.tx_max_age_secs
4611 );
4612 }
4613
4614 #[test]
4617 #[serial]
4618 fn checkpoint_config_rejects_equal_tx_thresholds() {
4619 let default = CheckpointConfig::default();
4620 std::env::set_var("KHIVE_TX_WARN_SECS", "60");
4621 std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "60");
4622
4623 let cfg = CheckpointConfig::from_env();
4624
4625 std::env::remove_var("KHIVE_TX_WARN_SECS");
4626 std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
4627
4628 assert_eq!(
4629 cfg.tx_warn_secs, default.tx_warn_secs,
4630 "an equal pair must fall back tx_warn_secs to its default, got: {:?}",
4631 cfg.tx_warn_secs
4632 );
4633 assert_eq!(
4634 cfg.tx_max_age_secs, default.tx_max_age_secs,
4635 "an equal pair must fall back tx_max_age_secs to its default, got: {:?}",
4636 cfg.tx_max_age_secs
4637 );
4638 }
4639
4640 #[test]
4643 fn skipped_tick_does_not_reset_high_water_crossing_state() {
4644 let mut was_above = false;
4645
4646 assert!(
4648 crossing_warn(true, &mut was_above),
4649 "should fire on first crossing"
4650 );
4651 assert!(was_above);
4652
4653 assert!(was_above, "was_above must stay true across skipped ticks");
4660
4661 let fired = crossing_warn(true, &mut was_above);
4663 assert!(!fired, "WARN must not re-fire while still above threshold");
4664
4665 let fired = crossing_warn(false, &mut was_above);
4667 assert!(!fired);
4668 assert!(!was_above);
4669
4670 let fired = crossing_warn(true, &mut was_above);
4672 assert!(fired, "WARN must fire again on a new below→above crossing");
4673 }
4674
4675 #[test]
4682 fn warn_pages_fires_once_on_crossing_not_every_tick() {
4683 let mut was_above_warn = false;
4684
4685 let fired_1 = crossing_warn(true, &mut was_above_warn);
4687 let fired_2 = crossing_warn(true, &mut was_above_warn);
4688 let fired_3 = crossing_warn(true, &mut was_above_warn);
4689
4690 assert!(fired_1, "WARN must fire on the first in-band tick");
4691 assert!(
4692 !fired_2,
4693 "WARN must not fire on the second consecutive in-band tick"
4694 );
4695 assert!(
4696 !fired_3,
4697 "WARN must not fire on the third consecutive in-band tick"
4698 );
4699
4700 crossing_warn(false, &mut was_above_warn);
4702 assert!(!was_above_warn);
4703
4704 let fired_reentry = crossing_warn(true, &mut was_above_warn);
4706 assert!(
4707 fired_reentry,
4708 "WARN must fire again on re-entry into warn band"
4709 );
4710 }
4711
4712 #[test]
4718 #[serial(tx_registry, checkpoint_skip_metrics)]
4719 fn truncate_attempts_when_high_water_crossed_with_no_prior_attempt() {
4720 let dir = tempfile::tempdir().unwrap();
4721 let path = dir.path().join("truncate_trigger.db");
4722 let pool = file_pool(&path);
4723
4724 {
4725 let writer = pool.try_writer().unwrap();
4726 writer
4727 .conn()
4728 .execute_batch(
4729 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
4730 )
4731 .unwrap();
4732 }
4733
4734 let config = CheckpointConfig {
4735 truncate_high_water_pages: 0,
4739 truncate_min_interval: Duration::from_secs(300),
4740 ..CheckpointConfig::default()
4741 };
4742 let mut state = TruncateState::default();
4743
4744 assert!(
4745 state.last_attempt.is_none(),
4746 "precondition: no attempt has run yet"
4747 );
4748
4749 let conn = checkpoint_conn(&pool);
4750 checkpoint_once(&pool, &conn, &config, &mut state)
4751 .expect("checkpoint_once must succeed against a healthy dedicated connection");
4752 assert!(
4753 state.last_attempt.is_some(),
4754 "an attempt must be stamped once the high-water threshold is crossed"
4755 );
4756 }
4757
4758 #[test]
4761 #[serial(tx_registry, checkpoint_skip_metrics)]
4762 fn truncate_does_not_attempt_below_high_water() {
4763 let dir = tempfile::tempdir().unwrap();
4764 let path = dir.path().join("truncate_below_threshold.db");
4765 let pool = file_pool(&path);
4766
4767 {
4768 let writer = pool.try_writer().unwrap();
4769 writer
4770 .conn()
4771 .execute_batch(
4772 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
4773 )
4774 .unwrap();
4775 }
4776
4777 let config = CheckpointConfig {
4779 truncate_high_water_pages: u64::MAX,
4780 ..CheckpointConfig::default()
4781 };
4782 let mut state = TruncateState::default();
4783
4784 let conn = checkpoint_conn(&pool);
4785 checkpoint_once(&pool, &conn, &config, &mut state)
4786 .expect("checkpoint_once must succeed against a healthy dedicated connection");
4787
4788 assert!(
4789 state.last_attempt.is_none(),
4790 "a below-threshold tick must never stamp last_attempt"
4791 );
4792 }
4793
4794 #[test]
4798 #[serial(tx_registry, checkpoint_skip_metrics)]
4799 fn truncate_min_interval_skip_does_not_restamp_last_attempt() {
4800 let dir = tempfile::tempdir().unwrap();
4801 let path = dir.path().join("truncate_min_interval.db");
4802 let pool = file_pool(&path);
4803
4804 {
4805 let writer = pool.try_writer().unwrap();
4806 writer
4807 .conn()
4808 .execute_batch(
4809 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
4810 )
4811 .unwrap();
4812 }
4813
4814 let config = CheckpointConfig {
4815 truncate_high_water_pages: 0,
4816 truncate_min_interval: Duration::from_secs(300),
4817 ..CheckpointConfig::default()
4818 };
4819 let mut state = TruncateState::default();
4820 let conn = checkpoint_conn(&pool);
4821
4822 checkpoint_once(&pool, &conn, &config, &mut state)
4823 .expect("checkpoint_once must succeed against a healthy dedicated connection");
4824 let first_attempt = state.last_attempt.expect("first tick must attempt");
4825
4826 checkpoint_once(&pool, &conn, &config, &mut state)
4831 .expect("checkpoint_once must succeed against a healthy dedicated connection");
4832 let second_attempt = state.last_attempt.expect("attempt timestamp must persist");
4833
4834 assert_eq!(
4835 first_attempt, second_attempt,
4836 "a tick within truncate_min_interval must not re-stamp last_attempt"
4837 );
4838 }
4839
4840 #[test]
4852 #[serial(tx_registry, checkpoint_skip_metrics)]
4853 fn checkpoint_once_proceeds_and_can_attempt_truncate_while_pool_writer_held() {
4854 reset_checkpoint_metrics_for_tests();
4855
4856 let dir = tempfile::tempdir().unwrap();
4857 let path = dir.path().join("truncate_busy_skip.db");
4858 let pool = file_pool(&path);
4859
4860 {
4861 let writer = pool.try_writer().unwrap();
4862 writer
4863 .conn()
4864 .execute_batch(
4865 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
4866 )
4867 .unwrap();
4868 }
4869
4870 let conn = checkpoint_conn(&pool);
4871
4872 let _held = pool.try_writer().unwrap();
4876
4877 let config = CheckpointConfig {
4878 truncate_high_water_pages: 0,
4879 ..CheckpointConfig::default()
4880 };
4881 let mut state = TruncateState::default();
4882
4883 checkpoint_once(&pool, &conn, &config, &mut state).expect(
4884 "checkpoint_once must observe normally on its own dedicated connection even \
4885 while a concurrent caller holds the pool's writer mutex",
4886 );
4887
4888 assert!(
4889 state.last_attempt.is_some(),
4890 "a threshold-armed tick must still evaluate (and attempt) TRUNCATE even while \
4891 the pool writer is held — the dedicated connection is unaffected by it"
4892 );
4893 assert_eq!(
4894 checkpoint_skipped_ticks(),
4895 0,
4896 "a busy pool writer must no longer count as a skipped checkpoint tick"
4897 );
4898 assert_eq!(
4899 checkpoint_consecutive_skips(),
4900 0,
4901 "a busy pool writer must not bump the consecutive-skip run length"
4902 );
4903 }
4904
4905 #[test]
4921 #[serial(checkpoint_skip_metrics)]
4922 fn all_checkpoint_metrics_callers_are_serial_tagged() {
4923 const SELF_SRC: &str = include_str!("checkpoint.rs");
4924 let lines: Vec<&str> = SELF_SRC.lines().collect();
4925
4926 let attr_starts: Vec<usize> = lines
4927 .iter()
4928 .enumerate()
4929 .filter(|(_, l)| {
4930 let t = l.trim();
4931 t == "#[test]" || t.starts_with("#[tokio::test")
4932 })
4933 .map(|(i, _)| i)
4934 .collect();
4935
4936 let mut offenders = Vec::new();
4937
4938 for (idx, &start) in attr_starts.iter().enumerate() {
4939 let end = attr_starts.get(idx + 1).copied().unwrap_or(lines.len());
4940 let span = &lines[start..end];
4941
4942 let touches_shared_metrics = span.iter().any(|l| {
4943 l.contains("checkpoint_once(")
4944 || l.contains("checkpoint_once_core(")
4945 || l.contains("run_checkpoint_task(")
4946 });
4947 if !touches_shared_metrics {
4948 continue;
4949 }
4950
4951 let mut in_serial_attr = false;
4954 let has_group_tag = span.iter().any(|line| {
4955 let trimmed = line.trim();
4956 if !in_serial_attr {
4957 in_serial_attr = trimmed.starts_with("#[serial(");
4958 }
4959 if !in_serial_attr {
4960 return false;
4961 }
4962
4963 let has_group = trimmed
4964 .split(|ch: char| !ch.is_ascii_alphanumeric() && ch != '_')
4965 .any(|token| token == "checkpoint_skip_metrics");
4966 if trimmed.ends_with(")]") {
4967 in_serial_attr = false;
4968 }
4969 has_group
4970 });
4971
4972 if !has_group_tag {
4973 let name = span
4974 .iter()
4975 .find_map(|l| {
4976 let t = l.trim_start();
4977 let t = t.strip_prefix("pub(crate) ").unwrap_or(t);
4978 let t = t.strip_prefix("pub ").unwrap_or(t);
4979 let t = t.strip_prefix("async ").unwrap_or(t);
4980 t.strip_prefix("fn ")
4981 .map(|rest| rest.split(['(', '<']).next().unwrap_or("").trim())
4982 })
4983 .unwrap_or("<unknown test>");
4984 offenders.push(name.to_string());
4985 }
4986 }
4987
4988 assert!(
4989 offenders.is_empty(),
4990 "these tests call checkpoint_once/checkpoint_once_core/run_checkpoint_task (which write the \
4991 process-wide LAST_WAL_PAGES/CHECKPOINT_* atomics via query_wal_pages) but \
4992 are not tagged #[serial(checkpoint_skip_metrics)] (or a group including it); \
4993 an untagged caller running concurrently on cargo's default test thread pool \
4994 can clobber those atomics mid-assertion in another test (the #828/#845 race): \
4995 {offenders:?}"
4996 );
4997 }
4998
4999 #[test]
5008 #[serial(tx_registry, checkpoint_skip_metrics)]
5009 fn observed_tick_resets_consecutive_skips_but_not_lifetime_total() {
5010 reset_checkpoint_metrics_for_tests();
5011
5012 let dir = tempfile::tempdir().unwrap();
5013 let path = dir.path().join("skip_then_observe.db");
5014 let pool = file_pool(&path);
5015
5016 {
5017 let writer = pool.try_writer().unwrap();
5018 writer
5019 .conn()
5020 .execute_batch(
5021 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
5022 )
5023 .unwrap();
5024 }
5025
5026 note_checkpoint_skipped();
5028 note_checkpoint_skipped();
5029 assert_eq!(checkpoint_skipped_ticks(), 2);
5030 assert_eq!(checkpoint_consecutive_skips(), 2);
5031
5032 let conn = checkpoint_conn(&pool);
5035 let mut state = TruncateState::default();
5036 checkpoint_once(&pool, &conn, &CheckpointConfig::default(), &mut state)
5037 .expect("checkpoint_once must succeed against a healthy dedicated connection");
5038
5039 assert_eq!(
5040 checkpoint_skipped_ticks(),
5041 2,
5042 "an observed tick must not change the lifetime skipped-tick total"
5043 );
5044 assert_eq!(
5045 checkpoint_consecutive_skips(),
5046 0,
5047 "an observed tick must reset the consecutive-skip run length"
5048 );
5049 }
5050
5051 #[test]
5056 fn note_truncate_outcome_warns_once_at_third_consecutive_failure() {
5057 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5058 let subscriber = CaptureSubscriber {
5059 events: std::sync::Arc::clone(&buffer),
5060 };
5061
5062 let config = CheckpointConfig {
5063 warn_pages: 2000,
5064 ..CheckpointConfig::default()
5065 };
5066 let mut state = TruncateState::default();
5067
5068 tracing::subscriber::with_default(subscriber, || {
5069 note_truncate_outcome(&config, 5000, &mut state);
5071 note_truncate_outcome(&config, 5000, &mut state);
5072 note_truncate_outcome(&config, 5000, &mut state);
5073 note_truncate_outcome(&config, 5000, &mut state);
5075 });
5076
5077 assert_eq!(state.consecutive_failures, 4);
5078
5079 let events = buffer.lock().unwrap();
5080 let escalation_count = events
5081 .iter()
5082 .filter(|e| {
5083 e.message.as_deref()
5084 == Some(
5085 "WAL TRUNCATE has failed to clear WAL pressure for 3 consecutive attempts",
5086 )
5087 })
5088 .count();
5089 assert_eq!(
5090 escalation_count, 1,
5091 "escalation WARN must fire exactly once at the 3rd consecutive failure, got: {events:?}"
5092 );
5093
5094 note_truncate_outcome(&config, 100, &mut state);
5096 assert_eq!(
5097 state.consecutive_failures, 0,
5098 "an attempt that clears warn_pages must reset the consecutive-failure counter"
5099 );
5100 }
5101
5102 fn severity_test_config() -> CheckpointConfig {
5105 CheckpointConfig {
5106 warn_pages: 100,
5107 warn_sustained_cycles: 3,
5108 ..CheckpointConfig::default()
5109 }
5110 }
5111
5112 #[test]
5115 fn severity_ladder_info_on_first_crossing_no_warn() {
5116 let config = severity_test_config();
5117 let mut state = CheckpointSeverityState::default();
5118
5119 let below = state.observe_wal_pages(10, &config);
5120 assert!(below.is_empty(), "below-warn tick must emit nothing");
5121
5122 let above = state.observe_wal_pages(150, &config);
5123 assert_eq!(
5124 above,
5125 vec![CheckpointSeverityEmission {
5126 rung: CheckpointSeverityRung::Info,
5127 wal_pages: 150,
5128 threshold_pages: 100,
5129 consecutive_cycles: 1,
5130 }],
5131 "first below->above crossing must emit exactly one INFO and no WARN"
5132 );
5133 }
5134
5135 #[test]
5138 fn severity_ladder_warn_on_third_consecutive_cycle() {
5139 let config = severity_test_config();
5140 let mut state = CheckpointSeverityState::default();
5141
5142 let tick1 = state.observe_wal_pages(150, &config);
5143 assert_eq!(tick1.len(), 1);
5144 assert_eq!(tick1[0].rung, CheckpointSeverityRung::Info);
5145
5146 let tick2 = state.observe_wal_pages(150, &config);
5147 assert!(
5148 tick2.is_empty(),
5149 "second consecutive above-warn tick must emit nothing yet"
5150 );
5151
5152 let tick3 = state.observe_wal_pages(150, &config);
5153 assert_eq!(
5154 tick3,
5155 vec![CheckpointSeverityEmission {
5156 rung: CheckpointSeverityRung::Warn,
5157 wal_pages: 150,
5158 threshold_pages: 100,
5159 consecutive_cycles: 3,
5160 }],
5161 "WARN must fire exactly on the third consecutive above-warn tick"
5162 );
5163
5164 let tick4 = state.observe_wal_pages(150, &config);
5165 assert!(
5166 tick4.is_empty(),
5167 "WARN must not repeat on a fourth consecutive above-warn tick"
5168 );
5169 }
5170
5171 #[test]
5174 fn severity_ladder_rearms_warn_after_drain() {
5175 let config = severity_test_config();
5176 let mut state = CheckpointSeverityState::default();
5177
5178 for _ in 0..3 {
5180 state.observe_wal_pages(150, &config);
5181 }
5182 assert!(state.warn_emitted_for_episode);
5183
5184 let drain = state.observe_wal_pages(10, &config);
5186 assert!(drain.is_empty(), "a draining tick must emit nothing");
5187
5188 let reentry = state.observe_wal_pages(150, &config);
5190 assert_eq!(reentry.len(), 1);
5191 assert_eq!(reentry[0].rung, CheckpointSeverityRung::Info);
5192
5193 let mid = state.observe_wal_pages(150, &config);
5194 assert!(mid.is_empty());
5195
5196 let second_warn = state.observe_wal_pages(150, &config);
5197 assert_eq!(
5198 second_warn,
5199 vec![CheckpointSeverityEmission {
5200 rung: CheckpointSeverityRung::Warn,
5201 wal_pages: 150,
5202 threshold_pages: 100,
5203 consecutive_cycles: 3,
5204 }],
5205 "a fresh elevation episode after a drain must WARN again"
5206 );
5207 }
5208
5209 #[test]
5212 fn severity_ladder_isolated_crossings_never_warn() {
5213 let config = severity_test_config();
5214 let mut state = CheckpointSeverityState::default();
5215
5216 for _ in 0..3 {
5217 let crossing = state.observe_wal_pages(150, &config);
5218 assert_eq!(
5219 crossing.len(),
5220 1,
5221 "each isolated crossing must emit exactly one INFO"
5222 );
5223 assert_eq!(crossing[0].rung, CheckpointSeverityRung::Info);
5224
5225 let drain = state.observe_wal_pages(10, &config);
5226 assert!(drain.is_empty(), "the drain tick must emit nothing");
5227 }
5228
5229 assert!(
5230 !state.warn_emitted_for_episode,
5231 "isolated single-tick crossings must never accumulate into a WARN"
5232 );
5233 }
5234
5235 #[test]
5240 fn severity_ladder_never_emits_alarm() {
5241 let config = CheckpointConfig {
5242 warn_pages: 100,
5243 warn_sustained_cycles: 1,
5244 ..CheckpointConfig::default()
5245 };
5246 let mut state = CheckpointSeverityState::default();
5247
5248 for wal_pages in [150, 200, 250, u64::MAX] {
5249 let emissions = state.observe_wal_pages(wal_pages, &config);
5250 assert!(
5251 emissions
5252 .iter()
5253 .all(|e| e.rung != CheckpointSeverityRung::Alarm),
5254 "observe_wal_pages must never emit the ALARM rung, got: {emissions:?}"
5255 );
5256 }
5257 }
5258
5259 fn tx_age_test_config() -> CheckpointConfig {
5263 CheckpointConfig {
5264 tx_warn_secs: Duration::from_secs(30),
5265 tx_max_age_secs: Duration::from_secs(120),
5266 ..CheckpointConfig::default()
5267 }
5268 }
5269
5270 fn tx_id(n: u64) -> khive_storage::tx_registry::TxId {
5275 khive_storage::tx_registry::TxId(n)
5276 }
5277
5278 #[test]
5280 fn tx_age_sweep_empty_registry_emits_nothing() {
5281 let config = tx_age_test_config();
5282 let mut state = TxAgeSweepState::default();
5283
5284 let emissions = state.observe(None, config.tx_warn_secs, config.tx_max_age_secs);
5285 assert!(emissions.is_empty(), "no open entry must emit nothing");
5286 }
5287
5288 #[test]
5290 fn tx_age_sweep_fresh_entry_emits_nothing() {
5291 let config = tx_age_test_config();
5292 let mut state = TxAgeSweepState::default();
5293
5294 let emissions = state.observe(
5295 Some((
5296 tx_id(1),
5297 Duration::from_secs(5),
5298 Some("fresh_span".to_string()),
5299 )),
5300 config.tx_warn_secs,
5301 config.tx_max_age_secs,
5302 );
5303 assert!(emissions.is_empty(), "a fresh entry must emit nothing");
5304 }
5305
5306 #[test]
5310 fn tx_age_sweep_warn_fires_once_on_crossing() {
5311 let config = tx_age_test_config();
5312 let mut state = TxAgeSweepState::default();
5313
5314 let tick1 = state.observe(
5315 Some((
5316 tx_id(1),
5317 Duration::from_secs(45),
5318 Some("stale_span".to_string()),
5319 )),
5320 config.tx_warn_secs,
5321 config.tx_max_age_secs,
5322 );
5323 assert_eq!(
5324 tick1,
5325 vec![TxAgeEmission {
5326 rung: TxAgeRung::Warn,
5327 age: Duration::from_secs(45),
5328 label: Some("stale_span".to_string()),
5329 }],
5330 "crossing tx_warn_secs must emit exactly one Warn"
5331 );
5332
5333 let tick2 = state.observe(
5334 Some((
5335 tx_id(1),
5336 Duration::from_secs(50),
5337 Some("stale_span".to_string()),
5338 )),
5339 config.tx_warn_secs,
5340 config.tx_max_age_secs,
5341 );
5342 assert!(
5343 tick2.is_empty(),
5344 "Warn must not repeat while the entry stays in the warn band"
5345 );
5346 }
5347
5348 #[test]
5351 fn tx_age_sweep_stale_fires_once_on_crossing() {
5352 let config = tx_age_test_config();
5353 let mut state = TxAgeSweepState::default();
5354
5355 state.observe(
5358 Some((
5359 tx_id(1),
5360 Duration::from_secs(45),
5361 Some("stuck_writer_task_tx".to_string()),
5362 )),
5363 config.tx_warn_secs,
5364 config.tx_max_age_secs,
5365 );
5366
5367 let tick = state.observe(
5368 Some((
5369 tx_id(1),
5370 Duration::from_secs(130),
5371 Some("stuck_writer_task_tx".to_string()),
5372 )),
5373 config.tx_warn_secs,
5374 config.tx_max_age_secs,
5375 );
5376 assert_eq!(
5377 tick,
5378 vec![TxAgeEmission {
5379 rung: TxAgeRung::Stale,
5380 age: Duration::from_secs(130),
5381 label: Some("stuck_writer_task_tx".to_string()),
5382 }],
5383 "crossing tx_max_age_secs must emit exactly one Stale"
5384 );
5385
5386 let tick_repeat = state.observe(
5387 Some((
5388 tx_id(1),
5389 Duration::from_secs(200),
5390 Some("stuck_writer_task_tx".to_string()),
5391 )),
5392 config.tx_warn_secs,
5393 config.tx_max_age_secs,
5394 );
5395 assert!(
5396 tick_repeat.is_empty(),
5397 "Stale must not repeat while the entry stays above tx_max_age_secs"
5398 );
5399 }
5400
5401 #[test]
5405 fn tx_age_sweep_already_stale_entry_emits_both_rungs_same_tick() {
5406 let config = tx_age_test_config();
5407 let mut state = TxAgeSweepState::default();
5408
5409 let tick = state.observe(
5410 Some((
5411 tx_id(1),
5412 Duration::from_secs(300),
5413 Some("ancient_tx".to_string()),
5414 )),
5415 config.tx_warn_secs,
5416 config.tx_max_age_secs,
5417 );
5418 assert_eq!(
5419 tick,
5420 vec![
5421 TxAgeEmission {
5422 rung: TxAgeRung::Warn,
5423 age: Duration::from_secs(300),
5424 label: Some("ancient_tx".to_string()),
5425 },
5426 TxAgeEmission {
5427 rung: TxAgeRung::Stale,
5428 age: Duration::from_secs(300),
5429 label: Some("ancient_tx".to_string()),
5430 },
5431 ],
5432 "an already-stale entry must cross both rungs on its first observed tick"
5433 );
5434 }
5435
5436 #[test]
5439 fn tx_age_sweep_rearms_after_entry_clears() {
5440 let config = tx_age_test_config();
5441 let mut state = TxAgeSweepState::default();
5442
5443 state.observe(
5444 Some((
5445 tx_id(1),
5446 Duration::from_secs(150),
5447 Some("first_span".to_string()),
5448 )),
5449 config.tx_warn_secs,
5450 config.tx_max_age_secs,
5451 );
5452
5453 let cleared = state.observe(None, config.tx_warn_secs, config.tx_max_age_secs);
5455 assert!(cleared.is_empty(), "a clearing tick must emit nothing");
5456
5457 let fresh = state.observe(
5459 Some((
5460 tx_id(2),
5461 Duration::from_secs(2),
5462 Some("second_span".to_string()),
5463 )),
5464 config.tx_warn_secs,
5465 config.tx_max_age_secs,
5466 );
5467 assert!(fresh.is_empty(), "a fresh oldest entry must emit nothing");
5468
5469 let rewarn = state.observe(
5471 Some((
5472 tx_id(2),
5473 Duration::from_secs(35),
5474 Some("second_span".to_string()),
5475 )),
5476 config.tx_warn_secs,
5477 config.tx_max_age_secs,
5478 );
5479 assert_eq!(
5480 rewarn,
5481 vec![TxAgeEmission {
5482 rung: TxAgeRung::Warn,
5483 age: Duration::from_secs(35),
5484 label: Some("second_span".to_string()),
5485 }],
5486 "a fresh stale episode after a clear must Warn again"
5487 );
5488 }
5489
5490 #[test]
5494 fn tx_age_sweep_stale_replacement_without_intervening_clear_still_names_new_entry() {
5495 let config = tx_age_test_config();
5496 let mut state = TxAgeSweepState::default();
5497
5498 let tick_a = state.observe(
5499 Some((
5500 tx_id(1),
5501 Duration::from_secs(300),
5502 Some("stale_entry_a".to_string()),
5503 )),
5504 config.tx_warn_secs,
5505 config.tx_max_age_secs,
5506 );
5507 assert_eq!(
5508 tick_a.len(),
5509 2,
5510 "entry A must cross both rungs on its first observed tick, got: {tick_a:?}"
5511 );
5512
5513 let tick_b = state.observe(
5516 Some((
5517 tx_id(2),
5518 Duration::from_secs(400),
5519 Some("stale_entry_b".to_string()),
5520 )),
5521 config.tx_warn_secs,
5522 config.tx_max_age_secs,
5523 );
5524 assert_eq!(
5525 tick_b,
5526 vec![
5527 TxAgeEmission {
5528 rung: TxAgeRung::Warn,
5529 age: Duration::from_secs(400),
5530 label: Some("stale_entry_b".to_string()),
5531 },
5532 TxAgeEmission {
5533 rung: TxAgeRung::Stale,
5534 age: Duration::from_secs(400),
5535 label: Some("stale_entry_b".to_string()),
5536 },
5537 ],
5538 "a same-tick identity change to an already-stale successor must re-emit both \
5539 rungs naming the NEW entry, got: {tick_b:?}"
5540 );
5541 }
5542
5543 #[test]
5546 fn tx_age_sweep_uses_configured_thresholds_not_hardcoded_defaults() {
5547 let config = CheckpointConfig {
5548 tx_warn_secs: Duration::from_millis(1),
5549 tx_max_age_secs: Duration::from_millis(2),
5550 ..CheckpointConfig::default()
5551 };
5552 let mut state = TxAgeSweepState::default();
5553
5554 let tick = state.observe(
5555 Some((
5556 tx_id(1),
5557 Duration::from_millis(5),
5558 Some("fast_cap_span".to_string()),
5559 )),
5560 config.tx_warn_secs,
5561 config.tx_max_age_secs,
5562 );
5563 assert_eq!(
5564 tick.len(),
5565 2,
5566 "a millisecond-scale cap must cross both rungs immediately, got: {tick:?}"
5567 );
5568 }
5569
5570 #[test]
5573 #[serial(tx_registry, checkpoint_skip_metrics)]
5574 fn tx_age_sweep_names_long_lived_reader_pinning_wal_past_high_water() {
5575 let dir = tempfile::tempdir().unwrap();
5576 let path = dir.path().join("tx_age_sweep_reader_pin.db");
5577 let pool = file_pool(&path);
5578
5579 {
5580 let writer = pool.try_writer().unwrap();
5581 writer
5582 .conn()
5583 .execute_batch(
5584 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
5585 )
5586 .unwrap();
5587 }
5588
5589 let reader = pool.reader().expect("reader");
5594 reader
5595 .execute_batch("BEGIN DEFERRED; SELECT * FROM t;")
5596 .expect("begin read tx");
5597 let _tx_handle =
5598 khive_storage::tx_registry::register(Some("tx_age_sweep_reader_pin_test".to_string()));
5599
5600 let config = CheckpointConfig {
5603 high_water_pages: 1,
5604 tx_warn_secs: Duration::from_millis(1),
5605 tx_max_age_secs: Duration::from_millis(1),
5606 ..CheckpointConfig::default()
5607 };
5608 {
5609 let writer = pool.try_writer().unwrap();
5610 for i in 0..50 {
5611 writer
5612 .conn()
5613 .execute_batch(&format!("INSERT INTO t VALUES ({i});"))
5614 .unwrap();
5615 }
5616 }
5617
5618 let conn = checkpoint_conn(&pool);
5619 let wal_pages = checkpoint_once(&pool, &conn, &config, &mut TruncateState::default())
5620 .expect("checkpoint_once must succeed against a healthy dedicated connection");
5621 assert!(
5622 wal_pages >= config.high_water_pages,
5623 "test setup must actually drive wal_pages ({wal_pages}) past high_water_pages \
5624 ({}) for this regression to mean anything",
5625 config.high_water_pages
5626 );
5627
5628 std::thread::sleep(Duration::from_millis(5));
5635 let our_entry = khive_storage::tx_registry::snapshot()
5646 .into_iter()
5647 .find(|(_, label)| label.as_deref() == Some("tx_age_sweep_reader_pin_test"))
5648 .expect("this test's own tx_registry entry must still be open");
5649 let mut tx_age_state = TxAgeSweepState::default();
5650 let emissions = tx_age_state.observe(
5651 Some((tx_id(1), our_entry.0, our_entry.1)),
5652 config.tx_warn_secs,
5653 config.tx_max_age_secs,
5654 );
5655 assert!(
5656 emissions.iter().any(|e| e.rung == TxAgeRung::Stale
5657 && e.label.as_deref() == Some("tx_age_sweep_reader_pin_test")),
5658 "expected a Stale emission naming the pinning reader, got: {emissions:?}"
5659 );
5660
5661 reader.execute_batch("COMMIT;").ok();
5662 drop(reader);
5663 drop(_tx_handle);
5664 }
5665
5666 #[test]
5669 #[serial(tx_registry, checkpoint_skip_metrics)]
5670 fn tx_age_sweep_own_entry_survives_concurrent_older_registration() {
5671 let _decoy = khive_storage::tx_registry::register(Some("decoy_unrelated_span".to_string()));
5672 std::thread::sleep(Duration::from_millis(2));
5673 let _own = khive_storage::tx_registry::register(Some("this_test_own_span".to_string()));
5674 std::thread::sleep(Duration::from_millis(5));
5675
5676 let global_oldest = khive_storage::tx_registry::oldest().expect("registry not empty");
5682 assert_ne!(
5683 global_oldest.2.as_deref(),
5684 Some("this_test_own_span"),
5685 "test setup must reproduce the race: an older, unrelated entry must be \
5686 the current global oldest, got: {global_oldest:?}"
5687 );
5688
5689 let our_entry = khive_storage::tx_registry::snapshot()
5690 .into_iter()
5691 .find(|(_, label)| label.as_deref() == Some("this_test_own_span"))
5692 .expect("this test's own tx_registry entry must still be open");
5693
5694 let config = CheckpointConfig {
5695 tx_warn_secs: Duration::from_millis(1),
5696 tx_max_age_secs: Duration::from_millis(1),
5697 ..CheckpointConfig::default()
5698 };
5699 let mut state = TxAgeSweepState::default();
5700 let emissions = state.observe(
5701 Some((tx_id(2), our_entry.0, our_entry.1)),
5702 config.tx_warn_secs,
5703 config.tx_max_age_secs,
5704 );
5705 assert!(
5706 emissions
5707 .iter()
5708 .any(|e| e.rung == TxAgeRung::Stale
5709 && e.label.as_deref() == Some("this_test_own_span")),
5710 "expected a Stale emission naming this test's own span despite an older, \
5711 unrelated concurrent registration, got: {emissions:?}"
5712 );
5713 }
5714
5715 #[test]
5717 #[serial]
5718 fn checkpoint_config_warn_sustained_cycles_env_override() {
5719 let default = CheckpointConfig::default();
5720 assert_eq!(default.warn_sustained_cycles, DEFAULT_WARN_SUSTAINED_CYCLES);
5721
5722 std::env::set_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES", "5");
5723 let cfg = CheckpointConfig::from_env();
5724 std::env::remove_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES");
5725 assert_eq!(cfg.warn_sustained_cycles, 5);
5726
5727 std::env::set_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES", "0");
5728 let cfg_zero = CheckpointConfig::from_env();
5729 std::env::remove_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES");
5730 assert_eq!(
5731 cfg_zero.warn_sustained_cycles, DEFAULT_WARN_SUSTAINED_CYCLES,
5732 "zero must fall back to the default"
5733 );
5734
5735 std::env::set_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES", "not_a_number");
5736 let cfg_invalid = CheckpointConfig::from_env();
5737 std::env::remove_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES");
5738 assert_eq!(
5739 cfg_invalid.warn_sustained_cycles,
5740 DEFAULT_WARN_SUSTAINED_CYCLES
5741 );
5742 }
5743
5744 #[derive(Clone, Copy)]
5747 enum FakeAppendBehavior {
5748 Record,
5749 Fail,
5750 }
5751
5752 struct FakeEventStore {
5753 events: std::sync::Mutex<Vec<khive_storage::Event>>,
5754 append_attempts: std::sync::atomic::AtomicUsize,
5755 append_behavior: FakeAppendBehavior,
5756 }
5757
5758 impl Default for FakeEventStore {
5759 fn default() -> Self {
5760 Self {
5761 events: std::sync::Mutex::new(Vec::new()),
5762 append_attempts: std::sync::atomic::AtomicUsize::new(0),
5763 append_behavior: FakeAppendBehavior::Record,
5764 }
5765 }
5766 }
5767
5768 impl FakeEventStore {
5769 fn failing() -> Self {
5770 Self {
5771 append_behavior: FakeAppendBehavior::Fail,
5772 ..Self::default()
5773 }
5774 }
5775 }
5776
5777 #[async_trait::async_trait]
5778 impl khive_storage::EventStore for FakeEventStore {
5779 async fn append_event(
5780 &self,
5781 event: khive_storage::Event,
5782 ) -> khive_storage::StorageResult<()> {
5783 self.append_attempts
5784 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
5785 match self.append_behavior {
5786 FakeAppendBehavior::Record => {
5787 self.events.lock().unwrap().push(event);
5788 Ok(())
5789 }
5790 FakeAppendBehavior::Fail => Err(khive_storage::StorageError::Internal(
5791 "synthetic checkpoint lifecycle append failure".to_string(),
5792 )),
5793 }
5794 }
5795
5796 async fn append_events(
5797 &self,
5798 events: Vec<khive_storage::Event>,
5799 ) -> khive_storage::StorageResult<khive_storage::BatchWriteSummary> {
5800 let count = events.len() as u64;
5801 self.events.lock().unwrap().extend(events);
5802 Ok(khive_storage::BatchWriteSummary {
5803 attempted: count,
5804 affected: count,
5805 failed: 0,
5806 first_error: String::new(),
5807 })
5808 }
5809
5810 async fn get_event(
5811 &self,
5812 id: uuid::Uuid,
5813 ) -> khive_storage::StorageResult<Option<khive_storage::Event>> {
5814 Ok(self
5815 .events
5816 .lock()
5817 .unwrap()
5818 .iter()
5819 .find(|e| e.id == id)
5820 .cloned())
5821 }
5822
5823 async fn query_events(
5824 &self,
5825 _filter: khive_storage::EventFilter,
5826 _page: khive_storage::PageRequest,
5827 ) -> khive_storage::StorageResult<khive_storage::Page<khive_storage::Event>> {
5828 unimplemented!("not exercised by the checkpoint lifecycle-event tests")
5829 }
5830
5831 async fn count_events(
5832 &self,
5833 _filter: khive_storage::EventFilter,
5834 ) -> khive_storage::StorageResult<u64> {
5835 Ok(self.events.lock().unwrap().len() as u64)
5836 }
5837 }
5838
5839 #[test]
5843 fn checkpoint_outcome_should_emit_covers_all_transitions() {
5844 assert!(
5845 checkpoint_outcome_should_emit(true, false),
5846 "first elevated tick must emit"
5847 );
5848 assert!(
5849 !checkpoint_outcome_should_emit(true, true),
5850 "sustained elevated ticks must aggregate in memory instead of writing the WAL"
5851 );
5852 assert!(
5853 checkpoint_outcome_should_emit(false, true),
5854 "the single drain row (elevated -> healthy) must emit"
5855 );
5856 assert!(
5857 !checkpoint_outcome_should_emit(false, false),
5858 "an ordinary below-warn tick must not emit"
5859 );
5860 }
5861
5862 #[test]
5865 fn persistent_pressure_lifecycle_rows_are_o_state_transitions() {
5866 let observations = [true; 128].into_iter().chain([false]).chain([false; 128]);
5867 let mut was_elevated = false;
5868 let writes = observations
5869 .filter(|above_warn| {
5870 let emit = checkpoint_outcome_should_emit(*above_warn, was_elevated);
5871 if emit {
5872 was_elevated = *above_warn;
5873 }
5874 emit
5875 })
5876 .count();
5877
5878 assert_eq!(
5879 writes, 2,
5880 "one elevation row plus one recovery summary must cover any number of attempts"
5881 );
5882 }
5883
5884 #[test]
5885 fn checkpoint_pressure_episode_retains_recovery_summary() {
5886 let mut episode = CheckpointPressureEpisode::start(2_500);
5887 episode.observe(2_300);
5888 episode.observe(8_100);
5889 episode.observe(4_000);
5890
5891 assert_eq!(episode.elevated_ticks, 4);
5892 assert_eq!(episode.peak_wal_pages, 8_100);
5893 }
5894
5895 fn drive_pressure_ticks(
5901 config: &CheckpointConfig,
5902 ticks: &[(bool, u64)],
5903 mut fail_on: impl FnMut(usize) -> bool,
5904 ) -> Vec<khive_storage::CheckpointOutcomeRecordedPayload> {
5905 let mut event_elevation_open = false;
5906 let mut pressure_episode: Option<CheckpointPressureEpisode> = None;
5907 let mut pending_recovery: Option<khive_storage::CheckpointOutcomeRecordedPayload> = None;
5908 let mut delivered = Vec::new();
5909 let mut call_index = 0usize;
5910 for &(above_warn, wal_pages) in ticks {
5911 observe_checkpoint_pressure_tick(
5912 above_warn,
5913 wal_pages,
5914 false,
5915 false,
5916 config,
5917 &mut event_elevation_open,
5918 &mut pressure_episode,
5919 &mut pending_recovery,
5920 |payload| {
5921 let idx = call_index;
5922 call_index += 1;
5923 if fail_on(idx) {
5924 false
5925 } else {
5926 delivered.push(payload);
5927 true
5928 }
5929 },
5930 );
5931 }
5932 delivered
5933 }
5934
5935 #[test]
5949 fn dropped_recovery_handoff_does_not_merge_pressure_episodes() {
5950 let config = CheckpointConfig {
5951 warn_pages: 1_000,
5952 ..CheckpointConfig::default()
5953 };
5954 let ticks = [
5955 (true, 1_500), (true, 1_800), (true, 2_000), (false, 500), (false, 400), (true, 3_000), (true, 3_500), (false, 300), ];
5964
5965 let delivered = drive_pressure_ticks(&config, &ticks, |idx| matches!(idx, 1..=3));
5966
5967 assert_eq!(
5968 delivered.len(),
5969 4,
5970 "expected episode-1 open, episode-1 delayed recovery, episode-2 open, \
5971 episode-2 recovery: {delivered:?}"
5972 );
5973
5974 let ep1_open = &delivered[0];
5975 assert!(ep1_open.above_warn);
5976 assert_eq!(ep1_open.episode_elevated_ticks, Some(1));
5977 assert_eq!(ep1_open.episode_peak_wal_pages, Some(1_500));
5978
5979 let ep1_recovery = &delivered[1];
5980 assert!(
5981 !ep1_recovery.above_warn,
5982 "episode 1's recovery must be delivered BEFORE episode 2's opening; \
5983 an opening in this slot means the barrier failed: {delivered:?}"
5984 );
5985 assert_eq!(
5986 ep1_recovery.episode_elevated_ticks,
5987 Some(3),
5988 "episode 1's delayed recovery must report only its own 3 elevated ticks, \
5989 not ticks absorbed from episode 2"
5990 );
5991 assert_eq!(ep1_recovery.episode_peak_wal_pages, Some(2_000));
5992
5993 let ep2_open = &delivered[2];
5994 assert!(ep2_open.above_warn);
5995 assert_eq!(
5996 ep2_open.episode_elevated_ticks,
5997 Some(2),
5998 "episode 2 opens fresh (never continuing episode 1's count), deferred one \
5999 tick by the barrier, so its opening reports 2 elevated ticks"
6000 );
6001 assert_eq!(ep2_open.episode_peak_wal_pages, Some(3_500));
6002
6003 let ep2_recovery = &delivered[3];
6004 assert!(!ep2_recovery.above_warn);
6005 assert_eq!(
6006 ep2_recovery.episode_elevated_ticks,
6007 Some(2),
6008 "episode 2's recovery must report only its own 2 elevated ticks"
6009 );
6010 assert_eq!(ep2_recovery.episode_peak_wal_pages, Some(3_500));
6011 }
6012
6013 #[test]
6020 fn episode_elapsed_entirely_behind_barrier_is_discarded_not_reordered() {
6021 let config = CheckpointConfig {
6022 warn_pages: 1_000,
6023 ..CheckpointConfig::default()
6024 };
6025 let ticks = [
6026 (true, 1_500), (false, 500), (true, 9_000), (false, 400), (false, 300), (true, 2_500), (false, 200), ];
6034
6035 let delivered = drive_pressure_ticks(&config, &ticks, |idx| matches!(idx, 1..=3));
6036
6037 let peaks: Vec<_> = delivered
6038 .iter()
6039 .map(|payload| (payload.above_warn, payload.episode_peak_wal_pages))
6040 .collect();
6041 assert_eq!(
6042 peaks,
6043 vec![
6044 (true, Some(1_500)), (false, Some(1_500)), (true, Some(2_500)), (false, Some(2_500)), ],
6049 "an episode elapsed entirely behind the barrier must not surface late or \
6050 out of order: {delivered:?}"
6051 );
6052 }
6053
6054 #[test]
6059 fn no_dropped_handoff_reports_two_separate_episodes() {
6060 let config = CheckpointConfig {
6061 warn_pages: 1_000,
6062 ..CheckpointConfig::default()
6063 };
6064 let ticks = [
6065 (true, 1_500),
6066 (true, 1_800),
6067 (true, 2_000),
6068 (false, 500),
6069 (false, 400),
6070 (true, 3_000),
6071 (true, 3_500),
6072 (false, 300),
6073 ];
6074
6075 let delivered = drive_pressure_ticks(&config, &ticks, |_idx| false);
6076
6077 assert_eq!(delivered.len(), 4, "{delivered:?}");
6078 assert_eq!(
6079 (
6080 delivered[0].above_warn,
6081 delivered[0].episode_elevated_ticks,
6082 delivered[0].episode_peak_wal_pages
6083 ),
6084 (true, Some(1), Some(1_500)),
6085 "episode 1 open"
6086 );
6087 assert_eq!(
6088 (
6089 delivered[1].above_warn,
6090 delivered[1].episode_elevated_ticks,
6091 delivered[1].episode_peak_wal_pages
6092 ),
6093 (false, Some(3), Some(2_000)),
6094 "episode 1 recovery"
6095 );
6096 assert_eq!(
6097 (
6098 delivered[2].above_warn,
6099 delivered[2].episode_elevated_ticks,
6100 delivered[2].episode_peak_wal_pages
6101 ),
6102 (true, Some(1), Some(3_000)),
6103 "episode 2 open"
6104 );
6105 assert_eq!(
6106 (
6107 delivered[3].above_warn,
6108 delivered[3].episode_elevated_ticks,
6109 delivered[3].episode_peak_wal_pages
6110 ),
6111 (false, Some(2), Some(3_500)),
6112 "episode 2 recovery"
6113 );
6114 }
6115
6116 #[test]
6117 #[serial(checkpoint_skip_metrics)]
6118 fn pressure_diagnostics_count_observations_and_transitions_separately() {
6119 reset_checkpoint_metrics_for_tests();
6120
6121 note_checkpoint_pressure_observation(true, false);
6122 note_checkpoint_pressure_observation(true, true);
6123 note_checkpoint_pressure_observation(true, true);
6124 note_checkpoint_pressure_observation(false, true);
6125 note_checkpoint_pressure_observation(false, false);
6126
6127 assert_eq!(checkpoint_pressure_elevated_ticks(), 3);
6128 assert_eq!(checkpoint_pressure_episodes_started(), 1);
6129 assert_eq!(checkpoint_pressure_episodes_recovered(), 1);
6130 assert_eq!(checkpoint_lifecycle_append_attempts(), 0);
6131 }
6132
6133 #[tokio::test]
6134 #[serial(checkpoint_skip_metrics)]
6135 async fn checkpoint_task_emits_one_opening_for_persistent_pressure() {
6136 reset_checkpoint_metrics_for_tests();
6137 let dir = tempfile::tempdir().unwrap();
6138 let path = dir.path().join("outcome_emit.db");
6139 let pool = file_pool(&path);
6140
6141 let cfg = CheckpointConfig {
6144 interval: Duration::from_millis(10),
6145 warn_pages: 0,
6146 ..CheckpointConfig::default()
6147 };
6148 let store = Arc::new(FakeEventStore::default());
6149 let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
6150
6151 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
6152 let handle = tokio::spawn(run_checkpoint_task(
6153 pool,
6154 cfg,
6155 Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
6156 shutdown_rx,
6157 true,
6158 ));
6159
6160 let progressed = wait_for(Duration::from_secs(10), || {
6161 checkpoint_pressure_elevated_ticks() >= 10
6162 })
6163 .await;
6164 let emitted = wait_for(Duration::from_secs(10), || {
6165 !store.events.lock().unwrap().is_empty()
6166 })
6167 .await;
6168 shutdown_tx.send(()).expect("send shutdown signal");
6169 tokio::time::timeout(Duration::from_secs(1), handle)
6170 .await
6171 .expect("checkpoint task should exit within 1s")
6172 .expect("checkpoint task panicked");
6173
6174 let events = store.events.lock().unwrap();
6175 assert!(
6176 progressed,
6177 "the simulated persistent-pressure episode must span at least ten checkpoint ticks"
6178 );
6179 assert!(
6180 emitted,
6181 "an always-elevated config must append one CheckpointOutcomeRecorded event \
6182 within the poll deadline"
6183 );
6184 assert_eq!(
6185 checkpoint_lifecycle_append_attempts(),
6186 1,
6187 "primary-store lifecycle writes must stay O(state transitions), not O(attempts)"
6188 );
6189 assert_eq!(events.len(), 1);
6190 assert_eq!(events[0].payload_schema_version, 2);
6191 assert_eq!(events[0].payload["episode_elevated_ticks"], 1);
6192 assert_eq!(
6193 events[0].payload["episode_peak_wal_pages"],
6194 events[0].payload["wal_pages"]
6195 );
6196 assert!(
6197 events
6198 .iter()
6199 .all(|e| e.kind == khive_types::EventKind::CheckpointOutcomeRecorded),
6200 "every appended event must be CheckpointOutcomeRecorded, got: {events:?}"
6201 );
6202 assert!(
6203 events.iter().all(|e| e.namespace == "local"),
6204 "events must be stamped with the namespace passed to run_checkpoint_task"
6205 );
6206 }
6207
6208 #[tokio::test]
6212 #[serial(checkpoint_skip_metrics)]
6213 async fn checkpoint_cycles_and_task_shutdown_do_not_wait_for_a_contended_lifecycle_writer() {
6214 reset_checkpoint_metrics_for_tests();
6215 let dir = tempfile::tempdir().unwrap();
6216 let path = dir.path().join("outcome_contended_sink.db");
6217 let checkpoint_pool = file_pool(&path);
6218
6219 let event_pool = Arc::new(
6220 ConnectionPool::new(PoolConfig {
6221 path: None,
6222 checkout_timeout: Duration::from_secs(5),
6223 write_queue_enabled: Some(false),
6224 ..PoolConfig::default()
6225 })
6226 .expect("event pool"),
6227 );
6228 {
6229 let writer = event_pool.try_writer().expect("initialize event schema");
6230 crate::stores::event::ensure_events_schema(writer.conn())
6231 .expect("initialize event schema");
6232 }
6233 let event_store: Arc<dyn khive_storage::EventStore> =
6234 Arc::new(crate::stores::event::SqlEventStore::new_scoped(
6235 Arc::clone(&event_pool),
6236 false,
6237 "local",
6238 ));
6239 let held_event_writer = event_pool
6240 .try_writer()
6241 .expect("hold the event-store writer");
6242
6243 let cfg = CheckpointConfig {
6244 interval: Duration::from_millis(10),
6245 warn_pages: 0,
6246 ..CheckpointConfig::default()
6247 };
6248 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
6249 let handle = tokio::spawn(run_checkpoint_task(
6250 checkpoint_pool,
6251 cfg,
6252 Some(CheckpointLifecycleOwner::new(event_store, "local")),
6253 shutdown_rx,
6254 true,
6255 ));
6256
6257 let progressed = wait_for(Duration::from_secs(2), || {
6258 checkpoint_pressure_elevated_ticks() >= 10
6259 })
6260 .await;
6261 assert!(
6262 progressed,
6263 "checkpoint observations must continue while the lifecycle append is contended"
6264 );
6265 assert_eq!(checkpoint_lifecycle_append_attempts(), 1);
6266 assert_eq!(checkpoint_lifecycle_enqueue_drops(), 0);
6267
6268 shutdown_tx.send(()).expect("send shutdown signal");
6269 tokio::time::timeout(Duration::from_secs(1), handle)
6270 .await
6271 .expect(
6272 "the run_checkpoint_task handle must not wait for the event store's \
6273 five-second writer checkout",
6274 )
6275 .expect("checkpoint task panicked");
6276
6277 drop(held_event_writer);
6282 }
6283
6284 #[tokio::test]
6287 #[serial(checkpoint_skip_metrics)]
6288 async fn checkpoint_task_continues_after_lifecycle_append_failure() {
6289 reset_checkpoint_metrics_for_tests();
6290 let dir = tempfile::tempdir().unwrap();
6291 let path = dir.path().join("outcome_failing_sink.db");
6292 let pool = file_pool(&path);
6293 let store = Arc::new(FakeEventStore::failing());
6294 let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
6295
6296 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6297 let subscriber = CaptureSubscriber {
6298 events: std::sync::Arc::clone(&buffer),
6299 };
6300 let _tracing_guard = tracing::subscriber::set_default(subscriber);
6301
6302 let cfg = CheckpointConfig {
6303 interval: Duration::from_millis(10),
6304 warn_pages: 0,
6305 ..CheckpointConfig::default()
6306 };
6307 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
6308 let handle = tokio::spawn(run_checkpoint_task(
6309 pool,
6310 cfg,
6311 Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
6312 shutdown_rx,
6313 true,
6314 ));
6315
6316 let progressed = wait_for(Duration::from_secs(2), || {
6317 checkpoint_pressure_elevated_ticks() >= 10
6318 })
6319 .await;
6320 shutdown_tx.send(()).expect("send shutdown signal");
6321 tokio::time::timeout(Duration::from_secs(1), handle)
6322 .await
6323 .expect("checkpoint task should remain responsive after sink failure")
6324 .expect("checkpoint task panicked");
6325
6326 assert!(
6327 progressed,
6328 "a failed append must not terminate or stall the checkpoint task"
6329 );
6330 assert_eq!(
6331 store
6332 .append_attempts
6333 .load(std::sync::atomic::Ordering::Relaxed),
6334 1,
6335 "a persistent pressure state must not retry one primary-store append per tick"
6336 );
6337 assert_eq!(checkpoint_lifecycle_append_attempts(), 1);
6338 assert_eq!(checkpoint_lifecycle_append_failures(), 1);
6339 let captured = buffer.lock().unwrap().clone();
6340 assert!(
6341 captured.iter().any(|event| event.message.as_deref()
6342 == Some("checkpoint lifecycle event append failed")),
6343 "lifecycle append failures must remain observable; got: {:?}",
6344 captured
6345 );
6346 }
6347
6348 #[tokio::test]
6349 #[serial(checkpoint_skip_metrics)]
6350 async fn secondary_checkpoint_task_with_lifecycle_ownership_emits_outcome_events() {
6351 let dir = tempfile::tempdir().unwrap();
6352 let path = dir.path().join("secondary_outcome.db");
6353 let pool = file_pool(&path);
6354 let cfg = CheckpointConfig {
6355 interval: Duration::from_millis(10),
6356 warn_pages: 0,
6357 ..CheckpointConfig::default()
6358 };
6359 let store = Arc::new(FakeEventStore::default());
6360 let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
6361
6362 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
6363 let handle = tokio::spawn(run_checkpoint_task(
6364 pool,
6365 cfg,
6366 Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
6367 shutdown_rx,
6368 false,
6369 ));
6370
6371 let emitted = wait_for(Duration::from_secs(10), || {
6374 !store.events.lock().unwrap().is_empty()
6375 })
6376 .await;
6377 shutdown_tx.send(()).expect("send shutdown signal");
6378 tokio::time::timeout(Duration::from_secs(1), handle)
6379 .await
6380 .expect("checkpoint task should exit within 1s")
6381 .expect("checkpoint task panicked");
6382
6383 assert!(
6384 emitted,
6385 "a designated secondary lifecycle owner must append outcome events within the poll \
6386 deadline"
6387 );
6388 }
6389
6390 #[tokio::test]
6391 #[serial(checkpoint_skip_metrics)]
6392 async fn checkpoint_task_emits_nothing_while_healthy() {
6393 let dir = tempfile::tempdir().unwrap();
6394 let path = dir.path().join("outcome_no_emit.db");
6395 let pool = file_pool(&path);
6396
6397 let cfg = CheckpointConfig {
6400 interval: Duration::from_millis(10),
6401 warn_pages: u64::MAX,
6402 ..CheckpointConfig::default()
6403 };
6404 let store = Arc::new(FakeEventStore::default());
6405 let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
6406
6407 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
6408 let handle = tokio::spawn(run_checkpoint_task(
6409 pool,
6410 cfg,
6411 Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
6412 shutdown_rx,
6413 true,
6414 ));
6415
6416 tokio::time::sleep(Duration::from_millis(60)).await;
6417 shutdown_tx.send(()).expect("send shutdown signal");
6418 tokio::time::timeout(Duration::from_secs(1), handle)
6419 .await
6420 .expect("checkpoint task should exit within 1s")
6421 .expect("checkpoint task panicked");
6422
6423 assert!(
6424 store.events.lock().unwrap().is_empty(),
6425 "a config that never crosses warn_pages must never append a lifecycle event"
6426 );
6427 }
6428
6429 #[tokio::test]
6430 #[serial(checkpoint_skip_metrics)]
6431 async fn checkpoint_task_with_no_event_store_does_not_panic() {
6432 let dir = tempfile::tempdir().unwrap();
6433 let path = dir.path().join("outcome_none_store.db");
6434 let pool = file_pool(&path);
6435
6436 let cfg = CheckpointConfig {
6437 interval: Duration::from_millis(10),
6438 warn_pages: 0,
6439 ..CheckpointConfig::default()
6440 };
6441
6442 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
6443 let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
6444
6445 tokio::time::sleep(Duration::from_millis(40)).await;
6446 shutdown_tx.send(()).expect("send shutdown signal");
6447 tokio::time::timeout(Duration::from_secs(1), handle)
6448 .await
6449 .expect("checkpoint task should exit within 1s")
6450 .expect("checkpoint task panicked");
6451 }
6452
6453 #[tokio::test]
6470 #[serial(tx_registry, checkpoint_skip_metrics)]
6471 async fn checkpoint_task_sweeps_stale_registry_entry_while_wal_is_healthy() {
6472 let dir = tempfile::tempdir().unwrap();
6473 let path = dir.path().join("tx_age_sweep_task_healthy_wal.db");
6474 let pool = file_pool(&path);
6475
6476 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6477 let subscriber = CaptureSubscriber {
6478 events: std::sync::Arc::clone(&buffer),
6479 };
6480 let _tracing_guard = tracing::subscriber::set_default(subscriber);
6481
6482 let _tx_handle = khive_storage::tx_registry::register(Some(
6483 "checkpoint_task_healthy_wal_sweep_test".to_string(),
6484 ));
6485
6486 let cfg = CheckpointConfig {
6487 interval: Duration::from_millis(10),
6488 warn_pages: u64::MAX,
6489 high_water_pages: u64::MAX,
6490 truncate_high_water_pages: u64::MAX,
6491 tx_warn_secs: Duration::from_millis(1),
6492 tx_max_age_secs: Duration::from_millis(1),
6493 ..CheckpointConfig::default()
6494 };
6495
6496 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
6497 let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
6498
6499 let swept = wait_for(Duration::from_secs(10), || {
6506 buffer.lock().unwrap().iter().any(|e| {
6507 e.tx_label.as_deref() == Some("checkpoint_task_healthy_wal_sweep_test")
6508 && e.message
6509 .as_deref()
6510 .is_some_and(|m| m.contains("stale-op cap"))
6511 })
6512 })
6513 .await;
6514
6515 shutdown_tx.send(()).expect("send shutdown signal");
6516 tokio::time::timeout(Duration::from_secs(1), handle)
6517 .await
6518 .expect("checkpoint task should exit within 1s")
6519 .expect("checkpoint task panicked");
6520
6521 drop(_tx_handle);
6522
6523 let events = buffer.lock().unwrap();
6524 assert!(
6525 swept,
6526 "expected the spawned task to sweep and escalate the stale registry entry \
6527 to Stale on its own within the poll deadline, got: {events:?}"
6528 );
6529 }
6530
6531 #[tokio::test]
6535 #[serial(tx_registry, checkpoint_skip_metrics)]
6536 async fn checkpoint_task_emits_no_age_alert_for_an_empty_registry() {
6537 let dir = tempfile::tempdir().unwrap();
6538 let path = dir.path().join("tx_age_sweep_task_empty_registry.db");
6539 let pool = file_pool(&path);
6540
6541 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6542 let subscriber = CaptureSubscriber {
6543 events: std::sync::Arc::clone(&buffer),
6544 };
6545 let _tracing_guard = tracing::subscriber::set_default(subscriber);
6546
6547 let cfg = CheckpointConfig {
6548 interval: Duration::from_millis(10),
6549 tx_warn_secs: Duration::from_millis(1),
6550 tx_max_age_secs: Duration::from_millis(1),
6551 ..CheckpointConfig::default()
6552 };
6553
6554 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
6555 let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
6556
6557 tokio::time::sleep(Duration::from_millis(40)).await;
6558 shutdown_tx.send(()).expect("send shutdown signal");
6559 tokio::time::timeout(Duration::from_secs(1), handle)
6560 .await
6561 .expect("checkpoint task should exit within 1s")
6562 .expect("checkpoint task panicked");
6563
6564 let events = buffer.lock().unwrap();
6565 assert!(
6566 events.iter().all(|e| e
6567 .message
6568 .as_deref()
6569 .is_none_or(|m| !m.contains("ADR-091 Plank 1"))),
6570 "an empty registry must never produce a Plank 1 age emission, got: {events:?}"
6571 );
6572 }
6573
6574 #[tokio::test]
6587 #[serial(tx_registry, checkpoint_skip_metrics)]
6588 async fn checkpoint_task_sweeps_stale_entry_even_when_dedicated_connection_is_unavailable_every_tick(
6589 ) {
6590 reset_checkpoint_metrics_for_tests();
6591
6592 let dir = tempfile::tempdir().unwrap();
6593 let path = dir.path().join("tx_age_sweep_task_conn_unavailable.db");
6594 {
6595 let seed_pool = file_pool(&path);
6599 let writer = seed_pool.try_writer().unwrap();
6600 writer
6601 .conn()
6602 .execute_batch("CREATE TABLE IF NOT EXISTS t (x INTEGER);")
6603 .unwrap();
6604 }
6605
6606 #[cfg(unix)]
6607 {
6608 use std::os::unix::fs::PermissionsExt;
6609 for suffix in ["-wal", "-shm"] {
6610 let mut name = path.file_name().expect("db file name").to_os_string();
6611 name.push(suffix);
6612 let sidecar = path.parent().expect("db parent dir").join(name);
6613 if sidecar.exists() {
6614 let mut permissions = std::fs::metadata(&sidecar)
6615 .expect("sidecar metadata")
6616 .permissions();
6617 permissions.set_mode(0o444);
6618 std::fs::set_permissions(&sidecar, permissions).expect("freeze sidecar");
6619 }
6620 }
6621 }
6622
6623 let pool = Arc::new(
6624 ConnectionPool::new(PoolConfig {
6625 path: Some(path.clone()),
6626 read_only: true,
6627 ..PoolConfig::default()
6628 })
6629 .expect("read-only pool open"),
6630 );
6631 assert!(
6632 pool.open_standalone_writer().is_err(),
6633 "test precondition: a read-only pool must never be able to open a dedicated \
6634 checkpoint connection"
6635 );
6636
6637 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6638 let subscriber = CaptureSubscriber {
6639 events: std::sync::Arc::clone(&buffer),
6640 };
6641 let _tracing_guard = tracing::subscriber::set_default(subscriber);
6642
6643 let _tx_handle = khive_storage::tx_registry::register(Some(
6644 "checkpoint_task_conn_unavailable_sweep_test".to_string(),
6645 ));
6646
6647 let cfg = CheckpointConfig {
6648 interval: Duration::from_millis(10),
6649 tx_warn_secs: Duration::from_millis(1),
6650 tx_max_age_secs: Duration::from_millis(1),
6651 ..CheckpointConfig::default()
6652 };
6653
6654 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
6655 let handle = tokio::spawn(run_checkpoint_task(
6656 Arc::clone(&pool),
6657 cfg,
6658 None,
6659 shutdown_rx,
6660 true,
6661 ));
6662
6663 let deadline = tokio::time::Instant::now() + Duration::from_secs(10);
6669 while checkpoint_skipped_ticks() == 0 {
6670 assert!(
6671 tokio::time::Instant::now() < deadline,
6672 "test setup must actually drive at least one Skipped tick for this \
6673 regression to mean anything (none within 10s)"
6674 );
6675 tokio::time::sleep(Duration::from_millis(10)).await;
6676 }
6677
6678 loop {
6679 let events = buffer.lock().unwrap().clone();
6680 if events.iter().any(|e| {
6681 e.tx_label.as_deref() == Some("checkpoint_task_conn_unavailable_sweep_test")
6682 && e.message
6683 .as_deref()
6684 .is_some_and(|m| m.contains("stale-op cap"))
6685 }) {
6686 break;
6687 }
6688 assert!(
6689 tokio::time::Instant::now() < deadline,
6690 "expected the age sweep to fire even though every tick's dedicated \
6691 connection was unavailable within 10s, got: {events:?}"
6692 );
6693 tokio::time::sleep(Duration::from_millis(10)).await;
6694 }
6695
6696 shutdown_tx.send(()).expect("send shutdown signal");
6697 tokio::time::timeout(Duration::from_secs(1), handle)
6698 .await
6699 .expect("checkpoint task should exit within 1s")
6700 .expect("checkpoint task panicked");
6701 drop(_tx_handle);
6702 }
6703
6704 #[tokio::test]
6708 async fn session_sweep_task_exits_on_shutdown_signal() {
6709 let cfg = SessionSweepConfig {
6710 interval: Duration::from_millis(10),
6711 ..SessionSweepConfig::default()
6712 };
6713 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
6714 let handle = tokio::spawn(run_session_sweep_task(Vec::new(), cfg, shutdown_rx));
6715
6716 shutdown_tx.send(()).expect("send shutdown signal");
6717
6718 tokio::time::timeout(Duration::from_secs(1), handle)
6719 .await
6720 .expect("session sweep task should exit within 1s")
6721 .expect("session sweep task panicked");
6722 }
6723
6724 async fn wait_for(deadline: Duration, mut cond: impl FnMut() -> bool) -> bool {
6728 let start = std::time::Instant::now();
6729 while start.elapsed() < deadline {
6730 if cond() {
6731 return true;
6732 }
6733 tokio::time::sleep(Duration::from_millis(5)).await;
6734 }
6735 cond()
6736 }
6737
6738 #[tokio::test]
6739 #[serial(khive_walpin_sidecar_env)]
6740 async fn walpin_observe_drops_beacon_when_heartbeat_write_fails() {
6741 let dir = tempfile::tempdir().unwrap();
6742 let db_path = dir.path().join("observe_gate.db");
6743 let sidecar_dir = crate::walpin::sidecar_dir_for(&db_path);
6744 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
6745 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
6746
6747 let mut state = WalpinSidecarState::new(
6748 Some(db_path.as_path()),
6749 true,
6750 "session",
6751 Duration::from_millis(500),
6752 )
6753 .expect("sidecar enabled for a file-backed path");
6754 let pid = std::process::id();
6755 state.register_beacon().await;
6756 let beacon_path = sidecar_dir.join(format!("{pid}.beacon"));
6757 let before = std::fs::metadata(&beacon_path)
6758 .expect("register_beacon must create the beacon file")
6759 .modified()
6760 .unwrap();
6761
6762 let obstruction = sidecar_dir.join(format!(".{pid}.json.tmp"));
6767 std::fs::create_dir(&obstruction).unwrap();
6768
6769 tokio::time::sleep(Duration::from_millis(20)).await;
6770 let over_threshold = Some(khive_storage::tx_registry::OldestSpan {
6771 id: khive_storage::tx_registry::TxId(1),
6772 age: Duration::from_secs(60),
6773 label: None,
6774 origin: khive_storage::tx_registry::TxOrigin::Unscoped,
6775 });
6776 state
6777 .observe(over_threshold.clone(), Duration::from_secs(30))
6778 .await;
6779
6780 assert!(
6781 !sidecar_dir.join(format!("{pid}.json")).exists(),
6782 "heartbeat write must have failed"
6783 );
6784 assert!(
6787 !beacon_path.exists(),
6788 "a failed heartbeat write must remove the beacon — a still-fresh \
6789 beacon with no heartbeat would classify registered-silent \
6790 (before-mtime {before:?})"
6791 );
6792
6793 std::fs::remove_dir(&obstruction).unwrap();
6796 state.observe(over_threshold, Duration::from_secs(30)).await;
6797 assert!(
6798 sidecar_dir.join(format!("{pid}.json")).exists(),
6799 "heartbeat must land once the write path recovers"
6800 );
6801 assert!(
6802 beacon_path.exists(),
6803 "beacon must re-register on the first healthy tick after removal"
6804 );
6805 }
6806
6807 #[tokio::test]
6808 #[serial(khive_walpin_sidecar_env)]
6809 async fn walpin_observe_touches_mtime_without_rewriting_body_when_content_unchanged() {
6810 let dir = tempfile::tempdir().unwrap();
6811 let db_path = dir.path().join("observe_touch.db");
6812 let sidecar_dir = crate::walpin::sidecar_dir_for(&db_path);
6813 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
6814 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
6815
6816 let mut state = WalpinSidecarState::new(
6817 Some(db_path.as_path()),
6818 true,
6819 "session",
6820 Duration::from_millis(500),
6821 )
6822 .expect("sidecar enabled for a file-backed path");
6823 let pid = std::process::id();
6824 let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
6825 let span = khive_storage::tx_registry::OldestSpan {
6826 id: khive_storage::tx_registry::TxId(1),
6827 age: Duration::from_secs(60),
6828 label: None,
6829 origin: khive_storage::tx_registry::TxOrigin::Unscoped,
6830 };
6831
6832 state
6833 .observe(Some(span.clone()), Duration::from_secs(30))
6834 .await;
6835 let body_after_create = std::fs::read(&heartbeat_path).expect("heartbeat written");
6836
6837 let backdated = std::time::SystemTime::now() - Duration::from_secs(120);
6843 std::fs::OpenOptions::new()
6847 .write(true)
6848 .open(&heartbeat_path)
6849 .unwrap()
6850 .set_modified(backdated)
6851 .unwrap();
6852
6853 state.observe(Some(span), Duration::from_secs(30)).await;
6854
6855 let body_after_second_observe =
6856 std::fs::read(&heartbeat_path).expect("heartbeat still present");
6857 assert_eq!(
6858 body_after_create, body_after_second_observe,
6859 "unchanged oldest-span identity/label/attribution/cadence must touch mtime, \
6860 not rewrite the body"
6861 );
6862 let mtime_after = std::fs::metadata(&heartbeat_path)
6863 .unwrap()
6864 .modified()
6865 .unwrap();
6866 assert!(
6867 mtime_after > backdated,
6868 "the touch must advance mtime past the backdated value"
6869 );
6870 }
6871
6872 #[tokio::test]
6873 #[serial(khive_walpin_sidecar_env)]
6874 async fn walpin_observe_recreates_heartbeat_after_it_is_deleted_while_span_still_live() {
6875 let dir = tempfile::tempdir().unwrap();
6876 let db_path = dir.path().join("observe_recreate.db");
6877 let sidecar_dir = crate::walpin::sidecar_dir_for(&db_path);
6878 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
6879 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
6880
6881 let mut state = WalpinSidecarState::new(
6882 Some(db_path.as_path()),
6883 true,
6884 "session",
6885 Duration::from_millis(500),
6886 )
6887 .expect("sidecar enabled for a file-backed path");
6888 let pid = std::process::id();
6889 let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
6890 let span = khive_storage::tx_registry::OldestSpan {
6891 id: khive_storage::tx_registry::TxId(1),
6892 age: Duration::from_secs(60),
6893 label: None,
6894 origin: khive_storage::tx_registry::TxOrigin::Unscoped,
6895 };
6896
6897 state
6898 .observe(Some(span.clone()), Duration::from_secs(30))
6899 .await;
6900 assert!(heartbeat_path.exists(), "heartbeat written on first tick");
6901
6902 std::fs::remove_file(&heartbeat_path).unwrap();
6908 assert!(!heartbeat_path.exists());
6909
6910 state.observe(Some(span), Duration::from_secs(30)).await;
6911
6912 assert!(
6913 heartbeat_path.exists(),
6914 "a touch failure against a deleted heartbeat must recreate it via a full write"
6915 );
6916 let recreated: crate::walpin::WalpinHeartbeat =
6917 serde_json::from_slice(&std::fs::read(&heartbeat_path).unwrap()).unwrap();
6918 assert_eq!(recreated.pid, pid);
6919 assert_eq!(recreated.oldest_tx_age_secs, 60.0);
6920 }
6921
6922 #[tokio::test]
6923 #[serial(tx_registry, khive_walpin_sidecar_env)]
6924 async fn session_sweep_task_writes_and_clears_walpin_heartbeat() {
6925 let dir = tempfile::tempdir().unwrap();
6926 let db_path = dir.path().join("session_sweep.db");
6927 let pool = file_pool(&db_path);
6928 let sidecar_dir =
6929 crate::walpin::sidecar_dir_for(pool.canonical_path().expect("file-backed pool"));
6930 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
6931 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
6932
6933 let cfg = SessionSweepConfig {
6934 interval: Duration::from_millis(10),
6935 tx_warn_secs: Duration::from_millis(20),
6936 tx_max_age_secs: Duration::from_millis(500),
6937 };
6938 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
6939 let handle = tokio::spawn(run_session_sweep_task(
6940 vec![SweepBackend {
6941 pool: Arc::clone(&pool),
6942 is_main: true,
6943 }],
6944 cfg,
6945 shutdown_rx,
6946 ));
6947
6948 let pid = std::process::id();
6955 let beacon = crate::walpin::beacon_path(&sidecar_dir, pid);
6956 let beacon_registered = wait_for(Duration::from_secs(2), || beacon.exists()).await;
6957 assert!(
6958 beacon_registered,
6959 "a quiet process must still register its one-time beacon"
6960 );
6961 assert!(
6962 !sidecar_dir.join(format!("{pid}.json")).exists(),
6963 "a quiet process must not write a walpin heartbeat"
6964 );
6965
6966 let tx_handle =
6967 khive_storage::tx_registry::register(Some("session_sweep_walpin_test".to_string()));
6968 let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
6969 assert!(
6970 wait_for(Duration::from_secs(2), || heartbeat_path.exists()).await,
6971 "expected a walpin heartbeat once the span crossed tx_warn_secs"
6972 );
6973 let body = std::fs::read_to_string(&heartbeat_path).unwrap();
6974 let hb: crate::walpin::WalpinHeartbeat = serde_json::from_str(&body).unwrap();
6975 assert_eq!(hb.pid, pid);
6976 assert_eq!(hb.process_role, "session");
6977 assert_eq!(
6978 hb.oldest_tx_label.as_deref(),
6979 Some("session_sweep_walpin_test")
6980 );
6981 assert_eq!(
6982 hb.attribution_basis.as_deref(),
6983 Some("fallback"),
6984 "an Unscoped span observed only through the main view's fallback \
6985 must carry attribution_basis=\"fallback\", never \"origin\""
6986 );
6987
6988 drop(tx_handle);
6989 assert!(
6990 wait_for(Duration::from_secs(2), || !heartbeat_path.exists()).await,
6991 "heartbeat must be removed once the stale span clears"
6992 );
6993
6994 shutdown_tx.send(()).expect("send shutdown signal");
6995 tokio::time::timeout(Duration::from_secs(1), handle)
6996 .await
6997 .expect("session sweep task should exit within 1s")
6998 .expect("session sweep task panicked");
6999 }
7000
7001 #[tokio::test]
7012 #[serial(tx_registry, khive_walpin_sidecar_env)]
7013 async fn session_sweep_fan_out_scopes_secondary_span_to_secondary_sidecar_only() {
7014 let main_dir = tempfile::tempdir().unwrap();
7015 let secondary_dir = tempfile::tempdir().unwrap();
7016 let main_pool = file_pool(&main_dir.path().join("main.db"));
7017 let secondary_pool = file_pool(&secondary_dir.path().join("secondary.db"));
7018 let main_sidecar =
7019 crate::walpin::sidecar_dir_for(main_pool.canonical_path().expect("file-backed"));
7020 let secondary_sidecar =
7021 crate::walpin::sidecar_dir_for(secondary_pool.canonical_path().expect("file-backed"));
7022 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
7023 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
7024
7025 let cfg = SessionSweepConfig {
7026 interval: Duration::from_millis(10),
7027 tx_warn_secs: Duration::from_millis(20),
7028 tx_max_age_secs: Duration::from_millis(500),
7029 };
7030 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
7031 let handle = tokio::spawn(run_session_sweep_task(
7032 vec![
7033 SweepBackend {
7034 pool: Arc::clone(&main_pool),
7035 is_main: true,
7036 },
7037 SweepBackend {
7038 pool: Arc::clone(&secondary_pool),
7039 is_main: false,
7040 },
7041 ],
7042 cfg,
7043 shutdown_rx,
7044 ));
7045
7046 let pid = std::process::id();
7047 let secondary_heartbeat = secondary_sidecar.join(format!("{pid}.json"));
7048 let main_heartbeat = main_sidecar.join(format!("{pid}.json"));
7049
7050 let tx_handle = khive_storage::tx_registry::register_scoped(
7051 Some("graph_traverse_read".to_string()),
7052 secondary_pool.origin(),
7053 );
7054 assert!(
7055 wait_for(Duration::from_secs(2), || secondary_heartbeat.exists()).await,
7056 "expected a walpin heartbeat in the secondary backend's own sidecar"
7057 );
7058 assert!(
7059 !main_heartbeat.exists(),
7060 "a span scoped to the secondary backend's origin must never produce \
7061 a heartbeat in the main backend's sidecar"
7062 );
7063
7064 let body = std::fs::read_to_string(&secondary_heartbeat).unwrap();
7065 let hb: crate::walpin::WalpinHeartbeat = serde_json::from_str(&body).unwrap();
7066 assert_eq!(hb.oldest_tx_label.as_deref(), Some("graph_traverse_read"));
7067 assert_eq!(
7068 hb.attribution_basis.as_deref(),
7069 Some("origin"),
7070 "a Secondary-view winner is always Database-origin-backed — never fallback"
7071 );
7072
7073 drop(tx_handle);
7074 assert!(
7075 wait_for(Duration::from_secs(2), || !secondary_heartbeat.exists()).await,
7076 "secondary heartbeat must be removed once its span clears"
7077 );
7078 assert!(
7079 !main_heartbeat.exists(),
7080 "the main sidecar must have stayed untouched for the whole tick sequence"
7081 );
7082
7083 shutdown_tx.send(()).expect("send shutdown signal");
7084 tokio::time::timeout(Duration::from_secs(1), handle)
7085 .await
7086 .expect("session sweep task should exit within 1s")
7087 .expect("session sweep task panicked");
7088 }
7089
7090 #[tokio::test]
7099 #[serial(tx_registry, checkpoint_skip_metrics, khive_walpin_sidecar_env)]
7100 async fn checkpoint_task_ignores_span_registered_against_other_backend_origin_and_unscoped() {
7101 let dir_a = tempfile::tempdir().unwrap();
7102 let dir_b = tempfile::tempdir().unwrap();
7103 let pool_a = file_pool(&dir_a.path().join("backend_a.db"));
7104 let pool_b = file_pool(&dir_b.path().join("backend_b.db"));
7107 let sidecar_a =
7108 crate::walpin::sidecar_dir_for(pool_a.canonical_path().expect("file-backed"));
7109 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
7110 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
7111
7112 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7113 let subscriber = CaptureSubscriber {
7114 events: std::sync::Arc::clone(&buffer),
7115 };
7116 let _tracing_guard = tracing::subscriber::set_default(subscriber);
7117
7118 let _b_origin_handle = khive_storage::tx_registry::register_scoped(
7119 Some("b_origin_span_ignored_by_a".to_string()),
7120 pool_b.origin(),
7121 );
7122 let _unscoped_handle = khive_storage::tx_registry::register(Some(
7123 "unscoped_span_ignored_by_secondary".to_string(),
7124 ));
7125
7126 let cfg = CheckpointConfig {
7127 interval: Duration::from_millis(10),
7128 tx_warn_secs: Duration::from_millis(1),
7129 tx_max_age_secs: Duration::from_millis(1),
7130 ..CheckpointConfig::default()
7131 };
7132 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
7133 let handle = tokio::spawn(run_checkpoint_task(
7134 pool_a,
7135 cfg,
7136 None,
7137 shutdown_rx,
7138 false, ));
7140
7141 tokio::time::sleep(Duration::from_millis(60)).await;
7146 shutdown_tx.send(()).expect("send shutdown signal");
7147 tokio::time::timeout(Duration::from_secs(1), handle)
7148 .await
7149 .expect("checkpoint task should exit within 1s")
7150 .expect("checkpoint task panicked");
7151
7152 let events = buffer.lock().unwrap();
7153 assert!(
7154 events.iter().all(|e| {
7155 e.tx_label.as_deref() != Some("b_origin_span_ignored_by_a")
7156 && e.tx_label.as_deref() != Some("unscoped_span_ignored_by_secondary")
7157 }),
7158 "backend A's Secondary filter must never emit an age alert naming a span \
7159 registered against a different backend's origin or an Unscoped span, got: \
7160 {events:?}"
7161 );
7162 assert!(
7163 !sidecar_a
7164 .join(format!("{}.json", std::process::id()))
7165 .exists(),
7166 "backend A's own sidecar must never gain a heartbeat from a span it does not own"
7167 );
7168 }
7169
7170 #[tokio::test]
7177 #[serial(tx_registry, checkpoint_skip_metrics, khive_walpin_sidecar_env)]
7178 async fn checkpoint_task_detects_and_enumerates_secondary_backend_stall() {
7179 let dir = tempfile::tempdir().unwrap();
7180 let pool = file_pool(&dir.path().join("secondary_stall.db"));
7181 let sidecar_dir =
7182 crate::walpin::sidecar_dir_for(pool.canonical_path().expect("file-backed"));
7183 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
7184 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
7185
7186 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7187 let subscriber = CaptureSubscriber {
7188 events: std::sync::Arc::clone(&buffer),
7189 };
7190 let _tracing_guard = tracing::subscriber::set_default(subscriber);
7191
7192 let tx_handle = khive_storage::tx_registry::register_scoped(
7193 Some("secondary_stall_test".to_string()),
7194 pool.origin(),
7195 );
7196
7197 let cfg = CheckpointConfig {
7198 interval: Duration::from_millis(10),
7199 tx_warn_secs: Duration::from_millis(5),
7200 tx_max_age_secs: Duration::from_millis(500),
7201 ..CheckpointConfig::default()
7202 };
7203 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
7204 let pid = std::process::id();
7205 let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
7206 let handle = tokio::spawn(run_checkpoint_task(
7207 pool,
7208 cfg,
7209 None,
7210 shutdown_rx,
7211 false, ));
7213
7214 assert!(
7215 wait_for(Duration::from_secs(2), || heartbeat_path.exists()).await,
7216 "expected a walpin heartbeat once the secondary backend's own span crossed \
7217 tx_warn_secs"
7218 );
7219 let body = std::fs::read_to_string(&heartbeat_path).unwrap();
7220 let hb: crate::walpin::WalpinHeartbeat = serde_json::from_str(&body).unwrap();
7221 assert_eq!(hb.oldest_tx_label.as_deref(), Some("secondary_stall_test"));
7222 assert_eq!(
7223 hb.attribution_basis.as_deref(),
7224 Some("origin"),
7225 "a Secondary-view winner is always Database-origin-backed — never fallback"
7226 );
7227 assert!(
7228 hb.oldest_tx_age_secs > 0.0,
7229 "the heartbeat must reflect a nonzero stale age for the secondary backend's own \
7230 span, got {hb:?}"
7231 );
7232
7233 shutdown_tx.send(()).expect("send shutdown signal");
7234 tokio::time::timeout(Duration::from_secs(1), handle)
7235 .await
7236 .expect("checkpoint task should exit within 1s")
7237 .expect("checkpoint task panicked");
7238
7239 drop(tx_handle);
7240
7241 let events = buffer.lock().unwrap();
7242 assert!(
7243 events.iter().any(|e| {
7244 e.tx_label.as_deref() == Some("secondary_stall_test")
7245 && e.message
7246 .as_deref()
7247 .is_some_and(|m| m.contains("ADR-091 Plank 1"))
7248 }),
7249 "expected the secondary backend's own checkpoint task to emit a Plank 1 age alert \
7250 for its own stalled span, got: {events:?}"
7251 );
7252 }
7253
7254 #[test]
7255 fn wal_pin_depth_arithmetic_against_real_connection() {
7256 let dir = tempfile::tempdir().unwrap();
7257 let path = dir.path().join("pin_depth.db");
7258 let pool = file_pool(&path);
7259 let writer = pool.try_writer().expect("acquire writer");
7260 let conn = writer.conn();
7261
7262 conn.execute_batch("CREATE TABLE t (v INTEGER)").unwrap();
7263 conn.execute_batch("INSERT INTO t (v) VALUES (1)").unwrap();
7264
7265 let (log, checkpointed) =
7266 query_wal_pin_depth(conn).expect("PRAGMA wal_checkpoint(PASSIVE) must succeed");
7267 assert!(
7271 log >= checkpointed,
7272 "checkpointed frames cannot exceed log frames"
7273 );
7274 assert_eq!(
7275 log - checkpointed,
7276 0,
7277 "an unpinned WAL must fully checkpoint under PASSIVE"
7278 );
7279 }
7280
7281 #[test]
7282 fn wal_pin_depth_arithmetic_on_in_memory_pool_errors_cleanly() {
7283 let cfg = PoolConfig {
7287 path: None,
7288 ..PoolConfig::default()
7289 };
7290 let pool = ConnectionPool::new(cfg).expect("in-memory pool");
7291 let writer = pool.try_writer().expect("acquire writer");
7292 let _ = query_wal_pin_depth(writer.conn());
7295 }
7296
7297 #[cfg(unix)]
7301 #[test]
7302 fn routine_wal_backend_key_preserves_non_utf8_path_bytes() {
7303 use std::ffi::OsString;
7304 use std::os::unix::ffi::OsStringExt;
7305
7306 let path_a = PathBuf::from(OsString::from_vec(b"/tmp/khive-wal-\x80.db".to_vec()));
7307 let path_b = PathBuf::from(OsString::from_vec(b"/tmp/khive-wal-\x81.db".to_vec()));
7308 assert_eq!(
7309 path_a.display().to_string(),
7310 path_b.display().to_string(),
7311 "fixture must reproduce the lossy display-label collision"
7312 );
7313 assert_ne!(
7314 checkpoint_db_key_from_path(Some(&path_a)),
7315 checkpoint_db_key_from_path(Some(&path_b)),
7316 "backend keys must retain the canonical path's exact OS bytes"
7317 );
7318 }
7319
7320 #[test]
7325 #[serial(checkpoint_skip_metrics)]
7326 fn routine_checkpoint_records_one_pass_logical_and_physical_wal_sample() {
7327 let dir = tempfile::tempdir().unwrap();
7328 let path = dir.path().join("routine_wal_sample.db");
7329 let pool = file_pool(&path);
7330
7331 {
7332 let writer = pool.try_writer().expect("writer");
7333 writer
7334 .conn()
7335 .execute_batch(
7336 "PRAGMA wal_autocheckpoint=0; \
7337 CREATE TABLE t (id INTEGER PRIMARY KEY, payload TEXT); \
7338 INSERT INTO t VALUES (0, 'seed');",
7339 )
7340 .unwrap();
7341 }
7342
7343 let reader = rusqlite::Connection::open_with_flags(
7344 &path,
7345 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY,
7346 )
7347 .unwrap();
7348 reader.execute_batch("BEGIN").unwrap();
7349 let _: i64 = reader
7350 .query_row("SELECT COUNT(*) FROM t", [], |row| row.get(0))
7351 .unwrap();
7352
7353 {
7354 let writer = pool.try_writer().expect("writer");
7355 writer.conn().execute_batch("BEGIN IMMEDIATE").unwrap();
7356 for id in 1..=256_i64 {
7357 writer
7358 .conn()
7359 .execute("INSERT INTO t VALUES (?1, printf('%.*c', 2048, 'x'))", [id])
7360 .unwrap();
7361 }
7362 writer.conn().execute_batch("COMMIT").unwrap();
7363 }
7364
7365 let checkpoint_conn = pool.open_standalone_writer().unwrap();
7366 let pragma_calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
7367 let pragma_calls_from_hook = Arc::clone(&pragma_calls);
7368 checkpoint_conn
7369 .authorizer(Some(move |context: rusqlite::hooks::AuthContext<'_>| {
7370 if matches!(
7371 context.action,
7372 AuthAction::Pragma { pragma_name, .. }
7373 if pragma_name.eq_ignore_ascii_case("wal_checkpoint")
7374 ) {
7375 pragma_calls_from_hook.fetch_add(1, Ordering::SeqCst);
7376 }
7377 Authorization::Allow
7378 }))
7379 .unwrap();
7380
7381 checkpoint_once(
7382 &pool,
7383 &checkpoint_conn,
7384 &CheckpointConfig::default(),
7385 &mut TruncateState::default(),
7386 )
7387 .unwrap();
7388 checkpoint_conn
7389 .authorizer(None::<fn(rusqlite::hooks::AuthContext<'_>) -> Authorization>)
7390 .unwrap();
7391
7392 assert_eq!(
7393 pragma_calls.load(Ordering::SeqCst),
7394 1,
7395 "one routine tick must issue exactly one PASSIVE checkpoint"
7396 );
7397 let pinned = routine_wal_observation(&pool).expect("routine sample");
7398 assert!(pinned.log_frames > 0, "the test must create WAL frames");
7399 assert!(
7400 pinned.pending_frames > 0,
7401 "the old reader must leave a logical backlog: {pinned:?}"
7402 );
7403 assert_eq!(
7404 pinned.pending_frames,
7405 pinned.log_frames.saturating_sub(pinned.checkpointed_frames)
7406 );
7407 assert!(
7408 pinned.physical_wal_bytes.is_some_and(|bytes| bytes > 0),
7409 "the physical sidecar high-water must be reported separately: {pinned:?}"
7410 );
7411
7412 reader.execute_batch("COMMIT").unwrap();
7413 checkpoint_once(
7414 &pool,
7415 &checkpoint_conn,
7416 &CheckpointConfig::default(),
7417 &mut TruncateState::default(),
7418 )
7419 .unwrap();
7420 let drained = routine_wal_observation(&pool).expect("drained routine sample");
7421 assert_eq!(drained.pending_frames, 0, "unpinned PASSIVE must drain");
7422 assert!(
7423 drained.physical_wal_bytes.is_some_and(|bytes| bytes > 0),
7424 "PASSIVE may reuse rather than shrink the physical WAL; the two gauges must remain \
7425 independently visible: {drained:?}"
7426 );
7427 }
7428}