1use std::collections::HashMap;
45use std::path::{Path, PathBuf};
46use std::sync::atomic::{AtomicU64, Ordering};
47use std::sync::{Arc, Mutex, OnceLock};
48use std::time::{Duration, Instant};
49
50use crate::pool::ConnectionPool;
51
52static LAST_WAL_PAGES: AtomicU64 = AtomicU64::new(u64::MAX);
61
62static TRUNCATE_ATTEMPTS: AtomicU64 = AtomicU64::new(0);
65
66static TRUNCATE_CONSECUTIVE_FAILURES: AtomicU64 = AtomicU64::new(0);
70
71static CHECKPOINT_SKIPPED_TICKS: AtomicU64 = AtomicU64::new(0);
77
78static CHECKPOINT_CONSECUTIVE_SKIPS: AtomicU64 = AtomicU64::new(0);
83
84static CHECKPOINT_LAST_SKIP_WAL_PAGES: AtomicU64 = AtomicU64::new(u64::MAX);
88
89static CHECKPOINT_PRESSURE_ELEVATED_TICKS: AtomicU64 = AtomicU64::new(0);
92
93static CHECKPOINT_PRESSURE_EPISODES_STARTED: AtomicU64 = AtomicU64::new(0);
95
96static CHECKPOINT_PRESSURE_EPISODES_RECOVERED: AtomicU64 = AtomicU64::new(0);
98
99static CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS: AtomicU64 = AtomicU64::new(0);
101
102static CHECKPOINT_LIFECYCLE_APPEND_FAILURES: AtomicU64 = AtomicU64::new(0);
104
105static CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS: AtomicU64 = AtomicU64::new(0);
108
109static READ_TX_MAX_AGE_EVICTIONS: AtomicU64 = AtomicU64::new(0);
115
116#[derive(Debug, Clone, PartialEq, Eq)]
121pub struct RoutineWalObservation {
122 pub busy: i64,
123 pub log_frames: u64,
124 pub checkpointed_frames: u64,
125 pub pending_frames: u64,
126 pub physical_wal_bytes: Option<u64>,
127 pub observed_at_unix_ms: u64,
128}
129
130#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
134#[serde(default)]
135pub struct CheckpointTiming {
136 pub ticks: u64,
137 pub elapsed_us_sum: u64,
138 pub elapsed_us_max: u64,
139 pub busy_ticks: u64,
140 pub error_ticks: u64,
141}
142
143static CHECKPOINT_TIMINGS: OnceLock<Mutex<HashMap<Option<PathBuf>, CheckpointTiming>>> =
144 OnceLock::new();
145
146fn checkpoint_timings() -> &'static Mutex<HashMap<Option<PathBuf>, CheckpointTiming>> {
147 CHECKPOINT_TIMINGS.get_or_init(|| Mutex::new(HashMap::new()))
148}
149
150fn record_checkpoint_timing(pool: &ConnectionPool, elapsed_us: u64, busy: Option<i64>) {
151 let mut timings = checkpoint_timings()
152 .lock()
153 .unwrap_or_else(std::sync::PoisonError::into_inner);
154 let timing = timings.entry(checkpoint_db_key(pool)).or_default();
155 timing.ticks = timing.ticks.saturating_add(1);
156 timing.elapsed_us_sum = timing.elapsed_us_sum.saturating_add(elapsed_us);
157 timing.elapsed_us_max = timing.elapsed_us_max.max(elapsed_us);
158 timing.busy_ticks = timing
159 .busy_ticks
160 .saturating_add(u64::from(busy.is_some_and(|value| value != 0)));
161 timing.error_ticks = timing.error_ticks.saturating_add(u64::from(busy.is_none()));
162}
163
164pub fn checkpoint_timing(pool: &ConnectionPool) -> CheckpointTiming {
166 checkpoint_timings()
167 .lock()
168 .unwrap_or_else(std::sync::PoisonError::into_inner)
169 .get(&checkpoint_db_key(pool))
170 .copied()
171 .unwrap_or_default()
172}
173
174static ROUTINE_WAL_OBSERVATIONS: OnceLock<Mutex<HashMap<Option<PathBuf>, RoutineWalObservation>>> =
178 OnceLock::new();
179
180fn routine_wal_observations() -> &'static Mutex<HashMap<Option<PathBuf>, RoutineWalObservation>> {
181 ROUTINE_WAL_OBSERVATIONS.get_or_init(|| Mutex::new(HashMap::new()))
182}
183
184fn checkpoint_db_key_from_path(path: Option<&Path>) -> Option<PathBuf> {
185 path.map(Path::to_path_buf)
186}
187
188fn checkpoint_db_key(pool: &ConnectionPool) -> Option<PathBuf> {
189 checkpoint_db_key_from_path(pool.canonical_path())
190}
191
192fn observed_at_unix_ms() -> u64 {
193 std::time::SystemTime::now()
194 .duration_since(std::time::UNIX_EPOCH)
195 .map(|duration| duration.as_millis() as u64)
196 .unwrap_or(0)
197}
198
199fn physical_wal_bytes(pool: &ConnectionPool) -> Option<u64> {
200 let path = pool.canonical_path()?;
201 let mut sidecar = path.as_os_str().to_os_string();
202 sidecar.push("-wal");
203 std::fs::metadata(PathBuf::from(sidecar))
204 .ok()
205 .map(|metadata| metadata.len())
206}
207
208fn record_routine_wal_observation(
209 pool: &ConnectionPool,
210 raw: RawCheckpointObservation,
211) -> RoutineWalObservation {
212 let log_frames = raw.log_frames.max(0) as u64;
213 let checkpointed_frames = raw.checkpointed_frames.max(0) as u64;
214 let observation = RoutineWalObservation {
215 busy: raw.busy,
216 log_frames,
217 checkpointed_frames,
218 pending_frames: log_frames.saturating_sub(checkpointed_frames),
219 physical_wal_bytes: physical_wal_bytes(pool),
220 observed_at_unix_ms: observed_at_unix_ms(),
221 };
222 routine_wal_observations()
223 .lock()
224 .unwrap_or_else(std::sync::PoisonError::into_inner)
225 .insert(checkpoint_db_key(pool), observation.clone());
226 observation
227}
228
229pub fn routine_wal_observation(pool: &ConnectionPool) -> Option<RoutineWalObservation> {
232 routine_wal_observations()
233 .lock()
234 .unwrap_or_else(std::sync::PoisonError::into_inner)
235 .get(&checkpoint_db_key(pool))
236 .cloned()
237}
238
239pub fn last_observed_wal_pages() -> Option<u64> {
242 match LAST_WAL_PAGES.load(Ordering::Relaxed) {
243 u64::MAX => None,
244 pages => Some(pages),
245 }
246}
247
248pub fn truncate_attempts() -> u64 {
250 TRUNCATE_ATTEMPTS.load(Ordering::Relaxed)
251}
252
253pub fn truncate_consecutive_failures() -> u64 {
255 TRUNCATE_CONSECUTIVE_FAILURES.load(Ordering::Relaxed)
256}
257
258pub fn checkpoint_skipped_ticks() -> u64 {
261 CHECKPOINT_SKIPPED_TICKS.load(Ordering::Relaxed)
262}
263
264pub fn checkpoint_consecutive_skips() -> u64 {
266 CHECKPOINT_CONSECUTIVE_SKIPS.load(Ordering::Relaxed)
267}
268
269pub fn checkpoint_last_skip_wal_pages() -> Option<u64> {
272 match CHECKPOINT_LAST_SKIP_WAL_PAGES.load(Ordering::Relaxed) {
273 u64::MAX => None,
274 pages => Some(pages),
275 }
276}
277
278pub fn checkpoint_pressure_elevated_ticks() -> u64 {
280 CHECKPOINT_PRESSURE_ELEVATED_TICKS.load(Ordering::Relaxed)
281}
282
283pub fn checkpoint_pressure_episodes_started() -> u64 {
285 CHECKPOINT_PRESSURE_EPISODES_STARTED.load(Ordering::Relaxed)
286}
287
288pub fn checkpoint_pressure_episodes_recovered() -> u64 {
290 CHECKPOINT_PRESSURE_EPISODES_RECOVERED.load(Ordering::Relaxed)
291}
292
293pub fn checkpoint_lifecycle_append_attempts() -> u64 {
295 CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.load(Ordering::Relaxed)
296}
297
298pub fn checkpoint_lifecycle_append_failures() -> u64 {
300 CHECKPOINT_LIFECYCLE_APPEND_FAILURES.load(Ordering::Relaxed)
301}
302
303pub fn read_tx_max_age_evictions() -> u64 {
306 READ_TX_MAX_AGE_EVICTIONS.load(Ordering::Relaxed)
307}
308
309pub(crate) fn note_read_tx_max_age_eviction() {
315 READ_TX_MAX_AGE_EVICTIONS.fetch_add(1, Ordering::Relaxed);
316}
317
318pub fn checkpoint_lifecycle_enqueue_drops() -> u64 {
320 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.load(Ordering::Relaxed)
321}
322
323fn note_checkpoint_skipped() {
328 CHECKPOINT_SKIPPED_TICKS.fetch_add(1, Ordering::Relaxed);
329 CHECKPOINT_CONSECUTIVE_SKIPS.fetch_add(1, Ordering::Relaxed);
330 if let Some(pages) = last_observed_wal_pages() {
331 CHECKPOINT_LAST_SKIP_WAL_PAGES.store(pages, Ordering::Relaxed);
332 }
333}
334
335fn note_checkpoint_observed(_wal_pages: u64) {
340 CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
341}
342
343fn note_checkpoint_pressure_observation(above_warn: bool, was_above_warn: bool) {
344 if above_warn {
345 CHECKPOINT_PRESSURE_ELEVATED_TICKS.fetch_add(1, Ordering::Relaxed);
346 if !was_above_warn {
347 CHECKPOINT_PRESSURE_EPISODES_STARTED.fetch_add(1, Ordering::Relaxed);
348 }
349 } else if was_above_warn {
350 CHECKPOINT_PRESSURE_EPISODES_RECOVERED.fetch_add(1, Ordering::Relaxed);
351 }
352}
353
354#[cfg(test)]
358pub(crate) fn reset_checkpoint_metrics_for_tests() {
359 CHECKPOINT_SKIPPED_TICKS.store(0, Ordering::Relaxed);
360 CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
361 CHECKPOINT_LAST_SKIP_WAL_PAGES.store(u64::MAX, Ordering::Relaxed);
362 CHECKPOINT_PRESSURE_ELEVATED_TICKS.store(0, Ordering::Relaxed);
363 CHECKPOINT_PRESSURE_EPISODES_STARTED.store(0, Ordering::Relaxed);
364 CHECKPOINT_PRESSURE_EPISODES_RECOVERED.store(0, Ordering::Relaxed);
365 CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.store(0, Ordering::Relaxed);
366 CHECKPOINT_LIFECYCLE_APPEND_FAILURES.store(0, Ordering::Relaxed);
367 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.store(0, Ordering::Relaxed);
368 READ_TX_MAX_AGE_EVICTIONS.store(0, Ordering::Relaxed);
369}
370
371#[derive(Debug, Clone, Copy, PartialEq, Eq)]
382pub enum CheckpointTick {
383 Skipped,
386 Observed(u64),
388}
389
390pub const DEFAULT_WARN_SUSTAINED_CYCLES: u8 = 3;
393
394#[derive(Clone, Debug)]
399pub struct CheckpointConfig {
400 pub interval: Duration,
405
406 pub warn_pages: u64,
411
412 pub warn_sustained_cycles: u8,
420
421 pub high_water_pages: u64,
434
435 pub truncate_high_water_pages: u64,
446
447 pub truncate_min_interval: Duration,
458
459 pub truncate_busy_timeout: Duration,
466
467 pub tx_warn_secs: Duration,
475
476 pub tx_max_age_secs: Duration,
491}
492
493impl Default for CheckpointConfig {
494 fn default() -> Self {
495 Self {
496 interval: Duration::from_millis(500),
497 warn_pages: 2000,
498 warn_sustained_cycles: DEFAULT_WARN_SUSTAINED_CYCLES,
499 high_water_pages: 6000,
500 truncate_high_water_pages: 20_000,
501 truncate_min_interval: Duration::from_secs(300),
502 truncate_busy_timeout: Duration::from_millis(2000),
503 tx_warn_secs: Duration::from_secs(30),
504 tx_max_age_secs: Duration::from_secs(120),
505 }
506 }
507}
508
509impl CheckpointConfig {
510 pub fn from_env() -> Self {
514 let mut cfg = Self::default();
515
516 if let Ok(ms) = std::env::var("KHIVE_CHECKPOINT_INTERVAL_MS") {
517 if let Ok(v) = ms.parse::<u64>() {
518 if v > 0 {
519 cfg.interval = Duration::from_millis(v);
520 }
521 }
522 }
523
524 if let Ok(v) = std::env::var("KHIVE_WAL_WARN_PAGES") {
525 if let Ok(n) = v.parse::<u64>() {
526 if n > 0 {
527 cfg.warn_pages = n;
528 }
529 }
530 }
531
532 if let Ok(v) = std::env::var("KHIVE_WAL_WARN_SUSTAINED_CYCLES") {
533 if let Ok(n) = v.parse::<u8>() {
534 if n > 0 {
535 cfg.warn_sustained_cycles = n;
536 }
537 }
538 }
539
540 if let Ok(v) = std::env::var("KHIVE_WAL_HIGH_WATER_PAGES") {
541 if let Ok(n) = v.parse::<u64>() {
542 if n > 0 {
543 cfg.high_water_pages = n;
544 }
545 }
546 }
547
548 if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES") {
549 if let Ok(n) = v.parse::<u64>() {
550 if n > 0 {
551 cfg.truncate_high_water_pages = n;
552 }
553 }
554 }
555
556 if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS") {
557 if let Ok(n) = v.parse::<u64>() {
558 if n > 0 {
559 cfg.truncate_min_interval = Duration::from_secs(n);
560 }
561 }
562 }
563
564 if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_BUSY_MS") {
565 if let Ok(n) = v.parse::<u64>() {
566 if n > 0 {
567 cfg.truncate_busy_timeout = Duration::from_millis(n);
568 }
569 }
570 }
571
572 (cfg.tx_warn_secs, cfg.tx_max_age_secs) =
573 tx_age_thresholds_from_env(cfg.tx_warn_secs, cfg.tx_max_age_secs);
574
575 cfg
576 }
577}
578
579pub(crate) fn tx_age_thresholds_from_env(
593 default_warn: Duration,
594 default_max: Duration,
595) -> (Duration, Duration) {
596 let mut warn_secs = default_warn;
597 let mut max_age_secs = default_max;
598
599 if let Ok(v) = std::env::var("KHIVE_TX_WARN_SECS") {
600 if let Ok(n) = v.parse::<u64>() {
601 if n > 0 {
602 warn_secs = Duration::from_secs(n);
603 }
604 }
605 }
606
607 if let Ok(v) = std::env::var("KHIVE_TX_MAX_AGE_SECS") {
608 if let Ok(n) = v.parse::<u64>() {
609 if n > 0 {
610 max_age_secs = Duration::from_secs(n);
611 }
612 }
613 }
614
615 if warn_secs >= max_age_secs {
616 tracing::warn!(
617 configured_tx_warn_secs = warn_secs.as_secs_f64(),
618 configured_tx_max_age_secs = max_age_secs.as_secs_f64(),
619 fallback_tx_warn_secs = default_warn.as_secs_f64(),
620 fallback_tx_max_age_secs = default_max.as_secs_f64(),
621 "KHIVE_TX_WARN_SECS must be strictly less than KHIVE_TX_MAX_AGE_SECS; \
622 both transaction-age thresholds were rejected and reset to their defaults"
623 );
624 return (default_warn, default_max);
625 }
626
627 (warn_secs, max_age_secs)
628}
629
630#[cfg(unix)]
631const DEFAULT_WALPIN_FULL_SCAN_INTERVAL: Duration = Duration::from_secs(30);
632
633#[cfg(unix)]
634#[derive(Debug, Clone)]
635struct CachedWalpinAttribution {
636 report: crate::walpin::WalpinReport,
637 census: Result<crate::walpin::CensusResult, String>,
638 captured_at: Instant,
639}
640
641#[cfg(unix)]
642#[derive(Debug)]
643enum WalpinFullScanPlan {
644 Refresh {
645 previous_last_attempt: Option<Instant>,
646 },
647 Cached(CachedWalpinAttribution),
648 Suppressed,
649}
650
651#[derive(Debug)]
658pub struct TruncateState {
659 last_attempt: Option<Instant>,
663 consecutive_failures: u32,
668 #[cfg(unix)]
673 legacy_walpin_fallback_interval: Duration,
674 #[cfg(unix)]
679 walpin_full_scan_interval: Duration,
680 #[cfg(unix)]
681 walpin_full_scan_last_attempt: Option<Instant>,
682 #[cfg(unix)]
683 walpin_cached_attribution: Option<CachedWalpinAttribution>,
684 #[cfg(unix)]
687 sidecar_attribution_attempted_this_tick: bool,
688}
689
690impl Default for TruncateState {
691 fn default() -> Self {
692 Self {
693 last_attempt: None,
694 consecutive_failures: 0,
695 #[cfg(unix)]
696 legacy_walpin_fallback_interval: DEFAULT_SESSION_SWEEP_INTERVAL,
697 #[cfg(unix)]
698 walpin_full_scan_interval: DEFAULT_WALPIN_FULL_SCAN_INTERVAL,
699 #[cfg(unix)]
700 walpin_full_scan_last_attempt: None,
701 #[cfg(unix)]
702 walpin_cached_attribution: None,
703 #[cfg(unix)]
704 sidecar_attribution_attempted_this_tick: false,
705 }
706 }
707}
708
709impl TruncateState {
710 #[cfg(unix)]
711 fn with_legacy_walpin_fallback(interval: Duration) -> Self {
712 Self {
713 legacy_walpin_fallback_interval: interval,
714 ..Self::default()
715 }
716 }
717
718 #[cfg(all(test, unix))]
719 fn with_walpin_full_scan_cadence(interval: Duration) -> Self {
720 Self {
721 walpin_full_scan_interval: interval,
722 ..Self::default()
723 }
724 }
725
726 #[cfg(unix)]
727 fn begin_tick(&mut self) {
728 self.sidecar_attribution_attempted_this_tick = false;
729 }
730
731 #[cfg(unix)]
732 fn housekeeping_due(&self) -> bool {
733 !self.sidecar_attribution_attempted_this_tick
734 && self.walpin_full_scan_due_at(Instant::now())
735 }
736
737 #[cfg(unix)]
738 fn walpin_full_scan_due_at(&self, now: Instant) -> bool {
739 self.walpin_full_scan_last_attempt.is_none_or(|last| {
740 now.saturating_duration_since(last) >= self.walpin_full_scan_interval
741 })
742 }
743
744 #[cfg(unix)]
745 fn claim_walpin_full_scan_at(&mut self, now: Instant) -> bool {
746 if !self.walpin_full_scan_due_at(now) {
747 return false;
748 }
749 self.walpin_full_scan_last_attempt = Some(now);
750 true
751 }
752
753 #[cfg(unix)]
754 fn plan_walpin_attribution_at(&mut self, now: Instant) -> WalpinFullScanPlan {
755 if self.walpin_full_scan_due_at(now) {
756 let previous_last_attempt = self.walpin_full_scan_last_attempt.replace(now);
757 WalpinFullScanPlan::Refresh {
758 previous_last_attempt,
759 }
760 } else if let Some(cached) = self.walpin_cached_attribution.clone() {
761 WalpinFullScanPlan::Cached(cached)
762 } else {
763 WalpinFullScanPlan::Suppressed
764 }
765 }
766
767 #[cfg(unix)]
768 fn restore_walpin_full_scan_reservation(&mut self, previous_last_attempt: Option<Instant>) {
769 self.walpin_full_scan_last_attempt = previous_last_attempt;
770 }
771
772 #[cfg(unix)]
773 fn cache_walpin_attribution(
774 &mut self,
775 report: crate::walpin::WalpinReport,
776 census: Result<crate::walpin::CensusResult, String>,
777 captured_at: Instant,
778 ) {
779 self.walpin_cached_attribution = Some(CachedWalpinAttribution {
780 report,
781 census,
782 captured_at,
783 });
784 }
785}
786
787#[derive(Debug, Clone, Copy, PartialEq, Eq)]
794pub enum CheckpointSeverityRung {
795 Info,
797 Warn,
800 Alarm,
803}
804
805#[derive(Debug, Default, Clone)]
809pub struct CheckpointSeverityState {
810 was_above_warn: bool,
813 consecutive_above_warn: u8,
816 warn_emitted_for_episode: bool,
820}
821
822#[derive(Debug, Clone, Copy, PartialEq, Eq)]
825pub struct CheckpointSeverityEmission {
826 pub rung: CheckpointSeverityRung,
830 pub wal_pages: u64,
832 pub threshold_pages: u64,
834 pub consecutive_cycles: u8,
837}
838
839impl CheckpointSeverityState {
840 pub fn observe_wal_pages(
851 &mut self,
852 wal_pages: u64,
853 config: &CheckpointConfig,
854 ) -> Vec<CheckpointSeverityEmission> {
855 let mut emissions = Vec::new();
856 let above_warn = wal_pages >= config.warn_pages;
857
858 if above_warn {
859 self.consecutive_above_warn = self.consecutive_above_warn.saturating_add(1);
860
861 if !self.was_above_warn {
862 emissions.push(CheckpointSeverityEmission {
863 rung: CheckpointSeverityRung::Info,
864 wal_pages,
865 threshold_pages: config.warn_pages,
866 consecutive_cycles: self.consecutive_above_warn,
867 });
868 }
869
870 if !self.warn_emitted_for_episode
871 && self.consecutive_above_warn >= config.warn_sustained_cycles
872 {
873 emissions.push(CheckpointSeverityEmission {
874 rung: CheckpointSeverityRung::Warn,
875 wal_pages,
876 threshold_pages: config.warn_pages,
877 consecutive_cycles: self.consecutive_above_warn,
878 });
879 self.warn_emitted_for_episode = true;
880 }
881 } else {
882 self.consecutive_above_warn = 0;
883 self.warn_emitted_for_episode = false;
884 }
885
886 self.was_above_warn = above_warn;
887 emissions
888 }
889}
890
891#[derive(Debug, Clone, Copy, PartialEq, Eq)]
895pub enum TxAgeRung {
896 Warn,
898 Stale,
903}
904
905#[derive(Debug, Clone, PartialEq, Eq)]
907pub struct TxAgeEmission {
908 pub rung: TxAgeRung,
909 pub age: Duration,
910 pub label: Option<String>,
911}
912
913#[derive(Debug, Default, Clone)]
924pub struct TxAgeSweepState {
925 was_above_warn: bool,
928 was_above_max_age: bool,
931 tracked_id: Option<khive_storage::tx_registry::TxId>,
936}
937
938impl TxAgeSweepState {
939 pub fn observe(
951 &mut self,
952 oldest: Option<(khive_storage::tx_registry::TxId, Duration, Option<String>)>,
953 tx_warn_secs: Duration,
954 tx_max_age_secs: Duration,
955 ) -> Vec<TxAgeEmission> {
956 let mut emissions = Vec::new();
957
958 let Some((id, age, label)) = oldest else {
959 self.was_above_warn = false;
960 self.was_above_max_age = false;
961 self.tracked_id = None;
962 return emissions;
963 };
964
965 if self.tracked_id != Some(id) {
966 self.was_above_warn = false;
967 self.was_above_max_age = false;
968 }
969 self.tracked_id = Some(id);
970
971 let above_warn = age >= tx_warn_secs;
972 let above_max_age = age >= tx_max_age_secs;
973
974 if above_warn && !self.was_above_warn {
975 emissions.push(TxAgeEmission {
976 rung: TxAgeRung::Warn,
977 age,
978 label: label.clone(),
979 });
980 }
981 if above_max_age && !self.was_above_max_age {
982 emissions.push(TxAgeEmission {
983 rung: TxAgeRung::Stale,
984 age,
985 label,
986 });
987 }
988
989 self.was_above_warn = above_warn;
990 self.was_above_max_age = above_max_age;
991 emissions
992 }
993}
994
995fn log_tx_age_emission(emission: &TxAgeEmission) {
1000 let label = emission.label.as_deref().unwrap_or("<unlabeled>");
1001 match emission.rung {
1002 TxAgeRung::Warn => {
1003 tracing::warn!(
1004 tx_age_secs = emission.age.as_secs_f64(),
1005 tx_label = label,
1006 "ADR-091 Plank 1: open transaction registry entry exceeded soft-cap age"
1007 );
1008 }
1009 TxAgeRung::Stale => {
1010 tracing::error!(
1011 tx_age_secs = emission.age.as_secs_f64(),
1012 tx_label = label,
1013 "ADR-091 Plank 1: open transaction registry entry exceeded the cooperative \
1014 stale-op cap; no in-process mechanism can force-close it — investigate the \
1015 labeled caller directly"
1016 );
1017 }
1018 }
1019}
1020
1021struct WalpinSidecarState {
1029 dir: PathBuf,
1030 pid: u32,
1031 role: &'static str,
1032 started_at: i64,
1033 sweep_interval_ms: u64,
1038 wrote: bool,
1039 beacon_registered: bool,
1045 last_heartbeat: Option<LastHeartbeatState>,
1052}
1053
1054struct LastHeartbeatState {
1060 span_id: khive_storage::tx_registry::TxId,
1061 label: Option<String>,
1062 attribution_basis: &'static str,
1063 sweep_interval_ms: u64,
1064 oldest_tx_started_at: i64,
1065}
1066
1067impl LastHeartbeatState {
1068 fn content_matches(
1074 &self,
1075 span_id: khive_storage::tx_registry::TxId,
1076 label: &Option<String>,
1077 attribution_basis: &str,
1078 sweep_interval_ms: u64,
1079 ) -> bool {
1080 self.span_id == span_id
1081 && self.label == *label
1082 && self.attribution_basis == attribution_basis
1083 && self.sweep_interval_ms == sweep_interval_ms
1084 }
1085}
1086
1087impl WalpinSidecarState {
1088 fn new(
1091 db_path: Option<&Path>,
1092 is_file_backed: bool,
1093 role: &'static str,
1094 interval: Duration,
1095 ) -> Option<Self> {
1096 let path = db_path?;
1097 if !crate::walpin::sidecar_enabled(is_file_backed) {
1098 return None;
1099 }
1100 let pid = std::process::id();
1101 Some(Self {
1102 dir: crate::walpin::sidecar_dir_for(path),
1103 pid,
1104 role,
1105 started_at: crate::walpin::process_start_time_secs(pid).unwrap_or(0),
1106 sweep_interval_ms: interval.as_millis().min(u64::MAX as u128) as u64,
1107 wrote: false,
1108 last_heartbeat: None,
1109 beacon_registered: false,
1110 })
1111 }
1112
1113 async fn register_beacon(&mut self) {
1122 let dir = self.dir.clone();
1123 let beacon = crate::walpin::WalpinBeacon {
1124 pid: self.pid,
1125 process_role: self.role.to_string(),
1126 started_at: self.started_at,
1127 sweep_interval_ms: self.sweep_interval_ms,
1128 };
1129 let result =
1130 tokio::task::spawn_blocking(move || crate::walpin::write_beacon(&dir, &beacon)).await;
1131 match result {
1132 Ok(Ok(())) => {
1133 self.beacon_registered = true;
1134 }
1135 Ok(Err(e)) => {
1136 tracing::warn!(
1137 error = %e,
1138 "ADR-091 Amendment 2: failed to write walpin registration beacon; \
1139 this process's sidecar health will read as unknown, not registered-silent"
1140 );
1141 }
1142 Err(join_err) => {
1143 tracing::warn!(
1144 error = %join_err,
1145 "ADR-091 Amendment 2: walpin beacon write task panicked"
1146 );
1147 }
1148 }
1149 }
1150
1151 #[cfg(unix)]
1157 async fn reap_dead_entries_bounded(
1158 &self,
1159 legacy_fallback_interval: Duration,
1160 ) -> Option<crate::walpin::WalpinReport> {
1161 let dir = self.dir.clone();
1162 let result = tokio::task::spawn_blocking(move || {
1163 crate::walpin::housekeep_live(&dir, legacy_fallback_interval)
1164 })
1165 .await;
1166 match result {
1167 Ok(Ok(report)) => Some(report),
1168 Ok(Err(e)) => {
1169 tracing::warn!(
1170 error = %e,
1171 "ADR-091 Amendment 6: bounded walpin sidecar cleanup failed"
1172 );
1173 None
1174 }
1175 Err(join_err) => {
1176 tracing::warn!(
1177 error = %join_err,
1178 "ADR-091 Amendment 6: walpin sidecar cleanup task panicked"
1179 );
1180 None
1181 }
1182 }
1183 }
1184
1185 async fn refresh_beacon(&mut self) {
1195 if !self.beacon_registered {
1196 self.register_beacon().await;
1197 return;
1198 }
1199 let dir = self.dir.clone();
1200 let pid = self.pid;
1201 let result =
1202 tokio::task::spawn_blocking(move || crate::walpin::touch_beacon(&dir, pid)).await;
1203 match result {
1204 Ok(Ok(())) => {}
1205 Ok(Err(e)) => {
1206 self.beacon_registered = false;
1207 tracing::warn!(
1208 error = %e,
1209 "ADR-091 Amendment 2: failed to refresh walpin registration beacon; \
1210 this process's sidecar health will read as unknown, not registered-silent"
1211 );
1212 }
1213 Err(join_err) => {
1214 self.beacon_registered = false;
1215 tracing::warn!(
1216 error = %join_err,
1217 "ADR-091 Amendment 2: walpin beacon refresh task panicked"
1218 );
1219 }
1220 }
1221 }
1222
1223 async fn drop_beacon_fail_closed(&mut self) {
1234 let dir = self.dir.clone();
1235 let pid = self.pid;
1236 self.beacon_registered = false;
1237 let result =
1238 tokio::task::spawn_blocking(move || crate::walpin::remove_beacon(&dir, pid)).await;
1239 match result {
1240 Ok(Ok(())) => {}
1241 Ok(Err(e)) => {
1242 tracing::warn!(
1243 error = %e,
1244 "ADR-091 Amendment 2: failed to remove walpin beacon after a failed \
1245 heartbeat write; beacon will age out of the freshness window instead"
1246 );
1247 }
1248 Err(join_err) => {
1249 tracing::warn!(
1250 error = %join_err,
1251 "ADR-091 Amendment 2: walpin beacon removal task panicked"
1252 );
1253 }
1254 }
1255 }
1256
1257 async fn observe(
1261 &mut self,
1262 oldest: Option<khive_storage::tx_registry::OldestSpan>,
1263 tx_warn_secs: Duration,
1264 ) {
1265 match oldest {
1266 Some(span) if span.age >= tx_warn_secs => {
1267 let attribution_basis = match span.origin {
1274 khive_storage::tx_registry::TxOrigin::Database(_) => "origin",
1275 khive_storage::tx_registry::TxOrigin::Unscoped
1276 | khive_storage::tx_registry::TxOrigin::Memory => "fallback",
1277 };
1278
1279 let content_unchanged = self.wrote
1285 && self.last_heartbeat.as_ref().is_some_and(|last| {
1286 last.content_matches(
1287 span.id,
1288 &span.label,
1289 attribution_basis,
1290 self.sweep_interval_ms,
1291 )
1292 });
1293
1294 if content_unchanged {
1295 let dir = self.dir.clone();
1296 let pid = self.pid;
1297 let touch_result = tokio::task::spawn_blocking(move || {
1298 crate::walpin::touch_heartbeat(&dir, pid)
1299 })
1300 .await;
1301 match touch_result {
1302 Ok(Ok(())) => {
1303 self.refresh_beacon().await;
1304 return;
1305 }
1306 Ok(Err(e)) => {
1307 tracing::warn!(
1308 error = %e,
1309 "ADR-091 Amendment 3 Plank F1: walpin heartbeat touch failed; \
1310 recreating with a full body write"
1311 );
1312 }
1313 Err(join_err) => {
1314 tracing::warn!(
1315 error = %join_err,
1316 "ADR-091 Amendment 3 Plank F1: walpin heartbeat touch task \
1317 panicked; recreating with a full body write"
1318 );
1319 }
1320 }
1321 }
1326
1327 let oldest_tx_started_at = self
1334 .last_heartbeat
1335 .as_ref()
1336 .filter(|last| last.span_id == span.id)
1337 .map(|last| last.oldest_tx_started_at)
1338 .unwrap_or_else(|| now_epoch_secs().saturating_sub(span.age.as_secs() as i64));
1339
1340 let heartbeat = crate::walpin::WalpinHeartbeat {
1341 pid: self.pid,
1342 process_role: self.role.to_string(),
1343 started_at: self.started_at,
1344 oldest_tx_age_secs: span.age.as_secs_f64(),
1345 oldest_tx_label: span.label.clone(),
1346 oldest_tx_started_at: Some(oldest_tx_started_at),
1347 updated_at: now_epoch_secs(),
1348 sweep_interval_ms: self.sweep_interval_ms,
1349 attribution_basis: Some(attribution_basis.to_string()),
1350 };
1351 let dir = self.dir.clone();
1352 let result = tokio::task::spawn_blocking(move || {
1353 crate::walpin::write_heartbeat(&dir, &heartbeat)
1354 })
1355 .await;
1356 match result {
1366 Ok(Ok(())) => {
1367 self.wrote = true;
1368 self.last_heartbeat = Some(LastHeartbeatState {
1369 span_id: span.id,
1370 label: span.label,
1371 attribution_basis,
1372 sweep_interval_ms: self.sweep_interval_ms,
1373 oldest_tx_started_at,
1374 });
1375 self.refresh_beacon().await;
1376 }
1377 Ok(Err(e)) => {
1378 tracing::warn!(
1379 error = %e,
1380 "ADR-091 Amendment 2 Plank B: failed to write walpin heartbeat; \
1381 removing beacon so this process cannot read as \
1382 registered-silent while over threshold"
1383 );
1384 self.last_heartbeat = None;
1388 self.drop_beacon_fail_closed().await;
1389 }
1390 Err(join_err) => {
1391 tracing::warn!(
1392 error = %join_err,
1393 "ADR-091 Amendment 2 Plank B: walpin heartbeat write task panicked"
1394 );
1395 self.last_heartbeat = None;
1396 self.drop_beacon_fail_closed().await;
1397 }
1398 }
1399 }
1400 _ => {
1401 self.refresh_beacon().await;
1402 if self.wrote {
1403 let dir = self.dir.clone();
1404 let pid = self.pid;
1405 let result = tokio::task::spawn_blocking(move || {
1406 crate::walpin::remove_heartbeat(&dir, pid)
1407 })
1408 .await;
1409 match result {
1410 Ok(Ok(())) => {}
1411 Ok(Err(e)) => tracing::warn!(
1412 error = %e,
1413 "ADR-091 Amendment 2 Plank B: failed to remove walpin heartbeat"
1414 ),
1415 Err(join_err) => tracing::warn!(
1416 error = %join_err,
1417 "ADR-091 Amendment 2 Plank B: walpin heartbeat removal task panicked"
1418 ),
1419 }
1420 self.wrote = false;
1421 self.last_heartbeat = None;
1422 }
1423 }
1424 }
1425 }
1426
1427 async fn shutdown(&mut self) {
1428 if self.wrote {
1429 let dir = self.dir.clone();
1430 let pid = self.pid;
1431 let _ = tokio::task::spawn_blocking(move || crate::walpin::remove_heartbeat(&dir, pid))
1432 .await;
1433 self.wrote = false;
1434 }
1435 }
1436}
1437
1438#[cfg(unix)]
1439async fn run_walpin_housekeeping_if_due(
1440 sidecar: &WalpinSidecarState,
1441 state: &mut TruncateState,
1442 legacy_fallback_interval: Duration,
1443) -> bool {
1444 if !state.housekeeping_due() || !state.claim_walpin_full_scan_at(Instant::now()) {
1445 return false;
1446 }
1447 if let Some(report) = sidecar
1448 .reap_dead_entries_bounded(legacy_fallback_interval)
1449 .await
1450 {
1451 state.cache_walpin_attribution(
1452 report,
1453 Err("OS holder census is unavailable for a housekeeping-only scan".to_string()),
1454 Instant::now(),
1455 );
1456 }
1457 true
1458}
1459
1460fn now_epoch_secs() -> i64 {
1461 std::time::SystemTime::now()
1462 .duration_since(std::time::UNIX_EPOCH)
1463 .map(|d| d.as_secs() as i64)
1464 .unwrap_or(0)
1465}
1466
1467const DEFAULT_SESSION_SWEEP_INTERVAL: Duration = Duration::from_secs(5);
1472
1473#[derive(Clone, Debug)]
1474pub struct SessionSweepConfig {
1475 pub interval: Duration,
1480 pub tx_warn_secs: Duration,
1482 pub tx_max_age_secs: Duration,
1484}
1485
1486impl Default for SessionSweepConfig {
1487 fn default() -> Self {
1488 Self {
1489 interval: DEFAULT_SESSION_SWEEP_INTERVAL,
1490 tx_warn_secs: Duration::from_secs(30),
1491 tx_max_age_secs: Duration::from_secs(120),
1492 }
1493 }
1494}
1495
1496impl SessionSweepConfig {
1497 pub fn from_env() -> Self {
1501 let mut cfg = Self {
1502 interval: session_sweep_interval_from_env(),
1503 ..Self::default()
1504 };
1505 (cfg.tx_warn_secs, cfg.tx_max_age_secs) =
1510 tx_age_thresholds_from_env(cfg.tx_warn_secs, cfg.tx_max_age_secs);
1511
1512 cfg
1513 }
1514}
1515
1516fn session_sweep_interval_from_env() -> Duration {
1517 std::env::var("KHIVE_SESSION_SWEEP_INTERVAL_MS")
1518 .ok()
1519 .and_then(|ms| ms.parse::<u64>().ok())
1520 .filter(|ms| *ms > 0)
1521 .map(Duration::from_millis)
1522 .unwrap_or(DEFAULT_SESSION_SWEEP_INTERVAL)
1523}
1524
1525pub struct SweepBackend {
1535 pub pool: Arc<ConnectionPool>,
1536 pub is_main: bool,
1537}
1538
1539struct BackendSweep {
1545 filter: khive_storage::tx_registry::TxOriginFilter,
1546 tx_age_state: TxAgeSweepState,
1547 sidecar: Option<WalpinSidecarState>,
1548}
1549
1550pub async fn run_session_sweep_task(
1564 backends: Vec<SweepBackend>,
1565 config: SessionSweepConfig,
1566 mut shutdown_rx: tokio::sync::watch::Receiver<()>,
1567) {
1568 let mut interval = tokio::time::interval(config.interval);
1569 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
1570
1571 let mut sweeps: Vec<BackendSweep> = Vec::with_capacity(backends.len());
1572 for backend in backends {
1573 let identity = match backend.pool.origin() {
1574 khive_storage::tx_registry::TxOrigin::Database(id) => id,
1575 khive_storage::tx_registry::TxOrigin::Memory
1578 | khive_storage::tx_registry::TxOrigin::Unscoped => continue,
1579 };
1580 let filter = if backend.is_main {
1581 khive_storage::tx_registry::TxOriginFilter::Main(identity)
1582 } else {
1583 khive_storage::tx_registry::TxOriginFilter::Secondary(identity)
1584 };
1585 let sidecar = WalpinSidecarState::new(
1586 backend.pool.canonical_path(),
1587 true,
1588 "session",
1589 config.interval,
1590 );
1591 sweeps.push(BackendSweep {
1592 filter,
1593 tx_age_state: TxAgeSweepState::default(),
1594 sidecar,
1595 });
1596 }
1597 for sweep in sweeps.iter_mut() {
1598 if let Some(sidecar) = sweep.sidecar.as_mut() {
1599 sidecar.register_beacon().await;
1600 }
1601 }
1602
1603 loop {
1604 tokio::select! {
1605 _ = interval.tick() => {}
1606 _ = shutdown_rx.changed() => break,
1607 }
1608
1609 for sweep in sweeps.iter_mut() {
1610 let oldest = khive_storage::tx_registry::oldest_for(&sweep.filter);
1611 for emission in sweep.tx_age_state.observe(
1612 oldest.as_ref().map(|s| (s.id, s.age, s.label.clone())),
1613 config.tx_warn_secs,
1614 config.tx_max_age_secs,
1615 ) {
1616 log_tx_age_emission(&emission);
1617 }
1618 if let Some(sidecar) = sweep.sidecar.as_mut() {
1619 sidecar.observe(oldest, config.tx_warn_secs).await;
1620 }
1621 }
1622 }
1623
1624 for sweep in sweeps.iter_mut() {
1625 if let Some(sidecar) = sweep.sidecar.as_mut() {
1626 sidecar.shutdown().await;
1627 }
1628 }
1629}
1630
1631#[derive(Clone)]
1636pub struct CheckpointLifecycleOwner {
1637 event_store: Arc<dyn khive_storage::EventStore>,
1638 namespace: String,
1639}
1640
1641impl CheckpointLifecycleOwner {
1642 pub fn new(
1644 event_store: Arc<dyn khive_storage::EventStore>,
1645 namespace: impl Into<String>,
1646 ) -> Self {
1647 Self {
1648 event_store,
1649 namespace: namespace.into(),
1650 }
1651 }
1652}
1653
1654const CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY: usize = 1;
1658
1659#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1660struct CheckpointPressureEpisode {
1661 elevated_ticks: u64,
1662 peak_wal_pages: u64,
1663}
1664
1665impl CheckpointPressureEpisode {
1666 fn start(wal_pages: u64) -> Self {
1667 Self {
1668 elevated_ticks: 1,
1669 peak_wal_pages: wal_pages,
1670 }
1671 }
1672
1673 fn observe(&mut self, wal_pages: u64) {
1674 self.elevated_ticks = self.elevated_ticks.saturating_add(1);
1675 self.peak_wal_pages = self.peak_wal_pages.max(wal_pages);
1676 }
1677}
1678
1679struct CheckpointLifecycleEmitter {
1688 namespace: Option<String>,
1689 sender: Option<tokio::sync::mpsc::Sender<khive_storage::Event>>,
1690 worker: Option<tokio::task::JoinHandle<()>>,
1691 busy_warning_emitted: bool,
1692}
1693
1694impl CheckpointLifecycleEmitter {
1695 fn new(owner: Option<CheckpointLifecycleOwner>) -> Self {
1696 let Some(owner) = owner else {
1697 return Self {
1698 namespace: None,
1699 sender: None,
1700 worker: None,
1701 busy_warning_emitted: false,
1702 };
1703 };
1704
1705 let namespace = owner.namespace.clone();
1706 let (sender, mut receiver) =
1707 tokio::sync::mpsc::channel::<khive_storage::Event>(CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY);
1708 let worker = tokio::spawn(async move {
1709 while let Some(event) = receiver.recv().await {
1710 let kind = event.kind;
1711 CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.fetch_add(1, Ordering::Relaxed);
1712 if let Err(err) = owner.event_store.append_event(event).await {
1713 CHECKPOINT_LIFECYCLE_APPEND_FAILURES.fetch_add(1, Ordering::Relaxed);
1714 tracing::warn!(
1715 error = %err,
1716 event_kind = %kind.name(),
1717 "checkpoint lifecycle event append failed"
1718 );
1719 }
1720 }
1721 });
1722
1723 Self {
1724 namespace: Some(namespace),
1725 sender: Some(sender),
1726 worker: Some(worker),
1727 busy_warning_emitted: false,
1728 }
1729 }
1730
1731 fn try_emit<P: serde::Serialize>(&mut self, kind: khive_types::EventKind, payload: P) -> bool {
1734 let (Some(namespace), Some(sender)) = (&self.namespace, &self.sender) else {
1735 return true;
1736 };
1737 let payload_value = match serde_json::to_value(&payload) {
1738 Ok(value) => value,
1739 Err(err) => {
1740 tracing::warn!(
1741 error = %err,
1742 event_kind = %kind.name(),
1743 "failed to serialize checkpoint lifecycle event payload"
1744 );
1745 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.fetch_add(1, Ordering::Relaxed);
1746 return false;
1747 }
1748 };
1749 let payload_schema_version = match kind {
1750 khive_types::EventKind::CheckpointOutcomeRecorded => 2,
1751 _ => 1,
1752 };
1753 let event = khive_storage::Event::new(
1754 namespace,
1755 "checkpoint.lifecycle",
1756 kind,
1757 khive_types::SubstrateKind::Event,
1758 "daemon:checkpoint_task",
1759 )
1760 .with_payload(payload_value)
1761 .with_payload_schema_version(payload_schema_version);
1762
1763 match sender.try_send(event) {
1764 Ok(()) => {
1765 self.busy_warning_emitted = false;
1766 true
1767 }
1768 Err(tokio::sync::mpsc::error::TrySendError::Full(event)) => {
1769 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.fetch_add(1, Ordering::Relaxed);
1770 if !self.busy_warning_emitted {
1771 tracing::warn!(
1772 event_kind = %event.kind.name(),
1773 queue_capacity = CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY,
1774 "checkpoint lifecycle event dropped because the append worker is busy"
1775 );
1776 self.busy_warning_emitted = true;
1777 }
1778 false
1779 }
1780 Err(tokio::sync::mpsc::error::TrySendError::Closed(event)) => {
1781 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.fetch_add(1, Ordering::Relaxed);
1782 tracing::warn!(
1783 event_kind = %event.kind.name(),
1784 "checkpoint lifecycle event dropped because the append worker stopped"
1785 );
1786 false
1787 }
1788 }
1789 }
1790
1791 async fn shutdown(mut self) {
1799 drop(self.sender.take());
1800 let Some(worker) = self.worker.take() else {
1801 return;
1802 };
1803 worker.abort();
1804 match worker.await {
1805 Ok(()) => {}
1806 Err(err) if err.is_cancelled() => {}
1807 Err(err) => tracing::warn!(
1808 error = %err,
1809 "checkpoint lifecycle event append worker terminated unexpectedly"
1810 ),
1811 }
1812 }
1813}
1814
1815impl Drop for CheckpointLifecycleEmitter {
1816 fn drop(&mut self) {
1817 if let Some(worker) = &self.worker {
1823 worker.abort();
1824 }
1825 }
1826}
1827
1828struct CheckpointConnection {
1858 conn: Option<rusqlite::Connection>,
1859 consecutive_open_failures: u32,
1867}
1868
1869impl CheckpointConnection {
1870 fn new() -> Self {
1871 Self {
1872 conn: None,
1873 consecutive_open_failures: 0,
1874 }
1875 }
1876
1877 fn ensure_open(&mut self, pool: &ConnectionPool) -> Option<&rusqlite::Connection> {
1893 if self.conn.is_none() {
1894 match pool.open_standalone_writer_untracked() {
1895 Ok(conn) => {
1896 if let Err(e) = conn.pragma_update(None, "wal_autocheckpoint", 0) {
1903 tracing::warn!(
1904 error = %e,
1905 "could not disable autocheckpoint on the dedicated checkpoint \
1906 connection"
1907 );
1908 }
1909 if self.consecutive_open_failures > 0 {
1910 tracing::info!(
1911 prior_consecutive_failures = self.consecutive_open_failures,
1912 "dedicated checkpoint connection opened successfully, ending a \
1913 failure streak"
1914 );
1915 }
1916 self.consecutive_open_failures = 0;
1917 self.conn = Some(conn);
1918 }
1919 Err(e) => {
1920 if self.consecutive_open_failures == 0 {
1921 tracing::warn!(
1922 error = %e,
1923 "failed to open the dedicated checkpoint connection; \
1924 this tick is skipped and the open retried next tick"
1925 );
1926 } else {
1927 tracing::debug!(
1928 error = %e,
1929 consecutive_failures = self.consecutive_open_failures,
1930 "dedicated checkpoint connection still unavailable; \
1931 this tick is skipped and the open retried next tick"
1932 );
1933 }
1934 self.consecutive_open_failures =
1935 self.consecutive_open_failures.saturating_add(1);
1936 return None;
1937 }
1938 }
1939 }
1940 self.conn.as_ref()
1941 }
1942}
1943
1944async fn run_fts_maintenance_off_worker(
1958 conn: rusqlite::Connection,
1959 config: crate::fts_maintenance::FtsMaintenanceConfig,
1960 mut state: crate::fts_maintenance::FtsMaintenanceState,
1961 now: Instant,
1962) -> Result<
1963 (
1964 rusqlite::Connection,
1965 crate::fts_maintenance::FtsMaintenanceState,
1966 Result<Option<crate::fts_maintenance::FtsMaintenanceStep>, String>,
1967 ),
1968 tokio::task::JoinError,
1969> {
1970 tokio::task::spawn_blocking(move || {
1971 let result = crate::fts_maintenance::run_if_due(&conn, &config, &mut state, now);
1972 (conn, state, result)
1973 })
1974 .await
1975}
1976
1977pub async fn run_checkpoint_task(
2011 pool: Arc<ConnectionPool>,
2012 config: CheckpointConfig,
2013 lifecycle_owner: Option<CheckpointLifecycleOwner>,
2014 mut shutdown_rx: tokio::sync::watch::Receiver<()>,
2015 is_main: bool,
2016) {
2017 match pool.claim_checkpoint_ownership() {
2024 Ok(()) => {
2025 if let Err(e) = pool.propagate_checkpoint_claim_to_writer_task().await {
2026 tracing::warn!(
2027 error = %e,
2028 "checkpoint task could not reach the writer task's connection; it keeps the \
2029 bounded autocheckpoint fallback"
2030 );
2031 }
2032 }
2033 Err(e) => {
2034 tracing::warn!(
2035 error = %e,
2036 "checkpoint task could not re-apply the ownership pragma on the pooled writer; \
2037 writer connections keep the bounded autocheckpoint fallback unless ownership is \
2038 claimed later"
2039 );
2040 }
2041 }
2042 let mut interval = tokio::time::interval(config.interval);
2043 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
2044 let mut severity_state = CheckpointSeverityState::default();
2045 let mut tx_age_state = TxAgeSweepState::default();
2046 let mut was_above_high_water = false;
2047 #[cfg(unix)]
2048 let legacy_walpin_fallback_interval = DEFAULT_SESSION_SWEEP_INTERVAL;
2049 #[cfg(unix)]
2050 let mut truncate_state =
2051 TruncateState::with_legacy_walpin_fallback(legacy_walpin_fallback_interval);
2052 #[cfg(not(unix))]
2053 let mut truncate_state = TruncateState::default();
2054 let mut lifecycle_emitter = CheckpointLifecycleEmitter::new(lifecycle_owner);
2055 let mut event_elevation_open = false;
2060 let mut pressure_episode: Option<CheckpointPressureEpisode> = None;
2061 let mut pending_recovery: Option<khive_storage::CheckpointOutcomeRecordedPayload> = None;
2066 let mut was_observed_above_warn = false;
2067 let tx_filter = match pool.origin() {
2080 khive_storage::tx_registry::TxOrigin::Database(id) => Some(if is_main {
2081 khive_storage::tx_registry::TxOriginFilter::Main(id)
2082 } else {
2083 khive_storage::tx_registry::TxOriginFilter::Secondary(id)
2084 }),
2085 khive_storage::tx_registry::TxOrigin::Memory
2086 | khive_storage::tx_registry::TxOrigin::Unscoped => None,
2087 };
2088 #[cfg(unix)]
2094 let mut walpin_state =
2095 WalpinSidecarState::new(pool.canonical_path(), true, "daemon", config.interval);
2096 #[cfg(unix)]
2097 if let Some(sidecar) = walpin_state.as_mut() {
2098 sidecar.register_beacon().await;
2099 }
2100
2101 let mut checkpoint_conn = CheckpointConnection::new();
2105 checkpoint_conn.ensure_open(&pool);
2106 let mut fts_maintenance_config = crate::fts_maintenance::FtsMaintenanceConfig::from_env();
2111 fts_maintenance_config.enabled &= is_main;
2115 let mut fts_maintenance_state =
2116 crate::fts_maintenance::FtsMaintenanceState::new(Instant::now());
2117
2118 loop {
2119 tokio::select! {
2124 _ = interval.tick() => {}
2125 _ = shutdown_rx.changed() => break,
2126 }
2127
2128 #[cfg(unix)]
2129 truncate_state.begin_tick();
2130
2131 #[cfg(unix)]
2132 let mut pending_sidecar_attribution = None;
2133
2134 let tick = if checkpoint_conn.ensure_open(&pool).is_none() {
2135 note_checkpoint_skipped();
2136 CheckpointTick::Skipped
2137 } else {
2138 let conn = checkpoint_conn
2144 .conn
2145 .take()
2146 .expect("ensure_open just confirmed a connection is open");
2147 match checkpoint_once_core(&pool, &conn, &config, &mut truncate_state) {
2148 Ok(outcome) => {
2149 #[cfg(unix)]
2150 {
2151 pending_sidecar_attribution = outcome.sidecar_attribution;
2152 }
2153 #[cfg(not(unix))]
2154 let _ = outcome.sidecar_attribution;
2155
2156 let fts_result = if !fts_maintenance_state
2160 .is_due(&fts_maintenance_config, Instant::now())
2161 {
2162 checkpoint_conn.conn = Some(conn);
2163 Ok(None)
2164 } else {
2165 match run_fts_maintenance_off_worker(
2166 conn,
2167 fts_maintenance_config.clone(),
2168 fts_maintenance_state,
2169 Instant::now(),
2170 )
2171 .await
2172 {
2173 Ok((conn, state, fts_result)) => {
2174 fts_maintenance_state = state;
2175 checkpoint_conn.conn = Some(conn);
2176 fts_result
2177 }
2178 Err(join_err) => {
2179 fts_maintenance_state =
2185 crate::fts_maintenance::FtsMaintenanceState::new(Instant::now());
2186 Err(format!(
2187 "bounded FTS5 segment maintenance task panicked: {join_err}"
2188 ))
2189 }
2190 }
2191 };
2192
2193 match fts_result {
2194 Ok(Some(step)) => match step.outcome {
2195 crate::fts_maintenance::FtsMaintenanceOutcome::Worked => {
2196 tracing::info!(
2197 table = step.table,
2198 requested_pages = step.requested_pages,
2199 segments_before = step.segments_before,
2200 segments_after = step.segments_after,
2201 "bounded FTS5 segment maintenance made progress"
2202 );
2203 }
2204 crate::fts_maintenance::FtsMaintenanceOutcome::Busy => {
2205 tracing::debug!(
2206 table = step.table,
2207 requested_pages = step.requested_pages,
2208 segments = step.segments_before,
2209 "bounded FTS5 segment maintenance skipped a busy writer"
2210 );
2211 }
2212 crate::fts_maintenance::FtsMaintenanceOutcome::Noop
2213 | crate::fts_maintenance::FtsMaintenanceOutcome::BelowThreshold => {
2214 tracing::debug!(
2215 table = step.table,
2216 outcome = ?step.outcome,
2217 segments = step.segments_before,
2218 "bounded FTS5 segment maintenance had no work"
2219 );
2220 }
2221 },
2222 Ok(None) => {}
2223 Err(error) => {
2224 tracing::warn!(
2229 error = %error,
2230 "bounded FTS5 segment maintenance failed"
2231 );
2232 }
2233 }
2234 CheckpointTick::Observed(outcome.wal_pages)
2235 }
2236 Err(e) => {
2237 tracing::warn!(
2238 error = %e,
2239 "dedicated checkpoint connection failed a pragma; \
2240 dropping it for a fresh reopen next tick"
2241 );
2242 note_checkpoint_skipped();
2243 CheckpointTick::Skipped
2244 }
2245 }
2246 };
2247
2248 #[cfg(unix)]
2257 if let Err(error) =
2258 complete_walpin_attribution(pending_sidecar_attribution, &mut truncate_state).await
2259 {
2260 tracing::warn!(
2261 error = %error,
2262 failure_kind = error.kind(),
2263 "ADR-091 Amendment 2 Plank B: no-progress sidecar attribution failed"
2264 );
2265 }
2266
2267 let oldest_tx = tx_filter
2282 .as_ref()
2283 .and_then(khive_storage::tx_registry::oldest_for);
2284 for emission in tx_age_state.observe(
2285 oldest_tx.as_ref().map(|s| (s.id, s.age, s.label.clone())),
2286 config.tx_warn_secs,
2287 config.tx_max_age_secs,
2288 ) {
2289 log_tx_age_emission(&emission);
2290 }
2291 #[cfg(unix)]
2295 if let Some(sidecar) = walpin_state.as_mut() {
2296 sidecar
2297 .observe(oldest_tx.clone(), config.tx_warn_secs)
2298 .await;
2299 let _ = run_walpin_housekeeping_if_due(
2300 sidecar,
2301 &mut truncate_state,
2302 legacy_walpin_fallback_interval,
2303 )
2304 .await;
2305 }
2306
2307 let wal_pages = match tick {
2310 CheckpointTick::Skipped => continue,
2311 CheckpointTick::Observed(n) => n,
2312 };
2313
2314 let above_warn = wal_pages >= config.warn_pages;
2315 let above_high_water = wal_pages >= config.high_water_pages;
2316 let above_truncate_high_water = wal_pages >= config.truncate_high_water_pages;
2317 note_checkpoint_pressure_observation(above_warn, was_observed_above_warn);
2318 was_observed_above_warn = above_warn;
2319
2320 log_tx_registry_oldest_debug(wal_pages, oldest_tx.as_ref());
2326
2327 for emission in severity_state.observe_wal_pages(wal_pages, &config) {
2332 match emission.rung {
2333 CheckpointSeverityRung::Info => {
2334 log_tx_registry_oldest_warn(wal_pages, oldest_tx.as_ref());
2335 tracing::info!(
2336 wal_pages = emission.wal_pages,
2337 warn_threshold = emission.threshold_pages,
2338 "WAL page count crossed warn threshold"
2339 );
2340 }
2341 CheckpointSeverityRung::Warn => {
2342 tracing::warn!(
2343 wal_pages = emission.wal_pages,
2344 warn_threshold = emission.threshold_pages,
2345 consecutive_cycles = emission.consecutive_cycles,
2346 "WAL page count failed to drain below warn threshold"
2347 );
2348 }
2349 CheckpointSeverityRung::Alarm => {
2350 }
2352 }
2353 }
2354
2355 let high_water_crossed = crossing_warn(above_high_water, &mut was_above_high_water);
2356 if high_water_crossed {
2357 log_tx_registry_snapshot_warn(wal_pages);
2358 log_wal_high_water_warn(
2359 wal_pages,
2360 config.high_water_pages,
2361 oldest_tx.as_ref(),
2362 config.tx_warn_secs,
2363 );
2364 }
2365
2366 observe_checkpoint_pressure_tick(
2372 above_warn,
2373 wal_pages,
2374 above_high_water,
2375 above_truncate_high_water,
2376 &config,
2377 &mut event_elevation_open,
2378 &mut pressure_episode,
2379 &mut pending_recovery,
2380 |payload| {
2381 lifecycle_emitter
2382 .try_emit(khive_types::EventKind::CheckpointOutcomeRecorded, payload)
2383 },
2384 );
2385 }
2386
2387 lifecycle_emitter.shutdown().await;
2388
2389 #[cfg(unix)]
2390 if let Some(sidecar) = walpin_state.as_mut() {
2391 sidecar.shutdown().await;
2392 }
2393}
2394
2395fn checkpoint_outcome_should_emit(above_warn: bool, was_elevated: bool) -> bool {
2399 above_warn != was_elevated
2400}
2401
2402#[allow(clippy::too_many_arguments)]
2415fn observe_checkpoint_pressure_tick(
2416 above_warn: bool,
2417 wal_pages: u64,
2418 above_high_water: bool,
2419 above_truncate_high_water: bool,
2420 config: &CheckpointConfig,
2421 event_elevation_open: &mut bool,
2422 pressure_episode: &mut Option<CheckpointPressureEpisode>,
2423 pending_recovery: &mut Option<khive_storage::CheckpointOutcomeRecordedPayload>,
2424 mut try_emit: impl FnMut(khive_storage::CheckpointOutcomeRecordedPayload) -> bool,
2425) {
2426 let pending_blocks_emission = if let Some(payload) = pending_recovery.clone() {
2434 if try_emit(payload) {
2435 *pending_recovery = None;
2436 false
2437 } else {
2438 true
2439 }
2440 } else {
2441 false
2442 };
2443
2444 if above_warn {
2445 match pressure_episode.as_mut() {
2446 Some(episode) => episode.observe(wal_pages),
2447 None => *pressure_episode = Some(CheckpointPressureEpisode::start(wal_pages)),
2448 }
2449 } else if !*event_elevation_open {
2450 if pending_blocks_emission && pressure_episode.is_some() {
2457 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.fetch_add(1, Ordering::Relaxed);
2458 tracing::warn!(
2459 wal_pages,
2460 "checkpoint pressure episode elapsed unreported behind an undelivered recovery summary"
2461 );
2462 }
2463 *pressure_episode = None;
2464 }
2465
2466 if pending_blocks_emission || !checkpoint_outcome_should_emit(above_warn, *event_elevation_open)
2467 {
2468 return;
2469 }
2470 let Some(episode) = *pressure_episode else {
2471 tracing::warn!(
2472 above_warn,
2473 event_elevation_open = *event_elevation_open,
2474 "checkpoint pressure transition has no episode aggregate"
2475 );
2476 return;
2477 };
2478 let payload = khive_storage::CheckpointOutcomeRecordedPayload {
2479 wal_pages,
2480 warn_pages: config.warn_pages,
2481 high_water_pages: config.high_water_pages,
2482 truncate_high_water_pages: config.truncate_high_water_pages,
2483 above_warn,
2484 above_high_water,
2485 above_truncate_high_water,
2486 episode_elevated_ticks: Some(episode.elevated_ticks),
2487 episode_peak_wal_pages: Some(episode.peak_wal_pages),
2488 };
2489 if try_emit(payload.clone()) {
2490 *event_elevation_open = above_warn;
2491 if !above_warn {
2492 *pressure_episode = None;
2493 }
2494 } else if !above_warn {
2495 debug_assert!(
2506 pending_recovery.is_none(),
2507 "recovery emission attempted while an earlier summary was still pending"
2508 );
2509 *event_elevation_open = false;
2510 *pressure_episode = None;
2511 *pending_recovery = Some(payload);
2512 }
2513}
2514
2515fn log_tx_registry_oldest_debug(
2524 wal_pages: u64,
2525 oldest: Option<&khive_storage::tx_registry::OldestSpan>,
2526) {
2527 if let Some(span) = oldest {
2528 tracing::debug!(
2529 wal_pages,
2530 oldest_tx_age_secs = span.age.as_secs_f64(),
2531 oldest_tx_label = span.label.as_deref().unwrap_or("<unlabeled>"),
2532 "WAL checkpoint tick: oldest open transaction registry entry"
2533 );
2534 }
2535}
2536
2537fn log_tx_registry_oldest_warn(
2541 wal_pages: u64,
2542 oldest: Option<&khive_storage::tx_registry::OldestSpan>,
2543) {
2544 if let Some(span) = oldest {
2545 tracing::warn!(
2546 wal_pages,
2547 oldest_tx_age_secs = span.age.as_secs_f64(),
2548 oldest_tx_label = span.label.as_deref().unwrap_or("<unlabeled>"),
2549 "WAL checkpoint tick: oldest open transaction registry entry"
2550 );
2551 }
2552}
2553
2554fn log_tx_registry_snapshot_warn(wal_pages: u64) {
2558 log_tx_registry_entries_warn(wal_pages, &khive_storage::tx_registry::snapshot());
2559}
2560
2561fn log_tx_registry_entries_warn(wal_pages: u64, snapshot: &[(Duration, Option<String>)]) {
2562 for (age, label) in snapshot {
2563 tracing::warn!(
2564 wal_pages,
2565 tx_age_secs = age.as_secs_f64(),
2566 tx_label = label.as_deref().unwrap_or("<unlabeled>"),
2567 "WAL high-water: open transaction registry entry"
2568 );
2569 }
2570}
2571
2572fn log_truncate_no_progress_warn(
2573 wal_pages_before: u64,
2574 wal_pages_after: u64,
2575 snapshot: &[(Duration, Option<String>)],
2576) {
2577 let open_tx_count = snapshot.len();
2578 let oldest_tx_age_secs = snapshot
2579 .iter()
2580 .map(|(age, _)| *age)
2581 .max()
2582 .map(|age| age.as_secs_f64());
2583 if snapshot.is_empty() {
2584 tracing::warn!(
2585 wal_pages_before,
2586 wal_pages_after,
2587 open_tx_count,
2588 oldest_tx_age_secs = ?oldest_tx_age_secs,
2589 "WAL TRUNCATE attempt made no progress; no open transaction in this process's registry"
2590 );
2591 } else {
2592 tracing::warn!(
2593 wal_pages_before,
2594 wal_pages_after,
2595 open_tx_count,
2596 oldest_tx_age_secs = ?oldest_tx_age_secs,
2597 "WAL TRUNCATE attempt made no progress; open transactions observed in this process's registry"
2598 );
2599 }
2600 log_tx_registry_entries_warn(wal_pages_after, snapshot);
2601}
2602
2603fn log_wal_high_water_warn(
2628 wal_pages: u64,
2629 high_water: u64,
2630 oldest: Option<&khive_storage::tx_registry::OldestSpan>,
2631 warn_after: Duration,
2632) {
2633 match oldest.filter(|span| span.age >= warn_after) {
2634 Some(span) => tracing::warn!(
2635 wal_pages,
2636 high_water,
2637 oldest_tx_age_secs = span.age.as_secs_f64(),
2638 oldest_tx_label = span.label.as_deref().unwrap_or("<unlabeled>"),
2639 "WAL high-water mark exceeded; an open transaction older than the age \
2640 threshold is pinning a snapshot PASSIVE cannot reclaim"
2641 ),
2642 None => tracing::warn!(
2643 wal_pages,
2644 high_water,
2645 oldest_tx_age_secs = ?oldest.map(|span| span.age.as_secs_f64()),
2646 oldest_tx_label = oldest
2647 .and_then(|span| span.label.as_deref())
2648 .unwrap_or("<none>"),
2649 "WAL high-water mark exceeded with no open transaction old enough to pin \
2650 a snapshot; the WAL is growing faster than PASSIVE checkpoints reclaim it"
2651 ),
2652 }
2653}
2654
2655#[derive(Debug)]
2660#[must_use]
2661struct CheckpointCoreOutcome {
2662 wal_pages: u64,
2663 sidecar_attribution: Option<WalpinAttributionRequest>,
2664}
2665
2666pub fn checkpoint_once(
2690 pool: &ConnectionPool,
2691 conn: &rusqlite::Connection,
2692 config: &CheckpointConfig,
2693 truncate_state: &mut TruncateState,
2694) -> Result<u64, rusqlite::Error> {
2695 checkpoint_once_core(pool, conn, config, truncate_state).map(|outcome| outcome.wal_pages)
2696}
2697
2698fn checkpoint_once_core(
2702 pool: &ConnectionPool,
2703 conn: &rusqlite::Connection,
2704 config: &CheckpointConfig,
2705 truncate_state: &mut TruncateState,
2706) -> Result<CheckpointCoreOutcome, rusqlite::Error> {
2707 #[cfg(unix)]
2708 truncate_state.begin_tick();
2709 let started = Instant::now();
2710 let checkpoint_result = query_checkpoint_observation(conn);
2711 let elapsed_us = started.elapsed().as_micros().min(u128::from(u64::MAX)) as u64;
2712 record_checkpoint_timing(
2713 pool,
2714 elapsed_us,
2715 checkpoint_result
2716 .as_ref()
2717 .ok()
2718 .map(|observation| observation.busy),
2719 );
2720 let raw_observation = match checkpoint_result {
2721 Ok(observation) => observation,
2722 Err(e) => {
2723 tracing::warn!(error = %e, elapsed_us, "WAL checkpoint failed");
2724 return Err(e);
2725 }
2726 };
2727 let observation = record_routine_wal_observation(pool, raw_observation);
2728 let wal_pages = observation.log_frames;
2729 LAST_WAL_PAGES.store(wal_pages, Ordering::Relaxed);
2730 note_checkpoint_observed(wal_pages);
2731
2732 if raw_observation.busy != 0 {
2733 tracing::debug!(
2734 busy = raw_observation.busy,
2735 wal_log_frames = raw_observation.log_frames,
2736 wal_checkpointed_frames = raw_observation.checkpointed_frames,
2737 wal_pending_frames = observation.pending_frames,
2738 wal_physical_bytes = ?observation.physical_wal_bytes,
2739 "WAL PASSIVE checkpoint reported incomplete progress"
2740 );
2741 }
2742 tracing::debug!(
2743 wal_pages,
2744 elapsed_us,
2745 busy = raw_observation.busy,
2746 wal_checkpointed_frames = observation.checkpointed_frames,
2747 wal_pending_frames = observation.pending_frames,
2748 wal_physical_bytes = ?observation.physical_wal_bytes,
2749 "WAL checkpoint issued"
2750 );
2751
2752 let sidecar_attribution = maybe_truncate(pool, conn, config, wal_pages, truncate_state);
2753
2754 Ok(CheckpointCoreOutcome {
2755 wal_pages,
2756 sidecar_attribution,
2757 })
2758}
2759
2760fn maybe_truncate(
2766 pool: &ConnectionPool,
2767 conn: &rusqlite::Connection,
2768 config: &CheckpointConfig,
2769 wal_pages_before: u64,
2770 truncate_state: &mut TruncateState,
2771) -> Option<WalpinAttributionRequest> {
2772 if wal_pages_before < config.truncate_high_water_pages {
2773 return None;
2774 }
2775
2776 if let Some(last) = truncate_state.last_attempt {
2777 if last.elapsed() < config.truncate_min_interval {
2778 return None;
2779 }
2780 }
2781
2782 log_tx_registry_snapshot_warn(wal_pages_before);
2785
2786 let original_busy_timeout = pool.config().busy_timeout;
2787
2788 if let Err(e) = conn.busy_timeout(config.truncate_busy_timeout) {
2789 tracing::warn!(error = %e, "failed to lower busy_timeout for TRUNCATE attempt; skipping");
2795 return None;
2796 }
2797
2798 #[cfg(unix)]
2799 let mut holder_attribution = capture_walpin_attribution_request(pool, truncate_state);
2800 #[cfg(unix)]
2801 let mut sidecar_attribution = None;
2802 #[cfg(not(unix))]
2803 let sidecar_attribution = None;
2804
2805 truncate_state.last_attempt = Some(Instant::now());
2809
2810 let start = Instant::now();
2811 let outcome = conn.execute_batch("PRAGMA wal_checkpoint(TRUNCATE)");
2812 let elapsed = start.elapsed();
2813
2814 if let Err(e) = conn.busy_timeout(original_busy_timeout) {
2817 tracing::warn!(error = %e, "failed to restore busy_timeout after TRUNCATE attempt");
2818 }
2819
2820 match outcome {
2821 Ok(()) => {
2822 let wal_pages_after = query_wal_pages(conn);
2823 tracing::info!(
2824 wal_pages_before,
2825 wal_pages_after,
2826 elapsed_ms = elapsed.as_millis() as u64,
2827 "WAL TRUNCATE checkpoint attempted"
2828 );
2829
2830 let made_progress = wal_pages_after < wal_pages_before;
2831 if !made_progress {
2832 let snapshot = khive_storage::tx_registry::snapshot();
2833 log_truncate_no_progress_warn(wal_pages_before, wal_pages_after, &snapshot);
2834 #[cfg(test)]
2835 if let Some(path) = pool.canonical_path() {
2836 truncate_report_test_sync::after_no_progress_before_report(path);
2837 }
2838 #[cfg(unix)]
2839 {
2840 sidecar_attribution = holder_attribution.take();
2847 }
2848 log_wal_pin_depth(conn);
2849 }
2850
2851 note_truncate_outcome(config, wal_pages_after, truncate_state);
2852 }
2853 Err(e) => {
2854 tracing::warn!(error = %e, wal_pages_before, "WAL TRUNCATE attempt failed");
2855 log_tx_registry_snapshot_warn(wal_pages_before);
2856 note_truncate_outcome(config, wal_pages_before, truncate_state);
2857 }
2858 }
2859 #[cfg(unix)]
2860 if let Some(WalpinAttributionRequest::Fresh {
2861 previous_last_attempt,
2862 ..
2863 }) = holder_attribution.as_ref()
2864 {
2865 truncate_state.restore_walpin_full_scan_reservation(*previous_last_attempt);
2866 }
2867 sidecar_attribution
2868}
2869
2870#[cfg(test)]
2871mod truncate_report_test_sync {
2872 use std::path::{Path, PathBuf};
2873 use std::sync::mpsc::{sync_channel, Receiver, SyncSender};
2874 use std::sync::Mutex;
2875
2876 struct Hook {
2877 db_path: PathBuf,
2878 reached_tx: SyncSender<()>,
2879 proceed_rx: Receiver<()>,
2880 }
2881
2882 static HOOK: Mutex<Option<Hook>> = Mutex::new(None);
2883
2884 pub(crate) fn install(db_path: PathBuf) -> (Receiver<()>, SyncSender<()>) {
2885 let (reached_tx, reached_rx) = sync_channel(0);
2886 let (proceed_tx, proceed_rx) = sync_channel(0);
2887 let replaced = HOOK
2888 .lock()
2889 .unwrap_or_else(|poisoned| poisoned.into_inner())
2890 .replace(Hook {
2891 db_path,
2892 reached_tx,
2893 proceed_rx,
2894 });
2895 assert!(replaced.is_none(), "truncate report hook already installed");
2896 (reached_rx, proceed_tx)
2897 }
2898
2899 pub(crate) fn uninstall() {
2900 *HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
2901 }
2902
2903 pub(crate) fn after_no_progress_before_report(db_path: &Path) {
2904 let hook = {
2905 let mut guard = HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
2906 match guard.as_ref() {
2907 Some(hook) if hook.db_path == db_path => guard.take(),
2908 _ => None,
2909 }
2910 };
2911 let Some(hook) = hook else {
2912 return;
2913 };
2914 let _ = hook.reached_tx.send(());
2915 let _ = hook.proceed_rx.recv();
2916 }
2917}
2918
2919#[cfg(all(test, unix))]
2924mod walpin_attribution_test_sync {
2925 use std::path::{Path, PathBuf};
2926 use std::sync::atomic::{AtomicUsize, Ordering};
2927 use std::sync::mpsc::{sync_channel, Receiver, SyncSender};
2928 use std::sync::{Arc, Mutex};
2929
2930 enum Behavior {
2931 Pause {
2932 reached_tx: tokio::sync::oneshot::Sender<std::thread::ThreadId>,
2933 proceed_rx: Receiver<()>,
2934 },
2935 Panic,
2936 }
2937
2938 struct Hook {
2939 dir: PathBuf,
2940 behavior: Behavior,
2941 }
2942
2943 static HOOK: Mutex<Option<Hook>> = Mutex::new(None);
2944 static REPORT_COUNTER: Mutex<Option<Arc<AtomicUsize>>> = Mutex::new(None);
2945
2946 pub(crate) fn install_pause(
2947 dir: PathBuf,
2948 ) -> (
2949 tokio::sync::oneshot::Receiver<std::thread::ThreadId>,
2950 SyncSender<()>,
2951 Arc<AtomicUsize>,
2952 ) {
2953 let (reached_tx, reached_rx) = tokio::sync::oneshot::channel();
2954 let (proceed_tx, proceed_rx) = sync_channel(0);
2955 let report_counter = Arc::new(AtomicUsize::new(0));
2956 let replaced = HOOK
2957 .lock()
2958 .unwrap_or_else(|poisoned| poisoned.into_inner())
2959 .replace(Hook {
2960 dir,
2961 behavior: Behavior::Pause {
2962 reached_tx,
2963 proceed_rx,
2964 },
2965 });
2966 assert!(
2967 replaced.is_none(),
2968 "walpin attribution hook already installed"
2969 );
2970 *REPORT_COUNTER
2971 .lock()
2972 .unwrap_or_else(|poisoned| poisoned.into_inner()) = Some(Arc::clone(&report_counter));
2973 (reached_rx, proceed_tx, report_counter)
2974 }
2975
2976 pub(crate) fn install_panic(dir: PathBuf) {
2977 let replaced = HOOK
2978 .lock()
2979 .unwrap_or_else(|poisoned| poisoned.into_inner())
2980 .replace(Hook {
2981 dir,
2982 behavior: Behavior::Panic,
2983 });
2984 assert!(
2985 replaced.is_none(),
2986 "walpin attribution hook already installed"
2987 );
2988 }
2989
2990 pub(crate) fn uninstall() {
2991 *HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
2992 *REPORT_COUNTER
2993 .lock()
2994 .unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
2995 }
2996
2997 pub(crate) fn before_enumeration(dir: &Path) {
2998 let hook = {
2999 let mut guard = HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
3000 match guard.as_ref() {
3001 Some(hook) if hook.dir == dir => guard.take(),
3002 _ => None,
3003 }
3004 };
3005 let Some(hook) = hook else {
3006 return;
3007 };
3008 match hook.behavior {
3009 Behavior::Pause {
3010 reached_tx,
3011 proceed_rx,
3012 } => {
3013 if reached_tx.send(std::thread::current().id()).is_ok() {
3014 let _ = proceed_rx.recv();
3015 }
3016 }
3017 Behavior::Panic => panic!("injected walpin attribution worker panic"),
3018 }
3019 }
3020
3021 pub(crate) fn report_used() {
3022 if let Some(counter) = REPORT_COUNTER
3023 .lock()
3024 .unwrap_or_else(|poisoned| poisoned.into_inner())
3025 .as_ref()
3026 {
3027 counter.fetch_add(1, Ordering::SeqCst);
3028 }
3029 }
3030}
3031
3032fn note_truncate_outcome(
3038 config: &CheckpointConfig,
3039 wal_pages_after: u64,
3040 state: &mut TruncateState,
3041) {
3042 TRUNCATE_ATTEMPTS.fetch_add(1, Ordering::Relaxed);
3047
3048 if wal_pages_after >= config.warn_pages {
3049 state.consecutive_failures = state.consecutive_failures.saturating_add(1);
3050 if state.consecutive_failures == 3 {
3051 tracing::warn!(
3052 wal_pages_after,
3053 warn_threshold = config.warn_pages,
3054 "WAL TRUNCATE has failed to clear WAL pressure for 3 consecutive attempts"
3055 );
3056 }
3057 } else {
3058 state.consecutive_failures = 0;
3059 }
3060
3061 TRUNCATE_CONSECUTIVE_FAILURES.store(state.consecutive_failures as u64, Ordering::Relaxed);
3062}
3063
3064#[cfg(unix)]
3072#[derive(Debug)]
3073enum WalpinAttributionRequest {
3074 Fresh {
3075 dir: PathBuf,
3076 census: Result<crate::walpin::CensusResult, String>,
3077 legacy_fallback_interval: Duration,
3078 previous_last_attempt: Option<Instant>,
3079 },
3080 Cached(CachedWalpinAttribution),
3081 Suppressed,
3082}
3083
3084#[cfg(unix)]
3085#[derive(Debug, Clone, Copy, PartialEq, Eq)]
3086enum WalpinReportFreshness {
3087 Fresh,
3088 Cached { age: Duration },
3089}
3090
3091#[cfg(unix)]
3092impl WalpinReportFreshness {
3093 fn is_fresh(self) -> bool {
3094 self == Self::Fresh
3095 }
3096}
3097
3098#[cfg(not(unix))]
3101type WalpinAttributionRequest = ();
3102
3103#[cfg(unix)]
3108#[derive(Debug, Clone, PartialEq, Eq)]
3109enum WalpinAttributionFailure {
3110 Enumeration(String),
3111 Worker(String),
3112}
3113
3114#[cfg(unix)]
3115impl WalpinAttributionFailure {
3116 fn kind(&self) -> &'static str {
3117 match self {
3118 Self::Enumeration(_) => "enumeration",
3119 Self::Worker(_) => "blocking_worker",
3120 }
3121 }
3122}
3123
3124#[cfg(unix)]
3125impl std::fmt::Display for WalpinAttributionFailure {
3126 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
3127 match self {
3128 Self::Enumeration(error) => write!(
3129 formatter,
3130 "sidecar directory failed the trust-boundary enumeration; cross-process \
3131 WAL-pin attribution is unestablished for this tick: {error}"
3132 ),
3133 Self::Worker(error) => write!(
3134 formatter,
3135 "sidecar attribution blocking worker failed; cross-process WAL-pin \
3136 attribution is unestablished for this tick: {error}"
3137 ),
3138 }
3139 }
3140}
3141
3142#[cfg(unix)]
3145fn capture_walpin_attribution_request(
3146 pool: &ConnectionPool,
3147 state: &mut TruncateState,
3148) -> Option<WalpinAttributionRequest> {
3149 let path = pool.canonical_path()?;
3150 if !crate::walpin::sidecar_enabled(true) {
3151 return None;
3152 }
3153 let legacy_fallback_interval = state.legacy_walpin_fallback_interval;
3154 Some(match state.plan_walpin_attribution_at(Instant::now()) {
3155 WalpinFullScanPlan::Refresh {
3156 previous_last_attempt,
3157 } => WalpinAttributionRequest::Fresh {
3158 dir: crate::walpin::sidecar_dir_for(path),
3159 census: crate::walpin::census_holders(path).map_err(|error| error.to_string()),
3160 legacy_fallback_interval,
3161 previous_last_attempt,
3162 },
3163 WalpinFullScanPlan::Cached(cached) => WalpinAttributionRequest::Cached(cached),
3164 WalpinFullScanPlan::Suppressed => WalpinAttributionRequest::Suppressed,
3165 })
3166}
3167
3168#[cfg(unix)]
3174async fn complete_walpin_attribution(
3175 request: Option<WalpinAttributionRequest>,
3176 state: &mut TruncateState,
3177) -> Result<bool, WalpinAttributionFailure> {
3178 let Some(request) = request else {
3179 return Ok(false);
3180 };
3181 match request {
3182 WalpinAttributionRequest::Suppressed => Ok(false),
3183 WalpinAttributionRequest::Cached(cached) => {
3184 log_walpin_sidecar_report(
3185 &cached.report,
3186 cached.census,
3187 WalpinReportFreshness::Cached {
3188 age: Instant::now().saturating_duration_since(cached.captured_at),
3189 },
3190 );
3191 Ok(true)
3192 }
3193 WalpinAttributionRequest::Fresh {
3194 dir,
3195 census,
3196 legacy_fallback_interval,
3197 previous_last_attempt: _,
3198 } => {
3199 state.sidecar_attribution_attempted_this_tick = true;
3200 if state.walpin_full_scan_last_attempt.is_none() {
3201 state.walpin_full_scan_last_attempt = Some(Instant::now());
3202 }
3203 let fallback = state.walpin_cached_attribution.clone();
3204 let result = tokio::task::spawn_blocking(move || {
3205 #[cfg(test)]
3206 walpin_attribution_test_sync::before_enumeration(&dir);
3207 crate::walpin::enumerate_live(&dir, legacy_fallback_interval)
3208 })
3209 .await
3210 .map_err(|error| WalpinAttributionFailure::Worker(error.to_string()))
3211 .and_then(|result| {
3212 result.map_err(|error| WalpinAttributionFailure::Enumeration(error.to_string()))
3213 });
3214
3215 match result {
3216 Ok(report) => {
3217 let captured_at = Instant::now();
3218 log_walpin_sidecar_report(
3219 &report,
3220 census.clone(),
3221 WalpinReportFreshness::Fresh,
3222 );
3223 state.cache_walpin_attribution(report, census, captured_at);
3224 Ok(true)
3225 }
3226 Err(error) => {
3227 if let Some(cached) = fallback {
3228 log_walpin_sidecar_report(
3229 &cached.report,
3230 cached.census,
3231 WalpinReportFreshness::Cached {
3232 age: Instant::now().saturating_duration_since(cached.captured_at),
3233 },
3234 );
3235 }
3236 Err(error)
3237 }
3238 }
3239 }
3240 }
3241}
3242
3243#[cfg(unix)]
3258fn log_walpin_sidecar_report(
3259 report: &crate::walpin::WalpinReport,
3260 census: Result<crate::walpin::CensusResult, String>,
3261 freshness: WalpinReportFreshness,
3262) {
3263 #[cfg(test)]
3264 walpin_attribution_test_sync::report_used();
3265 let now = now_epoch_secs();
3266 for hb in report.reporting() {
3267 tracing::warn!(
3273 walpin_pid = hb.pid,
3274 walpin_role = %hb.process_role,
3275 walpin_oldest_tx_age_secs = hb.current_oldest_tx_age_secs(now),
3276 walpin_oldest_tx_label = hb.oldest_tx_label.as_deref().unwrap_or("<unlabeled>"),
3277 walpin_attribution_basis = hb.attribution_basis.as_deref().unwrap_or("<unspecified>"),
3278 walpin_attribution_evidence_backed = hb.attribution_is_evidence_backed(),
3279 walpin_attribution_fresh = freshness.is_fresh(),
3280 walpin_health = "reporting",
3281 "ADR-091 Amendment 2 Plank B: live cross-process WAL-pin attribution report"
3282 );
3283 }
3284 for pid in report.registered_silent_pids() {
3285 tracing::debug!(
3286 walpin_pid = pid,
3287 walpin_health = "registered_silent",
3288 walpin_attribution_fresh = freshness.is_fresh(),
3289 "ADR-091 Amendment 2 Plank B: process affirmatively reports no over-threshold span"
3290 );
3291 }
3292 let mut unknown_pids: Vec<u32> = report.unknown_pids().collect();
3293 if let WalpinReportFreshness::Cached { age } = freshness {
3294 tracing::warn!(
3295 walpin_cache_age_ms = age.as_millis() as u64,
3296 "cached WAL-pin attribution is diagnostic-only; fully-attributed \
3297 conclusion is not licensed"
3298 );
3299 unknown_pids.push(0);
3300 }
3301
3302 match census {
3307 Ok(census) => {
3308 let sidecar_known: std::collections::HashSet<u32> = report
3309 .reporting()
3310 .map(|hb| hb.pid)
3311 .chain(report.registered_silent_pids())
3312 .chain(unknown_pids.iter().copied())
3313 .collect();
3314 let mut census_only: Vec<u32> =
3315 census.holders.difference(&sidecar_known).copied().collect();
3316 if !census_only.is_empty() {
3317 census_only.sort_unstable();
3318 tracing::warn!(
3319 ?census_only,
3320 "ADR-091 Amendment 2: these PIDs hold the database file open \
3321 at the OS level but have no sidecar data at all (pre-feature binary, \
3322 sidecar disabled, or wedged before its first write)"
3323 );
3324 unknown_pids.extend(census_only);
3325 }
3326 if !census.is_complete() {
3327 let mut uninspectable = census.uninspectable_pids.clone();
3328 uninspectable.sort_unstable();
3329 tracing::warn!(
3330 ?uninspectable,
3331 truncated = census.truncated,
3332 "ADR-091 Amendment 2: the OS-derived holder census is \
3333 INCOMPLETE — either specific PIDs' open file descriptors could not be \
3334 inspected (permission denied, or a listing race), or the enumeration walk \
3335 itself has positive evidence it did not see the full live-process universe \
3336 (namespace/visibility check, directory-iterator error, self-canary, or a \
3337 libproc buffer that stayed at capacity after bounded retries) — cannot \
3338 rule out an unregistered holder"
3339 );
3340 if uninspectable.is_empty() {
3341 unknown_pids.push(0);
3348 } else {
3349 unknown_pids.extend(uninspectable);
3350 }
3351 }
3352 }
3353 Err(e) => {
3354 tracing::warn!(
3355 error = %e,
3356 "ADR-091 Amendment 2: OS-derived holder census failed; \
3357 attribution cannot rule out an unregistered database holder this tick"
3358 );
3359 unknown_pids.push(0);
3363 }
3364 }
3365
3366 unknown_pids.sort_unstable();
3367 unknown_pids.dedup();
3368 if !unknown_pids.is_empty() {
3369 tracing::warn!(
3370 ?unknown_pids,
3371 "ADR-091 Amendment 2 Plank B: sidecar health unestablished for these PIDs; \
3372 attribution is inconclusive and the native/unregistered-mechanism conclusion \
3373 is NOT licensed this tick"
3374 );
3375 } else if report.reporting().next().is_none() {
3376 tracing::info!(
3377 "ADR-091 Amendment 2 Plank B: every live PID is reporting or registered-silent \
3378 with none pinning; the WAL pin is not attributable to any in-process registry \
3379 span this sidecar covers"
3380 );
3381 }
3382}
3383
3384fn log_wal_pin_depth(conn: &rusqlite::Connection) {
3390 match query_wal_pin_depth(conn) {
3391 Ok((log, checkpointed)) => {
3392 tracing::warn!(
3393 wal_log_frames = log,
3394 wal_checkpointed_frames = checkpointed,
3395 wal_pin_depth = (log - checkpointed).max(0),
3396 "ADR-091 Amendment 2 Plank C: WAL pin depth after TRUNCATE no-progress"
3397 );
3398 }
3399 Err(e) => {
3400 tracing::warn!(
3401 error = %e,
3402 "ADR-091 Amendment 2 Plank C: failed to query WAL pin depth"
3403 );
3404 }
3405 }
3406}
3407
3408fn query_wal_pin_depth(conn: &rusqlite::Connection) -> rusqlite::Result<(i64, i64)> {
3415 conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
3416 Ok((row.get::<_, i64>(1)?, row.get::<_, i64>(2)?))
3417 })
3418}
3419
3420fn crossing_warn(now_above: bool, was_above: &mut bool) -> bool {
3429 let fire = now_above && !*was_above;
3430 *was_above = now_above;
3431 fire
3432}
3433
3434#[derive(Debug, Clone, Copy)]
3435struct RawCheckpointObservation {
3436 busy: i64,
3437 log_frames: i64,
3438 checkpointed_frames: i64,
3439}
3440
3441fn query_checkpoint_observation(
3445 conn: &rusqlite::Connection,
3446) -> rusqlite::Result<RawCheckpointObservation> {
3447 conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
3448 Ok(RawCheckpointObservation {
3449 busy: row.get(0)?,
3450 log_frames: row.get(1)?,
3451 checkpointed_frames: row.get(2)?,
3452 })
3453 })
3454}
3455
3456fn query_wal_pages(conn: &rusqlite::Connection) -> u64 {
3462 let pages = query_checkpoint_observation(conn)
3463 .map(|observation| observation.log_frames)
3464 .unwrap_or(0)
3465 .max(0) as u64;
3466 LAST_WAL_PAGES.store(pages, Ordering::Relaxed);
3470 note_checkpoint_observed(pages);
3471 pages
3472}
3473
3474#[cfg(test)]
3475mod tests {
3476 use super::*;
3477 use crate::pool::PoolConfig;
3478 use crate::writer_task::WriterTaskHandle;
3479 use rusqlite::hooks::{AuthAction, Authorization};
3480 use serial_test::serial;
3481 use tracing::field::{Field, Visit};
3482
3483 #[derive(Clone, Debug, Default)]
3484 struct CapturedEvent {
3485 message: Option<String>,
3486 open_tx_count: Option<u64>,
3487 oldest_tx_age_secs: Option<String>,
3488 elapsed_us: Option<u64>,
3489 busy: Option<i64>,
3490 oldest_tx_label: Option<String>,
3491 tx_label: Option<String>,
3492 census_only: Option<String>,
3493 }
3494
3495 #[derive(Default)]
3496 struct CapturedEventVisitor(CapturedEvent);
3497
3498 impl Visit for CapturedEventVisitor {
3499 fn record_u64(&mut self, field: &Field, value: u64) {
3500 match field.name() {
3501 "open_tx_count" => self.0.open_tx_count = Some(value),
3502 "elapsed_us" => self.0.elapsed_us = Some(value),
3503 _ => {}
3504 }
3505 }
3506
3507 fn record_i64(&mut self, field: &Field, value: i64) {
3508 if field.name() == "busy" {
3509 self.0.busy = Some(value);
3510 }
3511 }
3512
3513 fn record_str(&mut self, field: &Field, value: &str) {
3514 match field.name() {
3515 "message" => self.0.message = Some(value.to_string()),
3516 "oldest_tx_label" => self.0.oldest_tx_label = Some(value.to_string()),
3517 "tx_label" => self.0.tx_label = Some(value.to_string()),
3518 _ => {}
3519 }
3520 }
3521
3522 fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
3523 let formatted = format!("{value:?}");
3524 let cleaned = formatted
3525 .trim_start_matches('"')
3526 .trim_end_matches('"')
3527 .to_string();
3528 match field.name() {
3529 "message" => self.0.message = Some(cleaned),
3530 "oldest_tx_label" => self.0.oldest_tx_label = Some(cleaned),
3531 "tx_label" => self.0.tx_label = Some(cleaned),
3532 "census_only" => self.0.census_only = Some(cleaned),
3533 "oldest_tx_age_secs" => self.0.oldest_tx_age_secs = Some(cleaned),
3534 _ => {}
3535 }
3536 }
3537 }
3538
3539 struct CaptureSubscriber {
3544 events: std::sync::Arc<std::sync::Mutex<Vec<CapturedEvent>>>,
3545 }
3546
3547 impl tracing::Subscriber for CaptureSubscriber {
3548 fn enabled(&self, _: &tracing::Metadata<'_>) -> bool {
3549 true
3550 }
3551 fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
3552 tracing::span::Id::from_u64(1)
3553 }
3554 fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
3555 fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
3556 fn event(&self, event: &tracing::Event<'_>) {
3557 let mut visitor = CapturedEventVisitor::default();
3558 event.record(&mut visitor);
3559 self.events.lock().unwrap().push(visitor.0);
3560 }
3561 fn enter(&self, _: &tracing::span::Id) {}
3562 fn exit(&self, _: &tracing::span::Id) {}
3563 }
3564
3565 fn spans_straddling(
3569 threshold: Duration,
3570 ) -> (
3571 khive_storage::tx_registry::OldestSpan,
3572 khive_storage::tx_registry::OldestSpan,
3573 ) {
3574 let _handle = khive_storage::tx_registry::register(Some("writer_task_tx".to_string()));
3575 let (id, _age, _label) =
3576 khive_storage::tx_registry::oldest().expect("a registration is open");
3577 let base = khive_storage::tx_registry::OldestSpan {
3578 id,
3579 age: Duration::ZERO,
3580 label: Some("writer_task_tx".to_string()),
3581 origin: khive_storage::tx_registry::TxOrigin::Unscoped,
3582 };
3583 let aged = khive_storage::tx_registry::OldestSpan {
3584 age: threshold + Duration::from_secs(1),
3585 ..base.clone()
3586 };
3587 let young = khive_storage::tx_registry::OldestSpan {
3588 age: Duration::from_micros(5_849),
3590 ..base
3591 };
3592 (aged, young)
3593 }
3594
3595 fn capture<F: FnOnce()>(f: F) -> Vec<CapturedEvent> {
3596 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
3597 let subscriber = CaptureSubscriber {
3598 events: std::sync::Arc::clone(&buffer),
3599 };
3600 tracing::subscriber::with_default(subscriber, f);
3601 let events = buffer.lock().unwrap();
3602 events.clone()
3603 }
3604
3605 #[test]
3606 fn truncate_no_progress_warn_reports_nonempty_registry_snapshot_facts() {
3607 let snapshot = [
3608 (Duration::from_secs(2), Some("younger-entry".to_string())),
3609 (Duration::from_secs(7), Some("older-entry".to_string())),
3610 ];
3611 let events = capture(|| log_truncate_no_progress_warn(6003, 6003, &snapshot));
3612 assert_eq!(events.len(), 3, "one summary and both captured entries");
3613 let summary = &events[0];
3614 assert_eq!(summary.open_tx_count, Some(2));
3615 assert_eq!(summary.oldest_tx_age_secs.as_deref(), Some("Some(7.0)"));
3616 let message = summary.message.as_deref().expect("summary message");
3617 assert_eq!(
3618 message,
3619 "WAL TRUNCATE attempt made no progress; open transactions observed in this process's registry"
3620 );
3621 assert!(!message.contains("pinning"));
3622 assert!(!message.contains("long-lived reader"));
3623 assert_eq!(events[1].tx_label.as_deref(), Some("younger-entry"));
3624 assert_eq!(events[2].tx_label.as_deref(), Some("older-entry"));
3625 assert!(events[1..].iter().all(|event| {
3626 event.message.as_deref() == Some("WAL high-water: open transaction registry entry")
3627 }));
3628 }
3629
3630 #[test]
3631 fn truncate_no_progress_warn_reports_empty_process_registry_without_pin_claim() {
3632 let events = capture(|| log_truncate_no_progress_warn(6003, 6003, &[]));
3633 assert_eq!(
3634 events.len(),
3635 1,
3636 "an empty snapshot has no entries to enumerate"
3637 );
3638 let summary = &events[0];
3639 assert_eq!(summary.open_tx_count, Some(0));
3640 assert_eq!(summary.oldest_tx_age_secs.as_deref(), Some("None"));
3641 let message = summary.message.as_deref().expect("summary message");
3642 assert_eq!(
3643 message,
3644 "WAL TRUNCATE attempt made no progress; no open transaction in this process's registry"
3645 );
3646 assert!(!message.contains("pinning"));
3647 assert!(!message.contains("long-lived reader"));
3648 let nonempty = capture(|| {
3649 log_truncate_no_progress_warn(6003, 6003, &[(Duration::ZERO, None)]);
3650 });
3651 assert_ne!(summary.message, nonempty[0].message);
3652 assert_eq!(nonempty[0].open_tx_count, Some(1));
3653 assert_eq!(nonempty[0].oldest_tx_age_secs.as_deref(), Some("Some(0.0)"));
3654 }
3655
3656 #[test]
3659 #[serial(tx_registry)]
3660 fn high_water_warn_names_the_pin_when_an_aged_entry_exists() {
3661 let threshold = Duration::from_secs(30);
3662 let (aged, _young) = spans_straddling(threshold);
3663
3664 let events = capture(|| log_wal_high_water_warn(6003, 6000, Some(&aged), threshold));
3665
3666 let message = events
3667 .iter()
3668 .find_map(|e| e.message.clone())
3669 .expect("one WARN is emitted");
3670 assert!(
3671 message.contains("is pinning a snapshot"),
3672 "an aged entry must produce the pin wording, got {message:?}"
3673 );
3674 assert_eq!(
3675 events.iter().find_map(|e| e.oldest_tx_label.clone()),
3676 Some("writer_task_tx".to_string()),
3677 "the named entry is the one handed in"
3678 );
3679 }
3680
3681 #[test]
3684 #[serial(tx_registry)]
3685 fn high_water_warn_rules_the_pin_out_when_the_oldest_entry_is_young() {
3686 let threshold = Duration::from_secs(30);
3687 let (aged, young) = spans_straddling(threshold);
3688
3689 let young_message =
3690 capture(|| log_wal_high_water_warn(6003, 6000, Some(&young), threshold))
3691 .iter()
3692 .find_map(|e| e.message.clone())
3693 .expect("one WARN is emitted");
3694 let aged_message = capture(|| log_wal_high_water_warn(6003, 6000, Some(&aged), threshold))
3695 .iter()
3696 .find_map(|e| e.message.clone())
3697 .expect("one WARN is emitted");
3698
3699 assert_ne!(
3701 young_message, aged_message,
3702 "the two registry states must produce different text"
3703 );
3704 assert!(
3705 young_message.contains("no open transaction old enough to pin"),
3706 "got {young_message:?}"
3707 );
3708 assert!(
3711 !young_message.contains("pinning"),
3712 "the young branch must not assert a pin, got {young_message:?}"
3713 );
3714 assert!(
3715 !young_message.contains("long-lived reader"),
3716 "the young branch must not name a reader, got {young_message:?}"
3717 );
3718 }
3719
3720 #[test]
3723 #[serial(tx_registry)]
3724 fn high_water_warn_with_an_empty_registry_takes_the_no_pin_branch() {
3725 let events = capture(|| log_wal_high_water_warn(6003, 6000, None, Duration::from_secs(30)));
3726 let message = events
3727 .iter()
3728 .find_map(|e| e.message.clone())
3729 .expect("one WARN is emitted");
3730 assert!(
3731 message.contains("no open transaction old enough to pin"),
3732 "got {message:?}"
3733 );
3734 assert_eq!(
3735 events.iter().find_map(|e| e.oldest_tx_label.clone()),
3736 Some("<none>".to_string()),
3737 "an absent entry is labelled as absent, never as unlabeled"
3738 );
3739 }
3740
3741 #[test]
3744 #[serial(tx_registry)]
3745 fn log_tx_registry_oldest_debug_reports_oldest_open_entry() {
3746 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
3747 let subscriber = CaptureSubscriber {
3748 events: std::sync::Arc::clone(&buffer),
3749 };
3750
3751 let _handle =
3752 khive_storage::tx_registry::register(Some("checkpoint_tick_test".to_string()));
3753
3754 let oldest = khive_storage::tx_registry::oldest().map(|(id, age, label)| {
3755 khive_storage::tx_registry::OldestSpan {
3756 id,
3757 age,
3758 label,
3759 origin: khive_storage::tx_registry::TxOrigin::Unscoped,
3760 }
3761 });
3762 let expected_label = oldest
3763 .as_ref()
3764 .and_then(|s| s.label.clone())
3765 .unwrap_or_else(|| "<unlabeled>".to_string());
3766
3767 tracing::subscriber::with_default(subscriber, || {
3768 log_tx_registry_oldest_debug(100, oldest.as_ref());
3769 });
3770
3771 let events = buffer.lock().unwrap();
3772 assert!(
3773 events.iter().any(|e| {
3774 e.message.as_deref()
3775 == Some("WAL checkpoint tick: oldest open transaction registry entry")
3776 && e.oldest_tx_label.as_deref() == Some(expected_label.as_str())
3777 }),
3778 "expected a log line naming the open registry entry's label, got: {events:?}"
3779 );
3780 }
3781
3782 #[test]
3788 #[serial(tx_registry)]
3789 fn registry_warns_fire_on_crossing_and_do_not_repeat() {
3790 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
3791 let subscriber = CaptureSubscriber {
3792 events: std::sync::Arc::clone(&buffer),
3793 };
3794
3795 let _handle =
3796 khive_storage::tx_registry::register(Some("registry_warn_crossing_test".to_string()));
3797 let oldest = khive_storage::tx_registry::oldest().map(|(id, age, label)| {
3798 khive_storage::tx_registry::OldestSpan {
3799 id,
3800 age,
3801 label,
3802 origin: khive_storage::tx_registry::TxOrigin::Unscoped,
3803 }
3804 });
3805
3806 let mut was_above_warn = false;
3807 let mut was_above_high_water = false;
3808
3809 tracing::subscriber::with_default(subscriber, || {
3810 if crossing_warn(true, &mut was_above_warn) {
3812 log_tx_registry_oldest_warn(6000, oldest.as_ref());
3813 }
3814 if crossing_warn(true, &mut was_above_high_water) {
3815 log_tx_registry_snapshot_warn(6000);
3816 }
3817
3818 if crossing_warn(true, &mut was_above_warn) {
3820 log_tx_registry_oldest_warn(6000, oldest.as_ref());
3821 }
3822 if crossing_warn(true, &mut was_above_high_water) {
3823 log_tx_registry_snapshot_warn(6000);
3824 }
3825 });
3826
3827 let events = buffer.lock().unwrap();
3828
3829 let oldest_warn_count = events
3839 .iter()
3840 .filter(|e| {
3841 e.message.as_deref()
3842 == Some("WAL checkpoint tick: oldest open transaction registry entry")
3843 })
3844 .count();
3845 assert_eq!(
3846 oldest_warn_count, 1,
3847 "oldest-entry WARN must fire exactly once across two above-threshold ticks, got: {events:?}"
3848 );
3849
3850 let snapshot_warn_count = events
3851 .iter()
3852 .filter(|e| {
3853 e.message.as_deref() == Some("WAL high-water: open transaction registry entry")
3854 && e.tx_label.as_deref() == Some("registry_warn_crossing_test")
3855 })
3856 .count();
3857 assert_eq!(
3858 snapshot_warn_count, 1,
3859 "high-water snapshot WARN must fire exactly once across two above-threshold ticks, got: {events:?}"
3860 );
3861 }
3862
3863 #[test]
3866 fn log_tx_age_emission_carries_label_for_both_rungs() {
3867 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
3868 let subscriber = CaptureSubscriber {
3869 events: std::sync::Arc::clone(&buffer),
3870 };
3871
3872 tracing::subscriber::with_default(subscriber, || {
3873 log_tx_age_emission(&TxAgeEmission {
3874 rung: TxAgeRung::Warn,
3875 age: Duration::from_secs(45),
3876 label: Some("plank1_warn_test".to_string()),
3877 });
3878 log_tx_age_emission(&TxAgeEmission {
3879 rung: TxAgeRung::Stale,
3880 age: Duration::from_secs(150),
3881 label: Some("plank1_stale_test".to_string()),
3882 });
3883 });
3884
3885 let events = buffer.lock().unwrap();
3886 assert!(
3887 events.iter().any(|e| {
3888 e.message.as_deref()
3889 == Some(
3890 "ADR-091 Plank 1: open transaction registry entry exceeded soft-cap age",
3891 )
3892 && e.tx_label.as_deref() == Some("plank1_warn_test")
3893 }),
3894 "expected a Warn-rung log line naming the entry, got: {events:?}"
3895 );
3896 assert!(
3897 events.iter().any(|e| {
3898 e.message.as_deref().is_some_and(|m| {
3899 m.starts_with(
3900 "ADR-091 Plank 1: open transaction registry entry exceeded the cooperative",
3901 )
3902 }) && e.tx_label.as_deref() == Some("plank1_stale_test")
3903 }),
3904 "expected a Stale-rung log line naming the entry, got: {events:?}"
3905 );
3906 }
3907
3908 fn file_pool(path: &std::path::Path) -> Arc<ConnectionPool> {
3909 let cfg = PoolConfig {
3910 path: Some(path.to_path_buf()),
3911 ..PoolConfig::for_test()
3912 };
3913 Arc::new(ConnectionPool::new(cfg).expect("pool open"))
3914 }
3915
3916 async fn writer_task_wal_autocheckpoint_pages(handle: &WriterTaskHandle) -> u32 {
3917 handle
3918 .send_top_level(|conn| {
3919 conn.pragma_query_value(None, "wal_autocheckpoint", |row| row.get::<_, u32>(0))
3920 .map_err(|error| khive_storage::error::StorageError::Pool {
3921 operation: "test_wal_autocheckpoint".into(),
3922 message: error.to_string(),
3923 })
3924 })
3925 .await
3926 .expect("query writer-task connection pragma")
3927 }
3928
3929 fn checkpoint_conn(pool: &ConnectionPool) -> rusqlite::Connection {
3933 pool.open_standalone_writer()
3934 .expect("open dedicated checkpoint connection")
3935 }
3936
3937 struct TruncateReportHookGuard;
3938
3939 impl Drop for TruncateReportHookGuard {
3940 fn drop(&mut self) {
3941 truncate_report_test_sync::uninstall();
3942 }
3943 }
3944
3945 #[cfg(unix)]
3946 struct WalpinAttributionHookGuard;
3947
3948 #[cfg(unix)]
3949 impl Drop for WalpinAttributionHookGuard {
3950 fn drop(&mut self) {
3951 walpin_attribution_test_sync::uninstall();
3952 }
3953 }
3954
3955 #[tokio::test(flavor = "current_thread")]
3956 #[cfg(unix)]
3957 #[serial(khive_walpin_sidecar_env)]
3958 async fn diagnostic_legacy_forecast_matches_housekeeping_with_distinct_cadences() {
3959 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
3960 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
3961 let root = tempfile::tempdir().unwrap();
3962 let pool = file_pool(&root.path().join("forecast.db"));
3963 let path = pool.canonical_path().unwrap();
3964 let checkpoint_interval = Duration::from_millis(500);
3965 let session_interval = SessionSweepConfig::default().interval;
3966 assert_eq!(session_interval, Duration::from_secs(5));
3967 let state = TruncateState::default();
3968 assert_eq!(state.legacy_walpin_fallback_interval, session_interval);
3969 let sidecar = WalpinSidecarState::new(Some(path), true, "daemon", checkpoint_interval)
3970 .expect("enabled fixture sidecar");
3971 crate::walpin::ensure_sidecar_dir(&sidecar.dir).unwrap();
3972 let mut paths = Vec::new();
3973 for (pid, age) in [(2_000_000_001, 5), (2_000_000_002, 40)] {
3974 assert!(!crate::walpin::is_process_alive(pid));
3975 let temp = sidecar.dir.join(format!(".{pid}.beacon.tmp"));
3976 std::fs::write(
3977 &temp,
3978 serde_json::to_vec(&serde_json::json!({
3979 "pid": pid, "process_role": "session", "started_at": 1
3980 }))
3981 .unwrap(),
3982 )
3983 .unwrap();
3984 std::fs::File::options()
3985 .write(true)
3986 .open(&temp)
3987 .unwrap()
3988 .set_modified(std::time::SystemTime::now() - Duration::from_secs(age))
3989 .unwrap();
3990 paths.push(temp);
3991 }
3992 let fast = crate::walpin::inspect_live(&sidecar.dir, checkpoint_interval).unwrap();
3993 assert_eq!(
3994 fast.cleanup_would_reap, 2,
3995 "control must distinguish the cadences"
3996 );
3997 let forecast = crate::diagnostics::wal_pin_attribution(path, session_interval);
3998 assert_eq!(forecast.sidecar_listing_truncated, Some(false));
3999 assert_eq!(forecast.sidecar_entries_cleanup_would_reap, Some(1));
4000 assert!(
4001 paths.iter().all(|path| path.exists()),
4002 "inspection retains evidence"
4003 );
4004
4005 let cleanup = sidecar
4006 .reap_dead_entries_bounded(state.legacy_walpin_fallback_interval)
4007 .await
4008 .expect("housekeeping report");
4009 assert_eq!(
4010 Some(cleanup.orphan_temps_reaped),
4011 forecast.sidecar_entries_cleanup_would_reap
4012 );
4013 assert!(paths[0].exists(), "the temp inside the 15s window remains");
4014 assert!(!paths[1].exists(), "the trusted older temp is reaped");
4015 }
4016
4017 #[test]
4018 #[cfg(unix)]
4019 fn walpin_full_scan_cadence_refreshes_first_then_reuses_until_boundary() {
4020 let cadence = Duration::from_secs(30);
4021 let started_at = Instant::now();
4022 let mut state = TruncateState::with_walpin_full_scan_cadence(cadence);
4023
4024 assert!(matches!(
4025 state.plan_walpin_attribution_at(started_at),
4026 WalpinFullScanPlan::Refresh { .. }
4027 ));
4028 state.cache_walpin_attribution(
4029 crate::walpin::WalpinReport::default(),
4030 Ok(crate::walpin::CensusResult::default()),
4031 started_at,
4032 );
4033
4034 assert!(matches!(
4035 state.plan_walpin_attribution_at(started_at + cadence - Duration::from_nanos(1)),
4036 WalpinFullScanPlan::Cached(_)
4037 ));
4038 assert!(matches!(
4039 state.plan_walpin_attribution_at(started_at + cadence),
4040 WalpinFullScanPlan::Refresh { .. }
4041 ));
4042 }
4043
4044 #[test]
4045 #[cfg(unix)]
4046 fn walpin_full_scan_failure_retries_only_after_cadence() {
4047 let cadence = Duration::from_secs(30);
4048 let started_at = Instant::now();
4049 let mut state = TruncateState::with_walpin_full_scan_cadence(cadence);
4050
4051 assert!(matches!(
4052 state.plan_walpin_attribution_at(started_at),
4053 WalpinFullScanPlan::Refresh { .. }
4054 ));
4055 assert!(matches!(
4058 state.plan_walpin_attribution_at(started_at + cadence - Duration::from_nanos(1)),
4059 WalpinFullScanPlan::Suppressed
4060 ));
4061 assert!(matches!(
4062 state.plan_walpin_attribution_at(started_at + cadence),
4063 WalpinFullScanPlan::Refresh { .. }
4064 ));
4065 }
4066
4067 #[test]
4068 #[cfg(unix)]
4069 fn cached_walpin_report_is_diagnostic_only_even_when_fully_attributed() {
4070 let report = crate::walpin::WalpinReport::default();
4071 assert!(
4072 report.fully_attributed(),
4073 "the fixture must otherwise license the sharp conclusion"
4074 );
4075 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
4076 let subscriber = CaptureSubscriber {
4077 events: std::sync::Arc::clone(&buffer),
4078 };
4079
4080 tracing::subscriber::with_default(subscriber, || {
4081 log_walpin_sidecar_report(
4082 &report,
4083 Ok(crate::walpin::CensusResult::default()),
4084 WalpinReportFreshness::Cached {
4085 age: Duration::from_secs(1),
4086 },
4087 );
4088 });
4089
4090 let events = buffer.lock().unwrap();
4091 assert!(
4092 events.iter().any(|event| {
4093 event.message.as_deref()
4094 == Some(
4095 "cached WAL-pin attribution is diagnostic-only; fully-attributed \
4096 conclusion is not licensed",
4097 )
4098 }),
4099 "cached attribution must declare its fail-closed status: {events:?}"
4100 );
4101 assert!(
4102 !events.iter().any(|event| {
4103 event.message.as_deref().is_some_and(|message| {
4104 message.starts_with("ADR-091 Amendment 2 Plank B: every live PID is reporting")
4105 })
4106 }),
4107 "cached attribution must never authorize the fully-attributed conclusion: {events:?}"
4108 );
4109 }
4110
4111 #[tokio::test(flavor = "current_thread")]
4112 #[cfg(unix)]
4113 #[serial(checkpoint_skip_metrics, khive_walpin_sidecar_env)]
4114 async fn progressing_truncate_releases_full_scan_reservation_to_housekeeping() {
4115 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
4116 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
4117 let dir = tempfile::tempdir().unwrap();
4118 let path = dir.path().join("walpin-progress-reservation.db");
4119 let pool = file_pool(&path);
4120 {
4121 let writer = pool.try_writer().unwrap();
4122 writer
4123 .conn()
4124 .execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
4125 .unwrap();
4126 }
4127 let conn = checkpoint_conn(&pool);
4128 let mut state = TruncateState::default();
4129 let config = CheckpointConfig {
4130 truncate_high_water_pages: 0,
4131 truncate_min_interval: Duration::ZERO,
4132 ..CheckpointConfig::default()
4133 };
4134
4135 assert!(
4136 maybe_truncate(&pool, &conn, &config, u64::MAX, &mut state).is_none(),
4137 "a progressing TRUNCATE must not schedule no-progress attribution"
4138 );
4139 assert!(
4140 state.housekeeping_due(),
4141 "unused pre-TRUNCATE reservation must be restored before housekeeping"
4142 );
4143 let sidecar =
4144 WalpinSidecarState::new(pool.canonical_path(), true, "daemon", config.interval)
4145 .expect("file-backed test sidecar");
4146 assert!(
4147 run_walpin_housekeeping_if_due(&sidecar, &mut state, DEFAULT_SESSION_SWEEP_INTERVAL,)
4148 .await,
4149 "the production housekeeping arm must consume one full scan"
4150 );
4151 assert!(state.walpin_cached_attribution.is_some());
4152 }
4153
4154 #[tokio::test(flavor = "current_thread")]
4155 #[cfg(unix)]
4156 #[serial(checkpoint_skip_metrics, khive_walpin_sidecar_env)]
4157 async fn erroring_truncate_releases_full_scan_reservation_to_housekeeping() {
4158 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
4159 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
4160 let dir = tempfile::tempdir().unwrap();
4161 let path = dir.path().join("walpin-error-reservation.db");
4162 let pool = file_pool(&path);
4163 {
4164 let writer = pool.try_writer().unwrap();
4165 writer
4166 .conn()
4167 .execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
4168 .unwrap();
4169 }
4170 let conn = checkpoint_conn(&pool);
4171 conn.authorizer(Some(
4172 |context: rusqlite::hooks::AuthContext<'_>| match context.action {
4173 AuthAction::Pragma { pragma_name, .. }
4174 if pragma_name.eq_ignore_ascii_case("wal_checkpoint") =>
4175 {
4176 Authorization::Deny
4177 }
4178 _ => Authorization::Allow,
4179 },
4180 ))
4181 .unwrap();
4182 let mut state = TruncateState::default();
4183 let config = CheckpointConfig {
4184 truncate_high_water_pages: 0,
4185 truncate_min_interval: Duration::ZERO,
4186 ..CheckpointConfig::default()
4187 };
4188
4189 assert!(
4190 maybe_truncate(&pool, &conn, &config, u64::MAX, &mut state).is_none(),
4191 "an erroring TRUNCATE must not schedule no-progress attribution"
4192 );
4193 conn.authorizer(None::<fn(rusqlite::hooks::AuthContext<'_>) -> Authorization>)
4194 .unwrap();
4195 assert!(
4196 state.housekeeping_due(),
4197 "failed TRUNCATE must restore its unused full-scan reservation"
4198 );
4199 let sidecar =
4200 WalpinSidecarState::new(pool.canonical_path(), true, "daemon", config.interval)
4201 .expect("file-backed test sidecar");
4202 assert!(
4203 run_walpin_housekeeping_if_due(&sidecar, &mut state, DEFAULT_SESSION_SWEEP_INTERVAL,)
4204 .await,
4205 "the production housekeeping arm must consume one full scan"
4206 );
4207 assert!(state.walpin_cached_attribution.is_some());
4208 }
4209
4210 struct ReaderProcess {
4211 child: std::process::Child,
4212 _stdout: std::io::BufReader<std::process::ChildStdout>,
4213 }
4214
4215 impl ReaderProcess {
4216 fn spawn(db_path: &std::path::Path) -> Self {
4217 use std::io::BufRead;
4218 use std::process::Stdio;
4219
4220 let mut child = std::process::Command::new(
4221 std::env::current_exe().expect("resolve current test executable"),
4222 )
4223 .args([
4224 "--exact",
4225 "checkpoint::tests::walpin_transient_reader_process_helper",
4226 "--nocapture",
4227 ])
4228 .env("KHIVE_CHECKPOINT_READER_HELPER_PATH", db_path)
4229 .stdin(Stdio::piped())
4230 .stdout(Stdio::piped())
4231 .spawn()
4232 .expect("spawn transient WAL reader helper");
4233
4234 let stdout = child.stdout.take().expect("capture helper stdout");
4235 let mut reader = std::io::BufReader::new(stdout);
4236 let mut line = String::new();
4237 loop {
4238 line.clear();
4239 let bytes = reader
4240 .read_line(&mut line)
4241 .expect("read transient reader readiness signal");
4242 assert!(bytes > 0, "reader helper exited before readiness signal");
4243 if line.contains("KHIVE_CHECKPOINT_READER_READY") {
4244 break;
4245 }
4246 }
4247 Self {
4248 child,
4249 _stdout: reader,
4250 }
4251 }
4252
4253 fn pid(&self) -> u32 {
4254 self.child.id()
4255 }
4256
4257 fn release(&mut self) {
4258 use std::io::Write;
4259
4260 let mut stdin = self.child.stdin.take().expect("helper stdin is available");
4261 stdin
4262 .write_all(b"release\n")
4263 .expect("release transient reader");
4264 drop(stdin);
4265 let status = self.child.wait().expect("wait for transient reader helper");
4266 assert!(status.success(), "transient reader helper failed: {status}");
4267 }
4268 }
4269
4270 impl Drop for ReaderProcess {
4271 fn drop(&mut self) {
4272 if self.child.try_wait().ok().flatten().is_none() {
4273 let _ = self.child.kill();
4274 let _ = self.child.wait();
4275 }
4276 }
4277 }
4278
4279 #[test]
4280 fn walpin_transient_reader_process_helper() {
4281 use std::io::Write;
4282
4283 let Some(path) = std::env::var_os("KHIVE_CHECKPOINT_READER_HELPER_PATH") else {
4284 return;
4285 };
4286 let conn = rusqlite::Connection::open(path).expect("helper opens database");
4287 conn.execute_batch("BEGIN DEFERRED; SELECT * FROM t;")
4288 .expect("helper pins a read snapshot");
4289 println!("KHIVE_CHECKPOINT_READER_READY");
4290 std::io::stdout().flush().expect("flush readiness signal");
4291 let mut release = String::new();
4292 std::io::stdin()
4293 .read_line(&mut release)
4294 .expect("wait for release signal");
4295 conn.execute_batch("COMMIT")
4296 .expect("helper releases read snapshot");
4297 }
4298
4299 #[tokio::test(flavor = "current_thread")]
4300 #[cfg(unix)]
4301 #[serial(
4302 checkpoint_skip_metrics,
4303 khive_walpin_sidecar_env,
4304 walpin_attribution_async,
4305 walpin_report_seam
4306 )]
4307 async fn no_progress_report_keeps_holder_released_after_truncate_timeout() {
4308 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
4309 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
4310 let dir = tempfile::tempdir().expect("tempdir");
4311 let path = dir.path().join("transient-reader.db");
4312 let pool = file_pool(&path);
4313 {
4314 let writer = pool.try_writer().expect("writer");
4315 writer
4316 .conn()
4317 .execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
4318 .expect("seed WAL before reader snapshot");
4319 }
4320
4321 let mut reader = ReaderProcess::spawn(&path);
4322 let reader_pid = reader.pid();
4323 {
4324 let writer = pool.try_writer().expect("writer");
4325 writer
4326 .conn()
4327 .execute_batch("INSERT INTO t VALUES (2);")
4328 .expect("append WAL behind reader snapshot");
4329 }
4330
4331 let canonical_path = pool
4332 .canonical_path()
4333 .expect("file-backed pool has canonical path")
4334 .to_path_buf();
4335 let (reached_rx, proceed_tx) = truncate_report_test_sync::install(canonical_path.clone());
4336 let _hook_guard = TruncateReportHookGuard;
4337 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
4338 let checkpoint_pool = Arc::clone(&pool);
4339 let dedicated_conn = checkpoint_conn(&checkpoint_pool);
4340 let checkpoint_events = Arc::clone(&buffer);
4341 let checkpoint = std::thread::spawn(move || {
4342 let subscriber = CaptureSubscriber {
4343 events: checkpoint_events,
4344 };
4345 let _subscriber_guard = tracing::subscriber::set_default(subscriber);
4346 let mut state = TruncateState::default();
4347 let result = checkpoint_once_core(
4348 &checkpoint_pool,
4349 &dedicated_conn,
4350 &CheckpointConfig {
4351 truncate_high_water_pages: 0,
4352 truncate_min_interval: Duration::ZERO,
4353 truncate_busy_timeout: Duration::from_millis(50),
4354 ..CheckpointConfig::default()
4355 },
4356 &mut state,
4357 );
4358 (result, state)
4359 });
4360
4361 reached_rx
4362 .recv_timeout(Duration::from_secs(5))
4363 .expect("TRUNCATE must report no progress while the reader is pinned");
4364 reader.release();
4365 let post_attempt_census =
4366 crate::walpin::census_holders(&canonical_path).expect("post-attempt holder census");
4367 assert!(
4368 !post_attempt_census.holders.contains(&reader_pid),
4369 "released reader PID must be absent from a post-attempt census"
4370 );
4371 proceed_tx
4372 .send(())
4373 .expect("allow no-progress reporting to continue");
4374 let (checkpoint_result, mut state) = checkpoint.join().expect("checkpoint thread");
4375 let outcome = checkpoint_result.expect("checkpoint succeeds");
4376 assert!(
4377 outcome.sidecar_attribution.is_some(),
4378 "the synchronous checkpoint result must carry a separate attribution request"
4379 );
4380 assert!(
4381 !state.sidecar_attribution_attempted_this_tick,
4382 "capturing a request is not the same as attempting its directory enumeration"
4383 );
4384
4385 let subscriber = CaptureSubscriber {
4386 events: std::sync::Arc::clone(&buffer),
4387 };
4388 let _subscriber_guard = tracing::subscriber::set_default(subscriber);
4389 let attribution_attempted =
4390 complete_walpin_attribution(outcome.sidecar_attribution, &mut state)
4391 .await
4392 .expect("deferred attribution succeeds");
4393 assert!(
4394 attribution_attempted,
4395 "a no-progress attribution pass must suppress the redundant healthy-housekeeping \
4396 pass for the same tick"
4397 );
4398
4399 let events = buffer.lock().expect("captured events");
4400 let summaries: Vec<_> = events
4401 .iter()
4402 .filter(|event| {
4403 event.message.as_deref().is_some_and(|message| {
4404 message.starts_with("WAL TRUNCATE attempt made no progress;")
4405 })
4406 })
4407 .collect();
4408 assert_eq!(
4409 summaries.len(),
4410 1,
4411 "the real no-progress path emits one summary"
4412 );
4413 let summary = summaries[0];
4414 assert!(summary.open_tx_count.is_some());
4415 assert!(summary.oldest_tx_age_secs.is_some());
4416 let message = summary.message.as_deref().expect("summary message");
4417 assert!(message.contains("in this process's registry"));
4418 assert!(!message.contains("pinning"));
4419 assert!(!message.contains("long-lived reader"));
4420 assert!(
4421 events.iter().any(|event| {
4422 event
4423 .census_only
4424 .as_deref()
4425 .is_some_and(|pids| pids.contains(&reader_pid.to_string()))
4426 }),
4427 "the no-progress report must retain PID {reader_pid} from the pre-attempt census: {events:?}"
4428 );
4429 }
4430
4431 #[tokio::test(flavor = "current_thread")]
4437 #[cfg(unix)]
4438 #[serial(walpin_attribution_async)]
4439 async fn no_progress_attribution_is_off_runtime_and_awaited_before_report_use() {
4440 use std::sync::atomic::Ordering;
4441
4442 let dir = tempfile::tempdir().expect("tempdir");
4443 let sidecar_dir = dir.path().join("checkpoint.db.walpin");
4444 let (reached_rx, proceed_tx, report_counter) =
4445 walpin_attribution_test_sync::install_pause(sidecar_dir.clone());
4446 let _hook_guard = WalpinAttributionHookGuard;
4447
4448 let runtime_thread = std::thread::current().id();
4449 let mut state = TruncateState::default();
4450 let request = Some(WalpinAttributionRequest::Fresh {
4451 dir: sidecar_dir,
4452 census: Ok(crate::walpin::CensusResult::default()),
4453 legacy_fallback_interval: DEFAULT_SESSION_SWEEP_INTERVAL,
4454 previous_last_attempt: None,
4455 });
4456 let completion = tokio::spawn(async move {
4457 let result = complete_walpin_attribution(request, &mut state).await;
4458 (result, state)
4459 });
4460
4461 let blocking_thread = reached_rx
4462 .await
4463 .expect("spawn_blocking attribution reached test seam");
4464 assert_ne!(
4465 blocking_thread, runtime_thread,
4466 "sidecar enumeration must not execute on the current-thread Tokio runtime worker"
4467 );
4468 assert!(
4469 !completion.is_finished(),
4470 "the async attribution owner must await the still-paused blocking enumeration"
4471 );
4472 assert_eq!(
4473 report_counter.load(Ordering::SeqCst),
4474 0,
4475 "the attribution report must not be consumed before enumeration completes"
4476 );
4477
4478 proceed_tx
4479 .send(())
4480 .expect("release blocking attribution enumeration");
4481 let (result, state) = completion.await.expect("attribution task joins");
4482 assert_eq!(result, Ok(true));
4483 assert_eq!(
4484 report_counter.load(Ordering::SeqCst),
4485 1,
4486 "the completed enumeration must feed exactly one report use"
4487 );
4488 assert!(state.sidecar_attribution_attempted_this_tick);
4489 assert!(
4490 !state.housekeeping_due(),
4491 "completed attribution must suppress same-tick housekeeping"
4492 );
4493 }
4494
4495 #[tokio::test(flavor = "current_thread")]
4500 #[cfg(unix)]
4501 #[serial(walpin_attribution_async)]
4502 async fn no_progress_attribution_join_failure_is_honest_and_suppresses_retry() {
4503 let dir = tempfile::tempdir().expect("tempdir");
4504 let sidecar_dir = dir.path().join("checkpoint.db.walpin");
4505 walpin_attribution_test_sync::install_panic(sidecar_dir.clone());
4506 let _hook_guard = WalpinAttributionHookGuard;
4507
4508 let mut state = TruncateState::default();
4509 let request = Some(WalpinAttributionRequest::Fresh {
4510 dir: sidecar_dir,
4511 census: Ok(crate::walpin::CensusResult::default()),
4512 legacy_fallback_interval: DEFAULT_SESSION_SWEEP_INTERVAL,
4513 previous_last_attempt: None,
4514 });
4515
4516 let error = complete_walpin_attribution(request, &mut state)
4517 .await
4518 .expect_err("injected worker panic must surface as failure");
4519 assert!(
4520 matches!(error, WalpinAttributionFailure::Worker(_)),
4521 "join failure must retain its worker classification: {error:?}"
4522 );
4523 assert!(state.sidecar_attribution_attempted_this_tick);
4524 assert!(
4525 !state.housekeeping_due(),
4526 "an indeterminate partial pass must not authorize a second scan"
4527 );
4528 }
4529
4530 #[test]
4535 #[cfg(unix)]
4536 #[serial(checkpoint_skip_metrics)]
4537 fn async_checkpoint_source_keeps_enumeration_behind_awaited_spawn_blocking() {
4538 fn section<'a>(source: &'a str, start: &str, end: &str) -> &'a str {
4539 source
4540 .split_once(start)
4541 .unwrap_or_else(|| panic!("missing source marker {start:?}"))
4542 .1
4543 .split_once(end)
4544 .unwrap_or_else(|| panic!("missing source marker {end:?}"))
4545 .0
4546 }
4547
4548 let source = include_str!("checkpoint.rs");
4549 let checkpoint_core = section(
4550 source,
4551 "fn checkpoint_once_core(",
4552 "/// Evaluate and, if due, attempt a TRUNCATE escalation",
4553 );
4554 let truncate_core = section(
4555 source,
4556 "fn maybe_truncate(",
4557 "#[cfg(test)]\nmod truncate_report_test_sync",
4558 );
4559 let report_logger = section(
4560 source,
4561 "fn log_walpin_sidecar_report(",
4562 "/// ADR-091 Amendment 2 Plank C",
4563 );
4564 for (name, body) in [
4565 ("checkpoint_once_core", checkpoint_core),
4566 ("maybe_truncate", truncate_core),
4567 ("log_walpin_sidecar_report", report_logger),
4568 ] {
4569 assert!(
4570 !body.contains("enumerate_live("),
4571 "{name} must not perform direct sidecar enumeration"
4572 );
4573 }
4574
4575 let async_completion = section(
4576 source,
4577 "async fn complete_walpin_attribution(",
4578 "/// When a TRUNCATE attempt makes no progress",
4579 );
4580 let spawn = async_completion
4581 .find("tokio::task::spawn_blocking")
4582 .expect("completion must spawn blocking work");
4583 let enumerate = async_completion
4584 .find("crate::walpin::enumerate_live")
4585 .expect("blocking closure must perform the attribution enumeration");
4586 let awaited = async_completion[enumerate..]
4587 .find(".await")
4588 .map(|offset| enumerate + offset)
4589 .expect("blocking worker must be awaited");
4590 assert!(spawn < enumerate && enumerate < awaited);
4591
4592 let task = section(
4593 source,
4594 "pub async fn run_checkpoint_task(",
4595 "/// Whether a `CheckpointOutcomeRecorded` transition should be enqueued",
4596 );
4597 let checkpoint = task
4598 .find("checkpoint_once_core(")
4599 .expect("checkpoint core call");
4600 let completion = task
4601 .find("complete_walpin_attribution(")
4602 .expect("awaited attribution completion");
4603 let housekeeping = task
4604 .find("run_walpin_housekeeping_if_due(")
4605 .expect("fallback housekeeping");
4606 let outcome = task
4607 .find("observe_checkpoint_pressure_tick(")
4608 .expect("lifecycle outcome use");
4609 assert!(
4610 checkpoint < completion && completion < housekeeping && housekeeping < outcome,
4611 "tick ordering must be checkpoint -> awaited attribution -> housekeeping decision -> outcome"
4612 );
4613 let housekeeping_helper = section(
4614 source,
4615 "async fn run_walpin_housekeeping_if_due(",
4616 "fn now_epoch_secs()",
4617 );
4618 assert!(
4619 housekeeping_helper.contains("reap_dead_entries_bounded(legacy_fallback_interval)"),
4620 "the ordered housekeeping arm must retain the bounded full scan"
4621 );
4622
4623 let pressure_tick = section(
4628 source,
4629 "fn observe_checkpoint_pressure_tick(",
4630 "/// ADR-091 Plank 0",
4631 );
4632 assert!(
4633 pressure_tick.contains("checkpoint_outcome_should_emit"),
4634 "extracted pressure tick helper must gate on the lifecycle emit decision"
4635 );
4636 }
4637
4638 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4643 async fn fts_maintenance_off_worker_matches_the_direct_call() {
4644 fn fragmented_fixture() -> (tempfile::TempDir, rusqlite::Connection) {
4645 let dir = tempfile::tempdir().expect("tempdir");
4646 let path = dir.path().join("fts-off-worker.db");
4647 let conn = rusqlite::Connection::open(&path).expect("open sqlite");
4648 conn.execute_batch(
4649 "CREATE VIRTUAL TABLE fts_entities USING fts5(namespace UNINDEXED, subject_id UNINDEXED, title, body, tokenize='trigram');
4650 CREATE VIRTUAL TABLE fts_notes USING fts5(namespace UNINDEXED, subject_id UNINDEXED, title, body, tokenize='trigram');
4651 INSERT INTO fts_entities(fts_entities, rank) VALUES('automerge', 0);
4652 INSERT INTO fts_notes(fts_notes, rank) VALUES('automerge', 0);",
4653 )
4654 .expect("create FTS fixtures");
4655 for index in 0..80 {
4658 let body = format!(
4659 "segment fixture {index} keeps enough repeated production recall text to span pages {}",
4660 "memory query corpus ".repeat(40)
4661 );
4662 conn.execute(
4663 "INSERT INTO fts_entities(namespace, subject_id, title, body) VALUES(?1, ?2, ?3, ?4)",
4664 rusqlite::params![
4665 "local",
4666 format!("id-{index}"),
4667 format!("title {index}"),
4668 body
4669 ],
4670 )
4671 .expect("one autocommit FTS write");
4672 }
4673 (dir, conn)
4674 }
4675
4676 let config = crate::fts_maintenance::FtsMaintenanceConfig {
4677 enabled: true,
4678 interval: Duration::ZERO,
4679 merge_pages: 8,
4680 minimum_segments: 2,
4681 };
4682
4683 let (_direct_dir, direct_conn) = fragmented_fixture();
4684 let mut direct_state = crate::fts_maintenance::FtsMaintenanceState::new(Instant::now());
4685 let direct_step = crate::fts_maintenance::run_if_due(
4686 &direct_conn,
4687 &config,
4688 &mut direct_state,
4689 Instant::now(),
4690 )
4691 .expect("direct maintenance step")
4692 .expect("fragmented fixture has a due step");
4693
4694 let (_wrapped_dir, wrapped_conn) = fragmented_fixture();
4695 let wrapped_state = crate::fts_maintenance::FtsMaintenanceState::new(Instant::now());
4696 let (_conn, _state, wrapped_result) =
4697 run_fts_maintenance_off_worker(wrapped_conn, config, wrapped_state, Instant::now())
4698 .await
4699 .expect("maintenance step did not panic");
4700 let wrapped_step = wrapped_result
4701 .expect("wrapped maintenance step")
4702 .expect("fragmented fixture has a due step");
4703
4704 assert_eq!(
4705 direct_step, wrapped_step,
4706 "the spawn_blocking wrapper must produce the same step outcome as calling \
4707 run_if_due directly"
4708 );
4709 assert_eq!(direct_step.table, "fts_entities");
4710 }
4711
4712 #[test]
4718 #[serial(checkpoint_skip_metrics)]
4719 fn checkpoint_once_succeeds_on_file_backed_pool() {
4720 let dir = tempfile::tempdir().unwrap();
4721 let path = dir.path().join("wal_test.db");
4722 let pool = file_pool(&path);
4723
4724 {
4726 let writer = pool.try_writer().unwrap();
4727 writer
4728 .conn()
4729 .execute_batch("CREATE TABLE IF NOT EXISTS t (x INTEGER);")
4730 .unwrap();
4731 writer
4732 .conn()
4733 .execute_batch("INSERT INTO t VALUES (1);")
4734 .unwrap();
4735 }
4736
4737 let conn = checkpoint_conn(&pool);
4738 checkpoint_once(
4739 &pool,
4740 &conn,
4741 &CheckpointConfig::default(),
4742 &mut TruncateState::default(),
4743 )
4744 .expect("checkpoint_once must succeed against a healthy dedicated connection");
4745 }
4746
4747 #[test]
4753 fn open_standalone_writer_fails_on_in_memory_pool() {
4754 let cfg = PoolConfig {
4755 path: None,
4756 ..PoolConfig::default()
4757 };
4758 let pool = Arc::new(ConnectionPool::new(cfg).expect("in-memory pool"));
4759 assert!(
4760 pool.open_standalone_writer().is_err(),
4761 "an in-memory pool must not be able to open a dedicated checkpoint connection"
4762 );
4763 }
4764
4765 #[test]
4771 fn ensure_open_does_not_move_writer_acquisition_counters() {
4772 let dir = tempfile::tempdir().unwrap();
4773 let path = dir.path().join("checkpoint_ensure_open.db");
4774 let pool = file_pool(&path);
4775
4776 let before_first_open = pool.writer_acquisition_snapshot();
4777 let mut checkpoint_conn = CheckpointConnection::new();
4778 checkpoint_conn
4779 .ensure_open(&pool)
4780 .expect("dedicated checkpoint connection must open against a file-backed pool");
4781 assert_eq!(
4782 pool.writer_acquisition_snapshot(),
4783 before_first_open,
4784 "the checkpoint connection's initial open must not count as a writer acquisition"
4785 );
4786
4787 checkpoint_conn.conn = None;
4788 let before_reopen = pool.writer_acquisition_snapshot();
4789 checkpoint_conn
4790 .ensure_open(&pool)
4791 .expect("dedicated checkpoint connection must reopen after invalidation");
4792 assert_eq!(
4793 pool.writer_acquisition_snapshot(),
4794 before_reopen,
4795 "reopening the checkpoint connection must not count as a writer acquisition either"
4796 );
4797 }
4798
4799 #[test]
4800 fn checkpoint_connection_disables_wal_autocheckpoint_on_open_and_reopen() {
4801 let dir = tempfile::tempdir().unwrap();
4802 let path = dir.path().join("checkpoint_autocheckpoint.db");
4803 let pool = file_pool(&path);
4804 let mut checkpoint_conn = CheckpointConnection::new();
4805
4806 let initial: u32 = checkpoint_conn
4807 .ensure_open(&pool)
4808 .expect("dedicated checkpoint connection must open")
4809 .pragma_query_value(None, "wal_autocheckpoint", |row| row.get(0))
4810 .expect("read initial autocheckpoint setting");
4811 assert_eq!(initial, 0);
4812
4813 checkpoint_conn.conn = None;
4814 let reopened: u32 = checkpoint_conn
4815 .ensure_open(&pool)
4816 .expect("dedicated checkpoint connection must reopen")
4817 .pragma_query_value(None, "wal_autocheckpoint", |row| row.get(0))
4818 .expect("read reopened autocheckpoint setting");
4819 assert_eq!(reopened, 0);
4820 }
4821
4822 #[tokio::test(flavor = "current_thread")]
4823 #[serial(checkpoint_skip_metrics)]
4824 async fn failed_checkpoint_claim_keeps_existing_writer_task_on_fallback() {
4825 let dir = tempfile::tempdir().unwrap();
4826 let path = dir.path().join("failed_claim_writer_task.db");
4827 let pool = Arc::new(
4828 ConnectionPool::new(PoolConfig {
4829 path: Some(path),
4830 checkout_timeout: Duration::from_millis(1),
4831 write_queue_enabled: Some(true),
4832 ..PoolConfig::for_test()
4833 })
4834 .expect("pool open"),
4835 );
4836 let writer_task = pool
4837 .writer_task_handle()
4838 .expect("writer-task resolution")
4839 .expect("writer task enabled");
4840 assert_eq!(
4841 writer_task_wal_autocheckpoint_pages(&writer_task).await,
4842 crate::pool::FALLBACK_WAL_AUTOCHECKPOINT_PAGES
4843 );
4844
4845 let legacy_conn = pool.legacy_conn();
4846 let (held_tx, held_rx) = tokio::sync::oneshot::channel();
4847 let (release_tx, release_rx) = std::sync::mpsc::channel();
4848 let holder = tokio::task::spawn_blocking(move || {
4849 let _held_writer = legacy_conn.lock();
4850 held_tx.send(()).expect("signal held pooled writer");
4851 release_rx.recv().expect("release held pooled writer");
4852 });
4853 held_rx.await.expect("pooled writer holder started");
4854 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4855 drop(shutdown_tx);
4856 run_checkpoint_task(
4857 Arc::clone(&pool),
4858 CheckpointConfig {
4859 interval: Duration::from_secs(60),
4860 ..CheckpointConfig::default()
4861 },
4862 None,
4863 shutdown_rx,
4864 true,
4865 )
4866 .await;
4867
4868 assert_eq!(pool.writer_acquisition_snapshot().timeouts, 1);
4869 assert_eq!(
4870 writer_task_wal_autocheckpoint_pages(&writer_task).await,
4871 crate::pool::FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
4872 "failed pooled-writer claim must not partially propagate ownership"
4873 );
4874 release_tx.send(()).expect("release pooled writer");
4875 holder.await.expect("pooled writer holder joined");
4876 }
4877
4878 #[tokio::test]
4879 #[serial(checkpoint_skip_metrics)]
4880 async fn checkpoint_task_exits_on_shutdown_signal() {
4881 let dir = tempfile::tempdir().unwrap();
4882 let path = dir.path().join("wal_task_shutdown.db");
4883 let pool = file_pool(&path);
4884
4885 let cfg = CheckpointConfig {
4887 interval: Duration::from_millis(10),
4888 ..Default::default()
4889 };
4890
4891 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4892 let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
4893
4894 shutdown_tx.send(()).expect("send shutdown signal");
4895
4896 tokio::time::timeout(Duration::from_secs(1), handle)
4897 .await
4898 .expect("checkpoint task should exit within 1s")
4899 .expect("checkpoint task panicked");
4900 }
4901
4902 #[cfg(unix)]
4903 #[tokio::test]
4904 #[serial(checkpoint_skip_metrics, khive_walpin_sidecar_env)]
4905 async fn healthy_checkpoint_tick_reaps_a_dead_walpin_beacon_without_truncate() {
4906 let dir = tempfile::tempdir().unwrap();
4907 let path = dir.path().join("healthy_sidecar_reap.db");
4908 let pool = file_pool(&path);
4909 let sidecar_dir =
4910 crate::walpin::sidecar_dir_for(pool.canonical_path().expect("file-backed pool"));
4911 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
4912 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
4913
4914 let dead_pid = 2_000_000_000;
4915 let dead_beacon = crate::walpin::WalpinBeacon {
4916 pid: dead_pid,
4917 process_role: "session".to_string(),
4918 started_at: 1,
4919 sweep_interval_ms: 5_000,
4920 };
4921 crate::walpin::write_beacon(&sidecar_dir, &dead_beacon)
4922 .expect("seed a crashed process's orphan beacon");
4923 let dead_beacon_path = crate::walpin::beacon_path(&sidecar_dir, dead_pid);
4924 assert!(dead_beacon_path.exists(), "orphan fixture must exist");
4925
4926 let cfg = CheckpointConfig {
4927 interval: Duration::from_millis(10),
4928 warn_pages: u64::MAX,
4929 high_water_pages: u64::MAX,
4930 truncate_high_water_pages: u64::MAX,
4931 ..CheckpointConfig::default()
4932 };
4933 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4934 let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
4935
4936 let reaped = wait_for(Duration::from_secs(2), || !dead_beacon_path.exists()).await;
4937 shutdown_tx.send(()).expect("send shutdown signal");
4938 tokio::time::timeout(Duration::from_secs(1), handle)
4939 .await
4940 .expect("checkpoint task should exit within 1s")
4941 .expect("checkpoint task panicked");
4942
4943 assert!(
4944 reaped,
4945 "the ordinary healthy tick must reap positively dead sidecar residue independently of \
4946 TRUNCATE diagnostics"
4947 );
4948 }
4949
4950 #[tokio::test]
4954 #[serial(checkpoint_skip_metrics)]
4955 async fn checkpoint_task_exits_via_shutdown_signal_with_live_event_store_pool_clone() {
4956 let dir = tempfile::tempdir().unwrap();
4957 let path = dir.path().join("wal_task_event_store.db");
4958 let pool = file_pool(&path);
4959
4960 let cfg = CheckpointConfig {
4961 interval: Duration::from_millis(10),
4962 ..Default::default()
4963 };
4964
4965 let event_store: Arc<dyn khive_storage::EventStore> =
4966 Arc::new(crate::stores::event::SqlEventStore::new_scoped(
4967 Arc::clone(&pool),
4968 true,
4969 "local".to_string(),
4970 ));
4971 let sibling_pool_clone = Arc::clone(&pool);
4976
4977 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4978 let handle = tokio::spawn(run_checkpoint_task(
4979 pool,
4980 cfg,
4981 Some(CheckpointLifecycleOwner::new(event_store, "local")),
4982 shutdown_rx,
4983 true,
4984 ));
4985
4986 assert!(
4990 Arc::strong_count(&sibling_pool_clone) > 1,
4991 "test setup must reproduce the multi-owner shape the bug depends on"
4992 );
4993
4994 shutdown_tx.send(()).expect("send shutdown signal");
4995
4996 tokio::time::timeout(Duration::from_secs(1), handle)
4997 .await
4998 .expect(
4999 "checkpoint task should exit within 1s via the watch signal, \
5000 even with a live sibling Arc<ConnectionPool> clone held by \
5001 the event store",
5002 )
5003 .expect("checkpoint task panicked");
5004 }
5005
5006 #[test]
5007 #[serial]
5008 fn checkpoint_config_env_override() {
5009 std::env::set_var("KHIVE_CHECKPOINT_INTERVAL_MS", "250");
5010 std::env::set_var("KHIVE_WAL_WARN_PAGES", "1500");
5011 std::env::set_var("KHIVE_WAL_HIGH_WATER_PAGES", "8000");
5012 std::env::set_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES", "12000");
5013 std::env::set_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS", "60");
5014 std::env::set_var("KHIVE_WAL_TRUNCATE_BUSY_MS", "500");
5015 std::env::set_var("KHIVE_TX_WARN_SECS", "15");
5016 std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "90");
5017
5018 let cfg = CheckpointConfig::from_env();
5019
5020 std::env::remove_var("KHIVE_CHECKPOINT_INTERVAL_MS");
5021 std::env::remove_var("KHIVE_WAL_WARN_PAGES");
5022 std::env::remove_var("KHIVE_WAL_HIGH_WATER_PAGES");
5023 std::env::remove_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES");
5024 std::env::remove_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS");
5025 std::env::remove_var("KHIVE_WAL_TRUNCATE_BUSY_MS");
5026 std::env::remove_var("KHIVE_TX_WARN_SECS");
5027 std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
5028
5029 assert_eq!(cfg.interval, Duration::from_millis(250));
5030 assert_eq!(cfg.warn_pages, 1500);
5031 assert_eq!(cfg.high_water_pages, 8000);
5032 assert_eq!(cfg.truncate_high_water_pages, 12000);
5033 assert_eq!(cfg.truncate_min_interval, Duration::from_secs(60));
5034 assert_eq!(cfg.truncate_busy_timeout, Duration::from_millis(500));
5035 assert_eq!(cfg.tx_warn_secs, Duration::from_secs(15));
5036 assert_eq!(cfg.tx_max_age_secs, Duration::from_secs(90));
5037 }
5038
5039 #[test]
5040 #[serial]
5041 fn checkpoint_config_defaults_on_invalid_env() {
5042 let default = CheckpointConfig::default();
5043
5044 std::env::set_var("KHIVE_CHECKPOINT_INTERVAL_MS", "not_a_number");
5045 std::env::set_var("KHIVE_WAL_WARN_PAGES", "");
5046 std::env::set_var("KHIVE_WAL_HIGH_WATER_PAGES", "0");
5047 std::env::set_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES", "not_a_number");
5048 std::env::set_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS", "");
5049 std::env::set_var("KHIVE_WAL_TRUNCATE_BUSY_MS", "0");
5050 std::env::set_var("KHIVE_TX_WARN_SECS", "not_a_number");
5051 std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "0");
5052
5053 let cfg = CheckpointConfig::from_env();
5054
5055 std::env::remove_var("KHIVE_CHECKPOINT_INTERVAL_MS");
5056 std::env::remove_var("KHIVE_WAL_WARN_PAGES");
5057 std::env::remove_var("KHIVE_WAL_HIGH_WATER_PAGES");
5058 std::env::remove_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES");
5059 std::env::remove_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS");
5060 std::env::remove_var("KHIVE_WAL_TRUNCATE_BUSY_MS");
5061 std::env::remove_var("KHIVE_TX_WARN_SECS");
5062 std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
5063
5064 assert_eq!(cfg.interval, default.interval);
5065 assert_eq!(cfg.warn_pages, default.warn_pages);
5066 assert_eq!(cfg.high_water_pages, default.high_water_pages);
5067 assert_eq!(
5068 cfg.truncate_high_water_pages,
5069 default.truncate_high_water_pages
5070 );
5071 assert_eq!(cfg.truncate_min_interval, default.truncate_min_interval);
5072 assert_eq!(cfg.truncate_busy_timeout, default.truncate_busy_timeout);
5073 assert_eq!(cfg.tx_warn_secs, default.tx_warn_secs);
5074 assert_eq!(cfg.tx_max_age_secs, default.tx_max_age_secs);
5075 }
5076
5077 #[test]
5082 #[serial(checkpoint_skip_metrics)]
5083 fn checkpoint_high_water_does_not_block_behind_reader() {
5084 let dir = tempfile::tempdir().unwrap();
5085 let path = dir.path().join("high_water_test.db");
5086
5087 let pool = Arc::new(
5091 ConnectionPool::new(PoolConfig {
5092 path: Some(path.clone()),
5093 busy_timeout: Duration::from_millis(2000),
5094 ..PoolConfig::for_test()
5095 })
5096 .expect("pool open"),
5097 );
5098
5099 {
5101 let writer = pool.try_writer().unwrap();
5102 writer
5103 .conn()
5104 .execute_batch(
5105 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
5106 )
5107 .unwrap();
5108 }
5109
5110 let reader = pool.reader().expect("reader");
5114 reader
5115 .conn()
5116 .execute_batch("BEGIN DEFERRED; SELECT * FROM t;")
5117 .expect("begin read tx");
5118
5119 {
5123 let writer = pool.try_writer().unwrap();
5124 writer
5125 .conn()
5126 .execute_batch("INSERT INTO t VALUES (2);")
5127 .unwrap();
5128 }
5129
5130 let checkpoint_config = CheckpointConfig::default();
5131 let conn = checkpoint_conn(&pool);
5132 let start = std::time::Instant::now();
5133 checkpoint_once(
5134 &pool,
5135 &conn,
5136 &checkpoint_config,
5137 &mut TruncateState::default(),
5138 )
5139 .expect("checkpoint_once must succeed against a healthy dedicated connection");
5140 let elapsed = start.elapsed();
5141
5142 reader.conn().execute_batch("COMMIT;").ok();
5144 drop(reader);
5145
5146 let max_elapsed = checkpoint_config.truncate_busy_timeout / 2;
5150 assert!(
5151 elapsed < max_elapsed,
5152 "checkpoint_once with active reader snapshot took {:?}; expected <{:?} \
5153 (PASSIVE must not block on readers; a TRUNCATE regression would block \
5154 for the configured {:?})",
5155 elapsed,
5156 max_elapsed,
5157 checkpoint_config.truncate_busy_timeout
5158 );
5159 }
5160
5161 #[test]
5162 #[serial]
5163 fn checkpoint_config_rejects_zero_for_all_fields() {
5164 let default = CheckpointConfig::default();
5165 std::env::set_var("KHIVE_CHECKPOINT_INTERVAL_MS", "0");
5166 std::env::set_var("KHIVE_WAL_WARN_PAGES", "0");
5167 std::env::set_var("KHIVE_WAL_HIGH_WATER_PAGES", "0");
5168 std::env::set_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES", "0");
5169 std::env::set_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS", "0");
5170 std::env::set_var("KHIVE_WAL_TRUNCATE_BUSY_MS", "0");
5171 std::env::set_var("KHIVE_TX_WARN_SECS", "0");
5172 std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "0");
5173
5174 let cfg = CheckpointConfig::from_env();
5175
5176 std::env::remove_var("KHIVE_CHECKPOINT_INTERVAL_MS");
5177 std::env::remove_var("KHIVE_WAL_WARN_PAGES");
5178 std::env::remove_var("KHIVE_WAL_HIGH_WATER_PAGES");
5179 std::env::remove_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES");
5180 std::env::remove_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS");
5181 std::env::remove_var("KHIVE_WAL_TRUNCATE_BUSY_MS");
5182 std::env::remove_var("KHIVE_TX_WARN_SECS");
5183 std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
5184
5185 assert_eq!(
5186 cfg.interval, default.interval,
5187 "zero interval must fall back to default"
5188 );
5189 assert_eq!(
5190 cfg.warn_pages, default.warn_pages,
5191 "zero warn_pages must fall back to default"
5192 );
5193 assert_eq!(
5194 cfg.high_water_pages, default.high_water_pages,
5195 "zero high_water_pages must fall back to default"
5196 );
5197 assert_eq!(
5198 cfg.truncate_high_water_pages, default.truncate_high_water_pages,
5199 "zero truncate_high_water_pages must fall back to default"
5200 );
5201 assert_eq!(
5202 cfg.truncate_min_interval, default.truncate_min_interval,
5203 "zero truncate_min_interval must fall back to default"
5204 );
5205 assert_eq!(
5206 cfg.truncate_busy_timeout, default.truncate_busy_timeout,
5207 "zero truncate_busy_timeout must fall back to default"
5208 );
5209 assert_eq!(
5210 cfg.tx_warn_secs, default.tx_warn_secs,
5211 "zero tx_warn_secs must fall back to default"
5212 );
5213 assert_eq!(
5214 cfg.tx_max_age_secs, default.tx_max_age_secs,
5215 "zero tx_max_age_secs must fall back to default"
5216 );
5217 }
5218
5219 #[test]
5222 #[serial]
5223 fn checkpoint_config_rejects_reversed_tx_thresholds() {
5224 let default = CheckpointConfig::default();
5225 std::env::set_var("KHIVE_TX_WARN_SECS", "120");
5226 std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "30");
5227
5228 let cfg = CheckpointConfig::from_env();
5229
5230 std::env::remove_var("KHIVE_TX_WARN_SECS");
5231 std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
5232
5233 assert_eq!(
5234 cfg.tx_warn_secs, default.tx_warn_secs,
5235 "a reversed pair must fall back tx_warn_secs to its default, got: {:?}",
5236 cfg.tx_warn_secs
5237 );
5238 assert_eq!(
5239 cfg.tx_max_age_secs, default.tx_max_age_secs,
5240 "a reversed pair must fall back tx_max_age_secs to its default, got: {:?}",
5241 cfg.tx_max_age_secs
5242 );
5243 }
5244
5245 #[test]
5248 #[serial]
5249 fn checkpoint_config_rejects_equal_tx_thresholds() {
5250 let default = CheckpointConfig::default();
5251 std::env::set_var("KHIVE_TX_WARN_SECS", "60");
5252 std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "60");
5253
5254 let cfg = CheckpointConfig::from_env();
5255
5256 std::env::remove_var("KHIVE_TX_WARN_SECS");
5257 std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
5258
5259 assert_eq!(
5260 cfg.tx_warn_secs, default.tx_warn_secs,
5261 "an equal pair must fall back tx_warn_secs to its default, got: {:?}",
5262 cfg.tx_warn_secs
5263 );
5264 assert_eq!(
5265 cfg.tx_max_age_secs, default.tx_max_age_secs,
5266 "an equal pair must fall back tx_max_age_secs to its default, got: {:?}",
5267 cfg.tx_max_age_secs
5268 );
5269 }
5270
5271 #[test]
5274 fn skipped_tick_does_not_reset_high_water_crossing_state() {
5275 let mut was_above = false;
5276
5277 assert!(
5279 crossing_warn(true, &mut was_above),
5280 "should fire on first crossing"
5281 );
5282 assert!(was_above);
5283
5284 assert!(was_above, "was_above must stay true across skipped ticks");
5291
5292 let fired = crossing_warn(true, &mut was_above);
5294 assert!(!fired, "WARN must not re-fire while still above threshold");
5295
5296 let fired = crossing_warn(false, &mut was_above);
5298 assert!(!fired);
5299 assert!(!was_above);
5300
5301 let fired = crossing_warn(true, &mut was_above);
5303 assert!(fired, "WARN must fire again on a new below→above crossing");
5304 }
5305
5306 #[test]
5313 fn warn_pages_fires_once_on_crossing_not_every_tick() {
5314 let mut was_above_warn = false;
5315
5316 let fired_1 = crossing_warn(true, &mut was_above_warn);
5318 let fired_2 = crossing_warn(true, &mut was_above_warn);
5319 let fired_3 = crossing_warn(true, &mut was_above_warn);
5320
5321 assert!(fired_1, "WARN must fire on the first in-band tick");
5322 assert!(
5323 !fired_2,
5324 "WARN must not fire on the second consecutive in-band tick"
5325 );
5326 assert!(
5327 !fired_3,
5328 "WARN must not fire on the third consecutive in-band tick"
5329 );
5330
5331 crossing_warn(false, &mut was_above_warn);
5333 assert!(!was_above_warn);
5334
5335 let fired_reentry = crossing_warn(true, &mut was_above_warn);
5337 assert!(
5338 fired_reentry,
5339 "WARN must fire again on re-entry into warn band"
5340 );
5341 }
5342
5343 #[test]
5349 #[serial(tx_registry, checkpoint_skip_metrics)]
5350 fn truncate_attempts_when_high_water_crossed_with_no_prior_attempt() {
5351 let dir = tempfile::tempdir().unwrap();
5352 let path = dir.path().join("truncate_trigger.db");
5353 let pool = file_pool(&path);
5354
5355 {
5356 let writer = pool.try_writer().unwrap();
5357 writer
5358 .conn()
5359 .execute_batch(
5360 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
5361 )
5362 .unwrap();
5363 }
5364
5365 let config = CheckpointConfig {
5366 truncate_high_water_pages: 0,
5370 truncate_min_interval: Duration::from_secs(300),
5371 ..CheckpointConfig::default()
5372 };
5373 let mut state = TruncateState::default();
5374
5375 assert!(
5376 state.last_attempt.is_none(),
5377 "precondition: no attempt has run yet"
5378 );
5379
5380 let conn = checkpoint_conn(&pool);
5381 checkpoint_once(&pool, &conn, &config, &mut state)
5382 .expect("checkpoint_once must succeed against a healthy dedicated connection");
5383 assert!(
5384 state.last_attempt.is_some(),
5385 "an attempt must be stamped once the high-water threshold is crossed"
5386 );
5387 }
5388
5389 #[test]
5392 #[serial(tx_registry, checkpoint_skip_metrics)]
5393 fn truncate_does_not_attempt_below_high_water() {
5394 let dir = tempfile::tempdir().unwrap();
5395 let path = dir.path().join("truncate_below_threshold.db");
5396 let pool = file_pool(&path);
5397
5398 {
5399 let writer = pool.try_writer().unwrap();
5400 writer
5401 .conn()
5402 .execute_batch(
5403 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
5404 )
5405 .unwrap();
5406 }
5407
5408 let config = CheckpointConfig {
5410 truncate_high_water_pages: u64::MAX,
5411 ..CheckpointConfig::default()
5412 };
5413 let mut state = TruncateState::default();
5414
5415 let conn = checkpoint_conn(&pool);
5416 checkpoint_once(&pool, &conn, &config, &mut state)
5417 .expect("checkpoint_once must succeed against a healthy dedicated connection");
5418
5419 assert!(
5420 state.last_attempt.is_none(),
5421 "a below-threshold tick must never stamp last_attempt"
5422 );
5423 }
5424
5425 #[test]
5429 #[serial(tx_registry, checkpoint_skip_metrics)]
5430 fn truncate_min_interval_skip_does_not_restamp_last_attempt() {
5431 let dir = tempfile::tempdir().unwrap();
5432 let path = dir.path().join("truncate_min_interval.db");
5433 let pool = file_pool(&path);
5434
5435 {
5436 let writer = pool.try_writer().unwrap();
5437 writer
5438 .conn()
5439 .execute_batch(
5440 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
5441 )
5442 .unwrap();
5443 }
5444
5445 let config = CheckpointConfig {
5446 truncate_high_water_pages: 0,
5447 truncate_min_interval: Duration::from_secs(300),
5448 ..CheckpointConfig::default()
5449 };
5450 let mut state = TruncateState::default();
5451 let conn = checkpoint_conn(&pool);
5452
5453 checkpoint_once(&pool, &conn, &config, &mut state)
5454 .expect("checkpoint_once must succeed against a healthy dedicated connection");
5455 let first_attempt = state.last_attempt.expect("first tick must attempt");
5456
5457 checkpoint_once(&pool, &conn, &config, &mut state)
5462 .expect("checkpoint_once must succeed against a healthy dedicated connection");
5463 let second_attempt = state.last_attempt.expect("attempt timestamp must persist");
5464
5465 assert_eq!(
5466 first_attempt, second_attempt,
5467 "a tick within truncate_min_interval must not re-stamp last_attempt"
5468 );
5469 }
5470
5471 #[test]
5483 #[serial(tx_registry, checkpoint_skip_metrics)]
5484 fn checkpoint_once_proceeds_and_can_attempt_truncate_while_pool_writer_held() {
5485 reset_checkpoint_metrics_for_tests();
5486
5487 let dir = tempfile::tempdir().unwrap();
5488 let path = dir.path().join("truncate_busy_skip.db");
5489 let pool = file_pool(&path);
5490
5491 {
5492 let writer = pool.try_writer().unwrap();
5493 writer
5494 .conn()
5495 .execute_batch(
5496 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
5497 )
5498 .unwrap();
5499 }
5500
5501 let conn = checkpoint_conn(&pool);
5502
5503 let _held = pool.try_writer().unwrap();
5507
5508 let config = CheckpointConfig {
5509 truncate_high_water_pages: 0,
5510 ..CheckpointConfig::default()
5511 };
5512 let mut state = TruncateState::default();
5513
5514 checkpoint_once(&pool, &conn, &config, &mut state).expect(
5515 "checkpoint_once must observe normally on its own dedicated connection even \
5516 while a concurrent caller holds the pool's writer mutex",
5517 );
5518
5519 assert!(
5520 state.last_attempt.is_some(),
5521 "a threshold-armed tick must still evaluate (and attempt) TRUNCATE even while \
5522 the pool writer is held — the dedicated connection is unaffected by it"
5523 );
5524 assert_eq!(
5525 checkpoint_skipped_ticks(),
5526 0,
5527 "a busy pool writer must no longer count as a skipped checkpoint tick"
5528 );
5529 assert_eq!(
5530 checkpoint_consecutive_skips(),
5531 0,
5532 "a busy pool writer must not bump the consecutive-skip run length"
5533 );
5534 }
5535
5536 #[test]
5552 #[serial(checkpoint_skip_metrics)]
5553 fn all_checkpoint_metrics_callers_are_serial_tagged() {
5554 const SELF_SRC: &str = include_str!("checkpoint.rs");
5555 let lines: Vec<&str> = SELF_SRC.lines().collect();
5556
5557 let attr_starts: Vec<usize> = lines
5558 .iter()
5559 .enumerate()
5560 .filter(|(_, l)| {
5561 let t = l.trim();
5562 t == "#[test]" || t.starts_with("#[tokio::test")
5563 })
5564 .map(|(i, _)| i)
5565 .collect();
5566
5567 let mut offenders = Vec::new();
5568
5569 for (idx, &start) in attr_starts.iter().enumerate() {
5570 let end = attr_starts.get(idx + 1).copied().unwrap_or(lines.len());
5571 let span = &lines[start..end];
5572
5573 let touches_shared_metrics = span.iter().any(|l| {
5574 l.contains("checkpoint_once(")
5575 || l.contains("checkpoint_once_core(")
5576 || l.contains("run_checkpoint_task(")
5577 });
5578 if !touches_shared_metrics {
5579 continue;
5580 }
5581
5582 let mut in_serial_attr = false;
5585 let has_group_tag = span.iter().any(|line| {
5586 let trimmed = line.trim();
5587 if !in_serial_attr {
5588 in_serial_attr = trimmed.starts_with("#[serial(");
5589 }
5590 if !in_serial_attr {
5591 return false;
5592 }
5593
5594 let has_group = trimmed
5595 .split(|ch: char| !ch.is_ascii_alphanumeric() && ch != '_')
5596 .any(|token| token == "checkpoint_skip_metrics");
5597 if trimmed.ends_with(")]") {
5598 in_serial_attr = false;
5599 }
5600 has_group
5601 });
5602
5603 if !has_group_tag {
5604 let name = span
5605 .iter()
5606 .find_map(|l| {
5607 let t = l.trim_start();
5608 let t = t.strip_prefix("pub(crate) ").unwrap_or(t);
5609 let t = t.strip_prefix("pub ").unwrap_or(t);
5610 let t = t.strip_prefix("async ").unwrap_or(t);
5611 t.strip_prefix("fn ")
5612 .map(|rest| rest.split(['(', '<']).next().unwrap_or("").trim())
5613 })
5614 .unwrap_or("<unknown test>");
5615 offenders.push(name.to_string());
5616 }
5617 }
5618
5619 assert!(
5620 offenders.is_empty(),
5621 "these tests call checkpoint_once/checkpoint_once_core/run_checkpoint_task (which write the \
5622 process-wide LAST_WAL_PAGES/CHECKPOINT_* atomics via query_wal_pages) but \
5623 are not tagged #[serial(checkpoint_skip_metrics)] (or a group including it); \
5624 an untagged caller running concurrently on cargo's default test thread pool \
5625 can clobber those atomics mid-assertion in another test (the #828/#845 race): \
5626 {offenders:?}"
5627 );
5628 }
5629
5630 #[test]
5639 #[serial(tx_registry, checkpoint_skip_metrics)]
5640 fn observed_tick_resets_consecutive_skips_but_not_lifetime_total() {
5641 reset_checkpoint_metrics_for_tests();
5642
5643 let dir = tempfile::tempdir().unwrap();
5644 let path = dir.path().join("skip_then_observe.db");
5645 let pool = file_pool(&path);
5646
5647 {
5648 let writer = pool.try_writer().unwrap();
5649 writer
5650 .conn()
5651 .execute_batch(
5652 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
5653 )
5654 .unwrap();
5655 }
5656
5657 note_checkpoint_skipped();
5659 note_checkpoint_skipped();
5660 assert_eq!(checkpoint_skipped_ticks(), 2);
5661 assert_eq!(checkpoint_consecutive_skips(), 2);
5662
5663 let conn = checkpoint_conn(&pool);
5666 let mut state = TruncateState::default();
5667 checkpoint_once(&pool, &conn, &CheckpointConfig::default(), &mut state)
5668 .expect("checkpoint_once must succeed against a healthy dedicated connection");
5669
5670 assert_eq!(
5671 checkpoint_skipped_ticks(),
5672 2,
5673 "an observed tick must not change the lifetime skipped-tick total"
5674 );
5675 assert_eq!(
5676 checkpoint_consecutive_skips(),
5677 0,
5678 "an observed tick must reset the consecutive-skip run length"
5679 );
5680 }
5681
5682 #[test]
5687 fn note_truncate_outcome_warns_once_at_third_consecutive_failure() {
5688 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5689 let subscriber = CaptureSubscriber {
5690 events: std::sync::Arc::clone(&buffer),
5691 };
5692
5693 let config = CheckpointConfig {
5694 warn_pages: 2000,
5695 ..CheckpointConfig::default()
5696 };
5697 let mut state = TruncateState::default();
5698
5699 tracing::subscriber::with_default(subscriber, || {
5700 note_truncate_outcome(&config, 5000, &mut state);
5702 note_truncate_outcome(&config, 5000, &mut state);
5703 note_truncate_outcome(&config, 5000, &mut state);
5704 note_truncate_outcome(&config, 5000, &mut state);
5706 });
5707
5708 assert_eq!(state.consecutive_failures, 4);
5709
5710 let events = buffer.lock().unwrap();
5711 let escalation_count = events
5712 .iter()
5713 .filter(|e| {
5714 e.message.as_deref()
5715 == Some(
5716 "WAL TRUNCATE has failed to clear WAL pressure for 3 consecutive attempts",
5717 )
5718 })
5719 .count();
5720 assert_eq!(
5721 escalation_count, 1,
5722 "escalation WARN must fire exactly once at the 3rd consecutive failure, got: {events:?}"
5723 );
5724
5725 note_truncate_outcome(&config, 100, &mut state);
5727 assert_eq!(
5728 state.consecutive_failures, 0,
5729 "an attempt that clears warn_pages must reset the consecutive-failure counter"
5730 );
5731 }
5732
5733 fn severity_test_config() -> CheckpointConfig {
5736 CheckpointConfig {
5737 warn_pages: 100,
5738 warn_sustained_cycles: 3,
5739 ..CheckpointConfig::default()
5740 }
5741 }
5742
5743 #[test]
5746 fn severity_ladder_info_on_first_crossing_no_warn() {
5747 let config = severity_test_config();
5748 let mut state = CheckpointSeverityState::default();
5749
5750 let below = state.observe_wal_pages(10, &config);
5751 assert!(below.is_empty(), "below-warn tick must emit nothing");
5752
5753 let above = state.observe_wal_pages(150, &config);
5754 assert_eq!(
5755 above,
5756 vec![CheckpointSeverityEmission {
5757 rung: CheckpointSeverityRung::Info,
5758 wal_pages: 150,
5759 threshold_pages: 100,
5760 consecutive_cycles: 1,
5761 }],
5762 "first below->above crossing must emit exactly one INFO and no WARN"
5763 );
5764 }
5765
5766 #[test]
5769 fn severity_ladder_warn_on_third_consecutive_cycle() {
5770 let config = severity_test_config();
5771 let mut state = CheckpointSeverityState::default();
5772
5773 let tick1 = state.observe_wal_pages(150, &config);
5774 assert_eq!(tick1.len(), 1);
5775 assert_eq!(tick1[0].rung, CheckpointSeverityRung::Info);
5776
5777 let tick2 = state.observe_wal_pages(150, &config);
5778 assert!(
5779 tick2.is_empty(),
5780 "second consecutive above-warn tick must emit nothing yet"
5781 );
5782
5783 let tick3 = state.observe_wal_pages(150, &config);
5784 assert_eq!(
5785 tick3,
5786 vec![CheckpointSeverityEmission {
5787 rung: CheckpointSeverityRung::Warn,
5788 wal_pages: 150,
5789 threshold_pages: 100,
5790 consecutive_cycles: 3,
5791 }],
5792 "WARN must fire exactly on the third consecutive above-warn tick"
5793 );
5794
5795 let tick4 = state.observe_wal_pages(150, &config);
5796 assert!(
5797 tick4.is_empty(),
5798 "WARN must not repeat on a fourth consecutive above-warn tick"
5799 );
5800 }
5801
5802 #[test]
5805 fn severity_ladder_rearms_warn_after_drain() {
5806 let config = severity_test_config();
5807 let mut state = CheckpointSeverityState::default();
5808
5809 for _ in 0..3 {
5811 state.observe_wal_pages(150, &config);
5812 }
5813 assert!(state.warn_emitted_for_episode);
5814
5815 let drain = state.observe_wal_pages(10, &config);
5817 assert!(drain.is_empty(), "a draining tick must emit nothing");
5818
5819 let reentry = state.observe_wal_pages(150, &config);
5821 assert_eq!(reentry.len(), 1);
5822 assert_eq!(reentry[0].rung, CheckpointSeverityRung::Info);
5823
5824 let mid = state.observe_wal_pages(150, &config);
5825 assert!(mid.is_empty());
5826
5827 let second_warn = state.observe_wal_pages(150, &config);
5828 assert_eq!(
5829 second_warn,
5830 vec![CheckpointSeverityEmission {
5831 rung: CheckpointSeverityRung::Warn,
5832 wal_pages: 150,
5833 threshold_pages: 100,
5834 consecutive_cycles: 3,
5835 }],
5836 "a fresh elevation episode after a drain must WARN again"
5837 );
5838 }
5839
5840 #[test]
5843 fn severity_ladder_isolated_crossings_never_warn() {
5844 let config = severity_test_config();
5845 let mut state = CheckpointSeverityState::default();
5846
5847 for _ in 0..3 {
5848 let crossing = state.observe_wal_pages(150, &config);
5849 assert_eq!(
5850 crossing.len(),
5851 1,
5852 "each isolated crossing must emit exactly one INFO"
5853 );
5854 assert_eq!(crossing[0].rung, CheckpointSeverityRung::Info);
5855
5856 let drain = state.observe_wal_pages(10, &config);
5857 assert!(drain.is_empty(), "the drain tick must emit nothing");
5858 }
5859
5860 assert!(
5861 !state.warn_emitted_for_episode,
5862 "isolated single-tick crossings must never accumulate into a WARN"
5863 );
5864 }
5865
5866 #[test]
5871 fn severity_ladder_never_emits_alarm() {
5872 let config = CheckpointConfig {
5873 warn_pages: 100,
5874 warn_sustained_cycles: 1,
5875 ..CheckpointConfig::default()
5876 };
5877 let mut state = CheckpointSeverityState::default();
5878
5879 for wal_pages in [150, 200, 250, u64::MAX] {
5880 let emissions = state.observe_wal_pages(wal_pages, &config);
5881 assert!(
5882 emissions
5883 .iter()
5884 .all(|e| e.rung != CheckpointSeverityRung::Alarm),
5885 "observe_wal_pages must never emit the ALARM rung, got: {emissions:?}"
5886 );
5887 }
5888 }
5889
5890 fn tx_age_test_config() -> CheckpointConfig {
5894 CheckpointConfig {
5895 tx_warn_secs: Duration::from_secs(30),
5896 tx_max_age_secs: Duration::from_secs(120),
5897 ..CheckpointConfig::default()
5898 }
5899 }
5900
5901 fn tx_id(n: u64) -> khive_storage::tx_registry::TxId {
5906 khive_storage::tx_registry::TxId(n)
5907 }
5908
5909 #[test]
5911 fn tx_age_sweep_empty_registry_emits_nothing() {
5912 let config = tx_age_test_config();
5913 let mut state = TxAgeSweepState::default();
5914
5915 let emissions = state.observe(None, config.tx_warn_secs, config.tx_max_age_secs);
5916 assert!(emissions.is_empty(), "no open entry must emit nothing");
5917 }
5918
5919 #[test]
5921 fn tx_age_sweep_fresh_entry_emits_nothing() {
5922 let config = tx_age_test_config();
5923 let mut state = TxAgeSweepState::default();
5924
5925 let emissions = state.observe(
5926 Some((
5927 tx_id(1),
5928 Duration::from_secs(5),
5929 Some("fresh_span".to_string()),
5930 )),
5931 config.tx_warn_secs,
5932 config.tx_max_age_secs,
5933 );
5934 assert!(emissions.is_empty(), "a fresh entry must emit nothing");
5935 }
5936
5937 #[test]
5941 fn tx_age_sweep_warn_fires_once_on_crossing() {
5942 let config = tx_age_test_config();
5943 let mut state = TxAgeSweepState::default();
5944
5945 let tick1 = state.observe(
5946 Some((
5947 tx_id(1),
5948 Duration::from_secs(45),
5949 Some("stale_span".to_string()),
5950 )),
5951 config.tx_warn_secs,
5952 config.tx_max_age_secs,
5953 );
5954 assert_eq!(
5955 tick1,
5956 vec![TxAgeEmission {
5957 rung: TxAgeRung::Warn,
5958 age: Duration::from_secs(45),
5959 label: Some("stale_span".to_string()),
5960 }],
5961 "crossing tx_warn_secs must emit exactly one Warn"
5962 );
5963
5964 let tick2 = state.observe(
5965 Some((
5966 tx_id(1),
5967 Duration::from_secs(50),
5968 Some("stale_span".to_string()),
5969 )),
5970 config.tx_warn_secs,
5971 config.tx_max_age_secs,
5972 );
5973 assert!(
5974 tick2.is_empty(),
5975 "Warn must not repeat while the entry stays in the warn band"
5976 );
5977 }
5978
5979 #[test]
5982 fn tx_age_sweep_stale_fires_once_on_crossing() {
5983 let config = tx_age_test_config();
5984 let mut state = TxAgeSweepState::default();
5985
5986 state.observe(
5989 Some((
5990 tx_id(1),
5991 Duration::from_secs(45),
5992 Some("stuck_writer_task_tx".to_string()),
5993 )),
5994 config.tx_warn_secs,
5995 config.tx_max_age_secs,
5996 );
5997
5998 let tick = state.observe(
5999 Some((
6000 tx_id(1),
6001 Duration::from_secs(130),
6002 Some("stuck_writer_task_tx".to_string()),
6003 )),
6004 config.tx_warn_secs,
6005 config.tx_max_age_secs,
6006 );
6007 assert_eq!(
6008 tick,
6009 vec![TxAgeEmission {
6010 rung: TxAgeRung::Stale,
6011 age: Duration::from_secs(130),
6012 label: Some("stuck_writer_task_tx".to_string()),
6013 }],
6014 "crossing tx_max_age_secs must emit exactly one Stale"
6015 );
6016
6017 let tick_repeat = state.observe(
6018 Some((
6019 tx_id(1),
6020 Duration::from_secs(200),
6021 Some("stuck_writer_task_tx".to_string()),
6022 )),
6023 config.tx_warn_secs,
6024 config.tx_max_age_secs,
6025 );
6026 assert!(
6027 tick_repeat.is_empty(),
6028 "Stale must not repeat while the entry stays above tx_max_age_secs"
6029 );
6030 }
6031
6032 #[test]
6036 fn tx_age_sweep_already_stale_entry_emits_both_rungs_same_tick() {
6037 let config = tx_age_test_config();
6038 let mut state = TxAgeSweepState::default();
6039
6040 let tick = state.observe(
6041 Some((
6042 tx_id(1),
6043 Duration::from_secs(300),
6044 Some("ancient_tx".to_string()),
6045 )),
6046 config.tx_warn_secs,
6047 config.tx_max_age_secs,
6048 );
6049 assert_eq!(
6050 tick,
6051 vec![
6052 TxAgeEmission {
6053 rung: TxAgeRung::Warn,
6054 age: Duration::from_secs(300),
6055 label: Some("ancient_tx".to_string()),
6056 },
6057 TxAgeEmission {
6058 rung: TxAgeRung::Stale,
6059 age: Duration::from_secs(300),
6060 label: Some("ancient_tx".to_string()),
6061 },
6062 ],
6063 "an already-stale entry must cross both rungs on its first observed tick"
6064 );
6065 }
6066
6067 #[test]
6070 fn tx_age_sweep_rearms_after_entry_clears() {
6071 let config = tx_age_test_config();
6072 let mut state = TxAgeSweepState::default();
6073
6074 state.observe(
6075 Some((
6076 tx_id(1),
6077 Duration::from_secs(150),
6078 Some("first_span".to_string()),
6079 )),
6080 config.tx_warn_secs,
6081 config.tx_max_age_secs,
6082 );
6083
6084 let cleared = state.observe(None, config.tx_warn_secs, config.tx_max_age_secs);
6086 assert!(cleared.is_empty(), "a clearing tick must emit nothing");
6087
6088 let fresh = state.observe(
6090 Some((
6091 tx_id(2),
6092 Duration::from_secs(2),
6093 Some("second_span".to_string()),
6094 )),
6095 config.tx_warn_secs,
6096 config.tx_max_age_secs,
6097 );
6098 assert!(fresh.is_empty(), "a fresh oldest entry must emit nothing");
6099
6100 let rewarn = state.observe(
6102 Some((
6103 tx_id(2),
6104 Duration::from_secs(35),
6105 Some("second_span".to_string()),
6106 )),
6107 config.tx_warn_secs,
6108 config.tx_max_age_secs,
6109 );
6110 assert_eq!(
6111 rewarn,
6112 vec![TxAgeEmission {
6113 rung: TxAgeRung::Warn,
6114 age: Duration::from_secs(35),
6115 label: Some("second_span".to_string()),
6116 }],
6117 "a fresh stale episode after a clear must Warn again"
6118 );
6119 }
6120
6121 #[test]
6125 fn tx_age_sweep_stale_replacement_without_intervening_clear_still_names_new_entry() {
6126 let config = tx_age_test_config();
6127 let mut state = TxAgeSweepState::default();
6128
6129 let tick_a = state.observe(
6130 Some((
6131 tx_id(1),
6132 Duration::from_secs(300),
6133 Some("stale_entry_a".to_string()),
6134 )),
6135 config.tx_warn_secs,
6136 config.tx_max_age_secs,
6137 );
6138 assert_eq!(
6139 tick_a.len(),
6140 2,
6141 "entry A must cross both rungs on its first observed tick, got: {tick_a:?}"
6142 );
6143
6144 let tick_b = state.observe(
6147 Some((
6148 tx_id(2),
6149 Duration::from_secs(400),
6150 Some("stale_entry_b".to_string()),
6151 )),
6152 config.tx_warn_secs,
6153 config.tx_max_age_secs,
6154 );
6155 assert_eq!(
6156 tick_b,
6157 vec![
6158 TxAgeEmission {
6159 rung: TxAgeRung::Warn,
6160 age: Duration::from_secs(400),
6161 label: Some("stale_entry_b".to_string()),
6162 },
6163 TxAgeEmission {
6164 rung: TxAgeRung::Stale,
6165 age: Duration::from_secs(400),
6166 label: Some("stale_entry_b".to_string()),
6167 },
6168 ],
6169 "a same-tick identity change to an already-stale successor must re-emit both \
6170 rungs naming the NEW entry, got: {tick_b:?}"
6171 );
6172 }
6173
6174 #[test]
6177 fn tx_age_sweep_uses_configured_thresholds_not_hardcoded_defaults() {
6178 let config = CheckpointConfig {
6179 tx_warn_secs: Duration::from_millis(1),
6180 tx_max_age_secs: Duration::from_millis(2),
6181 ..CheckpointConfig::default()
6182 };
6183 let mut state = TxAgeSweepState::default();
6184
6185 let tick = state.observe(
6186 Some((
6187 tx_id(1),
6188 Duration::from_millis(5),
6189 Some("fast_cap_span".to_string()),
6190 )),
6191 config.tx_warn_secs,
6192 config.tx_max_age_secs,
6193 );
6194 assert_eq!(
6195 tick.len(),
6196 2,
6197 "a millisecond-scale cap must cross both rungs immediately, got: {tick:?}"
6198 );
6199 }
6200
6201 #[test]
6204 #[serial(tx_registry, checkpoint_skip_metrics)]
6205 fn tx_age_sweep_names_long_lived_reader_pinning_wal_past_high_water() {
6206 let dir = tempfile::tempdir().unwrap();
6207 let path = dir.path().join("tx_age_sweep_reader_pin.db");
6208 let pool = file_pool(&path);
6209
6210 {
6211 let writer = pool.try_writer().unwrap();
6212 writer
6213 .conn()
6214 .execute_batch(
6215 "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
6216 )
6217 .unwrap();
6218 }
6219
6220 let reader = pool.reader().expect("reader");
6225 reader
6226 .conn()
6227 .execute_batch("BEGIN DEFERRED; SELECT * FROM t;")
6228 .expect("begin read tx");
6229 let _tx_handle =
6230 khive_storage::tx_registry::register(Some("tx_age_sweep_reader_pin_test".to_string()));
6231
6232 let config = CheckpointConfig {
6235 high_water_pages: 1,
6236 tx_warn_secs: Duration::from_millis(1),
6237 tx_max_age_secs: Duration::from_millis(1),
6238 ..CheckpointConfig::default()
6239 };
6240 {
6241 let writer = pool.try_writer().unwrap();
6242 for i in 0..50 {
6243 writer
6244 .conn()
6245 .execute_batch(&format!("INSERT INTO t VALUES ({i});"))
6246 .unwrap();
6247 }
6248 }
6249
6250 let conn = checkpoint_conn(&pool);
6251 let wal_pages = checkpoint_once(&pool, &conn, &config, &mut TruncateState::default())
6252 .expect("checkpoint_once must succeed against a healthy dedicated connection");
6253 assert!(
6254 wal_pages >= config.high_water_pages,
6255 "test setup must actually drive wal_pages ({wal_pages}) past high_water_pages \
6256 ({}) for this regression to mean anything",
6257 config.high_water_pages
6258 );
6259
6260 std::thread::sleep(Duration::from_millis(5));
6267 let our_entry = khive_storage::tx_registry::snapshot()
6278 .into_iter()
6279 .find(|(_, label)| label.as_deref() == Some("tx_age_sweep_reader_pin_test"))
6280 .expect("this test's own tx_registry entry must still be open");
6281 let mut tx_age_state = TxAgeSweepState::default();
6282 let emissions = tx_age_state.observe(
6283 Some((tx_id(1), our_entry.0, our_entry.1)),
6284 config.tx_warn_secs,
6285 config.tx_max_age_secs,
6286 );
6287 assert!(
6288 emissions.iter().any(|e| e.rung == TxAgeRung::Stale
6289 && e.label.as_deref() == Some("tx_age_sweep_reader_pin_test")),
6290 "expected a Stale emission naming the pinning reader, got: {emissions:?}"
6291 );
6292
6293 reader.conn().execute_batch("COMMIT;").ok();
6294 drop(reader);
6295 drop(_tx_handle);
6296 }
6297
6298 #[test]
6301 #[serial(tx_registry, checkpoint_skip_metrics)]
6302 fn tx_age_sweep_own_entry_survives_concurrent_older_registration() {
6303 let _decoy = khive_storage::tx_registry::register(Some("decoy_unrelated_span".to_string()));
6304 std::thread::sleep(Duration::from_millis(2));
6305 let _own = khive_storage::tx_registry::register(Some("this_test_own_span".to_string()));
6306 std::thread::sleep(Duration::from_millis(5));
6307
6308 let global_oldest = khive_storage::tx_registry::oldest().expect("registry not empty");
6314 assert_ne!(
6315 global_oldest.2.as_deref(),
6316 Some("this_test_own_span"),
6317 "test setup must reproduce the race: an older, unrelated entry must be \
6318 the current global oldest, got: {global_oldest:?}"
6319 );
6320
6321 let our_entry = khive_storage::tx_registry::snapshot()
6322 .into_iter()
6323 .find(|(_, label)| label.as_deref() == Some("this_test_own_span"))
6324 .expect("this test's own tx_registry entry must still be open");
6325
6326 let config = CheckpointConfig {
6327 tx_warn_secs: Duration::from_millis(1),
6328 tx_max_age_secs: Duration::from_millis(1),
6329 ..CheckpointConfig::default()
6330 };
6331 let mut state = TxAgeSweepState::default();
6332 let emissions = state.observe(
6333 Some((tx_id(2), our_entry.0, our_entry.1)),
6334 config.tx_warn_secs,
6335 config.tx_max_age_secs,
6336 );
6337 assert!(
6338 emissions
6339 .iter()
6340 .any(|e| e.rung == TxAgeRung::Stale
6341 && e.label.as_deref() == Some("this_test_own_span")),
6342 "expected a Stale emission naming this test's own span despite an older, \
6343 unrelated concurrent registration, got: {emissions:?}"
6344 );
6345 }
6346
6347 #[test]
6349 #[serial]
6350 fn checkpoint_config_warn_sustained_cycles_env_override() {
6351 let default = CheckpointConfig::default();
6352 assert_eq!(default.warn_sustained_cycles, DEFAULT_WARN_SUSTAINED_CYCLES);
6353
6354 std::env::set_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES", "5");
6355 let cfg = CheckpointConfig::from_env();
6356 std::env::remove_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES");
6357 assert_eq!(cfg.warn_sustained_cycles, 5);
6358
6359 std::env::set_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES", "0");
6360 let cfg_zero = CheckpointConfig::from_env();
6361 std::env::remove_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES");
6362 assert_eq!(
6363 cfg_zero.warn_sustained_cycles, DEFAULT_WARN_SUSTAINED_CYCLES,
6364 "zero must fall back to the default"
6365 );
6366
6367 std::env::set_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES", "not_a_number");
6368 let cfg_invalid = CheckpointConfig::from_env();
6369 std::env::remove_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES");
6370 assert_eq!(
6371 cfg_invalid.warn_sustained_cycles,
6372 DEFAULT_WARN_SUSTAINED_CYCLES
6373 );
6374 }
6375
6376 #[derive(Clone, Copy)]
6379 enum FakeAppendBehavior {
6380 Record,
6381 Fail,
6382 }
6383
6384 struct FakeEventStore {
6385 events: std::sync::Mutex<Vec<khive_storage::Event>>,
6386 append_attempts: std::sync::atomic::AtomicUsize,
6387 append_behavior: FakeAppendBehavior,
6388 }
6389
6390 impl Default for FakeEventStore {
6391 fn default() -> Self {
6392 Self {
6393 events: std::sync::Mutex::new(Vec::new()),
6394 append_attempts: std::sync::atomic::AtomicUsize::new(0),
6395 append_behavior: FakeAppendBehavior::Record,
6396 }
6397 }
6398 }
6399
6400 impl FakeEventStore {
6401 fn failing() -> Self {
6402 Self {
6403 append_behavior: FakeAppendBehavior::Fail,
6404 ..Self::default()
6405 }
6406 }
6407 }
6408
6409 #[async_trait::async_trait]
6410 impl khive_storage::EventStore for FakeEventStore {
6411 async fn append_event(
6412 &self,
6413 event: khive_storage::Event,
6414 ) -> khive_storage::StorageResult<()> {
6415 self.append_attempts
6416 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
6417 match self.append_behavior {
6418 FakeAppendBehavior::Record => {
6419 self.events.lock().unwrap().push(event);
6420 Ok(())
6421 }
6422 FakeAppendBehavior::Fail => Err(khive_storage::StorageError::Internal(
6423 "synthetic checkpoint lifecycle append failure".to_string(),
6424 )),
6425 }
6426 }
6427
6428 async fn append_events(
6429 &self,
6430 events: Vec<khive_storage::Event>,
6431 ) -> khive_storage::StorageResult<khive_storage::BatchWriteSummary> {
6432 let count = events.len() as u64;
6433 self.events.lock().unwrap().extend(events);
6434 Ok(khive_storage::BatchWriteSummary {
6435 attempted: count,
6436 affected: count,
6437 ..khive_storage::BatchWriteSummary::default()
6438 })
6439 }
6440
6441 async fn get_event(
6442 &self,
6443 id: uuid::Uuid,
6444 ) -> khive_storage::StorageResult<Option<khive_storage::Event>> {
6445 Ok(self
6446 .events
6447 .lock()
6448 .unwrap()
6449 .iter()
6450 .find(|e| e.id == id)
6451 .cloned())
6452 }
6453
6454 async fn query_events(
6455 &self,
6456 _filter: khive_storage::EventFilter,
6457 _page: khive_storage::PageRequest,
6458 ) -> khive_storage::StorageResult<khive_storage::Page<khive_storage::Event>> {
6459 unimplemented!("not exercised by the checkpoint lifecycle-event tests")
6460 }
6461
6462 async fn count_events(
6463 &self,
6464 _filter: khive_storage::EventFilter,
6465 ) -> khive_storage::StorageResult<u64> {
6466 Ok(self.events.lock().unwrap().len() as u64)
6467 }
6468 }
6469
6470 #[test]
6474 fn checkpoint_outcome_should_emit_covers_all_transitions() {
6475 assert!(
6476 checkpoint_outcome_should_emit(true, false),
6477 "first elevated tick must emit"
6478 );
6479 assert!(
6480 !checkpoint_outcome_should_emit(true, true),
6481 "sustained elevated ticks must aggregate in memory instead of writing the WAL"
6482 );
6483 assert!(
6484 checkpoint_outcome_should_emit(false, true),
6485 "the single drain row (elevated -> healthy) must emit"
6486 );
6487 assert!(
6488 !checkpoint_outcome_should_emit(false, false),
6489 "an ordinary below-warn tick must not emit"
6490 );
6491 }
6492
6493 #[test]
6496 fn persistent_pressure_lifecycle_rows_are_o_state_transitions() {
6497 let observations = [true; 128].into_iter().chain([false]).chain([false; 128]);
6498 let mut was_elevated = false;
6499 let writes = observations
6500 .filter(|above_warn| {
6501 let emit = checkpoint_outcome_should_emit(*above_warn, was_elevated);
6502 if emit {
6503 was_elevated = *above_warn;
6504 }
6505 emit
6506 })
6507 .count();
6508
6509 assert_eq!(
6510 writes, 2,
6511 "one elevation row plus one recovery summary must cover any number of attempts"
6512 );
6513 }
6514
6515 #[test]
6516 fn checkpoint_pressure_episode_retains_recovery_summary() {
6517 let mut episode = CheckpointPressureEpisode::start(2_500);
6518 episode.observe(2_300);
6519 episode.observe(8_100);
6520 episode.observe(4_000);
6521
6522 assert_eq!(episode.elevated_ticks, 4);
6523 assert_eq!(episode.peak_wal_pages, 8_100);
6524 }
6525
6526 fn drive_pressure_ticks(
6532 config: &CheckpointConfig,
6533 ticks: &[(bool, u64)],
6534 mut fail_on: impl FnMut(usize) -> bool,
6535 ) -> Vec<khive_storage::CheckpointOutcomeRecordedPayload> {
6536 let mut event_elevation_open = false;
6537 let mut pressure_episode: Option<CheckpointPressureEpisode> = None;
6538 let mut pending_recovery: Option<khive_storage::CheckpointOutcomeRecordedPayload> = None;
6539 let mut delivered = Vec::new();
6540 let mut call_index = 0usize;
6541 for &(above_warn, wal_pages) in ticks {
6542 observe_checkpoint_pressure_tick(
6543 above_warn,
6544 wal_pages,
6545 false,
6546 false,
6547 config,
6548 &mut event_elevation_open,
6549 &mut pressure_episode,
6550 &mut pending_recovery,
6551 |payload| {
6552 let idx = call_index;
6553 call_index += 1;
6554 if fail_on(idx) {
6555 false
6556 } else {
6557 delivered.push(payload);
6558 true
6559 }
6560 },
6561 );
6562 }
6563 delivered
6564 }
6565
6566 #[test]
6580 fn dropped_recovery_handoff_does_not_merge_pressure_episodes() {
6581 let config = CheckpointConfig {
6582 warn_pages: 1_000,
6583 ..CheckpointConfig::default()
6584 };
6585 let ticks = [
6586 (true, 1_500), (true, 1_800), (true, 2_000), (false, 500), (false, 400), (true, 3_000), (true, 3_500), (false, 300), ];
6595
6596 let delivered = drive_pressure_ticks(&config, &ticks, |idx| matches!(idx, 1..=3));
6597
6598 assert_eq!(
6599 delivered.len(),
6600 4,
6601 "expected episode-1 open, episode-1 delayed recovery, episode-2 open, \
6602 episode-2 recovery: {delivered:?}"
6603 );
6604
6605 let ep1_open = &delivered[0];
6606 assert!(ep1_open.above_warn);
6607 assert_eq!(ep1_open.episode_elevated_ticks, Some(1));
6608 assert_eq!(ep1_open.episode_peak_wal_pages, Some(1_500));
6609
6610 let ep1_recovery = &delivered[1];
6611 assert!(
6612 !ep1_recovery.above_warn,
6613 "episode 1's recovery must be delivered BEFORE episode 2's opening; \
6614 an opening in this slot means the barrier failed: {delivered:?}"
6615 );
6616 assert_eq!(
6617 ep1_recovery.episode_elevated_ticks,
6618 Some(3),
6619 "episode 1's delayed recovery must report only its own 3 elevated ticks, \
6620 not ticks absorbed from episode 2"
6621 );
6622 assert_eq!(ep1_recovery.episode_peak_wal_pages, Some(2_000));
6623
6624 let ep2_open = &delivered[2];
6625 assert!(ep2_open.above_warn);
6626 assert_eq!(
6627 ep2_open.episode_elevated_ticks,
6628 Some(2),
6629 "episode 2 opens fresh (never continuing episode 1's count), deferred one \
6630 tick by the barrier, so its opening reports 2 elevated ticks"
6631 );
6632 assert_eq!(ep2_open.episode_peak_wal_pages, Some(3_500));
6633
6634 let ep2_recovery = &delivered[3];
6635 assert!(!ep2_recovery.above_warn);
6636 assert_eq!(
6637 ep2_recovery.episode_elevated_ticks,
6638 Some(2),
6639 "episode 2's recovery must report only its own 2 elevated ticks"
6640 );
6641 assert_eq!(ep2_recovery.episode_peak_wal_pages, Some(3_500));
6642 }
6643
6644 #[test]
6651 fn episode_elapsed_entirely_behind_barrier_is_discarded_not_reordered() {
6652 let config = CheckpointConfig {
6653 warn_pages: 1_000,
6654 ..CheckpointConfig::default()
6655 };
6656 let ticks = [
6657 (true, 1_500), (false, 500), (true, 9_000), (false, 400), (false, 300), (true, 2_500), (false, 200), ];
6665
6666 let delivered = drive_pressure_ticks(&config, &ticks, |idx| matches!(idx, 1..=3));
6667
6668 let peaks: Vec<_> = delivered
6669 .iter()
6670 .map(|payload| (payload.above_warn, payload.episode_peak_wal_pages))
6671 .collect();
6672 assert_eq!(
6673 peaks,
6674 vec![
6675 (true, Some(1_500)), (false, Some(1_500)), (true, Some(2_500)), (false, Some(2_500)), ],
6680 "an episode elapsed entirely behind the barrier must not surface late or \
6681 out of order: {delivered:?}"
6682 );
6683 }
6684
6685 #[test]
6690 fn no_dropped_handoff_reports_two_separate_episodes() {
6691 let config = CheckpointConfig {
6692 warn_pages: 1_000,
6693 ..CheckpointConfig::default()
6694 };
6695 let ticks = [
6696 (true, 1_500),
6697 (true, 1_800),
6698 (true, 2_000),
6699 (false, 500),
6700 (false, 400),
6701 (true, 3_000),
6702 (true, 3_500),
6703 (false, 300),
6704 ];
6705
6706 let delivered = drive_pressure_ticks(&config, &ticks, |_idx| false);
6707
6708 assert_eq!(delivered.len(), 4, "{delivered:?}");
6709 assert_eq!(
6710 (
6711 delivered[0].above_warn,
6712 delivered[0].episode_elevated_ticks,
6713 delivered[0].episode_peak_wal_pages
6714 ),
6715 (true, Some(1), Some(1_500)),
6716 "episode 1 open"
6717 );
6718 assert_eq!(
6719 (
6720 delivered[1].above_warn,
6721 delivered[1].episode_elevated_ticks,
6722 delivered[1].episode_peak_wal_pages
6723 ),
6724 (false, Some(3), Some(2_000)),
6725 "episode 1 recovery"
6726 );
6727 assert_eq!(
6728 (
6729 delivered[2].above_warn,
6730 delivered[2].episode_elevated_ticks,
6731 delivered[2].episode_peak_wal_pages
6732 ),
6733 (true, Some(1), Some(3_000)),
6734 "episode 2 open"
6735 );
6736 assert_eq!(
6737 (
6738 delivered[3].above_warn,
6739 delivered[3].episode_elevated_ticks,
6740 delivered[3].episode_peak_wal_pages
6741 ),
6742 (false, Some(2), Some(3_500)),
6743 "episode 2 recovery"
6744 );
6745 }
6746
6747 #[test]
6748 #[serial(checkpoint_skip_metrics)]
6749 fn pressure_diagnostics_count_observations_and_transitions_separately() {
6750 reset_checkpoint_metrics_for_tests();
6751
6752 note_checkpoint_pressure_observation(true, false);
6753 note_checkpoint_pressure_observation(true, true);
6754 note_checkpoint_pressure_observation(true, true);
6755 note_checkpoint_pressure_observation(false, true);
6756 note_checkpoint_pressure_observation(false, false);
6757
6758 assert_eq!(checkpoint_pressure_elevated_ticks(), 3);
6759 assert_eq!(checkpoint_pressure_episodes_started(), 1);
6760 assert_eq!(checkpoint_pressure_episodes_recovered(), 1);
6761 assert_eq!(checkpoint_lifecycle_append_attempts(), 0);
6762 }
6763
6764 #[tokio::test]
6765 #[serial(checkpoint_skip_metrics)]
6766 async fn checkpoint_task_emits_one_opening_for_persistent_pressure() {
6767 reset_checkpoint_metrics_for_tests();
6768 let dir = tempfile::tempdir().unwrap();
6769 let path = dir.path().join("outcome_emit.db");
6770 let pool = file_pool(&path);
6771
6772 let cfg = CheckpointConfig {
6775 interval: Duration::from_millis(10),
6776 warn_pages: 0,
6777 ..CheckpointConfig::default()
6778 };
6779 let store = Arc::new(FakeEventStore::default());
6780 let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
6781
6782 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
6783 let handle = tokio::spawn(run_checkpoint_task(
6784 pool,
6785 cfg,
6786 Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
6787 shutdown_rx,
6788 true,
6789 ));
6790
6791 let progressed = wait_for(Duration::from_secs(10), || {
6792 checkpoint_pressure_elevated_ticks() >= 10
6793 })
6794 .await;
6795 let emitted = wait_for(Duration::from_secs(10), || {
6796 !store.events.lock().unwrap().is_empty()
6797 })
6798 .await;
6799 shutdown_tx.send(()).expect("send shutdown signal");
6800 tokio::time::timeout(Duration::from_secs(1), handle)
6801 .await
6802 .expect("checkpoint task should exit within 1s")
6803 .expect("checkpoint task panicked");
6804
6805 let events = store.events.lock().unwrap();
6806 assert!(
6807 progressed,
6808 "the simulated persistent-pressure episode must span at least ten checkpoint ticks"
6809 );
6810 assert!(
6811 emitted,
6812 "an always-elevated config must append one CheckpointOutcomeRecorded event \
6813 within the poll deadline"
6814 );
6815 assert_eq!(
6816 checkpoint_lifecycle_append_attempts(),
6817 1,
6818 "primary-store lifecycle writes must stay O(state transitions), not O(attempts)"
6819 );
6820 assert_eq!(events.len(), 1);
6821 assert_eq!(events[0].payload_schema_version, 2);
6822 assert_eq!(events[0].payload["episode_elevated_ticks"], 1);
6823 assert_eq!(
6824 events[0].payload["episode_peak_wal_pages"],
6825 events[0].payload["wal_pages"]
6826 );
6827 assert!(
6828 events
6829 .iter()
6830 .all(|e| e.kind == khive_types::EventKind::CheckpointOutcomeRecorded),
6831 "every appended event must be CheckpointOutcomeRecorded, got: {events:?}"
6832 );
6833 assert!(
6834 events.iter().all(|e| e.namespace == "local"),
6835 "events must be stamped with the namespace passed to run_checkpoint_task"
6836 );
6837 }
6838
6839 #[tokio::test]
6843 #[serial(checkpoint_skip_metrics)]
6844 async fn checkpoint_cycles_and_task_shutdown_do_not_wait_for_a_contended_lifecycle_writer() {
6845 reset_checkpoint_metrics_for_tests();
6846 let dir = tempfile::tempdir().unwrap();
6847 let path = dir.path().join("outcome_contended_sink.db");
6848 let checkpoint_pool = file_pool(&path);
6849
6850 let event_pool = Arc::new(
6851 ConnectionPool::new(PoolConfig {
6852 path: None,
6853 checkout_timeout: Duration::from_secs(5),
6854 write_queue_enabled: Some(false),
6855 ..PoolConfig::default()
6856 })
6857 .expect("event pool"),
6858 );
6859 {
6860 let writer = event_pool.try_writer().expect("initialize event schema");
6861 crate::stores::event::ensure_events_schema(writer.conn())
6862 .expect("initialize event schema");
6863 }
6864 let event_store: Arc<dyn khive_storage::EventStore> =
6865 Arc::new(crate::stores::event::SqlEventStore::new_scoped(
6866 Arc::clone(&event_pool),
6867 false,
6868 "local",
6869 ));
6870 let held_event_writer = event_pool
6871 .try_writer()
6872 .expect("hold the event-store writer");
6873
6874 let cfg = CheckpointConfig {
6875 interval: Duration::from_millis(10),
6876 warn_pages: 0,
6877 ..CheckpointConfig::default()
6878 };
6879 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
6880 let handle = tokio::spawn(run_checkpoint_task(
6881 checkpoint_pool,
6882 cfg,
6883 Some(CheckpointLifecycleOwner::new(event_store, "local")),
6884 shutdown_rx,
6885 true,
6886 ));
6887
6888 let progressed = wait_for(Duration::from_secs(2), || {
6889 checkpoint_pressure_elevated_ticks() >= 10
6890 })
6891 .await;
6892 assert!(
6893 progressed,
6894 "checkpoint observations must continue while the lifecycle append is contended"
6895 );
6896 assert_eq!(checkpoint_lifecycle_append_attempts(), 1);
6897 assert_eq!(checkpoint_lifecycle_enqueue_drops(), 0);
6898
6899 shutdown_tx.send(()).expect("send shutdown signal");
6900 tokio::time::timeout(Duration::from_secs(1), handle)
6901 .await
6902 .expect(
6903 "the run_checkpoint_task handle must not wait for the event store's \
6904 five-second writer checkout",
6905 )
6906 .expect("checkpoint task panicked");
6907
6908 drop(held_event_writer);
6913 }
6914
6915 #[tokio::test]
6918 #[serial(checkpoint_skip_metrics)]
6919 async fn checkpoint_task_continues_after_lifecycle_append_failure() {
6920 reset_checkpoint_metrics_for_tests();
6921 let dir = tempfile::tempdir().unwrap();
6922 let path = dir.path().join("outcome_failing_sink.db");
6923 let pool = file_pool(&path);
6924 let store = Arc::new(FakeEventStore::failing());
6925 let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
6926
6927 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6928 let subscriber = CaptureSubscriber {
6929 events: std::sync::Arc::clone(&buffer),
6930 };
6931 let _tracing_guard = tracing::subscriber::set_default(subscriber);
6932
6933 let cfg = CheckpointConfig {
6934 interval: Duration::from_millis(10),
6935 warn_pages: 0,
6936 ..CheckpointConfig::default()
6937 };
6938 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
6939 let handle = tokio::spawn(run_checkpoint_task(
6940 pool,
6941 cfg,
6942 Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
6943 shutdown_rx,
6944 true,
6945 ));
6946
6947 let progressed = wait_for(Duration::from_secs(2), || {
6948 checkpoint_pressure_elevated_ticks() >= 10
6949 })
6950 .await;
6951 shutdown_tx.send(()).expect("send shutdown signal");
6952 tokio::time::timeout(Duration::from_secs(1), handle)
6953 .await
6954 .expect("checkpoint task should remain responsive after sink failure")
6955 .expect("checkpoint task panicked");
6956
6957 assert!(
6958 progressed,
6959 "a failed append must not terminate or stall the checkpoint task"
6960 );
6961 assert_eq!(
6962 store
6963 .append_attempts
6964 .load(std::sync::atomic::Ordering::Relaxed),
6965 1,
6966 "a persistent pressure state must not retry one primary-store append per tick"
6967 );
6968 assert_eq!(checkpoint_lifecycle_append_attempts(), 1);
6969 assert_eq!(checkpoint_lifecycle_append_failures(), 1);
6970 let captured = buffer.lock().unwrap().clone();
6971 assert!(
6972 captured.iter().any(|event| event.message.as_deref()
6973 == Some("checkpoint lifecycle event append failed")),
6974 "lifecycle append failures must remain observable; got: {:?}",
6975 captured
6976 );
6977 }
6978
6979 #[tokio::test]
6980 #[serial(checkpoint_skip_metrics)]
6981 async fn secondary_checkpoint_task_with_lifecycle_ownership_emits_outcome_events() {
6982 let dir = tempfile::tempdir().unwrap();
6983 let path = dir.path().join("secondary_outcome.db");
6984 let pool = file_pool(&path);
6985 let cfg = CheckpointConfig {
6986 interval: Duration::from_millis(10),
6987 warn_pages: 0,
6988 ..CheckpointConfig::default()
6989 };
6990 let store = Arc::new(FakeEventStore::default());
6991 let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
6992
6993 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
6994 let handle = tokio::spawn(run_checkpoint_task(
6995 pool,
6996 cfg,
6997 Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
6998 shutdown_rx,
6999 false,
7000 ));
7001
7002 let emitted = wait_for(Duration::from_secs(10), || {
7005 !store.events.lock().unwrap().is_empty()
7006 })
7007 .await;
7008 shutdown_tx.send(()).expect("send shutdown signal");
7009 tokio::time::timeout(Duration::from_secs(1), handle)
7010 .await
7011 .expect("checkpoint task should exit within 1s")
7012 .expect("checkpoint task panicked");
7013
7014 assert!(
7015 emitted,
7016 "a designated secondary lifecycle owner must append outcome events within the poll \
7017 deadline"
7018 );
7019 }
7020
7021 #[tokio::test]
7022 #[serial(checkpoint_skip_metrics)]
7023 async fn checkpoint_task_emits_nothing_while_healthy() {
7024 let dir = tempfile::tempdir().unwrap();
7025 let path = dir.path().join("outcome_no_emit.db");
7026 let pool = file_pool(&path);
7027
7028 let cfg = CheckpointConfig {
7031 interval: Duration::from_millis(10),
7032 warn_pages: u64::MAX,
7033 ..CheckpointConfig::default()
7034 };
7035 let store = Arc::new(FakeEventStore::default());
7036 let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
7037
7038 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
7039 let handle = tokio::spawn(run_checkpoint_task(
7040 pool,
7041 cfg,
7042 Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
7043 shutdown_rx,
7044 true,
7045 ));
7046
7047 tokio::time::sleep(Duration::from_millis(60)).await;
7048 shutdown_tx.send(()).expect("send shutdown signal");
7049 tokio::time::timeout(Duration::from_secs(1), handle)
7050 .await
7051 .expect("checkpoint task should exit within 1s")
7052 .expect("checkpoint task panicked");
7053
7054 assert!(
7055 store.events.lock().unwrap().is_empty(),
7056 "a config that never crosses warn_pages must never append a lifecycle event"
7057 );
7058 }
7059
7060 #[tokio::test]
7061 #[serial(checkpoint_skip_metrics)]
7062 async fn checkpoint_task_with_no_event_store_does_not_panic() {
7063 let dir = tempfile::tempdir().unwrap();
7064 let path = dir.path().join("outcome_none_store.db");
7065 let pool = file_pool(&path);
7066
7067 let cfg = CheckpointConfig {
7068 interval: Duration::from_millis(10),
7069 warn_pages: 0,
7070 ..CheckpointConfig::default()
7071 };
7072
7073 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
7074 let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
7075
7076 tokio::time::sleep(Duration::from_millis(40)).await;
7077 shutdown_tx.send(()).expect("send shutdown signal");
7078 tokio::time::timeout(Duration::from_secs(1), handle)
7079 .await
7080 .expect("checkpoint task should exit within 1s")
7081 .expect("checkpoint task panicked");
7082 }
7083
7084 #[tokio::test]
7101 #[serial(tx_registry, checkpoint_skip_metrics)]
7102 async fn checkpoint_task_sweeps_stale_registry_entry_while_wal_is_healthy() {
7103 let dir = tempfile::tempdir().unwrap();
7104 let path = dir.path().join("tx_age_sweep_task_healthy_wal.db");
7105 let pool = file_pool(&path);
7106
7107 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7108 let subscriber = CaptureSubscriber {
7109 events: std::sync::Arc::clone(&buffer),
7110 };
7111 let _tracing_guard = tracing::subscriber::set_default(subscriber);
7112
7113 let _tx_handle = khive_storage::tx_registry::register(Some(
7114 "checkpoint_task_healthy_wal_sweep_test".to_string(),
7115 ));
7116
7117 let cfg = CheckpointConfig {
7118 interval: Duration::from_millis(10),
7119 warn_pages: u64::MAX,
7120 high_water_pages: u64::MAX,
7121 truncate_high_water_pages: u64::MAX,
7122 tx_warn_secs: Duration::from_millis(1),
7123 tx_max_age_secs: Duration::from_millis(1),
7124 ..CheckpointConfig::default()
7125 };
7126
7127 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
7128 let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
7129
7130 let swept = wait_for(Duration::from_secs(10), || {
7137 buffer.lock().unwrap().iter().any(|e| {
7138 e.tx_label.as_deref() == Some("checkpoint_task_healthy_wal_sweep_test")
7139 && e.message
7140 .as_deref()
7141 .is_some_and(|m| m.contains("stale-op cap"))
7142 })
7143 })
7144 .await;
7145
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 drop(_tx_handle);
7153
7154 let events = buffer.lock().unwrap();
7155 assert!(
7156 swept,
7157 "expected the spawned task to sweep and escalate the stale registry entry \
7158 to Stale on its own within the poll deadline, got: {events:?}"
7159 );
7160 }
7161
7162 #[tokio::test]
7166 #[serial(tx_registry, checkpoint_skip_metrics)]
7167 async fn checkpoint_task_emits_no_age_alert_for_an_empty_registry() {
7168 let dir = tempfile::tempdir().unwrap();
7169 let path = dir.path().join("tx_age_sweep_task_empty_registry.db");
7170 let pool = file_pool(&path);
7171
7172 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7173 let subscriber = CaptureSubscriber {
7174 events: std::sync::Arc::clone(&buffer),
7175 };
7176 let _tracing_guard = tracing::subscriber::set_default(subscriber);
7177
7178 let cfg = CheckpointConfig {
7179 interval: Duration::from_millis(10),
7180 tx_warn_secs: Duration::from_millis(1),
7181 tx_max_age_secs: Duration::from_millis(1),
7182 ..CheckpointConfig::default()
7183 };
7184
7185 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
7186 let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
7187
7188 tokio::time::sleep(Duration::from_millis(40)).await;
7189 shutdown_tx.send(()).expect("send shutdown signal");
7190 tokio::time::timeout(Duration::from_secs(1), handle)
7191 .await
7192 .expect("checkpoint task should exit within 1s")
7193 .expect("checkpoint task panicked");
7194
7195 let events = buffer.lock().unwrap();
7196 assert!(
7197 events.iter().all(|e| e
7198 .message
7199 .as_deref()
7200 .is_none_or(|m| !m.contains("ADR-091 Plank 1"))),
7201 "an empty registry must never produce a Plank 1 age emission, got: {events:?}"
7202 );
7203 }
7204
7205 #[tokio::test]
7218 #[serial(tx_registry, checkpoint_skip_metrics)]
7219 async fn checkpoint_task_sweeps_stale_entry_even_when_dedicated_connection_is_unavailable_every_tick(
7220 ) {
7221 reset_checkpoint_metrics_for_tests();
7222
7223 let dir = tempfile::tempdir().unwrap();
7224 let path = dir.path().join("tx_age_sweep_task_conn_unavailable.db");
7225 {
7226 let seed_pool = file_pool(&path);
7230 let writer = seed_pool.try_writer().unwrap();
7231 writer
7232 .conn()
7233 .execute_batch("CREATE TABLE IF NOT EXISTS t (x INTEGER);")
7234 .unwrap();
7235 }
7236
7237 #[cfg(unix)]
7238 {
7239 khive_storage::test_support::freeze_snapshot_sidecars(&path);
7240 }
7241
7242 let pool = Arc::new(
7243 ConnectionPool::new(PoolConfig {
7244 path: Some(path.clone()),
7245 read_only: true,
7246 ..PoolConfig::for_test()
7247 })
7248 .expect("read-only pool open"),
7249 );
7250 assert!(
7251 pool.open_standalone_writer().is_err(),
7252 "test precondition: a read-only pool must never be able to open a dedicated \
7253 checkpoint connection"
7254 );
7255
7256 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7257 let subscriber = CaptureSubscriber {
7258 events: std::sync::Arc::clone(&buffer),
7259 };
7260 let _tracing_guard = tracing::subscriber::set_default(subscriber);
7261
7262 let _tx_handle = khive_storage::tx_registry::register(Some(
7263 "checkpoint_task_conn_unavailable_sweep_test".to_string(),
7264 ));
7265
7266 let cfg = CheckpointConfig {
7267 interval: Duration::from_millis(10),
7268 tx_warn_secs: Duration::from_millis(1),
7269 tx_max_age_secs: Duration::from_millis(1),
7270 ..CheckpointConfig::default()
7271 };
7272
7273 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
7274 let handle = tokio::spawn(run_checkpoint_task(
7275 Arc::clone(&pool),
7276 cfg,
7277 None,
7278 shutdown_rx,
7279 true,
7280 ));
7281
7282 let deadline = tokio::time::Instant::now() + Duration::from_secs(10);
7288 while checkpoint_skipped_ticks() == 0 {
7289 assert!(
7290 tokio::time::Instant::now() < deadline,
7291 "test setup must actually drive at least one Skipped tick for this \
7292 regression to mean anything (none within 10s)"
7293 );
7294 tokio::time::sleep(Duration::from_millis(10)).await;
7295 }
7296
7297 loop {
7298 let events = buffer.lock().unwrap().clone();
7299 if events.iter().any(|e| {
7300 e.tx_label.as_deref() == Some("checkpoint_task_conn_unavailable_sweep_test")
7301 && e.message
7302 .as_deref()
7303 .is_some_and(|m| m.contains("stale-op cap"))
7304 }) {
7305 break;
7306 }
7307 assert!(
7308 tokio::time::Instant::now() < deadline,
7309 "expected the age sweep to fire even though every tick's dedicated \
7310 connection was unavailable within 10s, got: {events:?}"
7311 );
7312 tokio::time::sleep(Duration::from_millis(10)).await;
7313 }
7314
7315 shutdown_tx.send(()).expect("send shutdown signal");
7316 tokio::time::timeout(Duration::from_secs(1), handle)
7317 .await
7318 .expect("checkpoint task should exit within 1s")
7319 .expect("checkpoint task panicked");
7320 drop(_tx_handle);
7321 }
7322
7323 #[tokio::test]
7327 async fn session_sweep_task_exits_on_shutdown_signal() {
7328 let cfg = SessionSweepConfig {
7329 interval: Duration::from_millis(10),
7330 ..SessionSweepConfig::default()
7331 };
7332 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
7333 let handle = tokio::spawn(run_session_sweep_task(Vec::new(), cfg, shutdown_rx));
7334
7335 shutdown_tx.send(()).expect("send shutdown signal");
7336
7337 tokio::time::timeout(Duration::from_secs(1), handle)
7338 .await
7339 .expect("session sweep task should exit within 1s")
7340 .expect("session sweep task panicked");
7341 }
7342
7343 async fn wait_for(deadline: Duration, mut cond: impl FnMut() -> bool) -> bool {
7347 let start = std::time::Instant::now();
7348 while start.elapsed() < deadline {
7349 if cond() {
7350 return true;
7351 }
7352 tokio::time::sleep(Duration::from_millis(5)).await;
7353 }
7354 cond()
7355 }
7356
7357 #[tokio::test]
7358 #[serial(khive_walpin_sidecar_env)]
7359 async fn walpin_observe_drops_beacon_when_heartbeat_write_fails() {
7360 let dir = tempfile::tempdir().unwrap();
7361 let db_path = dir.path().join("observe_gate.db");
7362 let sidecar_dir = crate::walpin::sidecar_dir_for(&db_path);
7363 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
7364 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
7365
7366 let mut state = WalpinSidecarState::new(
7367 Some(db_path.as_path()),
7368 true,
7369 "session",
7370 Duration::from_millis(500),
7371 )
7372 .expect("sidecar enabled for a file-backed path");
7373 let pid = std::process::id();
7374 state.register_beacon().await;
7375 let beacon_path = sidecar_dir.join(format!("{pid}.beacon"));
7376 let before = std::fs::metadata(&beacon_path)
7377 .expect("register_beacon must create the beacon file")
7378 .modified()
7379 .unwrap();
7380
7381 let obstruction = sidecar_dir.join(format!(".{pid}.json.tmp"));
7386 std::fs::create_dir(&obstruction).unwrap();
7387
7388 tokio::time::sleep(Duration::from_millis(20)).await;
7389 let over_threshold = Some(khive_storage::tx_registry::OldestSpan {
7390 id: khive_storage::tx_registry::TxId(1),
7391 age: Duration::from_secs(60),
7392 label: None,
7393 origin: khive_storage::tx_registry::TxOrigin::Unscoped,
7394 });
7395 state
7396 .observe(over_threshold.clone(), Duration::from_secs(30))
7397 .await;
7398
7399 assert!(
7400 !sidecar_dir.join(format!("{pid}.json")).exists(),
7401 "heartbeat write must have failed"
7402 );
7403 assert!(
7406 !beacon_path.exists(),
7407 "a failed heartbeat write must remove the beacon — a still-fresh \
7408 beacon with no heartbeat would classify registered-silent \
7409 (before-mtime {before:?})"
7410 );
7411
7412 std::fs::remove_dir(&obstruction).unwrap();
7415 state.observe(over_threshold, Duration::from_secs(30)).await;
7416 assert!(
7417 sidecar_dir.join(format!("{pid}.json")).exists(),
7418 "heartbeat must land once the write path recovers"
7419 );
7420 assert!(
7421 beacon_path.exists(),
7422 "beacon must re-register on the first healthy tick after removal"
7423 );
7424 }
7425
7426 #[tokio::test]
7427 #[serial(khive_walpin_sidecar_env)]
7428 async fn walpin_observe_touches_mtime_without_rewriting_body_when_content_unchanged() {
7429 let dir = tempfile::tempdir().unwrap();
7430 let db_path = dir.path().join("observe_touch.db");
7431 let sidecar_dir = crate::walpin::sidecar_dir_for(&db_path);
7432 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
7433 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
7434
7435 let mut state = WalpinSidecarState::new(
7436 Some(db_path.as_path()),
7437 true,
7438 "session",
7439 Duration::from_millis(500),
7440 )
7441 .expect("sidecar enabled for a file-backed path");
7442 let pid = std::process::id();
7443 let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
7444 let span = khive_storage::tx_registry::OldestSpan {
7445 id: khive_storage::tx_registry::TxId(1),
7446 age: Duration::from_secs(60),
7447 label: None,
7448 origin: khive_storage::tx_registry::TxOrigin::Unscoped,
7449 };
7450
7451 state
7452 .observe(Some(span.clone()), Duration::from_secs(30))
7453 .await;
7454 let body_after_create = std::fs::read(&heartbeat_path).expect("heartbeat written");
7455
7456 let backdated = std::time::SystemTime::now() - Duration::from_secs(120);
7462 std::fs::OpenOptions::new()
7466 .write(true)
7467 .open(&heartbeat_path)
7468 .unwrap()
7469 .set_modified(backdated)
7470 .unwrap();
7471
7472 state.observe(Some(span), Duration::from_secs(30)).await;
7473
7474 let body_after_second_observe =
7475 std::fs::read(&heartbeat_path).expect("heartbeat still present");
7476 assert_eq!(
7477 body_after_create, body_after_second_observe,
7478 "unchanged oldest-span identity/label/attribution/cadence must touch mtime, \
7479 not rewrite the body"
7480 );
7481 let mtime_after = std::fs::metadata(&heartbeat_path)
7482 .unwrap()
7483 .modified()
7484 .unwrap();
7485 assert!(
7486 mtime_after > backdated,
7487 "the touch must advance mtime past the backdated value"
7488 );
7489 }
7490
7491 #[tokio::test]
7492 #[serial(khive_walpin_sidecar_env)]
7493 async fn walpin_observe_recreates_heartbeat_after_it_is_deleted_while_span_still_live() {
7494 let dir = tempfile::tempdir().unwrap();
7495 let db_path = dir.path().join("observe_recreate.db");
7496 let sidecar_dir = crate::walpin::sidecar_dir_for(&db_path);
7497 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
7498 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
7499
7500 let mut state = WalpinSidecarState::new(
7501 Some(db_path.as_path()),
7502 true,
7503 "session",
7504 Duration::from_millis(500),
7505 )
7506 .expect("sidecar enabled for a file-backed path");
7507 let pid = std::process::id();
7508 let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
7509 let span = khive_storage::tx_registry::OldestSpan {
7510 id: khive_storage::tx_registry::TxId(1),
7511 age: Duration::from_secs(60),
7512 label: None,
7513 origin: khive_storage::tx_registry::TxOrigin::Unscoped,
7514 };
7515
7516 state
7517 .observe(Some(span.clone()), Duration::from_secs(30))
7518 .await;
7519 assert!(heartbeat_path.exists(), "heartbeat written on first tick");
7520
7521 std::fs::remove_file(&heartbeat_path).unwrap();
7527 assert!(!heartbeat_path.exists());
7528
7529 state.observe(Some(span), Duration::from_secs(30)).await;
7530
7531 assert!(
7532 heartbeat_path.exists(),
7533 "a touch failure against a deleted heartbeat must recreate it via a full write"
7534 );
7535 let recreated: crate::walpin::WalpinHeartbeat =
7536 serde_json::from_slice(&std::fs::read(&heartbeat_path).unwrap()).unwrap();
7537 assert_eq!(recreated.pid, pid);
7538 assert_eq!(recreated.oldest_tx_age_secs, 60.0);
7539 }
7540
7541 #[tokio::test]
7542 #[serial(tx_registry, khive_walpin_sidecar_env)]
7543 async fn session_sweep_task_writes_and_clears_walpin_heartbeat() {
7544 let dir = tempfile::tempdir().unwrap();
7545 let db_path = dir.path().join("session_sweep.db");
7546 let pool = file_pool(&db_path);
7547 let sidecar_dir =
7548 crate::walpin::sidecar_dir_for(pool.canonical_path().expect("file-backed pool"));
7549 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
7550 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
7551
7552 let cfg = SessionSweepConfig {
7553 interval: Duration::from_millis(10),
7554 tx_warn_secs: Duration::from_millis(20),
7555 tx_max_age_secs: Duration::from_millis(500),
7556 };
7557 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
7558 let handle = tokio::spawn(run_session_sweep_task(
7559 vec![SweepBackend {
7560 pool: Arc::clone(&pool),
7561 is_main: true,
7562 }],
7563 cfg,
7564 shutdown_rx,
7565 ));
7566
7567 let pid = std::process::id();
7574 let beacon = crate::walpin::beacon_path(&sidecar_dir, pid);
7575 let beacon_registered = wait_for(Duration::from_secs(2), || beacon.exists()).await;
7576 assert!(
7577 beacon_registered,
7578 "a quiet process must still register its one-time beacon"
7579 );
7580 assert!(
7581 !sidecar_dir.join(format!("{pid}.json")).exists(),
7582 "a quiet process must not write a walpin heartbeat"
7583 );
7584
7585 let tx_handle =
7586 khive_storage::tx_registry::register(Some("session_sweep_walpin_test".to_string()));
7587 let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
7588 assert!(
7589 wait_for(Duration::from_secs(2), || heartbeat_path.exists()).await,
7590 "expected a walpin heartbeat once the span crossed tx_warn_secs"
7591 );
7592 let body = std::fs::read_to_string(&heartbeat_path).unwrap();
7593 let hb: crate::walpin::WalpinHeartbeat = serde_json::from_str(&body).unwrap();
7594 assert_eq!(hb.pid, pid);
7595 assert_eq!(hb.process_role, "session");
7596 assert_eq!(
7597 hb.oldest_tx_label.as_deref(),
7598 Some("session_sweep_walpin_test")
7599 );
7600 assert_eq!(
7601 hb.attribution_basis.as_deref(),
7602 Some("fallback"),
7603 "an Unscoped span observed only through the main view's fallback \
7604 must carry attribution_basis=\"fallback\", never \"origin\""
7605 );
7606
7607 drop(tx_handle);
7608 assert!(
7609 wait_for(Duration::from_secs(2), || !heartbeat_path.exists()).await,
7610 "heartbeat must be removed once the stale span clears"
7611 );
7612
7613 shutdown_tx.send(()).expect("send shutdown signal");
7614 tokio::time::timeout(Duration::from_secs(1), handle)
7615 .await
7616 .expect("session sweep task should exit within 1s")
7617 .expect("session sweep task panicked");
7618 }
7619
7620 #[tokio::test]
7631 #[serial(tx_registry, khive_walpin_sidecar_env)]
7632 async fn session_sweep_fan_out_scopes_secondary_span_to_secondary_sidecar_only() {
7633 let main_dir = tempfile::tempdir().unwrap();
7634 let secondary_dir = tempfile::tempdir().unwrap();
7635 let main_pool = file_pool(&main_dir.path().join("main.db"));
7636 let secondary_pool = file_pool(&secondary_dir.path().join("secondary.db"));
7637 let main_sidecar =
7638 crate::walpin::sidecar_dir_for(main_pool.canonical_path().expect("file-backed"));
7639 let secondary_sidecar =
7640 crate::walpin::sidecar_dir_for(secondary_pool.canonical_path().expect("file-backed"));
7641 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
7642 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
7643
7644 let cfg = SessionSweepConfig {
7645 interval: Duration::from_millis(10),
7646 tx_warn_secs: Duration::from_millis(20),
7647 tx_max_age_secs: Duration::from_millis(500),
7648 };
7649 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
7650 let handle = tokio::spawn(run_session_sweep_task(
7651 vec![
7652 SweepBackend {
7653 pool: Arc::clone(&main_pool),
7654 is_main: true,
7655 },
7656 SweepBackend {
7657 pool: Arc::clone(&secondary_pool),
7658 is_main: false,
7659 },
7660 ],
7661 cfg,
7662 shutdown_rx,
7663 ));
7664
7665 let pid = std::process::id();
7666 let secondary_heartbeat = secondary_sidecar.join(format!("{pid}.json"));
7667 let main_heartbeat = main_sidecar.join(format!("{pid}.json"));
7668
7669 let tx_handle = khive_storage::tx_registry::register_scoped(
7670 Some("graph_traverse_read".to_string()),
7671 secondary_pool.origin(),
7672 );
7673 assert!(
7674 wait_for(Duration::from_secs(2), || secondary_heartbeat.exists()).await,
7675 "expected a walpin heartbeat in the secondary backend's own sidecar"
7676 );
7677 assert!(
7678 !main_heartbeat.exists(),
7679 "a span scoped to the secondary backend's origin must never produce \
7680 a heartbeat in the main backend's sidecar"
7681 );
7682
7683 let body = std::fs::read_to_string(&secondary_heartbeat).unwrap();
7684 let hb: crate::walpin::WalpinHeartbeat = serde_json::from_str(&body).unwrap();
7685 assert_eq!(hb.oldest_tx_label.as_deref(), Some("graph_traverse_read"));
7686 assert_eq!(
7687 hb.attribution_basis.as_deref(),
7688 Some("origin"),
7689 "a Secondary-view winner is always Database-origin-backed — never fallback"
7690 );
7691
7692 drop(tx_handle);
7693 assert!(
7694 wait_for(Duration::from_secs(2), || !secondary_heartbeat.exists()).await,
7695 "secondary heartbeat must be removed once its span clears"
7696 );
7697 assert!(
7698 !main_heartbeat.exists(),
7699 "the main sidecar must have stayed untouched for the whole tick sequence"
7700 );
7701
7702 shutdown_tx.send(()).expect("send shutdown signal");
7703 tokio::time::timeout(Duration::from_secs(1), handle)
7704 .await
7705 .expect("session sweep task should exit within 1s")
7706 .expect("session sweep task panicked");
7707 }
7708
7709 #[tokio::test]
7718 #[serial(tx_registry, checkpoint_skip_metrics, khive_walpin_sidecar_env)]
7719 async fn checkpoint_task_ignores_span_registered_against_other_backend_origin_and_unscoped() {
7720 let dir_a = tempfile::tempdir().unwrap();
7721 let dir_b = tempfile::tempdir().unwrap();
7722 let pool_a = file_pool(&dir_a.path().join("backend_a.db"));
7723 let pool_b = file_pool(&dir_b.path().join("backend_b.db"));
7726 let sidecar_a =
7727 crate::walpin::sidecar_dir_for(pool_a.canonical_path().expect("file-backed"));
7728 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
7729 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
7730
7731 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7732 let subscriber = CaptureSubscriber {
7733 events: std::sync::Arc::clone(&buffer),
7734 };
7735 let _tracing_guard = tracing::subscriber::set_default(subscriber);
7736
7737 let _b_origin_handle = khive_storage::tx_registry::register_scoped(
7738 Some("b_origin_span_ignored_by_a".to_string()),
7739 pool_b.origin(),
7740 );
7741 let _unscoped_handle = khive_storage::tx_registry::register(Some(
7742 "unscoped_span_ignored_by_secondary".to_string(),
7743 ));
7744
7745 let cfg = CheckpointConfig {
7746 interval: Duration::from_millis(10),
7747 tx_warn_secs: Duration::from_millis(1),
7748 tx_max_age_secs: Duration::from_millis(1),
7749 ..CheckpointConfig::default()
7750 };
7751 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
7752 let handle = tokio::spawn(run_checkpoint_task(
7753 pool_a,
7754 cfg,
7755 None,
7756 shutdown_rx,
7757 false, ));
7759
7760 tokio::time::sleep(Duration::from_millis(60)).await;
7765 shutdown_tx.send(()).expect("send shutdown signal");
7766 tokio::time::timeout(Duration::from_secs(1), handle)
7767 .await
7768 .expect("checkpoint task should exit within 1s")
7769 .expect("checkpoint task panicked");
7770
7771 let events = buffer.lock().unwrap();
7772 assert!(
7773 events.iter().all(|e| {
7774 e.tx_label.as_deref() != Some("b_origin_span_ignored_by_a")
7775 && e.tx_label.as_deref() != Some("unscoped_span_ignored_by_secondary")
7776 }),
7777 "backend A's Secondary filter must never emit an age alert naming a span \
7778 registered against a different backend's origin or an Unscoped span, got: \
7779 {events:?}"
7780 );
7781 assert!(
7782 !sidecar_a
7783 .join(format!("{}.json", std::process::id()))
7784 .exists(),
7785 "backend A's own sidecar must never gain a heartbeat from a span it does not own"
7786 );
7787 }
7788
7789 #[tokio::test]
7796 #[serial(tx_registry, checkpoint_skip_metrics, khive_walpin_sidecar_env)]
7797 async fn checkpoint_task_detects_and_enumerates_secondary_backend_stall() {
7798 let dir = tempfile::tempdir().unwrap();
7799 let pool = file_pool(&dir.path().join("secondary_stall.db"));
7800 let sidecar_dir =
7801 crate::walpin::sidecar_dir_for(pool.canonical_path().expect("file-backed"));
7802 let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
7803 std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
7804
7805 let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7806 let subscriber = CaptureSubscriber {
7807 events: std::sync::Arc::clone(&buffer),
7808 };
7809 let _tracing_guard = tracing::subscriber::set_default(subscriber);
7810
7811 let tx_handle = khive_storage::tx_registry::register_scoped(
7812 Some("secondary_stall_test".to_string()),
7813 pool.origin(),
7814 );
7815
7816 let cfg = CheckpointConfig {
7817 interval: Duration::from_millis(10),
7818 tx_warn_secs: Duration::from_millis(5),
7819 tx_max_age_secs: Duration::from_millis(500),
7820 ..CheckpointConfig::default()
7821 };
7822 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
7823 let pid = std::process::id();
7824 let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
7825 let handle = tokio::spawn(run_checkpoint_task(
7826 pool,
7827 cfg,
7828 None,
7829 shutdown_rx,
7830 false, ));
7832
7833 assert!(
7834 wait_for(Duration::from_secs(2), || heartbeat_path.exists()).await,
7835 "expected a walpin heartbeat once the secondary backend's own span crossed \
7836 tx_warn_secs"
7837 );
7838 let body = std::fs::read_to_string(&heartbeat_path).unwrap();
7839 let hb: crate::walpin::WalpinHeartbeat = serde_json::from_str(&body).unwrap();
7840 assert_eq!(hb.oldest_tx_label.as_deref(), Some("secondary_stall_test"));
7841 assert_eq!(
7842 hb.attribution_basis.as_deref(),
7843 Some("origin"),
7844 "a Secondary-view winner is always Database-origin-backed — never fallback"
7845 );
7846 assert!(
7847 hb.oldest_tx_age_secs > 0.0,
7848 "the heartbeat must reflect a nonzero stale age for the secondary backend's own \
7849 span, got {hb:?}"
7850 );
7851
7852 shutdown_tx.send(()).expect("send shutdown signal");
7853 tokio::time::timeout(Duration::from_secs(1), handle)
7854 .await
7855 .expect("checkpoint task should exit within 1s")
7856 .expect("checkpoint task panicked");
7857
7858 drop(tx_handle);
7859
7860 let events = buffer.lock().unwrap();
7861 assert!(
7862 events.iter().any(|e| {
7863 e.tx_label.as_deref() == Some("secondary_stall_test")
7864 && e.message
7865 .as_deref()
7866 .is_some_and(|m| m.contains("ADR-091 Plank 1"))
7867 }),
7868 "expected the secondary backend's own checkpoint task to emit a Plank 1 age alert \
7869 for its own stalled span, got: {events:?}"
7870 );
7871 }
7872
7873 #[test]
7874 fn wal_pin_depth_arithmetic_against_real_connection() {
7875 let dir = tempfile::tempdir().unwrap();
7876 let path = dir.path().join("pin_depth.db");
7877 let pool = file_pool(&path);
7878 let writer = pool.try_writer().expect("acquire writer");
7879 let conn = writer.conn();
7880
7881 conn.execute_batch("CREATE TABLE t (v INTEGER)").unwrap();
7882 conn.execute_batch("INSERT INTO t (v) VALUES (1)").unwrap();
7883
7884 let (log, checkpointed) =
7885 query_wal_pin_depth(conn).expect("PRAGMA wal_checkpoint(PASSIVE) must succeed");
7886 assert!(
7890 log >= checkpointed,
7891 "checkpointed frames cannot exceed log frames"
7892 );
7893 assert_eq!(
7894 log - checkpointed,
7895 0,
7896 "an unpinned WAL must fully checkpoint under PASSIVE"
7897 );
7898 }
7899
7900 #[test]
7901 fn wal_pin_depth_arithmetic_on_in_memory_pool_errors_cleanly() {
7902 let cfg = PoolConfig {
7906 path: None,
7907 ..PoolConfig::default()
7908 };
7909 let pool = ConnectionPool::new(cfg).expect("in-memory pool");
7910 let writer = pool.try_writer().expect("acquire writer");
7911 let _ = query_wal_pin_depth(writer.conn());
7914 }
7915
7916 #[cfg(unix)]
7920 #[test]
7921 fn routine_wal_backend_key_preserves_non_utf8_path_bytes() {
7922 use std::ffi::OsString;
7923 use std::os::unix::ffi::OsStringExt;
7924
7925 let path_a = PathBuf::from(OsString::from_vec(b"/tmp/khive-wal-\x80.db".to_vec()));
7926 let path_b = PathBuf::from(OsString::from_vec(b"/tmp/khive-wal-\x81.db".to_vec()));
7927 assert_eq!(
7928 path_a.display().to_string(),
7929 path_b.display().to_string(),
7930 "fixture must reproduce the lossy display-label collision"
7931 );
7932 assert_ne!(
7933 checkpoint_db_key_from_path(Some(&path_a)),
7934 checkpoint_db_key_from_path(Some(&path_b)),
7935 "backend keys must retain the canonical path's exact OS bytes"
7936 );
7937 }
7938
7939 #[test]
7944 #[serial(checkpoint_skip_metrics)]
7945 fn routine_checkpoint_records_one_pass_logical_and_physical_wal_sample() {
7946 let dir = tempfile::tempdir().unwrap();
7947 let path = dir.path().join("routine_wal_sample.db");
7948 let pool = file_pool(&path);
7949
7950 {
7951 let writer = pool.try_writer().expect("writer");
7952 writer
7953 .conn()
7954 .execute_batch(
7955 "PRAGMA wal_autocheckpoint=0; \
7956 CREATE TABLE t (id INTEGER PRIMARY KEY, payload TEXT); \
7957 INSERT INTO t VALUES (0, 'seed');",
7958 )
7959 .unwrap();
7960 }
7961
7962 let reader = rusqlite::Connection::open_with_flags(
7963 &path,
7964 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY,
7965 )
7966 .unwrap();
7967 reader.execute_batch("BEGIN").unwrap();
7968 let _: i64 = reader
7969 .query_row("SELECT COUNT(*) FROM t", [], |row| row.get(0))
7970 .unwrap();
7971
7972 {
7973 let writer = pool.try_writer().expect("writer");
7974 writer.conn().execute_batch("BEGIN IMMEDIATE").unwrap();
7975 for id in 1..=256_i64 {
7976 writer
7977 .conn()
7978 .execute("INSERT INTO t VALUES (?1, printf('%.*c', 2048, 'x'))", [id])
7979 .unwrap();
7980 }
7981 writer.conn().execute_batch("COMMIT").unwrap();
7982 }
7983
7984 let checkpoint_conn = pool.open_standalone_writer().unwrap();
7985 let pragma_calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
7986 let pragma_calls_from_hook = Arc::clone(&pragma_calls);
7987 checkpoint_conn
7988 .authorizer(Some(move |context: rusqlite::hooks::AuthContext<'_>| {
7989 if matches!(
7990 context.action,
7991 AuthAction::Pragma { pragma_name, .. }
7992 if pragma_name.eq_ignore_ascii_case("wal_checkpoint")
7993 ) {
7994 pragma_calls_from_hook.fetch_add(1, Ordering::SeqCst);
7995 }
7996 Authorization::Allow
7997 }))
7998 .unwrap();
7999
8000 checkpoint_once(
8001 &pool,
8002 &checkpoint_conn,
8003 &CheckpointConfig::default(),
8004 &mut TruncateState::default(),
8005 )
8006 .unwrap();
8007 checkpoint_conn
8008 .authorizer(None::<fn(rusqlite::hooks::AuthContext<'_>) -> Authorization>)
8009 .unwrap();
8010
8011 assert_eq!(
8012 pragma_calls.load(Ordering::SeqCst),
8013 1,
8014 "one routine tick must issue exactly one PASSIVE checkpoint"
8015 );
8016 let pinned = routine_wal_observation(&pool).expect("routine sample");
8017 assert_eq!(
8018 pinned.busy, 0,
8019 "a pinned reader is not checkpoint-lock contention"
8020 );
8021 let first_timing = checkpoint_timing(&pool);
8022 assert_eq!(first_timing.ticks, 1);
8023 assert_eq!(
8024 first_timing.busy_ticks, 0,
8025 "pending frames must not count as busy"
8026 );
8027 assert!(pinned.log_frames > 0, "the test must create WAL frames");
8028 assert!(
8029 pinned.pending_frames > 0,
8030 "the old reader must leave a logical backlog: {pinned:?}"
8031 );
8032 assert_eq!(
8033 pinned.pending_frames,
8034 pinned.log_frames.saturating_sub(pinned.checkpointed_frames)
8035 );
8036 assert!(
8037 pinned.physical_wal_bytes.is_some_and(|bytes| bytes > 0),
8038 "the physical sidecar high-water must be reported separately: {pinned:?}"
8039 );
8040
8041 reader.execute_batch("COMMIT").unwrap();
8042 checkpoint_once(
8043 &pool,
8044 &checkpoint_conn,
8045 &CheckpointConfig::default(),
8046 &mut TruncateState::default(),
8047 )
8048 .unwrap();
8049 let drained = routine_wal_observation(&pool).expect("drained routine sample");
8050 let drained_timing = checkpoint_timing(&pool);
8051 assert_eq!(drained_timing.ticks, first_timing.ticks + 1);
8052 assert!(drained_timing.elapsed_us_sum >= first_timing.elapsed_us_sum);
8053 assert!(drained_timing.elapsed_us_max >= first_timing.elapsed_us_max);
8054 assert_eq!(drained_timing.busy_ticks, 0);
8055 assert_eq!(drained.pending_frames, 0, "unpinned PASSIVE must drain");
8056 assert!(
8057 drained.physical_wal_bytes.is_some_and(|bytes| bytes > 0),
8058 "PASSIVE may reuse rather than shrink the physical WAL; the two gauges must remain \
8059 independently visible: {drained:?}"
8060 );
8061 }
8062
8063 #[test]
8064 #[serial(checkpoint_skip_metrics)]
8065 fn routine_checkpoint_timing_records_real_call_and_debug_fields() {
8066 let dir = tempfile::tempdir().unwrap();
8067 let pool = file_pool(&dir.path().join("timed_tick.db"));
8068 let conn = checkpoint_conn(&pool);
8069 conn.execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
8070 .unwrap();
8071 let (entered_tx, entered_rx) = std::sync::mpsc::sync_channel(0);
8072 let (release_tx, release_rx) = std::sync::mpsc::sync_channel(0);
8073 let release_rx = Mutex::new(release_rx);
8074 conn.authorizer(Some(move |context: rusqlite::hooks::AuthContext<'_>| {
8075 if matches!(context.action, AuthAction::Pragma { pragma_name, .. }
8076 if pragma_name.eq_ignore_ascii_case("wal_checkpoint"))
8077 {
8078 entered_tx.send(()).unwrap();
8079 release_rx
8080 .lock()
8081 .unwrap()
8082 .recv_timeout(Duration::from_secs(5))
8083 .unwrap();
8084 }
8085 Authorization::Allow
8086 }))
8087 .unwrap();
8088 let tick_pool = Arc::clone(&pool);
8089 let tick = std::thread::spawn(move || {
8090 let mut result = None;
8091 let events = capture(|| {
8092 result = Some(checkpoint_once(
8093 &tick_pool,
8094 &conn,
8095 &CheckpointConfig {
8096 truncate_high_water_pages: u64::MAX,
8097 ..CheckpointConfig::default()
8098 },
8099 &mut TruncateState::default(),
8100 ));
8101 });
8102 (result.unwrap(), events)
8103 });
8104 entered_rx.recv_timeout(Duration::from_secs(5)).unwrap();
8105 let during = checkpoint_timing(&pool);
8106 release_tx.send(()).unwrap();
8107 let (result, events) = tick.join().unwrap();
8108 result.unwrap();
8109 assert_eq!(
8110 during.ticks, 0,
8111 "in-flight call must not publish partial counters"
8112 );
8113 let timing = checkpoint_timing(&pool);
8114 assert_eq!(
8115 timing.ticks, 1,
8116 "one actual PASSIVE call must advance the count"
8117 );
8118 assert!(
8119 timing.elapsed_us_sum > 0,
8120 "channel-held checkpoint call must record elapsed time"
8121 );
8122 assert_eq!(timing.elapsed_us_max, timing.elapsed_us_sum);
8123 assert_eq!(timing.busy_ticks, 0);
8124 assert_eq!(timing.error_ticks, 0);
8125 let issued: Vec<_> = events
8126 .iter()
8127 .filter(|event| event.message.as_deref() == Some("WAL checkpoint issued"))
8128 .collect();
8129 assert_eq!(issued.len(), 1);
8130 assert_eq!(issued[0].elapsed_us, Some(timing.elapsed_us_sum));
8131 assert_eq!(issued[0].busy, Some(0));
8132 }
8133
8134 #[test]
8135 #[serial(checkpoint_skip_metrics)]
8136 fn routine_checkpoint_timing_counts_sqlite_busy_from_competing_checkpoint() {
8137 struct BusyGate {
8138 entered: std::sync::mpsc::SyncSender<()>,
8139 release: std::sync::mpsc::Receiver<()>,
8140 }
8141 static BUSY_GATE: Mutex<Option<BusyGate>> = Mutex::new(None);
8142 fn hold_checkpoint_lock(_attempt: i32) -> bool {
8143 let gate = BUSY_GATE
8144 .lock()
8145 .unwrap()
8146 .take()
8147 .expect("armed busy handler");
8148 gate.entered.send(()).unwrap();
8149 gate.release.recv_timeout(Duration::from_secs(5)).unwrap();
8150 false
8151 }
8152
8153 let dir = tempfile::tempdir().unwrap();
8154 let path = dir.path().join("busy_checkpoint.db");
8155 let pool = file_pool(&path);
8156 let conn = checkpoint_conn(&pool);
8157 conn.execute_batch(
8158 "PRAGMA wal_autocheckpoint=0; CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);",
8159 )
8160 .unwrap();
8161 let reader = rusqlite::Connection::open_with_flags(
8162 &path,
8163 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY,
8164 )
8165 .unwrap();
8166 reader.execute_batch("BEGIN").unwrap();
8167 let _: i64 = reader
8168 .query_row("SELECT COUNT(*) FROM t", [], |row| row.get(0))
8169 .unwrap();
8170 conn.execute_batch("INSERT INTO t VALUES (2);").unwrap();
8171
8172 let (entered_tx, entered_rx) = std::sync::mpsc::sync_channel(0);
8173 let (release_tx, release_rx) = std::sync::mpsc::sync_channel(0);
8174 *BUSY_GATE.lock().unwrap() = Some(BusyGate {
8175 entered: entered_tx,
8176 release: release_rx,
8177 });
8178 let competing = std::thread::spawn(move || {
8179 let checkpoint = rusqlite::Connection::open(path).unwrap();
8180 checkpoint.busy_handler(Some(hold_checkpoint_lock)).unwrap();
8181 checkpoint.query_row("PRAGMA wal_checkpoint(FULL)", [], |row| {
8182 row.get::<_, i64>(0)
8183 })
8184 });
8185 entered_rx
8186 .recv_timeout(Duration::from_secs(5))
8187 .expect("FULL checkpoint holds CKPT lock while waiting on reader");
8188 let mut result = None;
8189 let events = capture(|| {
8190 result = Some(checkpoint_once(
8191 &pool,
8192 &conn,
8193 &CheckpointConfig {
8194 truncate_high_water_pages: u64::MAX,
8195 ..CheckpointConfig::default()
8196 },
8197 &mut TruncateState::default(),
8198 ));
8199 });
8200 release_tx.send(()).unwrap();
8201 let competing_busy = competing.join().unwrap().unwrap();
8202 reader.execute_batch("COMMIT").unwrap();
8203 result.unwrap().unwrap();
8204 assert_eq!(competing_busy, 1);
8205 assert_eq!(
8206 routine_wal_observation(&pool).unwrap().busy,
8207 1,
8208 "fixture must reach SQLite busy"
8209 );
8210 let timing = checkpoint_timing(&pool);
8211 assert_eq!(timing.ticks, 1);
8212 assert_eq!(
8213 timing.busy_ticks, 1,
8214 "SQLite busy result must increment busy ticks"
8215 );
8216 assert_eq!(timing.error_ticks, 0);
8217 let issued = events
8218 .iter()
8219 .find(|event| event.message.as_deref() == Some("WAL checkpoint issued"))
8220 .expect("existing tick record");
8221 assert_eq!(issued.busy, Some(1), "tick record must retain SQLite busy");
8222 assert_eq!(issued.elapsed_us, Some(timing.elapsed_us_sum));
8223 }
8224
8225 #[test]
8226 #[serial(checkpoint_skip_metrics)]
8227 fn routine_checkpoint_timing_counts_errors_and_excludes_post_truncate_probes() {
8228 let dir = tempfile::tempdir().unwrap();
8229 let pool = file_pool(&dir.path().join("failed_tick.db"));
8230 let conn = checkpoint_conn(&pool);
8231 conn.authorizer(Some(|context: rusqlite::hooks::AuthContext<'_>| {
8232 if matches!(context.action, AuthAction::Pragma { pragma_name, .. }
8233 if pragma_name.eq_ignore_ascii_case("wal_checkpoint"))
8234 {
8235 Authorization::Deny
8236 } else {
8237 Authorization::Allow
8238 }
8239 }))
8240 .unwrap();
8241 let result = checkpoint_once(
8242 &pool,
8243 &conn,
8244 &CheckpointConfig::default(),
8245 &mut TruncateState::default(),
8246 );
8247 conn.authorizer(None::<fn(rusqlite::hooks::AuthContext<'_>) -> Authorization>)
8248 .unwrap();
8249 assert!(result.is_err());
8250 let timing = checkpoint_timing(&pool);
8251 assert_eq!(timing.ticks, 1);
8252 assert_eq!(timing.error_ticks, 1);
8253 assert_eq!(timing.busy_ticks, 0);
8254 query_wal_pages(&conn);
8255 assert_eq!(
8256 checkpoint_timing(&pool),
8257 timing,
8258 "post-TRUNCATE observation is not a routine tick"
8259 );
8260 }
8261
8262 #[test]
8263 fn checkpoint_timing_accumulates_per_store_and_saturates() {
8264 let dir = tempfile::tempdir().unwrap();
8265 let a = file_pool(&dir.path().join("a.db"));
8266 let b = file_pool(&dir.path().join("b.db"));
8267 record_checkpoint_timing(&a, 17, Some(0));
8268 record_checkpoint_timing(&a, 31, Some(1));
8269 record_checkpoint_timing(&a, 7, None);
8270 record_checkpoint_timing(&b, 3, Some(0));
8271 assert_eq!(
8272 checkpoint_timing(&a),
8273 CheckpointTiming {
8274 ticks: 3,
8275 elapsed_us_sum: 55,
8276 elapsed_us_max: 31,
8277 busy_ticks: 1,
8278 error_ticks: 1,
8279 }
8280 );
8281 assert_eq!(
8282 checkpoint_timing(&b),
8283 CheckpointTiming {
8284 ticks: 1,
8285 elapsed_us_sum: 3,
8286 elapsed_us_max: 3,
8287 busy_ticks: 0,
8288 error_ticks: 0,
8289 }
8290 );
8291 checkpoint_timings().lock().unwrap().insert(
8292 checkpoint_db_key(&a),
8293 CheckpointTiming {
8294 ticks: u64::MAX,
8295 elapsed_us_sum: u64::MAX,
8296 elapsed_us_max: 31,
8297 busy_ticks: u64::MAX,
8298 error_ticks: u64::MAX,
8299 },
8300 );
8301 record_checkpoint_timing(&a, 1, Some(1));
8302 assert_eq!(
8303 checkpoint_timing(&a).elapsed_us_sum,
8304 u64::MAX,
8305 "elapsed sum must saturate on the first overflowing addition"
8306 );
8307 record_checkpoint_timing(&a, u64::MAX, None);
8308 assert_eq!(
8309 checkpoint_timing(&a),
8310 CheckpointTiming {
8311 ticks: u64::MAX,
8312 elapsed_us_sum: u64::MAX,
8313 elapsed_us_max: u64::MAX,
8314 busy_ticks: u64::MAX,
8315 error_ticks: u64::MAX,
8316 }
8317 );
8318 }
8319}