khive_db/checkpoint.rs
1//! Periodic WAL checkpoint task for the connection pool (ADR-091; dedicated
2//! checkpoint connection amendment, see below).
3//!
4//! Issues `PRAGMA wal_checkpoint(PASSIVE)` on every tick — non-blocking, never
5//! waits for readers. A rare, separately-gated escalation may additionally run
6//! `PRAGMA wal_checkpoint(TRUNCATE)` once WAL pressure crosses
7//! `truncate_high_water_pages` and `truncate_min_interval` has elapsed since
8//! the last attempt (Plank 2); both run on the task's own dedicated
9//! standalone connection (`CheckpointConnection`), opened once at task
10//! startup and reused for every tick — `checkpoint_once` never checks out the
11//! pool's writer mutex at all, so a concurrent `pool.writer()` checkout can
12//! never queue behind a checkpoint tick's ADMISSION. That guarantee is
13//! admission-only: PASSIVE takes SQLite's CKPT lock, not the WRITE lock, so
14//! it never blocks writers at the SQLite level either — but TRUNCATE
15//! additionally acquires SQLite's writer lock and can still block a
16//! concurrent write transaction, on any connection, for up to
17//! `truncate_busy_timeout` while it waits on a pinning reader, exactly as
18//! before this connection split.
19//!
20//! If the dedicated connection is unavailable (never opened yet, or dropped
21//! after a prior tick's connection-level pragma failure), the tick reports
22//! `CheckpointTick::Skipped` and the next tick lazily reopens it. A busy or
23//! inconsistent PASSIVE result also skips pressure decisions without replacing
24//! the last valid WAL sample; a busy pool writer does not cause either skip.
25//!
26//! `warn_pages` / `high_water_pages` WARNs fire at most once per below→above
27//! crossing; a skipped tick leaves crossing state unchanged. An age-based
28//! background sweep (Plank 1) additionally checks the oldest span in
29//! `khive_storage::tx_registry` against `tx_warn_secs`/`tx_max_age_secs` on
30//! every tick (Skipped or Observed) and escalates to `warn!`/`error!` on each
31//! below→above crossing — visibility only, nothing here force-closes a stale
32//! span.
33//!
34//! See crates/khive-db/docs/api/checkpoint.md#module-overview-adr-091-planks-012
35//! for full ADR-091 Plank 0/1/2 design rationale (why TRUNCATE is excluded
36//! from ordinary ticks, the dedicated-connection invariant, and why Plank 1
37//! is a sweep rather than the ADR's originally-described per-statement guard).
38//!
39//! The same long-lived standalone connection also owns a separate, five-minute
40//! FTS5 maintenance cadence. One due call gives one index at most 500 pages of
41//! incremental merge work and uses a zero busy timeout, so this best-effort
42//! derived-index maintenance cannot queue behind application writes.
43
44use std::collections::{BTreeMap, 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
52mod off_worker;
53#[cfg(test)]
54mod off_worker_tests;
55
56// ── metrics read-surface (load/perf harness) ─────────────────────────────
57// Read-only process-wide gauges (never reset outside #[cfg(test)]). See
58// crates/khive-db/docs/api/checkpoint.md#metrics-read-surface-loadperf-harness
59
60/// Last-observed WAL page count (the routine PASSIVE row's `log` value, or a
61/// rare post-TRUNCATE observation from `maybe_truncate`).
62/// `u64::MAX` is the "never observed" sentinel — no checkpoint tick has run
63/// yet in this process — distinct from a genuine zero-page WAL.
64static LAST_WAL_PAGES: AtomicU64 = AtomicU64::new(u64::MAX);
65
66/// Count of TRUNCATE attempts (`maybe_truncate`'s pragma actually invoked,
67/// win or lose) across this process's lifetime.
68static TRUNCATE_ATTEMPTS: AtomicU64 = AtomicU64::new(0);
69
70/// Current consecutive-failure count, mirrored from the caller-owned
71/// `TruncateState::consecutive_failures` field into a process-readable
72/// gauge every time `note_truncate_outcome` runs.
73static TRUNCATE_CONSECUTIVE_FAILURES: AtomicU64 = AtomicU64::new(0);
74
75/// Count of checkpoint ticks without a usable WAL frame observation because
76/// the dedicated connection was unavailable or SQLite returned a busy or
77/// inconsistent PASSIVE row.
78/// Never reset outside `#[cfg(test)]`.
79static CHECKPOINT_SKIPPED_TICKS: AtomicU64 = AtomicU64::new(0);
80
81/// Current run-length of consecutive skipped ticks. Reset to 0 the next time
82/// a tick has a valid WAL frame observation, so a
83/// sustained skip streak is visible even between two successful
84/// observations.
85static CHECKPOINT_CONSECUTIVE_SKIPS: AtomicU64 = AtomicU64::new(0);
86
87/// WAL page count as of the most recent *observed* tick, snapshotted at the
88/// moment a skip occurs. `u64::MAX` is the "no skip has recorded a snapshot
89/// yet" sentinel, mirroring `LAST_WAL_PAGES`.
90static CHECKPOINT_LAST_SKIP_WAL_PAGES: AtomicU64 = AtomicU64::new(u64::MAX);
91
92/// Elevated checkpoint observations aggregated in memory instead of written
93/// as one primary-store lifecycle row per tick (#1838).
94static CHECKPOINT_PRESSURE_ELEVATED_TICKS: AtomicU64 = AtomicU64::new(0);
95
96/// Below-to-above `warn_pages` transitions observed by checkpoint tasks.
97static CHECKPOINT_PRESSURE_EPISODES_STARTED: AtomicU64 = AtomicU64::new(0);
98
99/// Above-to-below `warn_pages` transitions observed by checkpoint tasks.
100static CHECKPOINT_PRESSURE_EPISODES_RECOVERED: AtomicU64 = AtomicU64::new(0);
101
102/// Primary-store append calls actually made by checkpoint lifecycle workers.
103static CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS: AtomicU64 = AtomicU64::new(0);
104
105/// Checkpoint lifecycle append calls that returned a storage error.
106static CHECKPOINT_LIFECYCLE_APPEND_FAILURES: AtomicU64 = AtomicU64::new(0);
107
108/// Lifecycle transitions rejected before append because the bounded handoff
109/// was full, closed, or could not serialize the payload.
110static CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS: AtomicU64 = AtomicU64::new(0);
111
112/// Count of cached-reader explicit read transactions rolled back on reuse
113/// for exceeding `read_tx_max_age` (#1846), across this process's lifetime.
114/// Unlike the Plank 1 sweep above, this is reclamation, not just visibility:
115/// each count here is a WAL snapshot that was actually released rather than
116/// merely logged as stale. See `sql_bridge.rs::execute_standalone_read`.
117static READ_TX_MAX_AGE_EVICTIONS: AtomicU64 = AtomicU64::new(0);
118
119/// One backend-scoped observation produced by the periodic checkpoint task's
120/// own PASSIVE pass. Logical frame counts and the physical `-wal` allocation
121/// are intentionally separate: SQLite may retain/reuse the sidecar after the
122/// logical backlog drains (#1849).
123#[derive(Debug, Clone, PartialEq, Eq)]
124pub struct RoutineWalObservation {
125 pub busy: i64,
126 pub log_frames: u64,
127 pub checkpointed_frames: u64,
128 pub pending_frames: u64,
129 pub physical_wal_bytes: Option<u64>,
130 pub observed_at_unix_ms: u64,
131}
132
133/// One backend's current checkpoint run: informative checkpoint results that
134/// stopped at the same frame while frames remained to backfill, allowing
135/// bounded neutral busy results between them. Reported by `db_diagnostics` as
136/// `oldest_pinned_frame_run`.
137#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
138pub struct CheckpointRun {
139 /// The checkpointed frame every result in the run stopped at.
140 pub frame: i64,
141 /// Unix time in milliseconds of the run's first checkpoint result.
142 pub first_observed_at_unix_ms: u64,
143}
144
145#[derive(Debug, Clone, Copy, PartialEq, Eq)]
146pub(crate) enum CheckpointRunStatus {
147 NoTask,
148 NoObservation,
149 Observed(CheckpointRun),
150}
151
152#[derive(Debug, Clone, Copy, PartialEq, Eq)]
153struct CheckpointRunEntry {
154 run: CheckpointRun,
155 first_observed_at: Instant,
156 last_log_frames: i64,
157 last_informative_at: Instant,
158 busy_since_last_informative: bool,
159}
160
161#[derive(Debug, Default)]
162struct CheckpointRunState {
163 active_tasks: usize,
164 checkpoint_interval_ms: u64,
165 owner_intervals_ms: BTreeMap<u64, usize>,
166 entry: Option<CheckpointRunEntry>,
167}
168
169static CHECKPOINT_RUNS: OnceLock<Mutex<HashMap<Option<PathBuf>, CheckpointRunState>>> =
170 OnceLock::new();
171
172fn checkpoint_runs() -> &'static Mutex<HashMap<Option<PathBuf>, CheckpointRunState>> {
173 CHECKPOINT_RUNS.get_or_init(|| Mutex::new(HashMap::new()))
174}
175
176pub(crate) struct CheckpointRunTaskGuard {
177 key: Option<PathBuf>,
178 interval_ms: u64,
179}
180
181impl CheckpointRunTaskGuard {
182 pub(crate) fn start(pool: &ConnectionPool, interval: Duration) -> Self {
183 let key = checkpoint_db_key(pool);
184 let interval_ms = interval.as_millis().min(u128::from(u64::MAX)) as u64;
185 let interval_ms = interval_ms.max(1);
186 let mut runs = checkpoint_runs()
187 .lock()
188 .unwrap_or_else(std::sync::PoisonError::into_inner);
189 let state = runs.entry(key.clone()).or_default();
190 if state.active_tasks == 0 {
191 state.entry = None;
192 }
193 *state.owner_intervals_ms.entry(interval_ms).or_default() += 1;
194 state.active_tasks = state.active_tasks.saturating_add(1);
195 state.checkpoint_interval_ms = *state
196 .owner_intervals_ms
197 .first_key_value()
198 .expect("active owner has an interval")
199 .0;
200 Self { key, interval_ms }
201 }
202}
203
204impl Drop for CheckpointRunTaskGuard {
205 fn drop(&mut self) {
206 let mut runs = checkpoint_runs()
207 .lock()
208 .unwrap_or_else(std::sync::PoisonError::into_inner);
209 let Some(state) = runs.get_mut(&self.key) else {
210 return;
211 };
212 let Some(count) = state.owner_intervals_ms.get_mut(&self.interval_ms) else {
213 return;
214 };
215 *count -= 1;
216 if *count == 0 {
217 state.owner_intervals_ms.remove(&self.interval_ms);
218 }
219 state.active_tasks = state.active_tasks.saturating_sub(1);
220 if state.active_tasks == 0 {
221 runs.remove(&self.key);
222 } else {
223 state.checkpoint_interval_ms = *state
224 .owner_intervals_ms
225 .first_key_value()
226 .expect("surviving owner has an interval")
227 .0;
228 }
229 }
230}
231
232fn advance_checkpoint_run_at(
233 entry: &mut Option<CheckpointRunEntry>,
234 observation: Option<(i64, i64, i64)>,
235 observed_at_unix_ms: u64,
236 observed_at: Instant,
237 checkpoint_interval_ms: u64,
238) {
239 let Some((busy, log_frames, checkpointed_frames)) = observation else {
240 *entry = None;
241 return;
242 };
243 // A busy row has no reliable reading of the ceiling, even if SQLite fills
244 // in the other columns. It cannot advance or end the current run.
245 if busy != 0 {
246 if let Some(current) = entry {
247 current.busy_since_last_informative = true;
248 }
249 return;
250 }
251 if log_frames < 0 || checkpointed_frames < 0 || checkpointed_frames >= log_frames {
252 *entry = None;
253 return;
254 }
255
256 match entry {
257 Some(current)
258 if current.run.frame == checkpointed_frames
259 && log_frames >= current.last_log_frames
260 && (!current.busy_since_last_informative
261 || observed_at
262 .checked_duration_since(current.last_informative_at)
263 .is_some_and(|elapsed| {
264 elapsed
265 <= Duration::from_millis(checkpoint_interval_ms.saturating_mul(2))
266 })) =>
267 {
268 current.last_log_frames = log_frames;
269 current.last_informative_at = observed_at;
270 current.busy_since_last_informative = false;
271 }
272 _ => {
273 *entry = Some(CheckpointRunEntry {
274 run: CheckpointRun {
275 frame: checkpointed_frames,
276 first_observed_at_unix_ms: observed_at_unix_ms,
277 },
278 first_observed_at: observed_at,
279 last_log_frames: log_frames,
280 last_informative_at: observed_at,
281 busy_since_last_informative: false,
282 });
283 }
284 }
285}
286
287#[cfg(test)]
288fn advance_checkpoint_run(
289 entry: &mut Option<CheckpointRunEntry>,
290 observation: Option<(i64, i64, i64)>,
291 observed_at_unix_ms: u64,
292 checkpoint_interval_ms: u64,
293) {
294 static TEST_ORIGIN: OnceLock<Instant> = OnceLock::new();
295 let observed_at = TEST_ORIGIN
296 .get_or_init(Instant::now)
297 .checked_add(Duration::from_millis(observed_at_unix_ms))
298 .expect("test monotonic timestamp");
299 advance_checkpoint_run_at(
300 entry,
301 observation,
302 observed_at_unix_ms,
303 observed_at,
304 checkpoint_interval_ms,
305 );
306}
307
308pub(crate) fn record_checkpoint_run_result(
309 pool: &ConnectionPool,
310 observation: Option<(i64, i64, i64)>,
311) -> CheckpointRunStatus {
312 let mut runs = checkpoint_runs()
313 .lock()
314 .unwrap_or_else(std::sync::PoisonError::into_inner);
315 let Some(state) = runs.get_mut(&checkpoint_db_key(pool)) else {
316 return CheckpointRunStatus::NoTask;
317 };
318 if state.active_tasks == 0 {
319 return CheckpointRunStatus::NoTask;
320 }
321 let checkpoint_interval_ms = state.checkpoint_interval_ms;
322 advance_checkpoint_run_at(
323 &mut state.entry,
324 observation,
325 observed_at_unix_ms(),
326 Instant::now(),
327 checkpoint_interval_ms,
328 );
329 state
330 .entry
331 .map_or(CheckpointRunStatus::NoObservation, |entry| {
332 CheckpointRunStatus::Observed(entry.run)
333 })
334}
335
336#[cfg(test)]
337pub(crate) fn checkpoint_run_status(pool: &ConnectionPool) -> CheckpointRunStatus {
338 checkpoint_run_snapshot(pool).0
339}
340
341/// A run and its monotonic age from the same state snapshot. The Unix timestamp
342/// on `CheckpointRun` is display metadata and must not gate pin detection.
343pub(crate) fn checkpoint_run_snapshot(
344 pool: &ConnectionPool,
345) -> (CheckpointRunStatus, Option<Duration>) {
346 let runs = checkpoint_runs()
347 .lock()
348 .unwrap_or_else(std::sync::PoisonError::into_inner);
349 let Some(state) = runs.get(&checkpoint_db_key(pool)) else {
350 return (CheckpointRunStatus::NoTask, None);
351 };
352 if state.active_tasks == 0 {
353 return (CheckpointRunStatus::NoTask, None);
354 }
355 state
356 .entry
357 .map_or((CheckpointRunStatus::NoObservation, None), |entry| {
358 (
359 CheckpointRunStatus::Observed(entry.run),
360 Some(entry.first_observed_at.elapsed()),
361 )
362 })
363}
364
365/// Process-lifetime totals for actual routine PASSIVE calls on one store.
366/// Ticks with no PASSIVE call and post-TRUNCATE probes are excluded. Busy counts only
367/// SQLite's returned busy flag, never an incomplete checkpoint's pending frames.
368#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
369#[serde(default)]
370pub struct CheckpointTiming {
371 pub ticks: u64,
372 pub elapsed_us_sum: u64,
373 pub elapsed_us_max: u64,
374 pub busy_ticks: u64,
375 pub error_ticks: u64,
376}
377
378static CHECKPOINT_TIMINGS: OnceLock<Mutex<HashMap<Option<PathBuf>, CheckpointTiming>>> =
379 OnceLock::new();
380
381fn checkpoint_timings() -> &'static Mutex<HashMap<Option<PathBuf>, CheckpointTiming>> {
382 CHECKPOINT_TIMINGS.get_or_init(|| Mutex::new(HashMap::new()))
383}
384
385fn record_checkpoint_timing(pool: &ConnectionPool, elapsed_us: u64, busy: Option<i64>) {
386 let mut timings = checkpoint_timings()
387 .lock()
388 .unwrap_or_else(std::sync::PoisonError::into_inner);
389 let timing = timings.entry(checkpoint_db_key(pool)).or_default();
390 timing.ticks = timing.ticks.saturating_add(1);
391 timing.elapsed_us_sum = timing.elapsed_us_sum.saturating_add(elapsed_us);
392 timing.elapsed_us_max = timing.elapsed_us_max.max(elapsed_us);
393 timing.busy_ticks = timing
394 .busy_ticks
395 .saturating_add(u64::from(busy.is_some_and(|value| value != 0)));
396 timing.error_ticks = timing.error_ticks.saturating_add(u64::from(busy.is_none()));
397}
398
399/// Pure in-memory read; zero means no routine call has been recorded for this key.
400pub fn checkpoint_timing(pool: &ConnectionPool) -> CheckpointTiming {
401 checkpoint_timings()
402 .lock()
403 .unwrap_or_else(std::sync::PoisonError::into_inner)
404 .get(&checkpoint_db_key(pool))
405 .copied()
406 .unwrap_or_default()
407}
408
409/// Latest routine observation by canonical database identity. Checkpoint
410/// tasks fan out per backend, so a single process-global "last task wins"
411/// gauge would misattribute a secondary backend to the main metrics frame.
412static ROUTINE_WAL_OBSERVATIONS: OnceLock<Mutex<HashMap<Option<PathBuf>, RoutineWalObservation>>> =
413 OnceLock::new();
414
415fn routine_wal_observations() -> &'static Mutex<HashMap<Option<PathBuf>, RoutineWalObservation>> {
416 ROUTINE_WAL_OBSERVATIONS.get_or_init(|| Mutex::new(HashMap::new()))
417}
418
419fn checkpoint_db_key_from_path(path: Option<&Path>) -> Option<PathBuf> {
420 path.map(Path::to_path_buf)
421}
422
423fn checkpoint_db_key(pool: &ConnectionPool) -> Option<PathBuf> {
424 checkpoint_db_key_from_path(pool.canonical_path())
425}
426
427fn observed_at_unix_ms() -> u64 {
428 std::time::SystemTime::now()
429 .duration_since(std::time::UNIX_EPOCH)
430 .map(|duration| duration.as_millis() as u64)
431 .unwrap_or(0)
432}
433
434fn physical_wal_bytes(pool: &ConnectionPool) -> Option<u64> {
435 let path = pool.canonical_path()?;
436 let mut sidecar = path.as_os_str().to_os_string();
437 sidecar.push("-wal");
438 std::fs::metadata(PathBuf::from(sidecar))
439 .ok()
440 .map(|metadata| metadata.len())
441}
442
443fn record_routine_wal_observation(
444 pool: &ConnectionPool,
445 raw: RawCheckpointObservation,
446) -> RoutineWalObservation {
447 let log_frames = raw.log_frames.max(0) as u64;
448 let checkpointed_frames = raw.checkpointed_frames.max(0) as u64;
449 let observation = RoutineWalObservation {
450 busy: raw.busy,
451 log_frames,
452 checkpointed_frames,
453 pending_frames: log_frames.saturating_sub(checkpointed_frames),
454 physical_wal_bytes: physical_wal_bytes(pool),
455 observed_at_unix_ms: observed_at_unix_ms(),
456 };
457 routine_wal_observations()
458 .lock()
459 .unwrap_or_else(std::sync::PoisonError::into_inner)
460 .insert(checkpoint_db_key(pool), observation.clone());
461 observation
462}
463
464/// Latest periodic checkpoint sample for this exact backend. This is a pure
465/// in-memory read: it never issues `wal_checkpoint` or stats the filesystem.
466pub fn routine_wal_observation(pool: &ConnectionPool) -> Option<RoutineWalObservation> {
467 routine_wal_observations()
468 .lock()
469 .unwrap_or_else(std::sync::PoisonError::into_inner)
470 .get(&checkpoint_db_key(pool))
471 .cloned()
472}
473
474/// Last-observed WAL page count, if any checkpoint tick has run yet in this
475/// process. Read surface for the daemon-frame metrics snapshot.
476pub fn last_observed_wal_pages() -> Option<u64> {
477 match LAST_WAL_PAGES.load(Ordering::Relaxed) {
478 u64::MAX => None,
479 pages => Some(pages),
480 }
481}
482
483/// Total WAL TRUNCATE attempts made in this process's lifetime.
484pub fn truncate_attempts() -> u64 {
485 TRUNCATE_ATTEMPTS.load(Ordering::Relaxed)
486}
487
488/// Current measured TRUNCATE-failure streak; unmeasured attempts leave it unchanged.
489pub fn truncate_consecutive_failures() -> u64 {
490 TRUNCATE_CONSECUTIVE_FAILURES.load(Ordering::Relaxed)
491}
492
493/// Total checkpoint ticks without an available WAL frame observation in this
494/// process's lifetime (dedicated connection unavailable or SQLite busy).
495pub fn checkpoint_skipped_ticks() -> u64 {
496 CHECKPOINT_SKIPPED_TICKS.load(Ordering::Relaxed)
497}
498
499/// Current consecutive-skip run length; 0 once the next tick is observed.
500pub fn checkpoint_consecutive_skips() -> u64 {
501 CHECKPOINT_CONSECUTIVE_SKIPS.load(Ordering::Relaxed)
502}
503
504/// WAL page count last known at the time of the most recent skip, if any
505/// skip has occurred yet in this process.
506pub fn checkpoint_last_skip_wal_pages() -> Option<u64> {
507 match CHECKPOINT_LAST_SKIP_WAL_PAGES.load(Ordering::Relaxed) {
508 u64::MAX => None,
509 pages => Some(pages),
510 }
511}
512
513/// Total at/above-`warn_pages` observations aggregated in memory.
514pub fn checkpoint_pressure_elevated_ticks() -> u64 {
515 CHECKPOINT_PRESSURE_ELEVATED_TICKS.load(Ordering::Relaxed)
516}
517
518/// Total pressure episodes observed to start in this process.
519pub fn checkpoint_pressure_episodes_started() -> u64 {
520 CHECKPOINT_PRESSURE_EPISODES_STARTED.load(Ordering::Relaxed)
521}
522
523/// Total pressure episodes observed to recover in this process.
524pub fn checkpoint_pressure_episodes_recovered() -> u64 {
525 CHECKPOINT_PRESSURE_EPISODES_RECOVERED.load(Ordering::Relaxed)
526}
527
528/// Total primary-store append calls made by checkpoint lifecycle workers.
529pub fn checkpoint_lifecycle_append_attempts() -> u64 {
530 CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.load(Ordering::Relaxed)
531}
532
533/// Total checkpoint lifecycle append calls that returned a storage error.
534pub fn checkpoint_lifecycle_append_failures() -> u64 {
535 CHECKPOINT_LIFECYCLE_APPEND_FAILURES.load(Ordering::Relaxed)
536}
537
538/// Total cached-reader read transactions rolled back on reuse for exceeding
539/// `read_tx_max_age` (#1846), across this process's lifetime.
540pub fn read_tx_max_age_evictions() -> u64 {
541 READ_TX_MAX_AGE_EVICTIONS.load(Ordering::Relaxed)
542}
543
544/// Records one cached-reader read transaction rolled back on reuse for
545/// exceeding `read_tx_max_age`. Called from `sql_bridge.rs` at the point the
546/// rollback is issued, regardless of whether the rollback itself succeeds —
547/// this counts the eviction *attempt*, matching the `truncate_attempts`
548/// naming convention above.
549pub(crate) fn note_read_tx_max_age_eviction() {
550 READ_TX_MAX_AGE_EVICTIONS.fetch_add(1, Ordering::Relaxed);
551}
552
553/// Total checkpoint lifecycle transitions rejected before append.
554pub fn checkpoint_lifecycle_enqueue_drops() -> u64 {
555 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.load(Ordering::Relaxed)
556}
557
558/// A tick had no usable WAL frame observation: bump the
559/// lifetime and consecutive-skip counters and snapshot the last-known WAL
560/// pressure so an operator can see how bad the WAL was heading into the skip
561/// streak.
562fn note_checkpoint_skipped() {
563 CHECKPOINT_SKIPPED_TICKS.fetch_add(1, Ordering::Relaxed);
564 CHECKPOINT_CONSECUTIVE_SKIPS.fetch_add(1, Ordering::Relaxed);
565 if let Some(pages) = last_observed_wal_pages() {
566 CHECKPOINT_LAST_SKIP_WAL_PAGES.store(pages, Ordering::Relaxed);
567 }
568}
569
570/// A tick had a valid WAL frame observation: close out any prior skip
571/// streak. `_wal_pages` is accepted for call-site symmetry with
572/// `note_checkpoint_skipped` and to leave room for a future observed-side
573/// gauge without changing this function's signature again.
574fn note_checkpoint_observed(_wal_pages: u64) {
575 CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
576}
577
578fn note_checkpoint_pressure_observation(above_warn: bool, was_above_warn: bool) {
579 if above_warn {
580 CHECKPOINT_PRESSURE_ELEVATED_TICKS.fetch_add(1, Ordering::Relaxed);
581 if !was_above_warn {
582 CHECKPOINT_PRESSURE_EPISODES_STARTED.fetch_add(1, Ordering::Relaxed);
583 }
584 } else if was_above_warn {
585 CHECKPOINT_PRESSURE_EPISODES_RECOVERED.fetch_add(1, Ordering::Relaxed);
586 }
587}
588
589/// Reset the checkpoint-pressure atomics between tests. Process-wide gauges
590/// are otherwise shared across every test in this binary; tests that assert
591/// on them must reset first and run under a shared `#[serial(...)]` group.
592#[cfg(test)]
593pub(crate) fn reset_checkpoint_metrics_for_tests() {
594 CHECKPOINT_SKIPPED_TICKS.store(0, Ordering::Relaxed);
595 CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
596 CHECKPOINT_LAST_SKIP_WAL_PAGES.store(u64::MAX, Ordering::Relaxed);
597 CHECKPOINT_PRESSURE_ELEVATED_TICKS.store(0, Ordering::Relaxed);
598 CHECKPOINT_PRESSURE_EPISODES_STARTED.store(0, Ordering::Relaxed);
599 CHECKPOINT_PRESSURE_EPISODES_RECOVERED.store(0, Ordering::Relaxed);
600 CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.store(0, Ordering::Relaxed);
601 CHECKPOINT_LIFECYCLE_APPEND_FAILURES.store(0, Ordering::Relaxed);
602 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.store(0, Ordering::Relaxed);
603 READ_TX_MAX_AGE_EVICTIONS.store(0, Ordering::Relaxed);
604}
605
606/// Outcome of a single checkpoint attempt.
607///
608/// `Skipped` is returned when the dedicated connection is unavailable or its
609/// PASSIVE result is busy or inconsistent and has no usable WAL frame
610/// observation. A concurrent pool writer does not cause a skip: the task never
611/// checks out its writer mutex. `Observed` carries the WAL page count read
612/// during the tick. A skipped tick leaves threshold-crossing state unchanged.
613#[derive(Debug, Clone, Copy, PartialEq, Eq)]
614pub enum CheckpointTick {
615 /// No usable WAL frame observation was available this tick.
616 Skipped,
617 /// A checkpoint was issued; the value is the observed WAL page count.
618 Observed(u64),
619}
620
621/// Default number of consecutive above-`warn_pages` observed ticks required
622/// to escalate from the INFO to the WARN rung of the ADR-091 severity ladder.
623pub const DEFAULT_WARN_SUSTAINED_CYCLES: u8 = 3;
624
625/// Configuration for the WAL checkpoint background task.
626///
627/// All fields default to conservative production values. Override via the
628/// environment variables documented on each field.
629#[derive(Clone, Debug)]
630pub struct CheckpointConfig {
631 /// How often to run a passive checkpoint when there is no active write.
632 ///
633 /// Overridable via `KHIVE_CHECKPOINT_INTERVAL_MS` (milliseconds).
634 /// Default: 500 ms.
635 pub interval: Duration,
636
637 /// WAL page count above which a warning is logged.
638 ///
639 /// Overridable via `KHIVE_WAL_WARN_PAGES`.
640 /// Default: 2000 pages (~8 MB at 4 KiB page size).
641 pub warn_pages: u64,
642
643 /// Number of consecutive observed ticks with `wal_pages >= warn_pages`
644 /// required before the ADR-091 severity ladder escalates from INFO
645 /// (first crossing) to WARN (sustained pressure). Edge-triggered once
646 /// per elevation episode — see [`CheckpointSeverityState`].
647 ///
648 /// Overridable via `KHIVE_WAL_WARN_SUSTAINED_CYCLES`.
649 /// Default: 3 cycles.
650 pub warn_sustained_cycles: u8,
651
652 /// WAL page count above which a high-pressure WARNING is logged.
653 ///
654 /// The periodic task always runs PASSIVE regardless; this threshold signals
655 /// only that the WAL is not draining. Whether an old snapshot is pinning it
656 /// is informed at the crossing by the in-process transaction registry,
657 /// against `tx_warn_secs` — see `log_wal_high_water_warn`. This registry
658 /// cannot exclude readers in another process. Either way an
659 /// operator can schedule a blocking TRUNCATE at a safe moment outside
660 /// normal write traffic; the two cases differ in what else is worth doing.
661 ///
662 /// Overridable via `KHIVE_WAL_HIGH_WATER_PAGES`.
663 /// Default: 6000 pages (~24 MB at 4 KiB page size).
664 pub high_water_pages: u64,
665
666 /// WAL page count above which a TRUNCATE escalation attempt is armed
667 /// (ADR-091 Plank 2).
668 ///
669 /// This is a separate, much higher threshold than `high_water_pages`:
670 /// crossing it does not itself attempt TRUNCATE — it only arms the
671 /// attempt, which additionally requires `truncate_min_interval` to have
672 /// elapsed since the last attempt.
673 ///
674 /// Overridable via `KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES`.
675 /// Default: 20000 pages.
676 pub truncate_high_water_pages: u64,
677
678 /// Minimum spacing between TRUNCATE *attempts* (not successes).
679 ///
680 /// A skipped tick (dedicated connection unavailable, below threshold, or
681 /// interval not yet elapsed) never advances the "last attempt" clock, so
682 /// the next tick where the connection is available and the threshold is
683 /// still crossed is immediately eligible rather than waiting out the
684 /// full interval again.
685 ///
686 /// Overridable via `KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS`.
687 /// Default: 300 seconds (5 minutes).
688 pub truncate_min_interval: Duration,
689
690 /// Temporary `busy_timeout` used only for the duration of a TRUNCATE
691 /// attempt, restored to the pool's configured busy timeout immediately
692 /// after the attempt completes (win or lose).
693 ///
694 /// Overridable via `KHIVE_WAL_TRUNCATE_BUSY_MS`.
695 /// Default: 2000 ms.
696 pub truncate_busy_timeout: Duration,
697
698 /// ADR-091 Plank 1 soft cap: age past which the oldest entry in the
699 /// shared open-transaction registry is surfaced at `tracing::warn!` on
700 /// every tick (Skipped or Observed), independent of WAL page pressure.
701 /// See `crates/khive-db/docs/api/checkpoint.md` for the Plank 1 rationale.
702 ///
703 /// Overridable via `KHIVE_TX_WARN_SECS`.
704 /// Default: 30 seconds.
705 pub tx_warn_secs: Duration,
706
707 /// ADR-091 Plank 1 hard cap: age past which the same sweep escalates the
708 /// oldest registry entry to `tracing::error!`. The sweep itself is
709 /// visibility only — nothing in `TxAgeSweepState` force-closes a stale
710 /// span. `sql_bridge.rs`'s cached-reader read-transaction path shares
711 /// this exact value (via `PoolConfig::read_tx_max_age`, #1846) to
712 /// actually roll back and evict an explicit read transaction the next
713 /// time its handle is reused past this age — reclamation for the
714 /// "reused periodically" case, not the "held idle with no further calls"
715 /// case the ADR named as its accepted gap; see
716 /// `crates/khive-db/docs/api/checkpoint.md`'s Plank 1 section for the
717 /// distinction and why the latter remains open design work.
718 ///
719 /// Overridable via `KHIVE_TX_MAX_AGE_SECS`.
720 /// Default: 120 seconds.
721 pub tx_max_age_secs: Duration,
722}
723
724impl Default for CheckpointConfig {
725 fn default() -> Self {
726 Self {
727 interval: Duration::from_millis(500),
728 warn_pages: 2000,
729 warn_sustained_cycles: DEFAULT_WARN_SUSTAINED_CYCLES,
730 high_water_pages: 6000,
731 truncate_high_water_pages: 20_000,
732 truncate_min_interval: Duration::from_secs(300),
733 truncate_busy_timeout: Duration::from_millis(2000),
734 tx_warn_secs: Duration::from_secs(30),
735 tx_max_age_secs: Duration::from_secs(120),
736 }
737 }
738}
739
740impl CheckpointConfig {
741 /// Build a `CheckpointConfig` from the environment.
742 ///
743 /// Unset or unparseable variables fall back to the compiled-in defaults.
744 pub fn from_env() -> Self {
745 let mut cfg = Self::default();
746
747 if let Ok(ms) = std::env::var("KHIVE_CHECKPOINT_INTERVAL_MS") {
748 if let Ok(v) = ms.parse::<u64>() {
749 if v > 0 {
750 cfg.interval = Duration::from_millis(v);
751 }
752 }
753 }
754
755 if let Ok(v) = std::env::var("KHIVE_WAL_WARN_PAGES") {
756 if let Ok(n) = v.parse::<u64>() {
757 if n > 0 {
758 cfg.warn_pages = n;
759 }
760 }
761 }
762
763 if let Ok(v) = std::env::var("KHIVE_WAL_WARN_SUSTAINED_CYCLES") {
764 if let Ok(n) = v.parse::<u8>() {
765 if n > 0 {
766 cfg.warn_sustained_cycles = n;
767 }
768 }
769 }
770
771 if let Ok(v) = std::env::var("KHIVE_WAL_HIGH_WATER_PAGES") {
772 if let Ok(n) = v.parse::<u64>() {
773 if n > 0 {
774 cfg.high_water_pages = n;
775 }
776 }
777 }
778
779 if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES") {
780 if let Ok(n) = v.parse::<u64>() {
781 if n > 0 {
782 cfg.truncate_high_water_pages = n;
783 }
784 }
785 }
786
787 if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS") {
788 if let Ok(n) = v.parse::<u64>() {
789 if n > 0 {
790 cfg.truncate_min_interval = Duration::from_secs(n);
791 }
792 }
793 }
794
795 if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_BUSY_MS") {
796 if let Ok(n) = v.parse::<u64>() {
797 if n > 0 {
798 cfg.truncate_busy_timeout = Duration::from_millis(n);
799 }
800 }
801 }
802
803 (cfg.tx_warn_secs, cfg.tx_max_age_secs) =
804 tx_age_thresholds_from_env(cfg.tx_warn_secs, cfg.tx_max_age_secs);
805
806 cfg
807 }
808}
809
810/// Parse `KHIVE_TX_WARN_SECS`/`KHIVE_TX_MAX_AGE_SECS` against the given
811/// defaults, applying the same ordering guard both [`CheckpointConfig`] and
812/// [`SessionSweepConfig`] need (minor, ADR-091 Amendment 2: this was
813/// previously duplicated verbatim in both `from_env` methods).
814///
815/// The severity ladder assumes `tx_warn_secs < tx_max_age_secs` (Warn fires
816/// before Stale as an entry ages). A reversed or equal pair — whether from
817/// one misconfigured var or the interaction of both — would invert or
818/// collapse that ordering (e.g. WARN_SECS=120, MAX_AGE_SECS=30 emits Stale at
819/// 30s and never reaches the Warn crossing until 120s), so both are rejected
820/// together rather than silently honored. Resetting both to the caller's
821/// defaults (rather than just clamping one) avoids guessing which of the two
822/// the operator actually meant to change.
823pub(crate) fn tx_age_thresholds_from_env(
824 default_warn: Duration,
825 default_max: Duration,
826) -> (Duration, Duration) {
827 let mut warn_secs = default_warn;
828 let mut max_age_secs = default_max;
829
830 if let Ok(v) = std::env::var("KHIVE_TX_WARN_SECS") {
831 if let Ok(n) = v.parse::<u64>() {
832 if n > 0 {
833 warn_secs = Duration::from_secs(n);
834 }
835 }
836 }
837
838 if let Ok(v) = std::env::var("KHIVE_TX_MAX_AGE_SECS") {
839 if let Ok(n) = v.parse::<u64>() {
840 if n > 0 {
841 max_age_secs = Duration::from_secs(n);
842 }
843 }
844 }
845
846 if warn_secs >= max_age_secs {
847 tracing::warn!(
848 configured_tx_warn_secs = warn_secs.as_secs_f64(),
849 configured_tx_max_age_secs = max_age_secs.as_secs_f64(),
850 fallback_tx_warn_secs = default_warn.as_secs_f64(),
851 fallback_tx_max_age_secs = default_max.as_secs_f64(),
852 "KHIVE_TX_WARN_SECS must be strictly less than KHIVE_TX_MAX_AGE_SECS; \
853 both transaction-age thresholds were rejected and reset to their defaults"
854 );
855 return (default_warn, default_max);
856 }
857
858 (warn_secs, max_age_secs)
859}
860
861#[cfg(unix)]
862const DEFAULT_WALPIN_FULL_SCAN_INTERVAL: Duration = Duration::from_secs(30);
863
864#[cfg(unix)]
865#[derive(Debug, Clone)]
866struct CachedWalpinAttribution {
867 report: crate::walpin::WalpinReport,
868 census: Result<crate::walpin::CensusResult, String>,
869 captured_at: Instant,
870}
871
872#[cfg(unix)]
873#[derive(Debug)]
874enum WalpinFullScanPlan {
875 Refresh {
876 previous_last_attempt: Option<Instant>,
877 },
878 Cached(CachedWalpinAttribution),
879 Suppressed,
880}
881
882/// Mutable escalation state carried across ticks by the caller (ADR-091 Plank 2).
883///
884/// Kept separate from [`CheckpointConfig`] because it is *state*, not
885/// configuration: `last_attempt` and `consecutive_failures` mutate every tick,
886/// while `CheckpointConfig` is parsed once and held immutable for the life of
887/// the task.
888#[derive(Debug)]
889pub struct TruncateState {
890 /// When the last TRUNCATE *attempt* ran (armed + writer held), regardless
891 /// of whether it succeeded in reclaiming pages. `None` means no attempt
892 /// has ever run, so the first armed tick is immediately eligible.
893 last_attempt: Option<Instant>,
894 /// Count of measured TRUNCATE outcomes that failed to bring `wal_pages`
895 /// below `warn_pages`, ignoring attempts whose post-TRUNCATE measurement
896 /// was unavailable. A measured clearing result resets it; a one-shot
897 /// escalated WARN fires at exactly 3 failures.
898 consecutive_failures: u32,
899 /// Fallback freshness cadence for legacy sidecar records that do not
900 /// declare their producer interval. Captured once when the daemon task
901 /// starts; this is ADR-091's compiled 5000 ms session-sweep default, never
902 /// the daemon's faster checkpoint cadence or a local environment override.
903 #[cfg(unix)]
904 legacy_walpin_fallback_interval: Duration,
905 /// Minimum spacing between full sidecar/OS-holder enumeration attempts.
906 /// The attempt timestamp advances before blocking work starts, so an I/O
907 /// failure or worker panic cannot turn sustained pressure into a hot retry
908 /// loop. A successful report is retained only for diagnostic reuse.
909 #[cfg(unix)]
910 walpin_full_scan_interval: Duration,
911 #[cfg(unix)]
912 walpin_full_scan_last_attempt: Option<Instant>,
913 #[cfg(unix)]
914 walpin_cached_attribution: Option<CachedWalpinAttribution>,
915 /// Whether the no-progress attribution arm already attempted the one
916 /// bounded sidecar enumeration allowed for this checkpoint tick.
917 #[cfg(unix)]
918 sidecar_attribution_attempted_this_tick: bool,
919}
920
921impl Default for TruncateState {
922 fn default() -> Self {
923 Self {
924 last_attempt: None,
925 consecutive_failures: 0,
926 #[cfg(unix)]
927 legacy_walpin_fallback_interval: DEFAULT_SESSION_SWEEP_INTERVAL,
928 #[cfg(unix)]
929 walpin_full_scan_interval: DEFAULT_WALPIN_FULL_SCAN_INTERVAL,
930 #[cfg(unix)]
931 walpin_full_scan_last_attempt: None,
932 #[cfg(unix)]
933 walpin_cached_attribution: None,
934 #[cfg(unix)]
935 sidecar_attribution_attempted_this_tick: false,
936 }
937 }
938}
939
940impl TruncateState {
941 #[cfg(unix)]
942 fn with_legacy_walpin_fallback(interval: Duration) -> Self {
943 Self {
944 legacy_walpin_fallback_interval: interval,
945 ..Self::default()
946 }
947 }
948
949 #[cfg(all(test, unix))]
950 fn with_walpin_full_scan_cadence(interval: Duration) -> Self {
951 Self {
952 walpin_full_scan_interval: interval,
953 ..Self::default()
954 }
955 }
956
957 #[cfg(unix)]
958 fn begin_tick(&mut self) {
959 self.sidecar_attribution_attempted_this_tick = false;
960 }
961
962 #[cfg(unix)]
963 fn housekeeping_due(&self) -> bool {
964 !self.sidecar_attribution_attempted_this_tick
965 && self.walpin_full_scan_due_at(Instant::now())
966 }
967
968 #[cfg(unix)]
969 fn walpin_full_scan_due_at(&self, now: Instant) -> bool {
970 self.walpin_full_scan_last_attempt.is_none_or(|last| {
971 now.saturating_duration_since(last) >= self.walpin_full_scan_interval
972 })
973 }
974
975 #[cfg(unix)]
976 fn claim_walpin_full_scan_at(&mut self, now: Instant) -> bool {
977 if !self.walpin_full_scan_due_at(now) {
978 return false;
979 }
980 self.walpin_full_scan_last_attempt = Some(now);
981 true
982 }
983
984 #[cfg(unix)]
985 fn plan_walpin_attribution_at(&mut self, now: Instant) -> WalpinFullScanPlan {
986 if self.walpin_full_scan_due_at(now) {
987 let previous_last_attempt = self.walpin_full_scan_last_attempt.replace(now);
988 WalpinFullScanPlan::Refresh {
989 previous_last_attempt,
990 }
991 } else if let Some(cached) = self.walpin_cached_attribution.clone() {
992 WalpinFullScanPlan::Cached(cached)
993 } else {
994 WalpinFullScanPlan::Suppressed
995 }
996 }
997
998 #[cfg(unix)]
999 fn restore_walpin_full_scan_reservation(&mut self, previous_last_attempt: Option<Instant>) {
1000 self.walpin_full_scan_last_attempt = previous_last_attempt;
1001 }
1002
1003 #[cfg(unix)]
1004 fn cache_walpin_attribution(
1005 &mut self,
1006 report: crate::walpin::WalpinReport,
1007 census: Result<crate::walpin::CensusResult, String>,
1008 captured_at: Instant,
1009 ) {
1010 self.walpin_cached_attribution = Some(CachedWalpinAttribution {
1011 report,
1012 census,
1013 captured_at,
1014 });
1015 }
1016}
1017
1018/// ADR-091 graduated severity rung for sustained WAL pressure.
1019///
1020/// `Alarm` is never produced by [`CheckpointSeverityState::observe_wal_pages`]
1021/// — it labels the existing TRUNCATE-escalation tier (`maybe_truncate`),
1022/// which is gated on its own threshold/interval state, not on this ladder.
1023/// It exists here so callers and tests can name all three rungs uniformly.
1024#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1025pub enum CheckpointSeverityRung {
1026 /// First observed tick crossing `warn_pages` after a below-warn tick.
1027 Info,
1028 /// `warn_sustained_cycles` consecutive observed ticks at/above
1029 /// `warn_pages`; edge-triggered once per elevation episode.
1030 Warn,
1031 /// The TRUNCATE-escalation tier (`checkpoint_high_water_pages` and
1032 /// above); never emitted by `observe_wal_pages`.
1033 Alarm,
1034}
1035
1036/// ADR-091 severity ladder state, carried across ticks by the caller
1037/// alongside [`TruncateState`]. Pure state machine: no I/O, no logging —
1038/// callers turn the returned emissions into `tracing` calls.
1039#[derive(Debug, Default, Clone)]
1040pub struct CheckpointSeverityState {
1041 /// Whether the previous observed tick was at/above `warn_pages`. Drives
1042 /// the below→above edge that fires INFO.
1043 was_above_warn: bool,
1044 /// Run-length of consecutive observed ticks at/above `warn_pages` in the
1045 /// current elevation episode. Resets to 0 on any below-warn tick.
1046 consecutive_above_warn: u8,
1047 /// Whether WARN has already fired for the current elevation episode, so
1048 /// sustained pressure logs WARN once per episode, not once per tick past
1049 /// the threshold.
1050 warn_emitted_for_episode: bool,
1051}
1052
1053/// One severity-ladder emission produced by a single
1054/// [`CheckpointSeverityState::observe_wal_pages`] call.
1055#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1056pub struct CheckpointSeverityEmission {
1057 /// Which rung this emission represents (`Info` or `Warn`; see
1058 /// [`CheckpointSeverityRung::Alarm`] doc for why `Alarm` never appears
1059 /// here).
1060 pub rung: CheckpointSeverityRung,
1061 /// The WAL page count observed on the tick that produced this emission.
1062 pub wal_pages: u64,
1063 /// The `warn_pages` threshold in effect for this tick.
1064 pub threshold_pages: u64,
1065 /// Consecutive above-warn cycle count as of this tick (1 on the INFO
1066 /// edge, `warn_sustained_cycles` on the WARN edge).
1067 pub consecutive_cycles: u8,
1068}
1069
1070impl CheckpointSeverityState {
1071 /// Advance the severity ladder by one observed tick and return every
1072 /// rung crossed on this tick (zero, one, or two emissions: a fresh
1073 /// elevation episode can produce INFO and, if `warn_sustained_cycles`
1074 /// is 1, WARN on the very same tick).
1075 ///
1076 /// A below-warn tick resets both the consecutive-cycle counter and the
1077 /// per-episode WARN latch, re-arming INFO/WARN for a later episode.
1078 /// Skipped ticks must not be passed here at all — the caller only calls
1079 /// this on `CheckpointTick::Observed`, matching the existing
1080 /// threshold-crossing WARN's skip-leaves-state-unchanged rule.
1081 pub fn observe_wal_pages(
1082 &mut self,
1083 wal_pages: u64,
1084 config: &CheckpointConfig,
1085 ) -> Vec<CheckpointSeverityEmission> {
1086 let mut emissions = Vec::new();
1087 let above_warn = wal_pages >= config.warn_pages;
1088
1089 if above_warn {
1090 self.consecutive_above_warn = self.consecutive_above_warn.saturating_add(1);
1091
1092 if !self.was_above_warn {
1093 emissions.push(CheckpointSeverityEmission {
1094 rung: CheckpointSeverityRung::Info,
1095 wal_pages,
1096 threshold_pages: config.warn_pages,
1097 consecutive_cycles: self.consecutive_above_warn,
1098 });
1099 }
1100
1101 if !self.warn_emitted_for_episode
1102 && self.consecutive_above_warn >= config.warn_sustained_cycles
1103 {
1104 emissions.push(CheckpointSeverityEmission {
1105 rung: CheckpointSeverityRung::Warn,
1106 wal_pages,
1107 threshold_pages: config.warn_pages,
1108 consecutive_cycles: self.consecutive_above_warn,
1109 });
1110 self.warn_emitted_for_episode = true;
1111 }
1112 } else {
1113 self.consecutive_above_warn = 0;
1114 self.warn_emitted_for_episode = false;
1115 }
1116
1117 self.was_above_warn = above_warn;
1118 emissions
1119 }
1120}
1121
1122/// ADR-091 Plank 1 rung for the open-transaction registry's background age
1123/// sweep: independent of the WAL-pressure ladder above, keyed purely off how
1124/// long the registry's oldest entry has been open.
1125#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1126pub enum TxAgeRung {
1127 /// The oldest registry entry's age crossed `tx_warn_secs`.
1128 Warn,
1129 /// The oldest registry entry's age crossed `tx_max_age_secs` — the ADR's
1130 /// "cooperative stale-op guard" cap. No in-process mechanism force-closes
1131 /// it (see [`CheckpointConfig::tx_max_age_secs`]); this rung is the
1132 /// sweep's strongest available signal.
1133 Stale,
1134}
1135
1136/// One emission produced by a single [`TxAgeSweepState::observe`] call.
1137#[derive(Debug, Clone, PartialEq, Eq)]
1138pub struct TxAgeEmission {
1139 pub rung: TxAgeRung,
1140 pub age: Duration,
1141 pub label: Option<String>,
1142}
1143
1144/// ADR-091 Plank 1 background-sweep state, carried across ticks by the
1145/// caller alongside [`CheckpointSeverityState`] and [`TruncateState`]. Pure
1146/// state machine: no I/O, no logging — callers turn the returned emissions
1147/// into `tracing` calls, mirroring [`CheckpointSeverityState`]'s shape.
1148///
1149/// Keyed off `khive_storage::tx_registry::oldest()` — the single oldest
1150/// entry across every registered span, regardless of which call site created
1151/// it. Deliberately a different signal from the WAL-pressure ladder: a span
1152/// can go stale under low WAL pressure, or vice versa. See
1153/// `crates/khive-db/docs/api/checkpoint.md` for the full rationale.
1154#[derive(Debug, Default, Clone)]
1155pub struct TxAgeSweepState {
1156 /// Whether the previous observed tick's oldest entry was at/above
1157 /// `tx_warn_secs`. Drives the below→above edge that fires `Warn`.
1158 was_above_warn: bool,
1159 /// Whether the previous observed tick's oldest entry was at/above
1160 /// `tx_max_age_secs`. Drives the below→above edge that fires `Stale`.
1161 was_above_max_age: bool,
1162 /// Identity of the entry the previous observed tick reported as oldest,
1163 /// or `None` if the registry was empty. Tracked separately from the two
1164 /// latches above so a change in *which span* is oldest can be detected
1165 /// even when both latches are already `true` (see [`Self::observe`]).
1166 tracked_id: Option<khive_storage::tx_registry::TxId>,
1167}
1168
1169impl TxAgeSweepState {
1170 /// Advance by one observed tick given the registry's current oldest
1171 /// entry (identity, age, label), or `None` if empty. Returns zero, one,
1172 /// or two emissions — an entry already stale the first time it's seen
1173 /// under a given identity crosses both rungs on the same tick.
1174 ///
1175 /// A below-threshold (or absent) oldest entry resets both latches. A
1176 /// change in the oldest entry's [`TxId`](khive_storage::tx_registry::TxId)
1177 /// also force-resets both latches before re-evaluating age, so a
1178 /// departed span's latched state cannot suppress the crossing for an
1179 /// already-stale successor. See `crates/khive-db/docs/api/checkpoint.md`
1180 /// for why identity tracking is required here, not just the age check.
1181 pub fn observe(
1182 &mut self,
1183 oldest: Option<(khive_storage::tx_registry::TxId, Duration, Option<String>)>,
1184 tx_warn_secs: Duration,
1185 tx_max_age_secs: Duration,
1186 ) -> Vec<TxAgeEmission> {
1187 let mut emissions = Vec::new();
1188
1189 let Some((id, age, label)) = oldest else {
1190 self.was_above_warn = false;
1191 self.was_above_max_age = false;
1192 self.tracked_id = None;
1193 return emissions;
1194 };
1195
1196 if self.tracked_id != Some(id) {
1197 self.was_above_warn = false;
1198 self.was_above_max_age = false;
1199 }
1200 self.tracked_id = Some(id);
1201
1202 let above_warn = age >= tx_warn_secs;
1203 let above_max_age = age >= tx_max_age_secs;
1204
1205 if above_warn && !self.was_above_warn {
1206 emissions.push(TxAgeEmission {
1207 rung: TxAgeRung::Warn,
1208 age,
1209 label: label.clone(),
1210 });
1211 }
1212 if above_max_age && !self.was_above_max_age {
1213 emissions.push(TxAgeEmission {
1214 rung: TxAgeRung::Stale,
1215 age,
1216 label,
1217 });
1218 }
1219
1220 self.was_above_warn = above_warn;
1221 self.was_above_max_age = above_max_age;
1222 emissions
1223 }
1224}
1225
1226/// ADR-091 Plank 1: turn a [`TxAgeEmission`] into the appropriate `tracing`
1227/// call. Extracted from `run_checkpoint_task` so tests can drive the same
1228/// logging path `CaptureSubscriber`-style without spinning up the async task
1229/// (mirrors [`log_tx_registry_oldest_warn`]/[`log_tx_registry_snapshot_warn`]).
1230fn log_tx_age_emission(emission: &TxAgeEmission) {
1231 let label = emission.label.as_deref().unwrap_or("<unlabeled>");
1232 match emission.rung {
1233 TxAgeRung::Warn => {
1234 tracing::warn!(
1235 tx_age_secs = emission.age.as_secs_f64(),
1236 tx_label = label,
1237 "ADR-091 Plank 1: open transaction registry entry exceeded soft-cap age"
1238 );
1239 }
1240 TxAgeRung::Stale => {
1241 tracing::error!(
1242 tx_age_secs = emission.age.as_secs_f64(),
1243 tx_label = label,
1244 "ADR-091 Plank 1: open transaction registry entry exceeded the cooperative \
1245 stale-op cap; no in-process mechanism can force-close it — investigate the \
1246 labeled caller directly"
1247 );
1248 }
1249 }
1250}
1251
1252/// ADR-091 Amendment 2 Plank B: per-process walpin sidecar state, carried
1253/// across ticks by whichever sweep owns it (the daemon's `run_checkpoint_task`
1254/// or a session's `run_session_sweep_task`). Once the registry's oldest span
1255/// exceeds `tx_warn_secs`, the first observation and each content change
1256/// rewrite the heartbeat body; unchanged ticks refresh only its mtime. The
1257/// heartbeat is removed once when the condition clears (and on shutdown), so
1258/// a process that never crosses the threshold writes no heartbeat body.
1259struct WalpinSidecarState {
1260 dir: PathBuf,
1261 pid: u32,
1262 role: &'static str,
1263 started_at: i64,
1264 /// This sweep's own tick cadence, recorded into every beacon and
1265 /// heartbeat so the enumerating daemon judges freshness against the
1266 /// PRODUCER's interval — a session on an independently slower configured
1267 /// cadence must not be misread as stale.
1268 sweep_interval_ms: u64,
1269 wrote: bool,
1270 /// Whether this process's registration beacon is believed present on
1271 /// disk. Cleared when a failed heartbeat write escalates to beacon
1272 /// removal (fail-closed — see `observe`) or a beacon touch fails; the
1273 /// next healthy tick then re-registers with a full write instead of a
1274 /// metadata touch.
1275 beacon_registered: bool,
1276 /// The content actually on disk in the last successful heartbeat body
1277 /// write, if any (ADR-091 Amendment 3 Plank F1). `None` whenever the
1278 /// next tick must go through a full write — no heartbeat written yet,
1279 /// the last write failed, or the threshold cleared. Compared against
1280 /// each new observation to decide touch (content unchanged) vs.
1281 /// rewrite (content changed).
1282 last_heartbeat: Option<LastHeartbeatState>,
1283}
1284
1285/// ADR-091 Amendment 3 Plank F1: the content signature of the heartbeat
1286/// body currently on disk, plus the `oldest_tx_started_at` value that body
1287/// carries — kept separate from the signature proper because it is derived
1288/// (fixed for as long as the same span stays oldest), not an independent
1289/// change signal.
1290struct LastHeartbeatState {
1291 span_id: khive_storage::tx_registry::TxId,
1292 label: Option<String>,
1293 attribution_basis: &'static str,
1294 sweep_interval_ms: u64,
1295 oldest_tx_started_at: i64,
1296}
1297
1298impl LastHeartbeatState {
1299 /// Whether a fresh observation carries exactly the content already on
1300 /// disk — the licensing condition for a metadata-only touch instead of
1301 /// a full body rewrite (the first over-threshold observation, a change
1302 /// of the oldest span's identity or label, a change of
1303 /// `attribution_basis`, or a change of the declared sweep cadence).
1304 fn content_matches(
1305 &self,
1306 span_id: khive_storage::tx_registry::TxId,
1307 label: &Option<String>,
1308 attribution_basis: &str,
1309 sweep_interval_ms: u64,
1310 ) -> bool {
1311 self.span_id == span_id
1312 && self.label == *label
1313 && self.attribution_basis == attribution_basis
1314 && self.sweep_interval_ms == sweep_interval_ms
1315 }
1316}
1317
1318impl WalpinSidecarState {
1319 /// `None` when the sidecar is disabled for this backend/env, or the
1320 /// backend has no on-disk path (in-memory).
1321 fn new(
1322 db_path: Option<&Path>,
1323 is_file_backed: bool,
1324 role: &'static str,
1325 interval: Duration,
1326 ) -> Option<Self> {
1327 let path = db_path?;
1328 if !crate::walpin::sidecar_enabled(is_file_backed) {
1329 return None;
1330 }
1331 let pid = std::process::id();
1332 Some(Self {
1333 dir: crate::walpin::sidecar_dir_for(path),
1334 pid,
1335 role,
1336 started_at: crate::walpin::process_start_time_secs(pid).unwrap_or(0),
1337 sweep_interval_ms: interval.as_millis().min(u64::MAX as u128) as u64,
1338 wrote: false,
1339 last_heartbeat: None,
1340 beacon_registered: false,
1341 })
1342 }
1343
1344 /// Write this process's registration beacon (ADR-091 Amendment 2
1345 /// sidecar-health attribution). Called once right after construction,
1346 /// before the sweep loop starts, and again only when a fail-closed
1347 /// removal or failed touch cleared `beacon_registered` — steady state
1348 /// stays metadata-touch-only with no data writes. The blocking fs I/O
1349 /// runs on `spawn_blocking` (perf, ADR-091 Amendment 2): this is
1350 /// invoked from an async context and must not run synchronous I/O
1351 /// inline on the async runtime's worker thread.
1352 async fn register_beacon(&mut self) {
1353 let dir = self.dir.clone();
1354 let beacon = crate::walpin::WalpinBeacon {
1355 pid: self.pid,
1356 process_role: self.role.to_string(),
1357 started_at: self.started_at,
1358 sweep_interval_ms: self.sweep_interval_ms,
1359 };
1360 let result =
1361 tokio::task::spawn_blocking(move || crate::walpin::write_beacon(&dir, &beacon)).await;
1362 match result {
1363 Ok(Ok(())) => {
1364 self.beacon_registered = true;
1365 }
1366 Ok(Err(e)) => {
1367 tracing::warn!(
1368 error = %e,
1369 "ADR-091 Amendment 2: failed to write walpin registration beacon; \
1370 this process's sidecar health will read as unknown, not registered-silent"
1371 );
1372 }
1373 Err(join_err) => {
1374 tracing::warn!(
1375 error = %join_err,
1376 "ADR-091 Amendment 2: walpin beacon write task panicked"
1377 );
1378 }
1379 }
1380 }
1381
1382 /// Run one bounded housekeeping pass independently of WAL pressure. The
1383 /// collector removes only positively dead/reused-PID residue; uncertain
1384 /// evidence remains for a no-progress attribution pass. Directory work
1385 /// and report memory are capped, and all blocking filesystem operations
1386 /// stay off the async runtime worker.
1387 #[cfg(unix)]
1388 async fn reap_dead_entries_bounded(
1389 &self,
1390 legacy_fallback_interval: Duration,
1391 ) -> Option<crate::walpin::WalpinReport> {
1392 let dir = self.dir.clone();
1393 let result = tokio::task::spawn_blocking(move || {
1394 crate::walpin::housekeep_live(&dir, legacy_fallback_interval)
1395 })
1396 .await;
1397 match result {
1398 Ok(Ok(report)) => Some(report),
1399 Ok(Err(e)) => {
1400 tracing::warn!(
1401 error = %e,
1402 "ADR-091 Amendment 6: bounded walpin sidecar cleanup failed"
1403 );
1404 None
1405 }
1406 Err(join_err) => {
1407 tracing::warn!(
1408 error = %join_err,
1409 "ADR-091 Amendment 6: walpin sidecar cleanup task panicked"
1410 );
1411 None
1412 }
1413 }
1414 }
1415
1416 /// ADR-091 Amendment 2 beacon refresh rule: a metadata-only mtime touch
1417 /// of this process's already-registered beacon, performed on every
1418 /// sweep tick except one where an over-threshold heartbeat write failed
1419 /// (see `observe`) — `registered-silent` classification requires this
1420 /// refresh to stay within the freshness window, not just the beacon's
1421 /// original write. After a fail-closed beacon removal (or a failed
1422 /// touch), the beacon is re-registered with a full write on the next
1423 /// healthy tick. Best-effort: a failure here degrades this process to
1424 /// `unknown` at the next enumeration, not a sweep-task error.
1425 async fn refresh_beacon(&mut self) {
1426 if !self.beacon_registered {
1427 self.register_beacon().await;
1428 return;
1429 }
1430 let dir = self.dir.clone();
1431 let pid = self.pid;
1432 let result =
1433 tokio::task::spawn_blocking(move || crate::walpin::touch_beacon(&dir, pid)).await;
1434 match result {
1435 Ok(Ok(())) => {}
1436 Ok(Err(e)) => {
1437 self.beacon_registered = false;
1438 tracing::warn!(
1439 error = %e,
1440 "ADR-091 Amendment 2: failed to refresh walpin registration beacon; \
1441 this process's sidecar health will read as unknown, not registered-silent"
1442 );
1443 }
1444 Err(join_err) => {
1445 self.beacon_registered = false;
1446 tracing::warn!(
1447 error = %join_err,
1448 "ADR-091 Amendment 2: walpin beacon refresh task panicked"
1449 );
1450 }
1451 }
1452 }
1453
1454 /// Fail-closed escalation for a failed heartbeat write: remove this
1455 /// process's beacon so enumeration cannot classify it
1456 /// `registered-silent` off the still-fresh prior refresh — skipping one
1457 /// touch alone leaves the previous mtime inside the freshness window
1458 /// for up to three producer ticks, an exoneration window. With the
1459 /// beacon gone the process either reports (once writes recover, the
1460 /// next tick re-registers and writes the heartbeat) or is caught by the
1461 /// OS-level holder census as an unattributed holder. If the removal
1462 /// itself fails, the beacon ages out over the freshness window — the
1463 /// narrowed fallback, not the contract.
1464 async fn drop_beacon_fail_closed(&mut self) {
1465 let dir = self.dir.clone();
1466 let pid = self.pid;
1467 self.beacon_registered = false;
1468 let result =
1469 tokio::task::spawn_blocking(move || crate::walpin::remove_beacon(&dir, pid)).await;
1470 match result {
1471 Ok(Ok(())) => {}
1472 Ok(Err(e)) => {
1473 tracing::warn!(
1474 error = %e,
1475 "ADR-091 Amendment 2: failed to remove walpin beacon after a failed \
1476 heartbeat write; beacon will age out of the freshness window instead"
1477 );
1478 }
1479 Err(join_err) => {
1480 tracing::warn!(
1481 error = %join_err,
1482 "ADR-091 Amendment 2: walpin beacon removal task panicked"
1483 );
1484 }
1485 }
1486 }
1487
1488 /// Blocking heartbeat write/removal runs on `spawn_blocking` (perf,
1489 /// ADR-091 Amendment 2) — this async sweep task must not block its
1490 /// executor thread on synchronous filesystem I/O.
1491 async fn observe(
1492 &mut self,
1493 oldest: Option<khive_storage::tx_registry::OldestSpan>,
1494 tx_warn_secs: Duration,
1495 ) {
1496 match oldest {
1497 Some(span) if span.age >= tx_warn_secs => {
1498 // ADR-091 Amendment 3 Plank F2: the caller's `TxOriginFilter`
1499 // guarantees a `Main` view's winner is either `Database` (this
1500 // backend's own identity) or `Unscoped` (the fallback), and a
1501 // `Secondary` view's winner is always `Database` — `Memory`
1502 // can never win a filtered query, so it degrades to
1503 // fallback-confidence rather than a reachability panic.
1504 let attribution_basis = match span.origin {
1505 khive_storage::tx_registry::TxOrigin::Database(_) => "origin",
1506 khive_storage::tx_registry::TxOrigin::Unscoped
1507 | khive_storage::tx_registry::TxOrigin::Memory => "fallback",
1508 };
1509
1510 // ADR-091 Amendment 3 Plank F1: a metadata-only mtime touch
1511 // advances freshness whenever nothing content-relevant has
1512 // changed since the last body write; a full rewrite happens
1513 // only on the first over-threshold observation or a genuine
1514 // content change.
1515 let content_unchanged = self.wrote
1516 && self.last_heartbeat.as_ref().is_some_and(|last| {
1517 last.content_matches(
1518 span.id,
1519 &span.label,
1520 attribution_basis,
1521 self.sweep_interval_ms,
1522 )
1523 });
1524
1525 if content_unchanged {
1526 let dir = self.dir.clone();
1527 let pid = self.pid;
1528 let touch_result = tokio::task::spawn_blocking(move || {
1529 crate::walpin::touch_heartbeat(&dir, pid)
1530 })
1531 .await;
1532 match touch_result {
1533 Ok(Ok(())) => {
1534 self.refresh_beacon().await;
1535 return;
1536 }
1537 Ok(Err(e)) => {
1538 tracing::warn!(
1539 error = %e,
1540 "ADR-091 Amendment 3 Plank F1: walpin heartbeat touch failed; \
1541 recreating with a full body write"
1542 );
1543 }
1544 Err(join_err) => {
1545 tracing::warn!(
1546 error = %join_err,
1547 "ADR-091 Amendment 3 Plank F1: walpin heartbeat touch task \
1548 panicked; recreating with a full body write"
1549 );
1550 }
1551 }
1552 // Recovery rule: the touch path must never assume the
1553 // target still exists — enumeration can delete a slow
1554 // writer's heartbeat while its span is still live. Fall
1555 // through to the full write below unconditionally.
1556 }
1557
1558 // The oldest span's registration instant is fixed for as
1559 // long as it stays the SAME span: reuse the previously
1560 // recorded value rather than re-deriving it from `now -
1561 // age`, which would drift by measurement noise across ticks
1562 // for no reason. A genuinely new oldest span (or the first
1563 // observation) derives it fresh.
1564 let oldest_tx_started_at = self
1565 .last_heartbeat
1566 .as_ref()
1567 .filter(|last| last.span_id == span.id)
1568 .map(|last| last.oldest_tx_started_at)
1569 .unwrap_or_else(|| now_epoch_secs().saturating_sub(span.age.as_secs() as i64));
1570
1571 let heartbeat = crate::walpin::WalpinHeartbeat {
1572 pid: self.pid,
1573 process_role: self.role.to_string(),
1574 started_at: self.started_at,
1575 oldest_tx_age_secs: span.age.as_secs_f64(),
1576 oldest_tx_label: span.label.clone(),
1577 oldest_tx_started_at: Some(oldest_tx_started_at),
1578 updated_at: now_epoch_secs(),
1579 sweep_interval_ms: self.sweep_interval_ms,
1580 attribution_basis: Some(attribution_basis.to_string()),
1581 };
1582 let dir = self.dir.clone();
1583 let result = tokio::task::spawn_blocking(move || {
1584 crate::walpin::write_heartbeat(&dir, &heartbeat)
1585 })
1586 .await;
1587 // The beacon refresh is gated on the heartbeat write
1588 // landing: a fresh beacon with no heartbeat file classifies
1589 // as `registered-silent` at enumeration, so a failed write
1590 // would exonerate a process that currently holds an
1591 // over-threshold transaction. Skipping the refresh alone is
1592 // not enough — the previous touch stays inside the freshness
1593 // window for up to three producer ticks — so the failure
1594 // path removes the beacon outright (`drop_beacon_fail_closed`);
1595 // the next successful tick re-registers it.
1596 match result {
1597 Ok(Ok(())) => {
1598 self.wrote = true;
1599 self.last_heartbeat = Some(LastHeartbeatState {
1600 span_id: span.id,
1601 label: span.label,
1602 attribution_basis,
1603 sweep_interval_ms: self.sweep_interval_ms,
1604 oldest_tx_started_at,
1605 });
1606 self.refresh_beacon().await;
1607 }
1608 Ok(Err(e)) => {
1609 tracing::warn!(
1610 error = %e,
1611 "ADR-091 Amendment 2 Plank B: failed to write walpin heartbeat; \
1612 removing beacon so this process cannot read as \
1613 registered-silent while over threshold"
1614 );
1615 // Unknown what (if anything) is on disk now — the
1616 // next tick must go through a full write, never a
1617 // touch, until a write actually lands.
1618 self.last_heartbeat = None;
1619 self.drop_beacon_fail_closed().await;
1620 }
1621 Err(join_err) => {
1622 tracing::warn!(
1623 error = %join_err,
1624 "ADR-091 Amendment 2 Plank B: walpin heartbeat write task panicked"
1625 );
1626 self.last_heartbeat = None;
1627 self.drop_beacon_fail_closed().await;
1628 }
1629 }
1630 }
1631 _ => {
1632 self.refresh_beacon().await;
1633 if self.wrote {
1634 let dir = self.dir.clone();
1635 let pid = self.pid;
1636 let result = tokio::task::spawn_blocking(move || {
1637 crate::walpin::remove_heartbeat(&dir, pid)
1638 })
1639 .await;
1640 match result {
1641 Ok(Ok(())) => {}
1642 Ok(Err(e)) => tracing::warn!(
1643 error = %e,
1644 "ADR-091 Amendment 2 Plank B: failed to remove walpin heartbeat"
1645 ),
1646 Err(join_err) => tracing::warn!(
1647 error = %join_err,
1648 "ADR-091 Amendment 2 Plank B: walpin heartbeat removal task panicked"
1649 ),
1650 }
1651 self.wrote = false;
1652 self.last_heartbeat = None;
1653 }
1654 }
1655 }
1656 }
1657
1658 async fn shutdown(&mut self) {
1659 if self.wrote {
1660 let dir = self.dir.clone();
1661 let pid = self.pid;
1662 let _ = tokio::task::spawn_blocking(move || crate::walpin::remove_heartbeat(&dir, pid))
1663 .await;
1664 self.wrote = false;
1665 }
1666 }
1667}
1668
1669#[cfg(unix)]
1670async fn run_walpin_housekeeping_if_due(
1671 sidecar: &WalpinSidecarState,
1672 state: &mut TruncateState,
1673 legacy_fallback_interval: Duration,
1674) -> bool {
1675 if !state.housekeeping_due() || !state.claim_walpin_full_scan_at(Instant::now()) {
1676 return false;
1677 }
1678 if let Some(report) = sidecar
1679 .reap_dead_entries_bounded(legacy_fallback_interval)
1680 .await
1681 {
1682 state.cache_walpin_attribution(
1683 report,
1684 Err("OS holder census is unavailable for a housekeeping-only scan".to_string()),
1685 Instant::now(),
1686 );
1687 }
1688 true
1689}
1690
1691fn now_epoch_secs() -> i64 {
1692 std::time::SystemTime::now()
1693 .duration_since(std::time::UNIX_EPOCH)
1694 .map(|d| d.as_secs() as i64)
1695 .unwrap_or(0)
1696}
1697
1698/// ADR-091 Amendment 2 Plank A: config for the observe-only per-session
1699/// sweep. Sessions never checkpoint — that stays daemon-owned so N session
1700/// processes never compete for the writer mutex — this only watches
1701/// `tx_registry` (and, Plank B, refreshes this process's walpin heartbeat).
1702const DEFAULT_SESSION_SWEEP_INTERVAL: Duration = Duration::from_secs(5);
1703
1704#[derive(Clone, Debug)]
1705pub struct SessionSweepConfig {
1706 /// How often a session polls the registry. Coarser than the daemon's
1707 /// tick: sessions do not need the daemon's 500ms checkpoint cadence.
1708 ///
1709 /// Overridable via `KHIVE_SESSION_SWEEP_INTERVAL_MS`. Default: 5000 ms.
1710 pub interval: Duration,
1711 /// Same semantics and default as [`CheckpointConfig::tx_warn_secs`].
1712 pub tx_warn_secs: Duration,
1713 /// Same semantics and default as [`CheckpointConfig::tx_max_age_secs`].
1714 pub tx_max_age_secs: Duration,
1715}
1716
1717impl Default for SessionSweepConfig {
1718 fn default() -> Self {
1719 Self {
1720 interval: DEFAULT_SESSION_SWEEP_INTERVAL,
1721 tx_warn_secs: Duration::from_secs(30),
1722 tx_max_age_secs: Duration::from_secs(120),
1723 }
1724 }
1725}
1726
1727impl SessionSweepConfig {
1728 /// Build from the environment. Reuses `KHIVE_TX_WARN_SECS` /
1729 /// `KHIVE_TX_MAX_AGE_SECS` (the same knobs the daemon's checkpoint task
1730 /// reads) so a session and the daemon agree on the same thresholds.
1731 pub fn from_env() -> Self {
1732 let mut cfg = Self {
1733 interval: session_sweep_interval_from_env(),
1734 ..Self::default()
1735 };
1736 // Shares `tx_age_thresholds_from_env` with `CheckpointConfig::from_env`
1737 // (minor, ADR-091 Amendment 2) so a session and the daemon
1738 // parse and validate `KHIVE_TX_WARN_SECS`/`KHIVE_TX_MAX_AGE_SECS`
1739 // identically from one source, not two hand-copied blocks.
1740 (cfg.tx_warn_secs, cfg.tx_max_age_secs) =
1741 tx_age_thresholds_from_env(cfg.tx_warn_secs, cfg.tx_max_age_secs);
1742
1743 cfg
1744 }
1745}
1746
1747fn session_sweep_interval_from_env() -> Duration {
1748 std::env::var("KHIVE_SESSION_SWEEP_INTERVAL_MS")
1749 .ok()
1750 .and_then(|ms| ms.parse::<u64>().ok())
1751 .filter(|ms| *ms > 0)
1752 .map(Duration::from_millis)
1753 .unwrap_or(DEFAULT_SESSION_SWEEP_INTERVAL)
1754}
1755
1756/// One file-backed backend the session sweep observes (ADR-091 Amendment 3
1757/// fan-out). `is_main` selects which [`khive_storage::tx_registry::TxOriginFilter`]
1758/// variant scopes this backend's view of the registry: the main backend's
1759/// `Main` filter additionally observes `Unscoped` spans (the
1760/// never-silently-drop fallback for call sites not yet threaded to an
1761/// origin); a secondary backend's `Secondary` filter is scoped to exactly
1762/// its own identity. A pool whose origin is `Memory` contributes no entry —
1763/// in-memory backends have no sidecar and nothing to attribute
1764/// cross-process.
1765pub struct SweepBackend {
1766 pub pool: Arc<ConnectionPool>,
1767 pub is_main: bool,
1768}
1769
1770/// Per-backend state the session sweep carries across ticks: this backend's
1771/// registry view, its own edge-triggered age-sweep state machine (so a
1772/// sustained stale span on one backend logs independently of the others),
1773/// and its own walpin sidecar (`None` if the sidecar is disabled or this
1774/// backend's origin is `Memory`).
1775struct BackendSweep {
1776 filter: khive_storage::tx_registry::TxOriginFilter,
1777 tx_age_state: TxAgeSweepState,
1778 sidecar: Option<WalpinSidecarState>,
1779}
1780
1781/// ADR-091 Amendment 2 Plank A (Amendment 3: per-backend fan-out): run the
1782/// observe-only per-session sweep.
1783///
1784/// Every non-daemon `kkernel mcp` process runs this instead of the daemon's
1785/// `run_checkpoint_task`: same `tx_registry` age check and Plank B heartbeat
1786/// refresh, but no PASSIVE/TRUNCATE checkpointing — checkpointing stays
1787/// daemon-owned. Stays ONE task for the whole process, but fans out
1788/// internally: each file-backed backend in `backends` gets its own
1789/// registry view, age-sweep state, and sidecar directory, so a long span on
1790/// a secondary backend is attributed (and heartbeats) only in that
1791/// backend's own sidecar — never the main backend's. Loops until
1792/// `shutdown_rx` observes a change (or its sender is dropped), removing
1793/// every written heartbeat on the way out.
1794pub async fn run_session_sweep_task(
1795 backends: Vec<SweepBackend>,
1796 config: SessionSweepConfig,
1797 mut shutdown_rx: tokio::sync::watch::Receiver<()>,
1798) {
1799 let mut interval = tokio::time::interval(config.interval);
1800 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
1801
1802 let mut sweeps: Vec<BackendSweep> = Vec::with_capacity(backends.len());
1803 for backend in backends {
1804 let identity = match backend.pool.origin() {
1805 khive_storage::tx_registry::TxOrigin::Database(id) => id,
1806 // No on-disk file, so no sidecar and no cross-process
1807 // attribution surface — nothing for this sweep to fan out to.
1808 khive_storage::tx_registry::TxOrigin::Memory
1809 | khive_storage::tx_registry::TxOrigin::Unscoped => continue,
1810 };
1811 let filter = if backend.is_main {
1812 khive_storage::tx_registry::TxOriginFilter::Main(identity)
1813 } else {
1814 khive_storage::tx_registry::TxOriginFilter::Secondary(identity)
1815 };
1816 let sidecar = WalpinSidecarState::new(
1817 backend.pool.canonical_path(),
1818 true,
1819 "session",
1820 config.interval,
1821 );
1822 sweeps.push(BackendSweep {
1823 filter,
1824 tx_age_state: TxAgeSweepState::default(),
1825 sidecar,
1826 });
1827 }
1828 for sweep in sweeps.iter_mut() {
1829 if let Some(sidecar) = sweep.sidecar.as_mut() {
1830 sidecar.register_beacon().await;
1831 }
1832 }
1833
1834 loop {
1835 tokio::select! {
1836 _ = interval.tick() => {}
1837 _ = shutdown_rx.changed() => break,
1838 }
1839
1840 for sweep in sweeps.iter_mut() {
1841 let oldest = khive_storage::tx_registry::oldest_for(&sweep.filter);
1842 for emission in sweep.tx_age_state.observe(
1843 oldest.as_ref().map(|s| (s.id, s.age, s.label.clone())),
1844 config.tx_warn_secs,
1845 config.tx_max_age_secs,
1846 ) {
1847 log_tx_age_emission(&emission);
1848 }
1849 if let Some(sidecar) = sweep.sidecar.as_mut() {
1850 sidecar.observe(oldest, config.tx_warn_secs).await;
1851 }
1852 }
1853 }
1854
1855 for sweep in sweeps.iter_mut() {
1856 if let Some(sidecar) = sweep.sidecar.as_mut() {
1857 sidecar.shutdown().await;
1858 }
1859 }
1860}
1861
1862/// The event sink and namespace owned by one checkpoint task in a fan-out.
1863///
1864/// Backend role and lifecycle ownership are separate: a secondary task may
1865/// own lifecycle emission when the deployment's main backend is in-memory.
1866#[derive(Clone)]
1867pub struct CheckpointLifecycleOwner {
1868 event_store: Arc<dyn khive_storage::EventStore>,
1869 namespace: String,
1870}
1871
1872impl CheckpointLifecycleOwner {
1873 /// Designate `event_store` as the lifecycle sink for one checkpoint task.
1874 pub fn new(
1875 event_store: Arc<dyn khive_storage::EventStore>,
1876 namespace: impl Into<String>,
1877 ) -> Self {
1878 Self {
1879 event_store,
1880 namespace: namespace.into(),
1881 }
1882 }
1883}
1884
1885/// Maximum number of checkpoint lifecycle events waiting behind the append
1886/// currently owned by the worker. One queued row preserves a recent outcome
1887/// without allowing sustained writer contention to grow memory without bound.
1888const CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY: usize = 1;
1889
1890#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1891struct CheckpointPressureEpisode {
1892 elevated_ticks: u64,
1893 peak_wal_pages: u64,
1894}
1895
1896impl CheckpointPressureEpisode {
1897 fn start(wal_pages: u64) -> Self {
1898 Self {
1899 elevated_ticks: 1,
1900 peak_wal_pages: wal_pages,
1901 }
1902 }
1903
1904 fn observe(&mut self, wal_pages: u64) {
1905 self.elevated_ticks = self.elevated_ticks.saturating_add(1);
1906 self.peak_wal_pages = self.peak_wal_pages.max(wal_pages);
1907 }
1908}
1909
1910/// Zero-wait handoff from the checkpoint scheduler to its lifecycle sink.
1911///
1912/// The worker serializes appends, preserving the order of every event that is
1913/// accepted. The scheduler only calls [`tokio::sync::mpsc::Sender::try_send`]:
1914/// if the worker and its single queue slot are both occupied, telemetry is
1915/// dropped rather than delaying the next checkpoint cycle. The first drop in
1916/// each uninterrupted full-queue episode warns; a successful enqueue re-arms
1917/// that warning without producing per-tick log spam.
1918struct CheckpointLifecycleEmitter {
1919 namespace: Option<String>,
1920 sender: Option<tokio::sync::mpsc::Sender<khive_storage::Event>>,
1921 worker: Option<tokio::task::JoinHandle<()>>,
1922 busy_warning_emitted: bool,
1923}
1924
1925impl CheckpointLifecycleEmitter {
1926 fn new(owner: Option<CheckpointLifecycleOwner>) -> Self {
1927 let Some(owner) = owner else {
1928 return Self {
1929 namespace: None,
1930 sender: None,
1931 worker: None,
1932 busy_warning_emitted: false,
1933 };
1934 };
1935
1936 let namespace = owner.namespace.clone();
1937 let (sender, mut receiver) =
1938 tokio::sync::mpsc::channel::<khive_storage::Event>(CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY);
1939 let worker = tokio::spawn(async move {
1940 while let Some(event) = receiver.recv().await {
1941 let kind = event.kind;
1942 CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.fetch_add(1, Ordering::Relaxed);
1943 if let Err(err) = owner.event_store.append_event(event).await {
1944 CHECKPOINT_LIFECYCLE_APPEND_FAILURES.fetch_add(1, Ordering::Relaxed);
1945 tracing::warn!(
1946 error = %err,
1947 event_kind = %kind.name(),
1948 "checkpoint lifecycle event append failed"
1949 );
1950 }
1951 }
1952 });
1953
1954 Self {
1955 namespace: Some(namespace),
1956 sender: Some(sender),
1957 worker: Some(worker),
1958 busy_warning_emitted: false,
1959 }
1960 }
1961
1962 /// Serialize and enqueue one lifecycle event without awaiting sink I/O.
1963 /// Returns whether the row was accepted for delivery (or no sink exists).
1964 fn try_emit<P: serde::Serialize>(&mut self, kind: khive_types::EventKind, payload: P) -> bool {
1965 let (Some(namespace), Some(sender)) = (&self.namespace, &self.sender) else {
1966 return true;
1967 };
1968 let payload_value = match serde_json::to_value(&payload) {
1969 Ok(value) => value,
1970 Err(err) => {
1971 tracing::warn!(
1972 error = %err,
1973 event_kind = %kind.name(),
1974 "failed to serialize checkpoint lifecycle event payload"
1975 );
1976 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.fetch_add(1, Ordering::Relaxed);
1977 return false;
1978 }
1979 };
1980 let payload_schema_version = match kind {
1981 khive_types::EventKind::CheckpointOutcomeRecorded => 2,
1982 _ => 1,
1983 };
1984 let event = khive_storage::Event::new(
1985 namespace,
1986 "checkpoint.lifecycle",
1987 kind,
1988 khive_types::SubstrateKind::Event,
1989 "daemon:checkpoint_task",
1990 )
1991 .with_payload(payload_value)
1992 .with_payload_schema_version(payload_schema_version);
1993
1994 match sender.try_send(event) {
1995 Ok(()) => {
1996 self.busy_warning_emitted = false;
1997 true
1998 }
1999 Err(tokio::sync::mpsc::error::TrySendError::Full(event)) => {
2000 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.fetch_add(1, Ordering::Relaxed);
2001 if !self.busy_warning_emitted {
2002 tracing::warn!(
2003 event_kind = %event.kind.name(),
2004 queue_capacity = CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY,
2005 "checkpoint lifecycle event dropped because the append worker is busy"
2006 );
2007 self.busy_warning_emitted = true;
2008 }
2009 false
2010 }
2011 Err(tokio::sync::mpsc::error::TrySendError::Closed(event)) => {
2012 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.fetch_add(1, Ordering::Relaxed);
2013 tracing::warn!(
2014 event_kind = %event.kind.name(),
2015 "checkpoint lifecycle event dropped because the append worker stopped"
2016 );
2017 false
2018 }
2019 }
2020 }
2021
2022 /// Stop the scheduler-owned async worker without making
2023 /// [`run_checkpoint_task`] wait for its current append future.
2024 ///
2025 /// This bounds checkpoint-task shutdown only. If the event store already
2026 /// admitted the append to `spawn_blocking` or a `WriterTask`, aborting this
2027 /// worker cannot cancel that downstream operation; at most one such sink
2028 /// operation may outlive the checkpoint task.
2029 async fn shutdown(mut self) {
2030 drop(self.sender.take());
2031 let Some(worker) = self.worker.take() else {
2032 return;
2033 };
2034 worker.abort();
2035 match worker.await {
2036 Ok(()) => {}
2037 Err(err) if err.is_cancelled() => {}
2038 Err(err) => tracing::warn!(
2039 error = %err,
2040 "checkpoint lifecycle event append worker terminated unexpectedly"
2041 ),
2042 }
2043 }
2044}
2045
2046impl Drop for CheckpointLifecycleEmitter {
2047 fn drop(&mut self) {
2048 // The normal watch-signal path calls `shutdown` and takes the handle
2049 // first. This fallback covers an externally-aborted or panicking
2050 // checkpoint task so the scheduler-owned async worker itself is never
2051 // detached. One already-admitted downstream sink operation may outlive
2052 // it; see `shutdown`'s contract above.
2053 if let Some(worker) = &self.worker {
2054 worker.abort();
2055 }
2056 }
2057}
2058
2059/// The checkpoint task's dedicated, long-lived standalone connection to the
2060/// same database file — opened once at task startup and reused for every
2061/// tick's PASSIVE (and, when armed, TRUNCATE) pragma. NEVER the pool's writer
2062/// mutex, which is what removes the pool-mutex ADMISSION path: a concurrent
2063/// `pool.writer()` checkout no longer queues behind a checkpoint tick.
2064///
2065/// That removal is scoped to admission, not to SQLite-level blocking in
2066/// general. `PRAGMA wal_checkpoint(PASSIVE)` takes only SQLite's CKPT lock,
2067/// not the WRITE lock, so a concurrent writer can commit while a PASSIVE pass
2068/// runs on this connection — true of PASSIVE specifically, not of TRUNCATE.
2069/// TRUNCATE inherits RESTART semantics and additionally acquires SQLite's
2070/// writer lock, so it can still block a concurrent write transaction, on any
2071/// connection, for up to `truncate_busy_timeout` while it waits on a pinning
2072/// reader — the same bounded cost that existed pre-fix, now paid on this
2073/// dedicated connection instead of the pool writer. Serializing checkpoint
2074/// admission behind the pool's writer mutex (the pre-fix design) imposed
2075/// contention SQLite itself does not require; TRUNCATE's own SQLite-level
2076/// write-blocking window is unaffected by that removal.
2077///
2078/// `None` between ticks means the connection is unavailable (never opened
2079/// yet, or dropped after a prior tick's connection-level pragma failure) —
2080/// the caller must report that tick `Skipped` and retry the open on the next
2081/// one. A busy or inconsistent PASSIVE result also skips a tick while keeping
2082/// this connection open.
2083///
2084/// ADR-136 D1 gate 5 classification: **checkpoint writer**. Explicitly
2085/// exempt from `WriterTask`/queue routing by design (see the admission-path
2086/// note above), never `SqlAccess`-reachable, never counted as a
2087/// `direct_route_violation` — see the classification table in
2088/// `writer_task`'s module doc.
2089struct CheckpointConnection {
2090 conn: Option<rusqlite::Connection>,
2091 /// Consecutive failed `open_standalone_writer` attempts since the last
2092 /// successful open (or since task startup). Drives the WARN-once /
2093 /// debug-thereafter log rate-limiting in `ensure_open`: a file-backed
2094 /// pool that transiently loses its dedicated connection would otherwise
2095 /// log a WARN on every tick (default 500ms) for as long as the outage
2096 /// lasts, which for a read-only or in-memory pool — where the open can
2097 /// never succeed — means permanent per-tick WARN spam.
2098 consecutive_open_failures: u32,
2099}
2100
2101impl CheckpointConnection {
2102 fn new() -> Self {
2103 Self {
2104 conn: None,
2105 consecutive_open_failures: 0,
2106 }
2107 }
2108
2109 /// Ensure a usable connection is open, lazily (re)opening from `pool`
2110 /// when the current one is absent. Reuses the crate's existing untracked
2111 /// standalone-connection open path (`ConnectionPool::open_standalone_writer_untracked`),
2112 /// which applies the same pragmas (including `busy_timeout` from the pool
2113 /// config) as any other standalone connection, without counting this
2114 /// infrastructure connection as write-operation traffic. Returns `None` if
2115 /// opening fails — an in-memory pool (no on-disk file to open a second
2116 /// connection against), a read-only pool, or a transient filesystem error.
2117 ///
2118 /// Logging is rate-limited across a failure streak: the FIRST failure of
2119 /// a streak logs at `warn!`, every subsequent identical failure (while
2120 /// still failing) logs at `debug!` instead, and a successful open that
2121 /// ends a streak logs one `info!` recovery line. Without this, a
2122 /// permanently-unopenable pool (read-only or in-memory, selected by
2123 /// `checkpoint_pool_for`) would WARN on every tick forever.
2124 fn ensure_open(&mut self, pool: &ConnectionPool) -> Option<&rusqlite::Connection> {
2125 if self.conn.is_none() {
2126 match pool.open_standalone_writer_untracked() {
2127 Ok(conn) => {
2128 // This is the dedicated owner's own connection: disable
2129 // autocheckpoint on it unconditionally, independent of
2130 // whether the pool-level ownership claim has landed yet
2131 // (the standalone open applies the claim-dependent
2132 // value; this connection must never run an implicit
2133 // checkpoint inside its own PASSIVE/TRUNCATE work).
2134 if let Err(e) = conn.pragma_update(None, "wal_autocheckpoint", 0) {
2135 tracing::warn!(
2136 error = %e,
2137 "could not disable autocheckpoint on the dedicated checkpoint \
2138 connection"
2139 );
2140 }
2141 if self.consecutive_open_failures > 0 {
2142 tracing::info!(
2143 prior_consecutive_failures = self.consecutive_open_failures,
2144 "dedicated checkpoint connection opened successfully, ending a \
2145 failure streak"
2146 );
2147 }
2148 self.consecutive_open_failures = 0;
2149 self.conn = Some(conn);
2150 }
2151 Err(e) => {
2152 if self.consecutive_open_failures == 0 {
2153 tracing::warn!(
2154 error = %e,
2155 "failed to open the dedicated checkpoint connection; \
2156 this tick is skipped and the open retried next tick"
2157 );
2158 } else {
2159 tracing::debug!(
2160 error = %e,
2161 consecutive_failures = self.consecutive_open_failures,
2162 "dedicated checkpoint connection still unavailable; \
2163 this tick is skipped and the open retried next tick"
2164 );
2165 }
2166 self.consecutive_open_failures =
2167 self.consecutive_open_failures.saturating_add(1);
2168 return None;
2169 }
2170 }
2171 }
2172 self.conn.as_ref()
2173 }
2174}
2175
2176/// Run one due FTS5 maintenance step off this task's Tokio worker thread.
2177///
2178/// A due step can issue up to `config.merge_pages` pages of synchronous
2179/// SQLite incremental-merge I/O against a trigram index over a corpus of
2180/// hundreds of thousands of rows — the same class of blocking work the
2181/// WAL-pin beacon writes above already move off the worker via
2182/// `tokio::task::spawn_blocking`. `conn` and `state` are moved into the
2183/// blocking closure and handed back to the caller whenever the step returns,
2184/// so the checkpoint task can restore its dedicated connection and scheduler
2185/// state on every non-panicking path. A panic inside the step surfaces as the
2186/// `JoinError`, like the beacon writes above: the connection and scheduler
2187/// state moved into the task are gone with it, and the caller reopens both
2188/// on the next tick instead of taking the checkpoint task down.
2189async fn run_fts_maintenance_off_worker(
2190 conn: rusqlite::Connection,
2191 config: crate::fts_maintenance::FtsMaintenanceConfig,
2192 mut state: crate::fts_maintenance::FtsMaintenanceState,
2193 now: Instant,
2194) -> Result<
2195 (
2196 rusqlite::Connection,
2197 crate::fts_maintenance::FtsMaintenanceState,
2198 Result<Option<crate::fts_maintenance::FtsMaintenanceStep>, String>,
2199 ),
2200 tokio::task::JoinError,
2201> {
2202 tokio::task::spawn_blocking(move || {
2203 let result = crate::fts_maintenance::run_if_due(&conn, &config, &mut state, now);
2204 (conn, state, result)
2205 })
2206 .await
2207}
2208
2209/// Run the WAL checkpoint background task.
2210///
2211/// Long-running async task — spawn with `tokio::spawn`. Loops until
2212/// `shutdown_rx` observes a change (or its sender is dropped). Callers MUST
2213/// hold the paired `tokio::sync::watch::Sender` for the daemon's run scope
2214/// and send on it to shut down — do NOT rely on `pool`'s `Arc` refcount
2215/// reaching zero; a sibling owner (e.g. `event_store`) holding its own clone
2216/// makes that check unreachable (issue #774).
2217///
2218/// Issues `PRAGMA wal_checkpoint(PASSIVE)` every tick on the task's dedicated
2219/// `CheckpointConnection` — never the pool's writer mutex, so a concurrent
2220/// `pool.writer()` checkout can never queue behind a checkpoint tick. That
2221/// guarantee is admission-only: an armed TRUNCATE still takes SQLite's writer
2222/// lock and can block new write transactions, on any connection, for up to
2223/// `truncate_busy_timeout` (see `CheckpointConnection`'s contract). The
2224/// checkpoint call itself runs on `spawn_blocking`, so that wait never holds
2225/// one of the runtime's worker threads. A tick is
2226/// `Skipped` when that connection is unavailable or SQLite returns a busy
2227/// PASSIVE row without a usable pressure observation. A
2228/// WARNING fires once per below→above threshold crossing, not every tick.
2229///
2230/// `lifecycle_owner` (ADR-094): exactly one task in a multi-backend fan-out
2231/// should receive `Some`. That task appends a best-effort
2232/// `CheckpointOutcomeRecorded` event on the elevation transition and one
2233/// recovery summary when pressure falls back below `warn_pages`. Sustained
2234/// elevated ticks aggregate in memory and in `db_diagnostics`; they never
2235/// write one primary-store row per checkpoint attempt. `None` explicitly
2236/// marks a non-owner. See `crates/khive-db/docs/api/checkpoint.md` for the
2237/// full shutdown-mechanism and event-emission design history.
2238///
2239/// `is_main` (ADR-091 Amendment 3): whether `pool` is the deployment's main
2240/// backend. A daemon owning several file-backed backends spawns one task per
2241/// backend, each with its own pool and shutdown-channel clone (the sender
2242/// broadcasts to every receiver clone alike). Lifecycle ownership is selected
2243/// independently through `lifecycle_owner`; `is_main` only controls registry
2244/// filtering. See the `tx_filter` construction below.
2245pub async fn run_checkpoint_task(
2246 pool: Arc<ConnectionPool>,
2247 config: CheckpointConfig,
2248 lifecycle_owner: Option<CheckpointLifecycleOwner>,
2249 mut shutdown_rx: tokio::sync::watch::Receiver<()>,
2250 is_main: bool,
2251) {
2252 let _checkpoint_run_guard = CheckpointRunTaskGuard::start(&pool, config.interval);
2253 // This task IS the dedicated checkpoint owner: claim the pool so writer
2254 // connections drop the bounded autocheckpoint fallback and routine
2255 // checkpoint I/O stays off application commit paths. Pools without a
2256 // running checkpoint task never claim and keep SQLite's bounded WAL
2257 // reclamation. A failed claim leaves connections on the bounded fallback
2258 // — safe, just not the low-latency posture — so it warns and continues.
2259 match pool.claim_checkpoint_ownership() {
2260 Ok(()) => {
2261 if let Err(e) = pool.propagate_checkpoint_claim_to_writer_task().await {
2262 tracing::warn!(
2263 error = %e,
2264 "checkpoint task could not reach the writer task's connection; it keeps the \
2265 bounded autocheckpoint fallback"
2266 );
2267 }
2268 }
2269 Err(e) => {
2270 tracing::warn!(
2271 error = %e,
2272 "checkpoint task could not re-apply the ownership pragma on the pooled writer; \
2273 writer connections keep the bounded autocheckpoint fallback unless ownership is \
2274 claimed later"
2275 );
2276 }
2277 }
2278 let mut interval = tokio::time::interval(config.interval);
2279 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
2280 let mut severity_state = CheckpointSeverityState::default();
2281 let mut tx_age_state = TxAgeSweepState::default();
2282 let mut was_above_high_water = false;
2283 #[cfg(unix)]
2284 let legacy_walpin_fallback_interval = DEFAULT_SESSION_SWEEP_INTERVAL;
2285 #[cfg(unix)]
2286 let mut truncate_state =
2287 TruncateState::with_legacy_walpin_fallback(legacy_walpin_fallback_interval);
2288 #[cfg(not(unix))]
2289 let mut truncate_state = TruncateState::default();
2290 let mut lifecycle_emitter = CheckpointLifecycleEmitter::new(lifecycle_owner);
2291 // Independent of `severity_state` (which owns the WARN ladder): this
2292 // tracks the lifecycle sink's accepted elevation state. A full queue
2293 // leaves it unchanged, so an opening or recovery transition is retried
2294 // without admitting more than one primary-store write for that edge.
2295 let mut event_elevation_open = false;
2296 let mut pressure_episode: Option<CheckpointPressureEpisode> = None;
2297 // A recovery row whose `try_emit` lost the race against a full queue.
2298 // Retried on later ticks (before that tick's own transition handling)
2299 // instead of leaving `pressure_episode` open for a stale episode to
2300 // absorb the next, genuinely separate, pressure incident (#1857).
2301 let mut pending_recovery: Option<khive_storage::CheckpointOutcomeRecordedPayload> = None;
2302 let mut was_observed_above_warn = false;
2303 // ADR-091 Amendment 3: this task's own backend-scoped view of the
2304 // registry. `is_main` selects which `TxOriginFilter` variant applies —
2305 // the caller passes `true` for exactly the one checkpoint task covering
2306 // the deployment's main backend, so only that task also observes legacy
2307 // `Unscoped` spans from any call site not yet threaded to an origin, the
2308 // designed never-silently-drop fallback. A secondary backend's task
2309 // never falls back to `Unscoped`: those spans belong to the main view or
2310 // to no view, never to a database they were never registered against.
2311 // `None` only when this pool's own origin isn't `Database` (an in-memory
2312 // checkpoint pool) — degrades to "no open span observed" for the tick
2313 // rather than panicking a long-running daemon loop on an
2314 // assumed-impossible state.
2315 let tx_filter = match pool.origin() {
2316 khive_storage::tx_registry::TxOrigin::Database(id) => Some(if is_main {
2317 khive_storage::tx_registry::TxOriginFilter::Main(id)
2318 } else {
2319 khive_storage::tx_registry::TxOriginFilter::Secondary(id)
2320 }),
2321 khive_storage::tx_registry::TxOrigin::Memory
2322 | khive_storage::tx_registry::TxOrigin::Unscoped => None,
2323 };
2324 // ADR-091 Amendment 2 Plank B: the checkpoint pool is only ever wired for
2325 // file-backed backends (`checkpoint_pool_for`), so `is_file_backed: true`
2326 // is always correct here. `canonical_path()` (not `pool.config().path`)
2327 // so the sidecar directory is keyed off the same minted identity every
2328 // alias of this backend's configured path converges to.
2329 #[cfg(unix)]
2330 let mut walpin_state =
2331 WalpinSidecarState::new(pool.canonical_path(), true, "daemon", config.interval);
2332 #[cfg(unix)]
2333 if let Some(sidecar) = walpin_state.as_mut() {
2334 sidecar.register_beacon().await;
2335 }
2336
2337 // Opened once here, at task startup; `ensure_open` is a no-op in steady
2338 // state and only reopens after a connection-level failure or (for an
2339 // in-memory/read-only pool) retries the open on every subsequent tick.
2340 let mut checkpoint_conn = CheckpointConnection::new();
2341 checkpoint_conn.ensure_open(&pool);
2342 // FTS5 segment maintenance shares this task's standalone connection, but
2343 // not its 500 ms cadence. Each due call performs at most one bounded
2344 // merge step on one table and refuses immediately when another writer
2345 // owns SQLite's write lock.
2346 let mut fts_maintenance_config = crate::fts_maintenance::FtsMaintenanceConfig::from_env();
2347 // Secondary backends have independent schemas (for example the code-map
2348 // database) and are not required to contain the substrate FTS tables.
2349 // Exactly the main backend owns this derived-index maintenance.
2350 fts_maintenance_config.enabled &= is_main;
2351 let mut fts_maintenance_state =
2352 crate::fts_maintenance::FtsMaintenanceState::new(Instant::now());
2353
2354 loop {
2355 // A closed sender (the daemon returning without an explicit send)
2356 // makes `changed()` resolve with `Err` immediately, which `select!`
2357 // treats as ready — so shutdown is observed either way, not just on
2358 // an explicit send.
2359 tokio::select! {
2360 _ = interval.tick() => {}
2361 _ = shutdown_rx.changed() => break,
2362 }
2363
2364 #[cfg(unix)]
2365 truncate_state.begin_tick();
2366
2367 #[cfg(unix)]
2368 let mut pending_sidecar_attribution = None;
2369
2370 let tick = if checkpoint_conn.ensure_open(&pool).is_none() {
2371 note_checkpoint_skipped();
2372 CheckpointTick::Skipped
2373 } else {
2374 // `ensure_open` above just confirmed a connection is open. Take
2375 // ownership of it so the checkpoint cycle and the FTS maintenance
2376 // step below can each move it onto a blocking thread; every path
2377 // either restores it to `checkpoint_conn` or lets it drop, which is
2378 // the moved-ownership equivalent of the former `drop_connection()`
2379 // call.
2380 let conn = checkpoint_conn
2381 .conn
2382 .take()
2383 .expect("ensure_open just confirmed a connection is open");
2384 match off_worker::run_checkpoint_core_off_worker(
2385 Arc::clone(&pool),
2386 conn,
2387 config.clone(),
2388 truncate_state,
2389 )
2390 .await
2391 {
2392 Ok((conn, state, Ok(outcome))) => {
2393 truncate_state = state;
2394 #[cfg(unix)]
2395 {
2396 pending_sidecar_attribution = outcome.sidecar_attribution;
2397 }
2398 #[cfg(not(unix))]
2399 let _ = outcome.sidecar_attribution;
2400
2401 // Only a due tick moves the connection onto a blocking
2402 // thread; an ordinary tick between maintenance intervals
2403 // keeps it here and records a no-op.
2404 let fts_result = if !fts_maintenance_state
2405 .is_due(&fts_maintenance_config, Instant::now())
2406 {
2407 checkpoint_conn.conn = Some(conn);
2408 Ok(None)
2409 } else {
2410 match run_fts_maintenance_off_worker(
2411 conn,
2412 fts_maintenance_config.clone(),
2413 fts_maintenance_state,
2414 Instant::now(),
2415 )
2416 .await
2417 {
2418 Ok((conn, state, fts_result)) => {
2419 fts_maintenance_state = state;
2420 checkpoint_conn.conn = Some(conn);
2421 fts_result
2422 }
2423 Err(join_err) => {
2424 // The step panicked on the blocking thread. The
2425 // connection and scheduler state moved into it
2426 // are gone; `ensure_open` reopens the connection
2427 // next tick and the schedule restarts from now.
2428 // The checkpoint pragma itself already succeeded.
2429 fts_maintenance_state =
2430 crate::fts_maintenance::FtsMaintenanceState::new(Instant::now());
2431 Err(format!(
2432 "bounded FTS5 segment maintenance task panicked: {join_err}"
2433 ))
2434 }
2435 }
2436 };
2437
2438 match fts_result {
2439 Ok(Some(step)) => match step.outcome {
2440 crate::fts_maintenance::FtsMaintenanceOutcome::Worked => {
2441 tracing::info!(
2442 table = step.table,
2443 requested_pages = step.requested_pages,
2444 segments_before = step.segments_before,
2445 segments_after = step.segments_after,
2446 "bounded FTS5 segment maintenance made progress"
2447 );
2448 }
2449 crate::fts_maintenance::FtsMaintenanceOutcome::Busy => {
2450 tracing::debug!(
2451 table = step.table,
2452 requested_pages = step.requested_pages,
2453 segments = step.segments_before,
2454 "bounded FTS5 segment maintenance skipped a busy writer"
2455 );
2456 }
2457 crate::fts_maintenance::FtsMaintenanceOutcome::Noop
2458 | crate::fts_maintenance::FtsMaintenanceOutcome::BelowThreshold => {
2459 tracing::debug!(
2460 table = step.table,
2461 outcome = ?step.outcome,
2462 segments = step.segments_before,
2463 "bounded FTS5 segment maintenance had no work"
2464 );
2465 }
2466 },
2467 Ok(None) => {}
2468 Err(error) => {
2469 // The checkpoint pragma already succeeded. An FTS
2470 // structure/read/merge error is an independent,
2471 // best-effort maintenance failure and must not make
2472 // the task discard an otherwise healthy connection.
2473 tracing::warn!(
2474 error = %error,
2475 "bounded FTS5 segment maintenance failed"
2476 );
2477 }
2478 }
2479 match outcome.wal_pages {
2480 Some(wal_pages) => CheckpointTick::Observed(wal_pages),
2481 None => {
2482 note_checkpoint_skipped();
2483 CheckpointTick::Skipped
2484 }
2485 }
2486 }
2487 Ok((_conn, state, Err(e))) => {
2488 truncate_state = state;
2489 tracing::warn!(
2490 error = %e,
2491 "dedicated checkpoint connection failed a pragma; \
2492 dropping it for a fresh reopen next tick"
2493 );
2494 note_checkpoint_skipped();
2495 CheckpointTick::Skipped
2496 }
2497 Err(panicked) => {
2498 // The cycle panicked: its connection is gone (reopened next tick) and the
2499 // escalation state is the one it was handed, counted as a TRUNCATE attempt.
2500 truncate_state = panicked.truncate_state;
2501 tracing::warn!(
2502 error = %panicked.join_error,
2503 "WAL checkpoint cycle panicked on its blocking thread; \
2504 dropping the connection for a fresh reopen next tick"
2505 );
2506 note_checkpoint_skipped();
2507 CheckpointTick::Skipped
2508 }
2509 }
2510 };
2511
2512 // A no-progress or unmeasured TRUNCATE returns a bounded attribution
2513 // request alongside the core outcome. Consume it before any
2514 // report-derived decision or ordinary housekeeping for this tick.
2515 // The await is intentional: enumeration may perform up to 512
2516 // filesystem reads/classifications, so none of that work is allowed
2517 // to run on this Tokio worker, while one-pass-per-tick ordering still
2518 // requires the result (or an honest worker/enumeration failure) before
2519 // the fallback housekeeping arm is considered.
2520 #[cfg(unix)]
2521 if let Err(error) =
2522 complete_walpin_attribution(pending_sidecar_attribution, &mut truncate_state).await
2523 {
2524 tracing::warn!(
2525 error = %error,
2526 failure_kind = error.kind(),
2527 "ADR-091 Amendment 2 Plank B: no-progress sidecar attribution failed"
2528 );
2529 }
2530
2531 // ADR-091 Plank 1: age-based sweep over the registry's oldest entry
2532 // MUST run on every tick, including a Skipped one — deliberately
2533 // BEFORE the Skipped early-continue below. Since the dedicated
2534 // checkpoint connection amendment, a `Skipped` tick means that
2535 // connection was unavailable or its PASSIVE result was busy, not that
2536 // some registered span held the pool's writer mutex (a checkpoint
2537 // tick no longer touches it at all) — but the sweep must not go blind for the
2538 // duration of that outage, and the two failure surfaces are
2539 // independent: a registry span can go stale
2540 // (KHIVE_TX_WARN_SECS / KHIVE_TX_MAX_AGE_SECS) while wal_pages sits
2541 // well under warn_pages, or while the checkpoint connection itself is
2542 // down. Edge-triggered per rung, same debounce idiom as the severity
2543 // ladder below, so a sustained stale span logs once per rung rather
2544 // than once per tick.
2545 let oldest_tx = tx_filter
2546 .as_ref()
2547 .and_then(khive_storage::tx_registry::oldest_for);
2548 for emission in tx_age_state.observe(
2549 oldest_tx.as_ref().map(|s| (s.id, s.age, s.label.clone())),
2550 config.tx_warn_secs,
2551 config.tx_max_age_secs,
2552 ) {
2553 log_tx_age_emission(&emission);
2554 }
2555 // ADR-091 Amendment 2 Plank B: refresh (or clear) this daemon
2556 // process's own walpin heartbeat on the same cadence, so its own
2557 // pin — if any — is attributable the same way a session's is.
2558 #[cfg(unix)]
2559 if let Some(sidecar) = walpin_state.as_mut() {
2560 sidecar
2561 .observe(oldest_tx.clone(), config.tx_warn_secs)
2562 .await;
2563 let _ = run_walpin_housekeeping_if_due(
2564 sidecar,
2565 &mut truncate_state,
2566 legacy_walpin_fallback_interval,
2567 )
2568 .await;
2569 }
2570
2571 // Skipped ticks leave crossing state unchanged — a busy tick must not
2572 // re-arm the rate limit while WAL pressure is still elevated.
2573 let wal_pages = match tick {
2574 CheckpointTick::Skipped => continue,
2575 CheckpointTick::Observed(n) => n,
2576 };
2577
2578 let above_warn = wal_pages >= config.warn_pages;
2579 let above_high_water = wal_pages >= config.high_water_pages;
2580 let above_truncate_high_water = wal_pages >= config.truncate_high_water_pages;
2581 note_checkpoint_pressure_observation(above_warn, was_observed_above_warn);
2582 was_observed_above_warn = above_warn;
2583
2584 // Per-tick debug for the oldest open entry always fires (cheap —
2585 // reuses this tick's already-computed `oldest_tx`); the two
2586 // `warn!`-level registry logs below are gated on the SAME crossing
2587 // state as the WAL-threshold WARNs above, so sustained pressure
2588 // logs once per crossing, not once per tick.
2589 log_tx_registry_oldest_debug(wal_pages, oldest_tx.as_ref());
2590
2591 // ADR-091 severity ladder: INFO on the first below→above crossing,
2592 // WARN once `warn_sustained_cycles` consecutive ticks stay elevated.
2593 // The oldest-entry registry WARN rides the same INFO edge the old
2594 // binary crossing_warn used to gate on.
2595 for emission in severity_state.observe_wal_pages(wal_pages, &config) {
2596 match emission.rung {
2597 CheckpointSeverityRung::Info => {
2598 log_tx_registry_oldest_warn(wal_pages, oldest_tx.as_ref());
2599 tracing::info!(
2600 wal_pages = emission.wal_pages,
2601 warn_threshold = emission.threshold_pages,
2602 "WAL page count crossed warn threshold"
2603 );
2604 }
2605 CheckpointSeverityRung::Warn => {
2606 tracing::warn!(
2607 wal_pages = emission.wal_pages,
2608 warn_threshold = emission.threshold_pages,
2609 consecutive_cycles = emission.consecutive_cycles,
2610 "WAL page count failed to drain below warn threshold"
2611 );
2612 }
2613 CheckpointSeverityRung::Alarm => {
2614 // Never produced by `observe_wal_pages`; see its doc.
2615 }
2616 }
2617 }
2618
2619 let high_water_crossed = crossing_warn(above_high_water, &mut was_above_high_water);
2620 if high_water_crossed {
2621 log_tx_registry_snapshot_warn(wal_pages);
2622 log_wal_high_water_warn(
2623 wal_pages,
2624 config.high_water_pages,
2625 oldest_tx.as_ref(),
2626 config.tx_warn_secs,
2627 );
2628 }
2629
2630 // ADR-094/#1838, #1857: one elevation row and one recovery summary
2631 // per genuinely continuous episode. Sustained elevated ticks update
2632 // only the bounded in-memory aggregate and process diagnostics
2633 // above; a dropped recovery handoff must not fold the next,
2634 // separate, pressure incident into this episode's aggregate.
2635 observe_checkpoint_pressure_tick(
2636 above_warn,
2637 wal_pages,
2638 above_high_water,
2639 above_truncate_high_water,
2640 &config,
2641 &mut event_elevation_open,
2642 &mut pressure_episode,
2643 &mut pending_recovery,
2644 |payload| {
2645 lifecycle_emitter
2646 .try_emit(khive_types::EventKind::CheckpointOutcomeRecorded, payload)
2647 },
2648 );
2649 }
2650
2651 lifecycle_emitter.shutdown().await;
2652
2653 #[cfg(unix)]
2654 if let Some(sidecar) = walpin_state.as_mut() {
2655 sidecar.shutdown().await;
2656 }
2657}
2658
2659/// Whether a `CheckpointOutcomeRecorded` transition should be enqueued for
2660/// this tick. Repeated observations in either state aggregate in memory;
2661/// only elevation and recovery edges reach the primary store.
2662fn checkpoint_outcome_should_emit(above_warn: bool, was_elevated: bool) -> bool {
2663 above_warn != was_elevated
2664}
2665
2666/// Advance the pressure-episode/lifecycle-emission state machine for one
2667/// observed tick. `try_emit` mirrors [`CheckpointLifecycleEmitter::try_emit`]
2668/// — `true` means the row was handed off, `false` means the queue was full
2669/// or closed.
2670///
2671/// #1857: on a dropped recovery handoff (`try_emit` returns `false` while
2672/// `above_warn` is `false`), the closed episode's summary is stashed in
2673/// `pending_recovery` for retry on later ticks — flushed here before this
2674/// tick's own transition is evaluated — instead of leaving
2675/// `event_elevation_open` and `pressure_episode` open for the next elevated
2676/// tick to silently extend, which would report two separate pressure
2677/// incidents as one merged episode.
2678#[allow(clippy::too_many_arguments)]
2679fn observe_checkpoint_pressure_tick(
2680 above_warn: bool,
2681 wal_pages: u64,
2682 above_high_water: bool,
2683 above_truncate_high_water: bool,
2684 config: &CheckpointConfig,
2685 event_elevation_open: &mut bool,
2686 pressure_episode: &mut Option<CheckpointPressureEpisode>,
2687 pending_recovery: &mut Option<khive_storage::CheckpointOutcomeRecordedPayload>,
2688 mut try_emit: impl FnMut(khive_storage::CheckpointOutcomeRecordedPayload) -> bool,
2689) {
2690 // An undelivered recovery summary is a BARRIER, not merely a retry:
2691 // lifecycle consumers assert on the ordered event history (ADR-094), so
2692 // a later episode's opening must never be appended ahead of an earlier
2693 // episode's recovery. If the retry fails, the in-memory aggregate still
2694 // advances below, but no other emission is attempted this tick — a
2695 // deferred opening or recovery re-derives from state on a later tick,
2696 // after the pending summary has been delivered in order.
2697 let pending_blocks_emission = if let Some(payload) = pending_recovery.clone() {
2698 if try_emit(payload) {
2699 *pending_recovery = None;
2700 false
2701 } else {
2702 true
2703 }
2704 } else {
2705 false
2706 };
2707
2708 if above_warn {
2709 match pressure_episode.as_mut() {
2710 Some(episode) => episode.observe(wal_pages),
2711 None => *pressure_episode = Some(CheckpointPressureEpisode::start(wal_pages)),
2712 }
2713 } else if !*event_elevation_open {
2714 // No elevation row reached the bounded handoff, so from any
2715 // consumer's view this episode never opened; discarding it keeps
2716 // the delivered history self-consistent. When the discard happens
2717 // because the barrier suppressed the opening attempt entirely, the
2718 // loss would otherwise be invisible even to the drop counters that
2719 // record failed attempts, so it is counted and logged here.
2720 if pending_blocks_emission && pressure_episode.is_some() {
2721 CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.fetch_add(1, Ordering::Relaxed);
2722 tracing::warn!(
2723 wal_pages,
2724 "checkpoint pressure episode elapsed unreported behind an undelivered recovery summary"
2725 );
2726 }
2727 *pressure_episode = None;
2728 }
2729
2730 if pending_blocks_emission || !checkpoint_outcome_should_emit(above_warn, *event_elevation_open)
2731 {
2732 return;
2733 }
2734 let Some(episode) = *pressure_episode else {
2735 tracing::warn!(
2736 above_warn,
2737 event_elevation_open = *event_elevation_open,
2738 "checkpoint pressure transition has no episode aggregate"
2739 );
2740 return;
2741 };
2742 let payload = khive_storage::CheckpointOutcomeRecordedPayload {
2743 wal_pages,
2744 warn_pages: config.warn_pages,
2745 high_water_pages: config.high_water_pages,
2746 truncate_high_water_pages: config.truncate_high_water_pages,
2747 above_warn,
2748 above_high_water,
2749 above_truncate_high_water,
2750 episode_elevated_ticks: Some(episode.elevated_ticks),
2751 episode_peak_wal_pages: Some(episode.peak_wal_pages),
2752 };
2753 if try_emit(payload.clone()) {
2754 *event_elevation_open = above_warn;
2755 if !above_warn {
2756 *pressure_episode = None;
2757 }
2758 } else if !above_warn {
2759 // The recovery handoff was dropped. Close this episode locally
2760 // anyway — `event_elevation_open` MUST NOT stay true, or the next
2761 // elevated tick would extend this (already finished) episode's
2762 // aggregate instead of starting a fresh one for what is genuinely a
2763 // new pressure incident. The dropped summary itself isn't thrown
2764 // away: it is stashed in `pending_recovery` and delivered on a
2765 // later tick, ahead of (and as a barrier to) every subsequent
2766 // emission, so lifecycle ordering survives the retry. The slot is
2767 // structurally empty here: a tick that entered with an undelivered
2768 // summary returned at the barrier above and never reached this arm.
2769 debug_assert!(
2770 pending_recovery.is_none(),
2771 "recovery emission attempted while an earlier summary was still pending"
2772 );
2773 *event_elevation_open = false;
2774 *pressure_episode = None;
2775 *pending_recovery = Some(payload);
2776 }
2777}
2778
2779/// ADR-091 Plank 0 (Amendment 3: takes the tick's already-computed,
2780/// backend-scoped oldest span instead of re-querying the process-wide
2781/// aggregate): log the oldest open transaction registry entry alongside the
2782/// WAL frame count at `debug!`, on EVERY tick regardless of threshold
2783/// state. This is the low-volume per-tick trace; the WARN-level escalations
2784/// live in [`log_tx_registry_oldest_warn`] and
2785/// debug-level, unconditional per-tick trace. See
2786/// crates/khive-db/docs/api/checkpoint.md#private-tx-registry-logging-helpers-plank-0
2787fn log_tx_registry_oldest_debug(
2788 wal_pages: u64,
2789 oldest: Option<&khive_storage::tx_registry::OldestSpan>,
2790) {
2791 if let Some(span) = oldest {
2792 tracing::debug!(
2793 wal_pages,
2794 oldest_tx_age_secs = span.age.as_secs_f64(),
2795 oldest_tx_label = span.label.as_deref().unwrap_or("<unlabeled>"),
2796 "WAL checkpoint tick: oldest open transaction registry entry"
2797 );
2798 }
2799}
2800
2801/// Escalates the oldest open registry entry to `warn!`. NOT internally
2802/// rate-limited — caller MUST gate on a below→above `warn_pages` crossing
2803/// (`crossing_warn`) or every tick reproduces the log-spam bug this fixes.
2804fn log_tx_registry_oldest_warn(
2805 wal_pages: u64,
2806 oldest: Option<&khive_storage::tx_registry::OldestSpan>,
2807) {
2808 if let Some(span) = oldest {
2809 tracing::warn!(
2810 wal_pages,
2811 oldest_tx_age_secs = span.age.as_secs_f64(),
2812 oldest_tx_label = span.label.as_deref().unwrap_or("<unlabeled>"),
2813 "WAL checkpoint tick: oldest open transaction registry entry"
2814 );
2815 }
2816}
2817
2818/// Enumerates every open registry entry at `warn!`. NOT internally
2819/// rate-limited — caller MUST gate on a below→above `high_water_pages`
2820/// crossing (`crossing_warn`) or every tick repeats the full enumeration.
2821fn log_tx_registry_snapshot_warn(wal_pages: u64) {
2822 log_tx_registry_entries_warn(wal_pages, &khive_storage::tx_registry::snapshot());
2823}
2824
2825fn log_tx_registry_entries_warn(wal_pages: u64, snapshot: &[(Duration, Option<String>)]) {
2826 for (age, label) in snapshot {
2827 tracing::warn!(
2828 wal_pages,
2829 tx_age_secs = age.as_secs_f64(),
2830 tx_label = label.as_deref().unwrap_or("<unlabeled>"),
2831 "WAL high-water: open transaction registry entry"
2832 );
2833 }
2834}
2835
2836fn log_truncate_no_progress_warn(
2837 wal_pages_before: u64,
2838 wal_pages_after: u64,
2839 snapshot: &[(Duration, Option<String>)],
2840) {
2841 let open_tx_count = snapshot.len();
2842 let oldest_tx_age_secs = snapshot
2843 .iter()
2844 .map(|(age, _)| *age)
2845 .max()
2846 .map(|age| age.as_secs_f64());
2847 if snapshot.is_empty() {
2848 tracing::warn!(
2849 wal_pages_before,
2850 wal_pages_after,
2851 open_tx_count,
2852 oldest_tx_age_secs = ?oldest_tx_age_secs,
2853 "WAL TRUNCATE attempt made no progress; no open transaction in this process's registry"
2854 );
2855 } else {
2856 tracing::warn!(
2857 wal_pages_before,
2858 wal_pages_after,
2859 open_tx_count,
2860 oldest_tx_age_secs = ?oldest_tx_age_secs,
2861 "WAL TRUNCATE attempt made no progress; open transactions observed in this process's registry"
2862 );
2863 }
2864 log_tx_registry_entries_warn(wal_pages_after, snapshot);
2865}
2866
2867/// Emits the high-water WARN, deciding its text from the registry entry this
2868/// tick already read instead of asserting a cause the evidence beside it can
2869/// refute.
2870///
2871/// The registry only covers this process. An old registered transaction may
2872/// hold a snapshot; a young or empty registry cannot rule out an external
2873/// reader. `warn_after` is the same threshold as the transaction-age ladder.
2874fn log_wal_high_water_warn(
2875 wal_pages: u64,
2876 high_water: u64,
2877 oldest: Option<&khive_storage::tx_registry::OldestSpan>,
2878 warn_after: Duration,
2879) {
2880 match oldest.filter(|span| span.age >= warn_after) {
2881 Some(span) => tracing::warn!(
2882 wal_pages,
2883 high_water,
2884 oldest_tx_age_secs = span.age.as_secs_f64(),
2885 oldest_tx_label = span.label.as_deref().unwrap_or("<unlabeled>"),
2886 "WAL high-water mark exceeded; an in-process registered transaction is older \
2887 than the age threshold and may hold a snapshot"
2888 ),
2889 None => tracing::warn!(
2890 wal_pages,
2891 high_water,
2892 oldest_tx_age_secs = ?oldest.map(|span| span.age.as_secs_f64()),
2893 oldest_tx_label = oldest
2894 .and_then(|span| span.label.as_deref())
2895 .unwrap_or("<none>"),
2896 "WAL high-water mark exceeded; no in-process transaction older than the age \
2897 threshold is visible; a reader in another process may hold the snapshot, \
2898 or writes may outpace PASSIVE checkpoints"
2899 ),
2900 }
2901}
2902
2903/// Internal result of the synchronous SQLite checkpoint core. Keeping the
2904/// no-progress attribution request next to (but distinct from) `wal_pages`
2905/// makes the async handoff explicit and gives deferred work one caller-owned
2906/// lifetime instead of leaving it in mutable cross-tick state.
2907#[derive(Debug)]
2908#[must_use]
2909struct CheckpointCoreOutcome {
2910 wal_pages: Option<u64>,
2911 unavailable_reason: Option<CheckpointUnavailableReason>,
2912 sidecar_attribution: Option<WalpinAttributionRequest>,
2913}
2914
2915/// Issue one checkpoint cycle against the task's dedicated checkpoint
2916/// connection (`conn` — see `CheckpointConnection`; NEVER the pool's writer
2917/// mutex).
2918///
2919/// Returns the observed WAL page count on success. A busy PASSIVE row has no
2920/// usable observation and the compatibility wrapper returns `SQLITE_BUSY`;
2921/// an inconsistent nonbusy frame pair instead returns `SQLITE_ERROR`.
2922/// A connection-level pragma error is also returned; the task caller drops
2923/// that connection and reopens next tick. TRUNCATE errors remain non-fatal.
2924///
2925/// The caller owns all threshold-crossing WARN logging so that warnings fire
2926/// at most once per crossing, not every tick.
2927///
2928/// ADR-091 Plank 2: after the PASSIVE pass, this is also the single point
2929/// that may escalate to TRUNCATE (`maybe_truncate`) — on the SAME dedicated
2930/// connection, never a second connection or a pool checkout. A no-progress
2931/// result produces a separate cross-process attribution request; the
2932/// synchronous core never walks the sidecar directory. Production's
2933/// [`run_checkpoint_task`] consumes that request through an awaited
2934/// `spawn_blocking` before continuing the tick. This compatibility wrapper
2935/// intentionally returns only the historical page-count surface; the daemon
2936/// calls `checkpoint_once_core` so it cannot discard the request.
2937pub fn checkpoint_once(
2938 pool: &ConnectionPool,
2939 conn: &rusqlite::Connection,
2940 config: &CheckpointConfig,
2941 truncate_state: &mut TruncateState,
2942) -> Result<u64, rusqlite::Error> {
2943 let outcome = checkpoint_once_core(pool, conn, config, truncate_state)?;
2944 outcome.wal_pages.ok_or_else(|| match outcome
2945 .unavailable_reason
2946 .expect("unavailable core outcome has a reason")
2947 {
2948 CheckpointUnavailableReason::Busy => rusqlite::Error::SqliteFailure(
2949 rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_BUSY),
2950 Some("PASSIVE checkpoint returned a busy row without a WAL frame observation".into()),
2951 ),
2952 CheckpointUnavailableReason::InconsistentFrames => rusqlite::Error::SqliteFailure(
2953 rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_ERROR),
2954 Some("PASSIVE checkpoint returned an inconsistent frame pair without a WAL frame observation".into()),
2955 ),
2956 })
2957}
2958
2959/// Synchronous PASSIVE/TRUNCATE core used by the async task. Unlike the
2960/// compatibility wrapper [`checkpoint_once`], this preserves the explicit
2961/// no-progress attribution request for the caller to complete off-runtime.
2962fn checkpoint_once_core(
2963 pool: &ConnectionPool,
2964 conn: &rusqlite::Connection,
2965 config: &CheckpointConfig,
2966 truncate_state: &mut TruncateState,
2967) -> Result<CheckpointCoreOutcome, rusqlite::Error> {
2968 #[cfg(unix)]
2969 truncate_state.begin_tick();
2970 let started = Instant::now();
2971 let checkpoint_result = query_routine_checkpoint_observation(pool, conn);
2972 let elapsed_us = started.elapsed().as_micros().min(u128::from(u64::MAX)) as u64;
2973 record_checkpoint_timing(
2974 pool,
2975 elapsed_us,
2976 checkpoint_result
2977 .as_ref()
2978 .ok()
2979 .map(|observation| observation.busy),
2980 );
2981 let raw_observation = match checkpoint_result {
2982 Ok(observation) => observation,
2983 Err(e) => {
2984 record_checkpoint_run_result(pool, None);
2985 tracing::warn!(error = %e, elapsed_us, "WAL checkpoint failed");
2986 return Err(e);
2987 }
2988 };
2989 record_checkpoint_run_result(
2990 pool,
2991 Some((
2992 raw_observation.busy,
2993 raw_observation.log_frames,
2994 raw_observation.checkpointed_frames,
2995 )),
2996 );
2997 let wal_pages = match observed_wal_pages(raw_observation) {
2998 Ok(wal_pages) => wal_pages,
2999 Err(reason) => {
3000 match reason {
3001 CheckpointUnavailableReason::Busy => tracing::debug!(
3002 busy = raw_observation.busy,
3003 wal_log_frames = raw_observation.log_frames,
3004 wal_checkpointed_frames = raw_observation.checkpointed_frames,
3005 elapsed_us,
3006 "WAL PASSIVE checkpoint returned a busy row; frame observation unavailable"
3007 ),
3008 CheckpointUnavailableReason::InconsistentFrames => tracing::warn!(
3009 busy = raw_observation.busy,
3010 wal_log_frames = raw_observation.log_frames,
3011 wal_checkpointed_frames = raw_observation.checkpointed_frames,
3012 elapsed_us,
3013 "WAL PASSIVE checkpoint returned an inconsistent frame pair; frame observation unavailable"
3014 ),
3015 }
3016 return Ok(CheckpointCoreOutcome {
3017 wal_pages: None,
3018 unavailable_reason: Some(reason),
3019 sidecar_attribution: None,
3020 });
3021 }
3022 };
3023 let observation = record_routine_wal_observation(pool, raw_observation);
3024 LAST_WAL_PAGES.store(wal_pages, Ordering::Relaxed);
3025 note_checkpoint_observed(wal_pages);
3026 tracing::debug!(
3027 wal_pages,
3028 elapsed_us,
3029 busy = raw_observation.busy,
3030 wal_checkpointed_frames = observation.checkpointed_frames,
3031 wal_pending_frames = observation.pending_frames,
3032 wal_physical_bytes = ?observation.physical_wal_bytes,
3033 "WAL checkpoint issued"
3034 );
3035
3036 let sidecar_attribution = maybe_truncate(pool, conn, config, wal_pages, truncate_state);
3037
3038 Ok(CheckpointCoreOutcome {
3039 wal_pages: Some(wal_pages),
3040 unavailable_reason: None,
3041 sidecar_attribution,
3042 })
3043}
3044
3045fn truncate_needs_attribution(wal_pages_before: u64, wal_pages_after: Option<u64>) -> bool {
3046 wal_pages_after.is_none_or(|pages| pages >= wal_pages_before)
3047}
3048
3049/// Evaluate and, if due, attempt a TRUNCATE escalation on the same dedicated
3050/// checkpoint connection the caller already holds (never its own checkout —
3051/// there is no pool writer involved on this path at all). `last_attempt`
3052/// is stamped ONLY on an actual attempt, never on a skip. See
3053/// crates/khive-db/docs/api/checkpoint.md#maybe_truncate--truncate-attempt-gating-plank-2
3054fn maybe_truncate(
3055 pool: &ConnectionPool,
3056 conn: &rusqlite::Connection,
3057 config: &CheckpointConfig,
3058 wal_pages_before: u64,
3059 truncate_state: &mut TruncateState,
3060) -> Option<WalpinAttributionRequest> {
3061 if wal_pages_before < config.truncate_high_water_pages {
3062 return None;
3063 }
3064
3065 if let Some(last) = truncate_state.last_attempt {
3066 if last.elapsed() < config.truncate_min_interval {
3067 return None;
3068 }
3069 }
3070
3071 // Which caller (if any) is pinning the WAL — logged before the attempt so
3072 // it is available even if the attempt itself succeeds.
3073 log_tx_registry_snapshot_warn(wal_pages_before);
3074
3075 let original_busy_timeout = pool.config().busy_timeout;
3076
3077 if let Err(e) = conn.busy_timeout(config.truncate_busy_timeout) {
3078 // Setup failed before the TRUNCATE pragma ever ran — this is a skip,
3079 // not an attempt. `last_attempt` must NOT advance here (ADR-091
3080 // §377-382): stamping now would suppress the next eligible attempt
3081 // for the full `truncate_min_interval` on a path that never touched
3082 // the WAL at all.
3083 tracing::warn!(error = %e, "failed to lower busy_timeout for TRUNCATE attempt; skipping");
3084 return None;
3085 }
3086
3087 #[cfg(unix)]
3088 let mut holder_attribution = capture_walpin_attribution_request(pool, truncate_state);
3089 #[cfg(unix)]
3090 let mut sidecar_attribution = None;
3091 #[cfg(not(unix))]
3092 let sidecar_attribution = None;
3093
3094 // Only now is this a genuine attempt: the writer is held, the threshold
3095 // and interval gates passed, and the busy_timeout override is in effect
3096 // immediately before the TRUNCATE pragma itself.
3097 truncate_state.last_attempt = Some(Instant::now());
3098 #[cfg(test)]
3099 off_worker::cycle_panic_seam::after_attempt_decided(pool.canonical_path());
3100
3101 let start = Instant::now();
3102 let outcome = query_truncate_observation(conn);
3103 record_checkpoint_run_result(
3104 pool,
3105 outcome.as_ref().ok().map(|observation| {
3106 (
3107 observation.busy,
3108 observation.log_frames,
3109 observation.checkpointed_frames,
3110 )
3111 }),
3112 );
3113 let elapsed = start.elapsed();
3114
3115 // Restore the pool's configured busy_timeout immediately after the
3116 // attempt, win or lose, before any other logging or bookkeeping.
3117 if let Err(e) = conn.busy_timeout(original_busy_timeout) {
3118 tracing::warn!(error = %e, "failed to restore busy_timeout after TRUNCATE attempt");
3119 }
3120
3121 match outcome {
3122 Ok(_) => {
3123 let wal_pages_after = query_wal_pages(pool, conn);
3124 if let Some(pages) = wal_pages_after {
3125 tracing::info!(
3126 wal_pages_before,
3127 wal_pages_after = pages,
3128 elapsed_ms = elapsed.as_millis() as u64,
3129 "WAL TRUNCATE checkpoint attempted"
3130 );
3131 } else {
3132 tracing::info!(
3133 wal_pages_before,
3134 wal_pages_after_unavailable = true,
3135 elapsed_ms = elapsed.as_millis() as u64,
3136 "WAL TRUNCATE checkpoint attempted"
3137 );
3138 }
3139
3140 if truncate_needs_attribution(wal_pages_before, wal_pages_after) {
3141 let snapshot = khive_storage::tx_registry::snapshot();
3142 if let Some(pages) = wal_pages_after {
3143 log_truncate_no_progress_warn(wal_pages_before, pages, &snapshot);
3144 } else {
3145 tracing::warn!(
3146 wal_pages_before,
3147 "WAL TRUNCATE progress unmeasured; checking possible holders in this process and others"
3148 );
3149 log_tx_registry_entries_warn(wal_pages_before, &snapshot);
3150 }
3151 #[cfg(test)]
3152 if let Some(path) = pool.canonical_path() {
3153 truncate_report_test_sync::after_no_progress_before_report(path);
3154 }
3155 #[cfg(unix)]
3156 {
3157 // The census above had to be captured before TRUNCATE so
3158 // a transient holder remains attributable. The bounded
3159 // sidecar walk itself must not run here: it keeps its own
3160 // awaited pass in the async owner, which orders it against
3161 // the tick's housekeeping. Hand the immutable request back
3162 // to that owner for its `spawn_blocking` pass.
3163 sidecar_attribution = holder_attribution.take();
3164 }
3165 log_backfill_gap(pool, conn);
3166 }
3167
3168 note_truncate_outcome(config, wal_pages_after, truncate_state);
3169 }
3170 Err(e) => {
3171 tracing::warn!(error = %e, wal_pages_before, "WAL TRUNCATE attempt failed");
3172 log_tx_registry_snapshot_warn(wal_pages_before);
3173 note_truncate_outcome(config, Some(wal_pages_before), truncate_state);
3174 }
3175 }
3176 #[cfg(unix)]
3177 if let Some(WalpinAttributionRequest::Fresh {
3178 previous_last_attempt,
3179 ..
3180 }) = holder_attribution.as_ref()
3181 {
3182 truncate_state.restore_walpin_full_scan_reservation(*previous_last_attempt);
3183 }
3184 sidecar_attribution
3185}
3186
3187#[cfg(test)]
3188mod truncate_report_test_sync {
3189 use std::path::{Path, PathBuf};
3190 use std::sync::mpsc::{sync_channel, Receiver, SyncSender};
3191 use std::sync::Mutex;
3192
3193 struct Hook {
3194 db_path: PathBuf,
3195 reached_tx: SyncSender<()>,
3196 proceed_rx: Receiver<()>,
3197 }
3198
3199 static HOOK: Mutex<Option<Hook>> = Mutex::new(None);
3200
3201 pub(crate) fn install(db_path: PathBuf) -> (Receiver<()>, SyncSender<()>) {
3202 let (reached_tx, reached_rx) = sync_channel(0);
3203 let (proceed_tx, proceed_rx) = sync_channel(0);
3204 let replaced = HOOK
3205 .lock()
3206 .unwrap_or_else(|poisoned| poisoned.into_inner())
3207 .replace(Hook {
3208 db_path,
3209 reached_tx,
3210 proceed_rx,
3211 });
3212 assert!(replaced.is_none(), "truncate report hook already installed");
3213 (reached_rx, proceed_tx)
3214 }
3215
3216 pub(crate) fn uninstall() {
3217 *HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
3218 }
3219
3220 pub(crate) fn after_no_progress_before_report(db_path: &Path) {
3221 let hook = {
3222 let mut guard = HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
3223 match guard.as_ref() {
3224 Some(hook) if hook.db_path == db_path => guard.take(),
3225 _ => None,
3226 }
3227 };
3228 let Some(hook) = hook else {
3229 return;
3230 };
3231 let _ = hook.reached_tx.send(());
3232 let _ = hook.proceed_rx.recv();
3233 }
3234}
3235
3236/// Deterministic seam for the async-attribution regressions below. The hook
3237/// executes inside the actual `spawn_blocking` closure, so a current-thread
3238/// Tokio test can prove both thread displacement and awaited ordering without
3239/// relying on sleeps or scheduler timing.
3240#[cfg(all(test, unix))]
3241mod walpin_attribution_test_sync {
3242 use std::path::{Path, PathBuf};
3243 use std::sync::atomic::{AtomicUsize, Ordering};
3244 use std::sync::mpsc::{sync_channel, Receiver, SyncSender};
3245 use std::sync::{Arc, Mutex};
3246
3247 enum Behavior {
3248 Pause {
3249 reached_tx: tokio::sync::oneshot::Sender<std::thread::ThreadId>,
3250 proceed_rx: Receiver<()>,
3251 },
3252 Panic,
3253 }
3254
3255 struct Hook {
3256 dir: PathBuf,
3257 behavior: Behavior,
3258 }
3259
3260 static HOOK: Mutex<Option<Hook>> = Mutex::new(None);
3261 static REPORT_COUNTER: Mutex<Option<Arc<AtomicUsize>>> = Mutex::new(None);
3262
3263 pub(crate) fn install_pause(
3264 dir: PathBuf,
3265 ) -> (
3266 tokio::sync::oneshot::Receiver<std::thread::ThreadId>,
3267 SyncSender<()>,
3268 Arc<AtomicUsize>,
3269 ) {
3270 let (reached_tx, reached_rx) = tokio::sync::oneshot::channel();
3271 let (proceed_tx, proceed_rx) = sync_channel(0);
3272 let report_counter = Arc::new(AtomicUsize::new(0));
3273 let replaced = HOOK
3274 .lock()
3275 .unwrap_or_else(|poisoned| poisoned.into_inner())
3276 .replace(Hook {
3277 dir,
3278 behavior: Behavior::Pause {
3279 reached_tx,
3280 proceed_rx,
3281 },
3282 });
3283 assert!(
3284 replaced.is_none(),
3285 "walpin attribution hook already installed"
3286 );
3287 *REPORT_COUNTER
3288 .lock()
3289 .unwrap_or_else(|poisoned| poisoned.into_inner()) = Some(Arc::clone(&report_counter));
3290 (reached_rx, proceed_tx, report_counter)
3291 }
3292
3293 pub(crate) fn install_panic(dir: PathBuf) {
3294 let replaced = HOOK
3295 .lock()
3296 .unwrap_or_else(|poisoned| poisoned.into_inner())
3297 .replace(Hook {
3298 dir,
3299 behavior: Behavior::Panic,
3300 });
3301 assert!(
3302 replaced.is_none(),
3303 "walpin attribution hook already installed"
3304 );
3305 }
3306
3307 pub(crate) fn uninstall() {
3308 *HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
3309 *REPORT_COUNTER
3310 .lock()
3311 .unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
3312 }
3313
3314 pub(crate) fn before_enumeration(dir: &Path) {
3315 let hook = {
3316 let mut guard = HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
3317 match guard.as_ref() {
3318 Some(hook) if hook.dir == dir => guard.take(),
3319 _ => None,
3320 }
3321 };
3322 let Some(hook) = hook else {
3323 return;
3324 };
3325 match hook.behavior {
3326 Behavior::Pause {
3327 reached_tx,
3328 proceed_rx,
3329 } => {
3330 if reached_tx.send(std::thread::current().id()).is_ok() {
3331 let _ = proceed_rx.recv();
3332 }
3333 }
3334 Behavior::Panic => panic!("injected walpin attribution worker panic"),
3335 }
3336 }
3337
3338 pub(crate) fn report_used() {
3339 if let Some(counter) = REPORT_COUNTER
3340 .lock()
3341 .unwrap_or_else(|poisoned| poisoned.into_inner())
3342 .as_ref()
3343 {
3344 counter.fetch_add(1, Ordering::SeqCst);
3345 }
3346 }
3347}
3348
3349/// ADR-091 Plank 2: track measured TRUNCATE outcomes that fail to bring
3350/// `wal_pages` below `warn_pages`, firing a one-shot escalated WARN at the
3351/// third such failure. An unmeasured attempt leaves the streak unchanged;
3352/// only a measured result below `warn_pages` resets it.
3353fn note_truncate_outcome(
3354 config: &CheckpointConfig,
3355 wal_pages_after: Option<u64>,
3356 state: &mut TruncateState,
3357) {
3358 // Metrics read-surface (load/perf harness): this function runs exactly
3359 // once per genuine TRUNCATE attempt (both the `Ok` and `Err` outcome
3360 // arms in `maybe_truncate` call it once each), so incrementing here
3361 // counts total attempts without a separate call site.
3362 TRUNCATE_ATTEMPTS.fetch_add(1, Ordering::Relaxed);
3363
3364 if let Some(wal_pages_after) = wal_pages_after {
3365 if wal_pages_after >= config.warn_pages {
3366 state.consecutive_failures = state.consecutive_failures.saturating_add(1);
3367 if state.consecutive_failures == 3 {
3368 tracing::warn!(
3369 wal_pages_after,
3370 warn_threshold = config.warn_pages,
3371 "WAL TRUNCATE has failed to clear WAL pressure for 3 consecutive attempts"
3372 );
3373 }
3374 } else {
3375 state.consecutive_failures = 0;
3376 }
3377 }
3378
3379 TRUNCATE_CONSECUTIVE_FAILURES.store(state.consecutive_failures as u64, Ordering::Relaxed);
3380}
3381
3382/// Immutable work captured around an armed TRUNCATE and consumed by the
3383/// async checkpoint owner only when that attempt makes no progress.
3384///
3385/// The holder census belongs here because it must precede the bounded
3386/// TRUNCATE wait. The sidecar directory walk does not: it remains deferred
3387/// until after the outcome is known and is executed through an awaited
3388/// `spawn_blocking` by [`complete_walpin_attribution`].
3389#[cfg(unix)]
3390#[derive(Debug)]
3391enum WalpinAttributionRequest {
3392 Fresh {
3393 dir: PathBuf,
3394 census: Result<crate::walpin::CensusResult, String>,
3395 legacy_fallback_interval: Duration,
3396 previous_last_attempt: Option<Instant>,
3397 },
3398 Cached(CachedWalpinAttribution),
3399 Suppressed,
3400}
3401
3402#[cfg(unix)]
3403#[derive(Debug, Clone, Copy, PartialEq, Eq)]
3404enum WalpinReportFreshness {
3405 Fresh,
3406 Cached { age: Duration },
3407}
3408
3409#[cfg(unix)]
3410impl WalpinReportFreshness {
3411 fn is_fresh(self) -> bool {
3412 self == Self::Fresh
3413 }
3414}
3415
3416/// Non-Unix placeholder keeps the synchronous core's outcome shape stable;
3417/// daemon sidecar attribution itself is Unix-only.
3418#[cfg(not(unix))]
3419type WalpinAttributionRequest = ();
3420
3421/// Honest failure surface for an attempted no-progress attribution pass.
3422/// Both variants suppress same-tick housekeeping because a panicked blocking
3423/// worker may already have partially enumerated the directory; retrying a
3424/// second pass would violate the one-pass-per-tick bound.
3425#[cfg(unix)]
3426#[derive(Debug, Clone, PartialEq, Eq)]
3427enum WalpinAttributionFailure {
3428 Enumeration(String),
3429 Worker(String),
3430}
3431
3432#[cfg(unix)]
3433impl WalpinAttributionFailure {
3434 fn kind(&self) -> &'static str {
3435 match self {
3436 Self::Enumeration(_) => "enumeration",
3437 Self::Worker(_) => "blocking_worker",
3438 }
3439 }
3440}
3441
3442#[cfg(unix)]
3443impl std::fmt::Display for WalpinAttributionFailure {
3444 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
3445 match self {
3446 Self::Enumeration(error) => write!(
3447 formatter,
3448 "sidecar directory failed the trust-boundary enumeration; cross-process \
3449 WAL-pin attribution is unestablished for this tick: {error}"
3450 ),
3451 Self::Worker(error) => write!(
3452 formatter,
3453 "sidecar attribution blocking worker failed; cross-process WAL-pin \
3454 attribution is unestablished for this tick: {error}"
3455 ),
3456 }
3457 }
3458}
3459
3460/// Capture the pre-TRUNCATE OS holder census and stable sidecar inputs. A
3461/// no-op if the sidecar is disabled or this backend has no on-disk path.
3462#[cfg(unix)]
3463fn capture_walpin_attribution_request(
3464 pool: &ConnectionPool,
3465 state: &mut TruncateState,
3466) -> Option<WalpinAttributionRequest> {
3467 let path = pool.canonical_path()?;
3468 if !crate::walpin::sidecar_enabled(true) {
3469 return None;
3470 }
3471 let legacy_fallback_interval = state.legacy_walpin_fallback_interval;
3472 Some(match state.plan_walpin_attribution_at(Instant::now()) {
3473 WalpinFullScanPlan::Refresh {
3474 previous_last_attempt,
3475 } => WalpinAttributionRequest::Fresh {
3476 dir: crate::walpin::sidecar_dir_for(path),
3477 census: crate::walpin::census_holders(path).map_err(|error| error.to_string()),
3478 legacy_fallback_interval,
3479 previous_last_attempt,
3480 },
3481 WalpinFullScanPlan::Cached(cached) => WalpinAttributionRequest::Cached(cached),
3482 WalpinFullScanPlan::Suppressed => WalpinAttributionRequest::Suppressed,
3483 })
3484}
3485
3486/// Consume this tick's no-progress attribution request off the async runtime
3487/// worker and await it before any report or fallback housekeeping is used.
3488/// Returns `Ok(false)` when no pass was requested. Once a request exists the
3489/// state is marked attempted before spawning, so worker panic/cancellation
3490/// cannot accidentally authorize a second directory scan in the same tick.
3491#[cfg(unix)]
3492async fn complete_walpin_attribution(
3493 request: Option<WalpinAttributionRequest>,
3494 state: &mut TruncateState,
3495) -> Result<bool, WalpinAttributionFailure> {
3496 let Some(request) = request else {
3497 return Ok(false);
3498 };
3499 match request {
3500 WalpinAttributionRequest::Suppressed => Ok(false),
3501 WalpinAttributionRequest::Cached(cached) => {
3502 log_walpin_sidecar_report(
3503 &cached.report,
3504 cached.census,
3505 WalpinReportFreshness::Cached {
3506 age: Instant::now().saturating_duration_since(cached.captured_at),
3507 },
3508 );
3509 Ok(true)
3510 }
3511 WalpinAttributionRequest::Fresh {
3512 dir,
3513 census,
3514 legacy_fallback_interval,
3515 previous_last_attempt: _,
3516 } => {
3517 state.sidecar_attribution_attempted_this_tick = true;
3518 if state.walpin_full_scan_last_attempt.is_none() {
3519 state.walpin_full_scan_last_attempt = Some(Instant::now());
3520 }
3521 let fallback = state.walpin_cached_attribution.clone();
3522 let result = tokio::task::spawn_blocking(move || {
3523 #[cfg(test)]
3524 walpin_attribution_test_sync::before_enumeration(&dir);
3525 crate::walpin::enumerate_live(&dir, legacy_fallback_interval)
3526 })
3527 .await
3528 .map_err(|error| WalpinAttributionFailure::Worker(error.to_string()))
3529 .and_then(|result| {
3530 result.map_err(|error| WalpinAttributionFailure::Enumeration(error.to_string()))
3531 });
3532
3533 match result {
3534 Ok(report) => {
3535 let captured_at = Instant::now();
3536 log_walpin_sidecar_report(
3537 &report,
3538 census.clone(),
3539 WalpinReportFreshness::Fresh,
3540 );
3541 state.cache_walpin_attribution(report, census, captured_at);
3542 Ok(true)
3543 }
3544 Err(error) => {
3545 if let Some(cached) = fallback {
3546 log_walpin_sidecar_report(
3547 &cached.report,
3548 cached.census,
3549 WalpinReportFreshness::Cached {
3550 age: Instant::now().saturating_duration_since(cached.captured_at),
3551 },
3552 );
3553 }
3554 Err(error)
3555 }
3556 }
3557 }
3558 }
3559}
3560
3561/// When a TRUNCATE attempt makes no progress, enumerate the walpin sidecar and
3562/// combine it with the holder census captured immediately before that attempt.
3563/// This pass consumes the classifications for attribution and returns whether
3564/// enumeration was attempted; the caller uses that marker to suppress the
3565/// ordinary housekeeping pass later in the same tick. Holder identity cannot
3566/// be deferred because a transient blocker may have released by then.
3567///
3568/// Sidecar-health attribution (ADR-091 Amendment 2):
3569/// the sharper "unregistered/native mechanism" conclusion is licensed only
3570/// when every discovered PID is `reporting` or `registered-silent`
3571/// (`WalpinReport::fully_attributed`); any `unknown` PID — including the
3572/// directory itself failing the trust-boundary check — makes attribution
3573/// inconclusive, and the WARN below names exactly which PIDs are unresolved
3574/// instead of silently exonerating them.
3575#[cfg(unix)]
3576fn log_walpin_sidecar_report(
3577 report: &crate::walpin::WalpinReport,
3578 census: Result<crate::walpin::CensusResult, String>,
3579 freshness: WalpinReportFreshness,
3580) {
3581 #[cfg(test)]
3582 walpin_attribution_test_sync::report_used();
3583 let now = now_epoch_secs();
3584 for hb in report.reporting() {
3585 // ADR-091 Amendment 3 Plank F2 fail-closed reading rule: the
3586 // logger must never let a fallback-confidence entry read as live
3587 // cross-process ground truth, so the confidence distinction is
3588 // always emitted alongside the raw field — never inferred by the
3589 // reader of this log line.
3590 tracing::warn!(
3591 walpin_pid = hb.pid,
3592 walpin_role = %hb.process_role,
3593 walpin_oldest_tx_age_secs = hb.current_oldest_tx_age_secs(now),
3594 walpin_oldest_tx_label = hb.oldest_tx_label.as_deref().unwrap_or("<unlabeled>"),
3595 walpin_attribution_basis = hb.attribution_basis.as_deref().unwrap_or("<unspecified>"),
3596 walpin_attribution_evidence_backed = hb.attribution_is_evidence_backed(),
3597 walpin_attribution_fresh = freshness.is_fresh(),
3598 walpin_health = "reporting",
3599 "ADR-091 Amendment 2 Plank B: live cross-process WAL-pin attribution report"
3600 );
3601 }
3602 for pid in report.registered_silent_pids() {
3603 tracing::debug!(
3604 walpin_pid = pid,
3605 walpin_health = "registered_silent",
3606 walpin_attribution_fresh = freshness.is_fresh(),
3607 "ADR-091 Amendment 2 Plank B: process affirmatively reports no over-threshold span"
3608 );
3609 }
3610 let mut unknown_pids: Vec<u32> = report.unknown_pids().collect();
3611 if let WalpinReportFreshness::Cached { age } = freshness {
3612 tracing::warn!(
3613 walpin_cache_age_ms = age.as_millis() as u64,
3614 "cached WAL-pin attribution is diagnostic-only; fully-attributed \
3615 conclusion is not licensed"
3616 );
3617 unknown_pids.push(0);
3618 }
3619
3620 // The sidecar directory alone can only speak for PIDs that wrote
3621 // something there. Widen the universe to every PID the OS reports as
3622 // holding the database immediately before the TRUNCATE attempt; any holder
3623 // absent from `report` is unknown.
3624 match census {
3625 Ok(census) => {
3626 let sidecar_known: std::collections::HashSet<u32> = report
3627 .reporting()
3628 .map(|hb| hb.pid)
3629 .chain(report.registered_silent_pids())
3630 .chain(unknown_pids.iter().copied())
3631 .collect();
3632 let mut census_only: Vec<u32> =
3633 census.holders.difference(&sidecar_known).copied().collect();
3634 if !census_only.is_empty() {
3635 census_only.sort_unstable();
3636 tracing::warn!(
3637 ?census_only,
3638 "ADR-091 Amendment 2: these PIDs hold the database file open \
3639 at the OS level but have no sidecar data at all (pre-feature binary, \
3640 sidecar disabled, or wedged before its first write)"
3641 );
3642 unknown_pids.extend(census_only);
3643 }
3644 if !census.is_complete() {
3645 let mut uninspectable = census.uninspectable_pids.clone();
3646 uninspectable.sort_unstable();
3647 tracing::warn!(
3648 ?uninspectable,
3649 truncated = census.truncated,
3650 "ADR-091 Amendment 2: the OS-derived holder census is \
3651 INCOMPLETE — either specific PIDs' open file descriptors could not be \
3652 inspected (permission denied, or a listing race), or the enumeration walk \
3653 itself has positive evidence it did not see the full live-process universe \
3654 (namespace/visibility check, directory-iterator error, self-canary, or a \
3655 libproc buffer that stayed at capacity after bounded retries) — cannot \
3656 rule out an unregistered holder"
3657 );
3658 if uninspectable.is_empty() {
3659 // `truncated` fired with no specific PID list (a
3660 // namespace/visibility or buffer-truncation signal, not
3661 // a per-PID inspection failure) — still makes
3662 // attribution inconclusive. Mirror the census-failure
3663 // arm below with the same non-PID sentinel rather than
3664 // silently trusting a walk we know was incomplete.
3665 unknown_pids.push(0);
3666 } else {
3667 unknown_pids.extend(uninspectable);
3668 }
3669 }
3670 }
3671 Err(e) => {
3672 tracing::warn!(
3673 error = %e,
3674 "ADR-091 Amendment 2: OS-derived holder census failed; \
3675 attribution cannot rule out an unregistered database holder this tick"
3676 );
3677 // A failed census is itself a health failure for the sharper
3678 // conclusion below — treat it as if at least one PID were
3679 // unresolved, without fabricating a specific PID number.
3680 unknown_pids.push(0);
3681 }
3682 }
3683
3684 unknown_pids.sort_unstable();
3685 unknown_pids.dedup();
3686 if !unknown_pids.is_empty() {
3687 tracing::warn!(
3688 ?unknown_pids,
3689 "ADR-091 Amendment 2 Plank B: sidecar health unestablished for these PIDs; \
3690 attribution is inconclusive and the native/unregistered-mechanism conclusion \
3691 is NOT licensed this tick"
3692 );
3693 } else if report.reporting().next().is_none() {
3694 tracing::info!(
3695 "ADR-091 Amendment 2 Plank B: every live PID is reporting or registered-silent \
3696 with none pinning; the WAL pin is not attributable to any in-process registry \
3697 span this sidecar covers"
3698 );
3699 }
3700}
3701
3702/// ADR-091 Amendment 2 Plank C: on a TRUNCATE no-progress event, run a fresh
3703/// `PRAGMA wal_checkpoint(PASSIVE)` (never blocks readers or writers) and
3704/// report the one-row backfill gap as `log` minus `checkpointed` from its
3705/// 3-column return row when that row is informative. A busy or malformed row
3706/// reports the gap as unavailable, never as zero. A gap alone does not
3707/// establish a reader pin. Zero
3708/// dependence on SQLite's shm WAL-index layout (ADR-091 Amendment 22).
3709fn log_backfill_gap(pool: &ConnectionPool, conn: &rusqlite::Connection) {
3710 match query_backfill_gap(conn) {
3711 Ok(observation) => {
3712 record_checkpoint_run_result(
3713 pool,
3714 Some((
3715 observation.busy,
3716 observation.log_frames,
3717 observation.checkpointed_frames,
3718 )),
3719 );
3720 if observed_wal_pages(observation).is_ok() {
3721 tracing::warn!(
3722 busy = observation.busy,
3723 wal_log_frames = observation.log_frames,
3724 wal_checkpointed_frames = observation.checkpointed_frames,
3725 backfill_gap_frames = observation
3726 .log_frames
3727 .saturating_sub(observation.checkpointed_frames)
3728 .max(0),
3729 "ADR-091 Plank C: WAL backfill gap after TRUNCATE no-progress"
3730 );
3731 } else {
3732 tracing::warn!(
3733 busy = observation.busy,
3734 wal_log_frames = observation.log_frames,
3735 wal_checkpointed_frames = observation.checkpointed_frames,
3736 "ADR-091 Plank C: WAL backfill gap unavailable after TRUNCATE"
3737 );
3738 }
3739 }
3740 Err(e) => {
3741 record_checkpoint_run_result(pool, None);
3742 tracing::warn!(
3743 error = %e,
3744 "ADR-091 Plank C: failed to query WAL backfill gap"
3745 );
3746 }
3747 }
3748}
3749
3750/// ADR-091 Amendment 2 Plank C: issue `PRAGMA wal_checkpoint(PASSIVE)` and
3751/// return its `(log, checkpointed)` columns (index 1 and 2 of the 3-column
3752/// return row). PASSIVE never blocks readers or writers. The backfill gap is
3753/// `log - checkpointed`; extracted as its own pure query so the arithmetic is
3754/// unit-testable against a real SQLite connection without depending on
3755/// `tracing` capture.
3756fn query_backfill_gap(conn: &rusqlite::Connection) -> rusqlite::Result<RawCheckpointObservation> {
3757 conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
3758 Ok(RawCheckpointObservation {
3759 busy: row.get(0)?,
3760 log_frames: row.get(1)?,
3761 checkpointed_frames: row.get(2)?,
3762 })
3763 })
3764}
3765
3766/// Evaluate whether a threshold-crossing WARN should fire and advance the
3767/// crossing-state flag.
3768///
3769/// Returns `true` on a false→true transition in `now_above` (first observed
3770/// above-threshold tick after a below-threshold tick), `false` on any other
3771/// tick. The `was_above` flag is updated in-place to track state across calls.
3772/// `run_checkpoint_task` uses this for the `high_water_pages` threshold;
3773/// `observe_wal_pages` owns the separate `warn_pages` severity ladder.
3774fn crossing_warn(now_above: bool, was_above: &mut bool) -> bool {
3775 let fire = now_above && !*was_above;
3776 *was_above = now_above;
3777 fire
3778}
3779
3780#[derive(Debug, Clone, Copy)]
3781struct RawCheckpointObservation {
3782 busy: i64,
3783 log_frames: i64,
3784 checkpointed_frames: i64,
3785}
3786
3787/// Why a syntactically valid SQLite checkpoint row has no usable frame count.
3788/// Keep the raw `busy` indication distinct from an inconsistent frame pair:
3789/// only SQLite's nonzero busy column may become a busy error or busy log.
3790#[derive(Debug, Clone, Copy, PartialEq, Eq)]
3791enum CheckpointUnavailableReason {
3792 Busy,
3793 InconsistentFrames,
3794}
3795
3796fn observed_wal_pages(
3797 observation: RawCheckpointObservation,
3798) -> Result<u64, CheckpointUnavailableReason> {
3799 if observation.busy != 0 {
3800 return Err(CheckpointUnavailableReason::Busy);
3801 }
3802 if observation.log_frames == -1 && observation.checkpointed_frames == -1 {
3803 // SQLite reports an absent WAL with two -1 frame columns.
3804 return Ok(0);
3805 }
3806 if observation.log_frames >= 0
3807 && observation.checkpointed_frames >= 0
3808 && observation.checkpointed_frames <= observation.log_frames
3809 {
3810 Ok(observation.log_frames as u64)
3811 } else {
3812 Err(CheckpointUnavailableReason::InconsistentFrames)
3813 }
3814}
3815
3816/// Issue one PASSIVE checkpoint and retain the complete SQLite result row.
3817/// This is the periodic task's one routine checkpoint call: the same row
3818/// drives thresholds and the logical-backlog monitoring sample (#1849).
3819fn query_checkpoint_observation(
3820 conn: &rusqlite::Connection,
3821) -> rusqlite::Result<RawCheckpointObservation> {
3822 conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
3823 Ok(RawCheckpointObservation {
3824 busy: row.get(0)?,
3825 log_frames: row.get(1)?,
3826 checkpointed_frames: row.get(2)?,
3827 })
3828 })
3829}
3830
3831/// The routine caller uses the real SQLite row. Unit tests can inject one
3832/// exact raw row for a uniquely keyed pool to exercise a frame combination
3833/// that SQLite does not normally emit without changing the production path.
3834fn query_routine_checkpoint_observation(
3835 pool: &ConnectionPool,
3836 conn: &rusqlite::Connection,
3837) -> rusqlite::Result<RawCheckpointObservation> {
3838 #[cfg(test)]
3839 if let Some(row) = test_take_passive_row(pool) {
3840 return Ok(row);
3841 }
3842 #[cfg(not(test))]
3843 let _ = pool;
3844 query_checkpoint_observation(conn)
3845}
3846
3847#[cfg(test)]
3848static TEST_PASSIVE_ROWS: OnceLock<Mutex<HashMap<Option<PathBuf>, RawCheckpointObservation>>> =
3849 OnceLock::new();
3850
3851#[cfg(test)]
3852fn test_passive_rows() -> &'static Mutex<HashMap<Option<PathBuf>, RawCheckpointObservation>> {
3853 TEST_PASSIVE_ROWS.get_or_init(|| Mutex::new(HashMap::new()))
3854}
3855
3856#[cfg(test)]
3857fn test_arm_passive_row(pool: &ConnectionPool, row: RawCheckpointObservation) {
3858 test_passive_rows()
3859 .lock()
3860 .unwrap_or_else(std::sync::PoisonError::into_inner)
3861 .insert(checkpoint_db_key(pool), row);
3862}
3863
3864#[cfg(test)]
3865fn test_take_passive_row(pool: &ConnectionPool) -> Option<RawCheckpointObservation> {
3866 test_passive_rows()
3867 .lock()
3868 .unwrap_or_else(std::sync::PoisonError::into_inner)
3869 .remove(&checkpoint_db_key(pool))
3870}
3871
3872fn query_truncate_observation(
3873 conn: &rusqlite::Connection,
3874) -> rusqlite::Result<RawCheckpointObservation> {
3875 conn.query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |row| {
3876 Ok(RawCheckpointObservation {
3877 busy: row.get(0)?,
3878 log_frames: row.get(1)?,
3879 checkpointed_frames: row.get(2)?,
3880 })
3881 })
3882}
3883
3884/// Query the current WAL frame count with one PASSIVE checkpoint.
3885///
3886/// Used only for rare post-TRUNCATE outcome measurement. The ordinary
3887/// periodic path calls [`query_checkpoint_observation`] directly and stores
3888/// its complete row, avoiding the former double-checkpoint pass.
3889fn query_wal_pages(pool: &ConnectionPool, conn: &rusqlite::Connection) -> Option<u64> {
3890 let observation = query_checkpoint_observation(conn);
3891 record_checkpoint_run_result(
3892 pool,
3893 observation.as_ref().ok().map(|observation| {
3894 (
3895 observation.busy,
3896 observation.log_frames,
3897 observation.checkpointed_frames,
3898 )
3899 }),
3900 );
3901 let pages = observation
3902 .ok()
3903 .and_then(|row| observed_wal_pages(row).ok());
3904 if let Some(pages) = pages {
3905 LAST_WAL_PAGES.store(pages, Ordering::Relaxed);
3906 note_checkpoint_observed(pages);
3907 }
3908 pages
3909}
3910
3911#[cfg(test)]
3912#[path = "checkpoint_tests.rs"]
3913mod tests;