Skip to main content

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;