Skip to main content

khive_db/
checkpoint.rs

1//! Periodic WAL checkpoint task for the connection pool (ADR-091).
2//!
3//! Issues `PRAGMA wal_checkpoint(PASSIVE)` on every tick — non-blocking, never
4//! waits for readers. A rare, separately-gated escalation may additionally run
5//! `PRAGMA wal_checkpoint(TRUNCATE)` once WAL pressure crosses
6//! `truncate_high_water_pages` and `truncate_min_interval` has elapsed since
7//! the last attempt (Plank 2); both run under the single writer checkout
8//! `checkpoint_once` holds for that tick. `checkpoint_once` uses
9//! `try_writer_nowait` (zero-wait `try_lock`) so a tick is skipped immediately
10//! when the writer mutex is held, rather than blocking — a skipped tick is
11//! always preferable to stalling write traffic.
12//!
13//! `warn_pages` / `high_water_pages` WARNs fire at most once per below→above
14//! crossing; a skipped tick leaves crossing state unchanged. An age-based
15//! background sweep (Plank 1) additionally checks the oldest span in
16//! `khive_storage::tx_registry` against `tx_warn_secs`/`tx_max_age_secs` on
17//! every tick (Skipped or Observed) and escalates to `warn!`/`error!` on each
18//! below→above crossing — visibility only, nothing here force-closes a stale
19//! span.
20//!
21//! See crates/khive-db/docs/api/checkpoint.md#module-overview-adr-091-planks-012
22//! for full ADR-091 Plank 0/1/2 design rationale (why TRUNCATE is excluded
23//! from ordinary ticks, the single-writer-checkout invariant, and why Plank 1
24//! is a sweep rather than the ADR's originally-described per-statement guard).
25
26use std::path::{Path, PathBuf};
27use std::sync::atomic::{AtomicU64, Ordering};
28use std::sync::Arc;
29use std::time::{Duration, Instant};
30
31use crate::pool::{ConnectionPool, WriterGuard};
32
33// ── metrics read-surface (load/perf harness) ─────────────────────────────
34// Read-only process-wide gauges (never reset outside #[cfg(test)]). See
35// crates/khive-db/docs/api/checkpoint.md#metrics-read-surface-loadperf-harness
36
37/// Last-observed WAL page count (`query_wal_pages`'s return value on its
38/// most recent call, from either `checkpoint_once` or `maybe_truncate`).
39/// `u64::MAX` is the "never observed" sentinel — no checkpoint tick has run
40/// yet in this process — distinct from a genuine zero-page WAL.
41static LAST_WAL_PAGES: AtomicU64 = AtomicU64::new(u64::MAX);
42
43/// Count of TRUNCATE attempts (`maybe_truncate`'s pragma actually invoked,
44/// win or lose) across this process's lifetime.
45static TRUNCATE_ATTEMPTS: AtomicU64 = AtomicU64::new(0);
46
47/// Current consecutive-failure count, mirrored from the caller-owned
48/// `TruncateState::consecutive_failures` field into a process-readable
49/// gauge every time `note_truncate_outcome` runs.
50static TRUNCATE_CONSECUTIVE_FAILURES: AtomicU64 = AtomicU64::new(0);
51
52/// Count of checkpoint ticks skipped because the writer mutex was already
53/// held (ADR-091 checkpoint-pressure telemetry), across this process's
54/// lifetime. Never reset outside `#[cfg(test)]`.
55static CHECKPOINT_SKIPPED_TICKS: AtomicU64 = AtomicU64::new(0);
56
57/// Current run-length of consecutive skipped ticks. Reset to 0 the next time
58/// a tick is actually observed (writer free), so a sustained skip streak is
59/// visible even between two successful observations.
60static CHECKPOINT_CONSECUTIVE_SKIPS: AtomicU64 = AtomicU64::new(0);
61
62/// WAL page count as of the most recent *observed* tick, snapshotted at the
63/// moment a skip occurs. `u64::MAX` is the "no skip has recorded a snapshot
64/// yet" sentinel, mirroring `LAST_WAL_PAGES`.
65static CHECKPOINT_LAST_SKIP_WAL_PAGES: AtomicU64 = AtomicU64::new(u64::MAX);
66
67/// Last-observed WAL page count, if any checkpoint tick has run yet in this
68/// process. Read surface for the daemon-frame metrics snapshot.
69pub fn last_observed_wal_pages() -> Option<u64> {
70    match LAST_WAL_PAGES.load(Ordering::Relaxed) {
71        u64::MAX => None,
72        pages => Some(pages),
73    }
74}
75
76/// Total WAL TRUNCATE attempts made in this process's lifetime.
77pub fn truncate_attempts() -> u64 {
78    TRUNCATE_ATTEMPTS.load(Ordering::Relaxed)
79}
80
81/// Current consecutive TRUNCATE-attempt failure count.
82pub fn truncate_consecutive_failures() -> u64 {
83    TRUNCATE_CONSECUTIVE_FAILURES.load(Ordering::Relaxed)
84}
85
86/// Total checkpoint ticks skipped (writer busy) in this process's lifetime.
87pub fn checkpoint_skipped_ticks() -> u64 {
88    CHECKPOINT_SKIPPED_TICKS.load(Ordering::Relaxed)
89}
90
91/// Current consecutive-skip run length; 0 once the next tick is observed.
92pub fn checkpoint_consecutive_skips() -> u64 {
93    CHECKPOINT_CONSECUTIVE_SKIPS.load(Ordering::Relaxed)
94}
95
96/// WAL page count last known at the time of the most recent skip, if any
97/// skip has occurred yet in this process.
98pub fn checkpoint_last_skip_wal_pages() -> Option<u64> {
99    match CHECKPOINT_LAST_SKIP_WAL_PAGES.load(Ordering::Relaxed) {
100        u64::MAX => None,
101        pages => Some(pages),
102    }
103}
104
105/// A tick's writer checkout was skipped (mutex busy): bump the lifetime and
106/// consecutive-skip counters and snapshot the last-known WAL pressure so an
107/// operator can see how bad the WAL was heading into the skip streak.
108fn note_checkpoint_skipped() {
109    CHECKPOINT_SKIPPED_TICKS.fetch_add(1, Ordering::Relaxed);
110    CHECKPOINT_CONSECUTIVE_SKIPS.fetch_add(1, Ordering::Relaxed);
111    if let Some(pages) = last_observed_wal_pages() {
112        CHECKPOINT_LAST_SKIP_WAL_PAGES.store(pages, Ordering::Relaxed);
113    }
114}
115
116/// A tick was actually observed (writer free): close out any prior skip
117/// streak. `_wal_pages` is accepted for call-site symmetry with
118/// `note_checkpoint_skipped` and to leave room for a future observed-side
119/// gauge without changing this function's signature again.
120fn note_checkpoint_observed(_wal_pages: u64) {
121    CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
122}
123
124/// Reset the checkpoint-pressure atomics between tests. Process-wide gauges
125/// are otherwise shared across every test in this binary; tests that assert
126/// on them must reset first and run under a shared `#[serial(...)]` group.
127#[cfg(test)]
128pub(crate) fn reset_checkpoint_metrics_for_tests() {
129    CHECKPOINT_SKIPPED_TICKS.store(0, Ordering::Relaxed);
130    CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
131    CHECKPOINT_LAST_SKIP_WAL_PAGES.store(u64::MAX, Ordering::Relaxed);
132}
133
134/// Outcome of a single checkpoint attempt.
135///
136/// `Skipped` is returned when the writer mutex is already held (the tick is a
137/// no-op). `Observed` carries the WAL page count read during the tick. The
138/// distinction matters for threshold-crossing WARN rate-limiting: a skipped tick
139/// must leave the above/below state unchanged so that a busy tick cannot
140/// spuriously re-arm the rate limit while WAL pressure is still elevated.
141#[derive(Debug, Clone, Copy, PartialEq, Eq)]
142pub enum CheckpointTick {
143    /// The writer mutex was busy; no checkpoint was issued this tick.
144    Skipped,
145    /// A checkpoint was issued; the value is the observed WAL page count.
146    Observed(u64),
147}
148
149/// Default number of consecutive above-`warn_pages` observed ticks required
150/// to escalate from the INFO to the WARN rung of the ADR-091 severity ladder.
151pub const DEFAULT_WARN_SUSTAINED_CYCLES: u8 = 3;
152
153/// Configuration for the WAL checkpoint background task.
154///
155/// All fields default to conservative production values. Override via the
156/// environment variables documented on each field.
157#[derive(Clone, Debug)]
158pub struct CheckpointConfig {
159    /// How often to run a passive checkpoint when there is no active write.
160    ///
161    /// Overridable via `KHIVE_CHECKPOINT_INTERVAL_MS` (milliseconds).
162    /// Default: 500 ms.
163    pub interval: Duration,
164
165    /// WAL page count above which a warning is logged.
166    ///
167    /// Overridable via `KHIVE_WAL_WARN_PAGES`.
168    /// Default: 2000 pages (~8 MB at 4 KiB page size).
169    pub warn_pages: u64,
170
171    /// Number of consecutive observed ticks with `wal_pages >= warn_pages`
172    /// required before the ADR-091 severity ladder escalates from INFO
173    /// (first crossing) to WARN (sustained pressure). Edge-triggered once
174    /// per elevation episode — see [`CheckpointSeverityState`].
175    ///
176    /// Overridable via `KHIVE_WAL_WARN_SUSTAINED_CYCLES`.
177    /// Default: 3 cycles.
178    pub warn_sustained_cycles: u8,
179
180    /// WAL page count above which a high-pressure WARNING is logged.
181    ///
182    /// The periodic task always runs PASSIVE regardless; this threshold signals
183    /// that a long-lived reader may be pinning an old WAL snapshot that PASSIVE
184    /// cannot reclaim. An operator can then schedule a blocking TRUNCATE at a
185    /// safe moment outside normal write traffic.
186    ///
187    /// Overridable via `KHIVE_WAL_HIGH_WATER_PAGES`.
188    /// Default: 6000 pages (~24 MB at 4 KiB page size).
189    pub high_water_pages: u64,
190
191    /// WAL page count above which a TRUNCATE escalation attempt is armed
192    /// (ADR-091 Plank 2).
193    ///
194    /// This is a separate, much higher threshold than `high_water_pages`:
195    /// crossing it does not itself attempt TRUNCATE — it only arms the
196    /// attempt, which additionally requires `truncate_min_interval` to have
197    /// elapsed since the last attempt.
198    ///
199    /// Overridable via `KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES`.
200    /// Default: 20000 pages.
201    pub truncate_high_water_pages: u64,
202
203    /// Minimum spacing between TRUNCATE *attempts* (not successes).
204    ///
205    /// A skipped tick (writer busy, below threshold, or interval not yet
206    /// elapsed) never advances the "last attempt" clock, so the next tick
207    /// where the writer is free and the threshold is still crossed is
208    /// immediately eligible rather than waiting out the full interval again.
209    ///
210    /// Overridable via `KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS`.
211    /// Default: 300 seconds (5 minutes).
212    pub truncate_min_interval: Duration,
213
214    /// Temporary `busy_timeout` used only for the duration of a TRUNCATE
215    /// attempt, restored to the pool's configured busy timeout immediately
216    /// after the attempt completes (win or lose).
217    ///
218    /// Overridable via `KHIVE_WAL_TRUNCATE_BUSY_MS`.
219    /// Default: 2000 ms.
220    pub truncate_busy_timeout: Duration,
221
222    /// ADR-091 Plank 1 soft cap: age past which the oldest entry in the
223    /// shared open-transaction registry is surfaced at `tracing::warn!` on
224    /// every tick (Skipped or Observed), independent of WAL page pressure.
225    /// See `crates/khive-db/docs/api/checkpoint.md` for the Plank 1 rationale.
226    ///
227    /// Overridable via `KHIVE_TX_WARN_SECS`.
228    /// Default: 30 seconds.
229    pub tx_warn_secs: Duration,
230
231    /// ADR-091 Plank 1 hard cap: age past which the same sweep escalates the
232    /// oldest registry entry to `tracing::error!`. This is visibility only —
233    /// nothing here can force-close a stale span; see
234    /// `crates/khive-db/docs/design.md` for why.
235    ///
236    /// Overridable via `KHIVE_TX_MAX_AGE_SECS`.
237    /// Default: 120 seconds.
238    pub tx_max_age_secs: Duration,
239}
240
241impl Default for CheckpointConfig {
242    fn default() -> Self {
243        Self {
244            interval: Duration::from_millis(500),
245            warn_pages: 2000,
246            warn_sustained_cycles: DEFAULT_WARN_SUSTAINED_CYCLES,
247            high_water_pages: 6000,
248            truncate_high_water_pages: 20_000,
249            truncate_min_interval: Duration::from_secs(300),
250            truncate_busy_timeout: Duration::from_millis(2000),
251            tx_warn_secs: Duration::from_secs(30),
252            tx_max_age_secs: Duration::from_secs(120),
253        }
254    }
255}
256
257impl CheckpointConfig {
258    /// Build a `CheckpointConfig` from the environment.
259    ///
260    /// Unset or unparseable variables fall back to the compiled-in defaults.
261    pub fn from_env() -> Self {
262        let mut cfg = Self::default();
263
264        if let Ok(ms) = std::env::var("KHIVE_CHECKPOINT_INTERVAL_MS") {
265            if let Ok(v) = ms.parse::<u64>() {
266                if v > 0 {
267                    cfg.interval = Duration::from_millis(v);
268                }
269            }
270        }
271
272        if let Ok(v) = std::env::var("KHIVE_WAL_WARN_PAGES") {
273            if let Ok(n) = v.parse::<u64>() {
274                if n > 0 {
275                    cfg.warn_pages = n;
276                }
277            }
278        }
279
280        if let Ok(v) = std::env::var("KHIVE_WAL_WARN_SUSTAINED_CYCLES") {
281            if let Ok(n) = v.parse::<u8>() {
282                if n > 0 {
283                    cfg.warn_sustained_cycles = n;
284                }
285            }
286        }
287
288        if let Ok(v) = std::env::var("KHIVE_WAL_HIGH_WATER_PAGES") {
289            if let Ok(n) = v.parse::<u64>() {
290                if n > 0 {
291                    cfg.high_water_pages = n;
292                }
293            }
294        }
295
296        if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES") {
297            if let Ok(n) = v.parse::<u64>() {
298                if n > 0 {
299                    cfg.truncate_high_water_pages = n;
300                }
301            }
302        }
303
304        if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS") {
305            if let Ok(n) = v.parse::<u64>() {
306                if n > 0 {
307                    cfg.truncate_min_interval = Duration::from_secs(n);
308                }
309            }
310        }
311
312        if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_BUSY_MS") {
313            if let Ok(n) = v.parse::<u64>() {
314                if n > 0 {
315                    cfg.truncate_busy_timeout = Duration::from_millis(n);
316                }
317            }
318        }
319
320        (cfg.tx_warn_secs, cfg.tx_max_age_secs) =
321            tx_age_thresholds_from_env(cfg.tx_warn_secs, cfg.tx_max_age_secs);
322
323        cfg
324    }
325}
326
327/// Parse `KHIVE_TX_WARN_SECS`/`KHIVE_TX_MAX_AGE_SECS` against the given
328/// defaults, applying the same ordering guard both [`CheckpointConfig`] and
329/// [`SessionSweepConfig`] need (minor, ADR-091 Amendment 2: this was
330/// previously duplicated verbatim in both `from_env` methods).
331///
332/// The severity ladder assumes `tx_warn_secs < tx_max_age_secs` (Warn fires
333/// before Stale as an entry ages). A reversed or equal pair — whether from
334/// one misconfigured var or the interaction of both — would invert or
335/// collapse that ordering (e.g. WARN_SECS=120, MAX_AGE_SECS=30 emits Stale at
336/// 30s and never reaches the Warn crossing until 120s), so both are rejected
337/// together rather than silently honored. Resetting both to the caller's
338/// defaults (rather than just clamping one) avoids guessing which of the two
339/// the operator actually meant to change.
340fn tx_age_thresholds_from_env(
341    default_warn: Duration,
342    default_max: Duration,
343) -> (Duration, Duration) {
344    let mut warn_secs = default_warn;
345    let mut max_age_secs = default_max;
346
347    if let Ok(v) = std::env::var("KHIVE_TX_WARN_SECS") {
348        if let Ok(n) = v.parse::<u64>() {
349            if n > 0 {
350                warn_secs = Duration::from_secs(n);
351            }
352        }
353    }
354
355    if let Ok(v) = std::env::var("KHIVE_TX_MAX_AGE_SECS") {
356        if let Ok(n) = v.parse::<u64>() {
357            if n > 0 {
358                max_age_secs = Duration::from_secs(n);
359            }
360        }
361    }
362
363    if warn_secs >= max_age_secs {
364        tracing::warn!(
365            configured_tx_warn_secs = warn_secs.as_secs_f64(),
366            configured_tx_max_age_secs = max_age_secs.as_secs_f64(),
367            fallback_tx_warn_secs = default_warn.as_secs_f64(),
368            fallback_tx_max_age_secs = default_max.as_secs_f64(),
369            "KHIVE_TX_WARN_SECS must be strictly less than KHIVE_TX_MAX_AGE_SECS; \
370             both transaction-age thresholds were rejected and reset to their defaults"
371        );
372        return (default_warn, default_max);
373    }
374
375    (warn_secs, max_age_secs)
376}
377
378/// Mutable escalation state carried across ticks by the caller (ADR-091 Plank 2).
379///
380/// Kept separate from [`CheckpointConfig`] because it is *state*, not
381/// configuration: `last_attempt` and `consecutive_failures` mutate every tick,
382/// while `CheckpointConfig` is parsed once and held immutable for the life of
383/// the task.
384#[derive(Debug, Default)]
385pub struct TruncateState {
386    /// When the last TRUNCATE *attempt* ran (armed + writer held), regardless
387    /// of whether it succeeded in reclaiming pages. `None` means no attempt
388    /// has ever run, so the first armed tick is immediately eligible.
389    last_attempt: Option<Instant>,
390    /// Count of consecutive TRUNCATE attempts that failed to bring `wal_pages`
391    /// back below `warn_pages`. Resets to 0 the first time an attempt clears
392    /// `warn_pages`; used to fire a one-shot escalated WARN at exactly 3
393    /// consecutive failures (does not repeat every subsequent attempt).
394    consecutive_failures: u32,
395}
396
397/// ADR-091 graduated severity rung for sustained WAL pressure.
398///
399/// `Alarm` is never produced by [`CheckpointSeverityState::observe_wal_pages`]
400/// — it labels the existing TRUNCATE-escalation tier (`maybe_truncate`),
401/// which is gated on its own threshold/interval state, not on this ladder.
402/// It exists here so callers and tests can name all three rungs uniformly.
403#[derive(Debug, Clone, Copy, PartialEq, Eq)]
404pub enum CheckpointSeverityRung {
405    /// First observed tick crossing `warn_pages` after a below-warn tick.
406    Info,
407    /// `warn_sustained_cycles` consecutive observed ticks at/above
408    /// `warn_pages`; edge-triggered once per elevation episode.
409    Warn,
410    /// The TRUNCATE-escalation tier (`checkpoint_high_water_pages` and
411    /// above); never emitted by `observe_wal_pages`.
412    Alarm,
413}
414
415/// ADR-091 severity ladder state, carried across ticks by the caller
416/// alongside [`TruncateState`]. Pure state machine: no I/O, no logging —
417/// callers turn the returned emissions into `tracing` calls.
418#[derive(Debug, Default, Clone)]
419pub struct CheckpointSeverityState {
420    /// Whether the previous observed tick was at/above `warn_pages`. Drives
421    /// the below→above edge that fires INFO.
422    was_above_warn: bool,
423    /// Run-length of consecutive observed ticks at/above `warn_pages` in the
424    /// current elevation episode. Resets to 0 on any below-warn tick.
425    consecutive_above_warn: u8,
426    /// Whether WARN has already fired for the current elevation episode, so
427    /// sustained pressure logs WARN once per episode, not once per tick past
428    /// the threshold.
429    warn_emitted_for_episode: bool,
430}
431
432/// One severity-ladder emission produced by a single
433/// [`CheckpointSeverityState::observe_wal_pages`] call.
434#[derive(Debug, Clone, Copy, PartialEq, Eq)]
435pub struct CheckpointSeverityEmission {
436    /// Which rung this emission represents (`Info` or `Warn`; see
437    /// [`CheckpointSeverityRung::Alarm`] doc for why `Alarm` never appears
438    /// here).
439    pub rung: CheckpointSeverityRung,
440    /// The WAL page count observed on the tick that produced this emission.
441    pub wal_pages: u64,
442    /// The `warn_pages` threshold in effect for this tick.
443    pub threshold_pages: u64,
444    /// Consecutive above-warn cycle count as of this tick (1 on the INFO
445    /// edge, `warn_sustained_cycles` on the WARN edge).
446    pub consecutive_cycles: u8,
447}
448
449impl CheckpointSeverityState {
450    /// Advance the severity ladder by one observed tick and return every
451    /// rung crossed on this tick (zero, one, or two emissions: a fresh
452    /// elevation episode can produce INFO and, if `warn_sustained_cycles`
453    /// is 1, WARN on the very same tick).
454    ///
455    /// A below-warn tick resets both the consecutive-cycle counter and the
456    /// per-episode WARN latch, re-arming INFO/WARN for a later episode.
457    /// Skipped ticks must not be passed here at all — the caller only calls
458    /// this on `CheckpointTick::Observed`, matching the existing
459    /// threshold-crossing WARN's skip-leaves-state-unchanged rule.
460    pub fn observe_wal_pages(
461        &mut self,
462        wal_pages: u64,
463        config: &CheckpointConfig,
464    ) -> Vec<CheckpointSeverityEmission> {
465        let mut emissions = Vec::new();
466        let above_warn = wal_pages >= config.warn_pages;
467
468        if above_warn {
469            self.consecutive_above_warn = self.consecutive_above_warn.saturating_add(1);
470
471            if !self.was_above_warn {
472                emissions.push(CheckpointSeverityEmission {
473                    rung: CheckpointSeverityRung::Info,
474                    wal_pages,
475                    threshold_pages: config.warn_pages,
476                    consecutive_cycles: self.consecutive_above_warn,
477                });
478            }
479
480            if !self.warn_emitted_for_episode
481                && self.consecutive_above_warn >= config.warn_sustained_cycles
482            {
483                emissions.push(CheckpointSeverityEmission {
484                    rung: CheckpointSeverityRung::Warn,
485                    wal_pages,
486                    threshold_pages: config.warn_pages,
487                    consecutive_cycles: self.consecutive_above_warn,
488                });
489                self.warn_emitted_for_episode = true;
490            }
491        } else {
492            self.consecutive_above_warn = 0;
493            self.warn_emitted_for_episode = false;
494        }
495
496        self.was_above_warn = above_warn;
497        emissions
498    }
499}
500
501/// ADR-091 Plank 1 rung for the open-transaction registry's background age
502/// sweep: independent of the WAL-pressure ladder above, keyed purely off how
503/// long the registry's oldest entry has been open.
504#[derive(Debug, Clone, Copy, PartialEq, Eq)]
505pub enum TxAgeRung {
506    /// The oldest registry entry's age crossed `tx_warn_secs`.
507    Warn,
508    /// The oldest registry entry's age crossed `tx_max_age_secs` — the ADR's
509    /// "cooperative stale-op guard" cap. No in-process mechanism force-closes
510    /// it (see [`CheckpointConfig::tx_max_age_secs`]); this rung is the
511    /// sweep's strongest available signal.
512    Stale,
513}
514
515/// One emission produced by a single [`TxAgeSweepState::observe`] call.
516#[derive(Debug, Clone, PartialEq, Eq)]
517pub struct TxAgeEmission {
518    pub rung: TxAgeRung,
519    pub age: Duration,
520    pub label: Option<String>,
521}
522
523/// ADR-091 Plank 1 background-sweep state, carried across ticks by the
524/// caller alongside [`CheckpointSeverityState`] and [`TruncateState`]. Pure
525/// state machine: no I/O, no logging — callers turn the returned emissions
526/// into `tracing` calls, mirroring [`CheckpointSeverityState`]'s shape.
527///
528/// Keyed off `khive_storage::tx_registry::oldest()` — the single oldest
529/// entry across every registered span, regardless of which call site created
530/// it. Deliberately a different signal from the WAL-pressure ladder: a span
531/// can go stale under low WAL pressure, or vice versa. See
532/// `crates/khive-db/docs/api/checkpoint.md` for the full rationale.
533#[derive(Debug, Default, Clone)]
534pub struct TxAgeSweepState {
535    /// Whether the previous observed tick's oldest entry was at/above
536    /// `tx_warn_secs`. Drives the below→above edge that fires `Warn`.
537    was_above_warn: bool,
538    /// Whether the previous observed tick's oldest entry was at/above
539    /// `tx_max_age_secs`. Drives the below→above edge that fires `Stale`.
540    was_above_max_age: bool,
541    /// Identity of the entry the previous observed tick reported as oldest,
542    /// or `None` if the registry was empty. Tracked separately from the two
543    /// latches above so a change in *which span* is oldest can be detected
544    /// even when both latches are already `true` (see [`Self::observe`]).
545    tracked_id: Option<khive_storage::tx_registry::TxId>,
546}
547
548impl TxAgeSweepState {
549    /// Advance by one observed tick given the registry's current oldest
550    /// entry (identity, age, label), or `None` if empty. Returns zero, one,
551    /// or two emissions — an entry already stale the first time it's seen
552    /// under a given identity crosses both rungs on the same tick.
553    ///
554    /// A below-threshold (or absent) oldest entry resets both latches. A
555    /// change in the oldest entry's [`TxId`](khive_storage::tx_registry::TxId)
556    /// also force-resets both latches before re-evaluating age, so a
557    /// departed span's latched state cannot suppress the crossing for an
558    /// already-stale successor. See `crates/khive-db/docs/api/checkpoint.md`
559    /// for why identity tracking is required here, not just the age check.
560    pub fn observe(
561        &mut self,
562        oldest: Option<(khive_storage::tx_registry::TxId, Duration, Option<String>)>,
563        tx_warn_secs: Duration,
564        tx_max_age_secs: Duration,
565    ) -> Vec<TxAgeEmission> {
566        let mut emissions = Vec::new();
567
568        let Some((id, age, label)) = oldest else {
569            self.was_above_warn = false;
570            self.was_above_max_age = false;
571            self.tracked_id = None;
572            return emissions;
573        };
574
575        if self.tracked_id != Some(id) {
576            self.was_above_warn = false;
577            self.was_above_max_age = false;
578        }
579        self.tracked_id = Some(id);
580
581        let above_warn = age >= tx_warn_secs;
582        let above_max_age = age >= tx_max_age_secs;
583
584        if above_warn && !self.was_above_warn {
585            emissions.push(TxAgeEmission {
586                rung: TxAgeRung::Warn,
587                age,
588                label: label.clone(),
589            });
590        }
591        if above_max_age && !self.was_above_max_age {
592            emissions.push(TxAgeEmission {
593                rung: TxAgeRung::Stale,
594                age,
595                label,
596            });
597        }
598
599        self.was_above_warn = above_warn;
600        self.was_above_max_age = above_max_age;
601        emissions
602    }
603}
604
605/// ADR-091 Plank 1: turn a [`TxAgeEmission`] into the appropriate `tracing`
606/// call. Extracted from `run_checkpoint_task` so tests can drive the same
607/// logging path `CaptureSubscriber`-style without spinning up the async task
608/// (mirrors [`log_tx_registry_oldest_warn`]/[`log_tx_registry_snapshot_warn`]).
609fn log_tx_age_emission(emission: &TxAgeEmission) {
610    let label = emission.label.as_deref().unwrap_or("<unlabeled>");
611    match emission.rung {
612        TxAgeRung::Warn => {
613            tracing::warn!(
614                tx_age_secs = emission.age.as_secs_f64(),
615                tx_label = label,
616                "ADR-091 Plank 1: open transaction registry entry exceeded soft-cap age"
617            );
618        }
619        TxAgeRung::Stale => {
620            tracing::error!(
621                tx_age_secs = emission.age.as_secs_f64(),
622                tx_label = label,
623                "ADR-091 Plank 1: open transaction registry entry exceeded the cooperative \
624                 stale-op cap; no in-process mechanism can force-close it — investigate the \
625                 labeled caller directly"
626            );
627        }
628    }
629}
630
631/// ADR-091 Amendment 2 Plank B: per-process walpin sidecar state, carried
632/// across ticks by whichever sweep owns it (the daemon's `run_checkpoint_task`
633/// or a session's `run_session_sweep_task`). Once the registry's oldest span
634/// exceeds `tx_warn_secs`, the first observation and each content change
635/// rewrite the heartbeat body; unchanged ticks refresh only its mtime. The
636/// heartbeat is removed once when the condition clears (and on shutdown), so
637/// a process that never crosses the threshold writes no heartbeat body.
638struct WalpinSidecarState {
639    dir: PathBuf,
640    pid: u32,
641    role: &'static str,
642    started_at: i64,
643    /// This sweep's own tick cadence, recorded into every beacon and
644    /// heartbeat so the enumerating daemon judges freshness against the
645    /// PRODUCER's interval — a session on an independently slower configured
646    /// cadence must not be misread as stale.
647    sweep_interval_ms: u64,
648    wrote: bool,
649    /// Whether this process's registration beacon is believed present on
650    /// disk. Cleared when a failed heartbeat write escalates to beacon
651    /// removal (fail-closed — see `observe`) or a beacon touch fails; the
652    /// next healthy tick then re-registers with a full write instead of a
653    /// metadata touch.
654    beacon_registered: bool,
655    /// The content actually on disk in the last successful heartbeat body
656    /// write, if any (ADR-091 Amendment 3 Plank F1). `None` whenever the
657    /// next tick must go through a full write — no heartbeat written yet,
658    /// the last write failed, or the threshold cleared. Compared against
659    /// each new observation to decide touch (content unchanged) vs.
660    /// rewrite (content changed).
661    last_heartbeat: Option<LastHeartbeatState>,
662}
663
664/// ADR-091 Amendment 3 Plank F1: the content signature of the heartbeat
665/// body currently on disk, plus the `oldest_tx_started_at` value that body
666/// carries — kept separate from the signature proper because it is derived
667/// (fixed for as long as the same span stays oldest), not an independent
668/// change signal.
669struct LastHeartbeatState {
670    span_id: khive_storage::tx_registry::TxId,
671    label: Option<String>,
672    attribution_basis: &'static str,
673    sweep_interval_ms: u64,
674    oldest_tx_started_at: i64,
675}
676
677impl LastHeartbeatState {
678    /// Whether a fresh observation carries exactly the content already on
679    /// disk — the licensing condition for a metadata-only touch instead of
680    /// a full body rewrite (the first over-threshold observation, a change
681    /// of the oldest span's identity or label, a change of
682    /// `attribution_basis`, or a change of the declared sweep cadence).
683    fn content_matches(
684        &self,
685        span_id: khive_storage::tx_registry::TxId,
686        label: &Option<String>,
687        attribution_basis: &str,
688        sweep_interval_ms: u64,
689    ) -> bool {
690        self.span_id == span_id
691            && self.label == *label
692            && self.attribution_basis == attribution_basis
693            && self.sweep_interval_ms == sweep_interval_ms
694    }
695}
696
697impl WalpinSidecarState {
698    /// `None` when the sidecar is disabled for this backend/env, or the
699    /// backend has no on-disk path (in-memory).
700    fn new(
701        db_path: Option<&Path>,
702        is_file_backed: bool,
703        role: &'static str,
704        interval: Duration,
705    ) -> Option<Self> {
706        let path = db_path?;
707        if !crate::walpin::sidecar_enabled(is_file_backed) {
708            return None;
709        }
710        let pid = std::process::id();
711        Some(Self {
712            dir: crate::walpin::sidecar_dir_for(path),
713            pid,
714            role,
715            started_at: crate::walpin::process_start_time_secs(pid).unwrap_or(0),
716            sweep_interval_ms: interval.as_millis().min(u64::MAX as u128) as u64,
717            wrote: false,
718            last_heartbeat: None,
719            beacon_registered: false,
720        })
721    }
722
723    /// Write this process's registration beacon (ADR-091 Amendment 2
724    /// sidecar-health attribution). Called once right after construction,
725    /// before the sweep loop starts, and again only when a fail-closed
726    /// removal or failed touch cleared `beacon_registered` — steady state
727    /// stays metadata-touch-only with no data writes. The blocking fs I/O
728    /// runs on `spawn_blocking` (perf, ADR-091 Amendment 2): this is
729    /// invoked from an async context and must not run synchronous I/O
730    /// inline on the async runtime's worker thread.
731    async fn register_beacon(&mut self) {
732        let dir = self.dir.clone();
733        let beacon = crate::walpin::WalpinBeacon {
734            pid: self.pid,
735            process_role: self.role.to_string(),
736            started_at: self.started_at,
737            sweep_interval_ms: self.sweep_interval_ms,
738        };
739        let result =
740            tokio::task::spawn_blocking(move || crate::walpin::write_beacon(&dir, &beacon)).await;
741        match result {
742            Ok(Ok(())) => {
743                self.beacon_registered = true;
744            }
745            Ok(Err(e)) => {
746                tracing::warn!(
747                    error = %e,
748                    "ADR-091 Amendment 2: failed to write walpin registration beacon; \
749                     this process's sidecar health will read as unknown, not registered-silent"
750                );
751            }
752            Err(join_err) => {
753                tracing::warn!(
754                    error = %join_err,
755                    "ADR-091 Amendment 2: walpin beacon write task panicked"
756                );
757            }
758        }
759    }
760
761    /// ADR-091 Amendment 2 beacon refresh rule: a metadata-only mtime touch
762    /// of this process's already-registered beacon, performed on every
763    /// sweep tick except one where an over-threshold heartbeat write failed
764    /// (see `observe`) — `registered-silent` classification requires this
765    /// refresh to stay within the freshness window, not just the beacon's
766    /// original write. After a fail-closed beacon removal (or a failed
767    /// touch), the beacon is re-registered with a full write on the next
768    /// healthy tick. Best-effort: a failure here degrades this process to
769    /// `unknown` at the next enumeration, not a sweep-task error.
770    async fn refresh_beacon(&mut self) {
771        if !self.beacon_registered {
772            self.register_beacon().await;
773            return;
774        }
775        let dir = self.dir.clone();
776        let pid = self.pid;
777        let result =
778            tokio::task::spawn_blocking(move || crate::walpin::touch_beacon(&dir, pid)).await;
779        match result {
780            Ok(Ok(())) => {}
781            Ok(Err(e)) => {
782                self.beacon_registered = false;
783                tracing::warn!(
784                    error = %e,
785                    "ADR-091 Amendment 2: failed to refresh walpin registration beacon; \
786                     this process's sidecar health will read as unknown, not registered-silent"
787                );
788            }
789            Err(join_err) => {
790                self.beacon_registered = false;
791                tracing::warn!(
792                    error = %join_err,
793                    "ADR-091 Amendment 2: walpin beacon refresh task panicked"
794                );
795            }
796        }
797    }
798
799    /// Fail-closed escalation for a failed heartbeat write: remove this
800    /// process's beacon so enumeration cannot classify it
801    /// `registered-silent` off the still-fresh prior refresh — skipping one
802    /// touch alone leaves the previous mtime inside the freshness window
803    /// for up to three producer ticks, an exoneration window. With the
804    /// beacon gone the process either reports (once writes recover, the
805    /// next tick re-registers and writes the heartbeat) or is caught by the
806    /// OS-level holder census as an unattributed holder. If the removal
807    /// itself fails, the beacon ages out over the freshness window — the
808    /// narrowed fallback, not the contract.
809    async fn drop_beacon_fail_closed(&mut self) {
810        let dir = self.dir.clone();
811        let pid = self.pid;
812        self.beacon_registered = false;
813        let result =
814            tokio::task::spawn_blocking(move || crate::walpin::remove_beacon(&dir, pid)).await;
815        match result {
816            Ok(Ok(())) => {}
817            Ok(Err(e)) => {
818                tracing::warn!(
819                    error = %e,
820                    "ADR-091 Amendment 2: failed to remove walpin beacon after a failed \
821                     heartbeat write; beacon will age out of the freshness window instead"
822                );
823            }
824            Err(join_err) => {
825                tracing::warn!(
826                    error = %join_err,
827                    "ADR-091 Amendment 2: walpin beacon removal task panicked"
828                );
829            }
830        }
831    }
832
833    /// Blocking heartbeat write/removal runs on `spawn_blocking` (perf,
834    /// ADR-091 Amendment 2) — this async sweep task must not block its
835    /// executor thread on synchronous filesystem I/O.
836    async fn observe(
837        &mut self,
838        oldest: Option<khive_storage::tx_registry::OldestSpan>,
839        tx_warn_secs: Duration,
840    ) {
841        match oldest {
842            Some(span) if span.age >= tx_warn_secs => {
843                // ADR-091 Amendment 3 Plank F2: the caller's `TxOriginFilter`
844                // guarantees a `Main` view's winner is either `Database` (this
845                // backend's own identity) or `Unscoped` (the fallback), and a
846                // `Secondary` view's winner is always `Database` — `Memory`
847                // can never win a filtered query, so it degrades to
848                // fallback-confidence rather than a reachability panic.
849                let attribution_basis = match span.origin {
850                    khive_storage::tx_registry::TxOrigin::Database(_) => "origin",
851                    khive_storage::tx_registry::TxOrigin::Unscoped
852                    | khive_storage::tx_registry::TxOrigin::Memory => "fallback",
853                };
854
855                // ADR-091 Amendment 3 Plank F1: a metadata-only mtime touch
856                // advances freshness whenever nothing content-relevant has
857                // changed since the last body write; a full rewrite happens
858                // only on the first over-threshold observation or a genuine
859                // content change.
860                let content_unchanged = self.wrote
861                    && self.last_heartbeat.as_ref().is_some_and(|last| {
862                        last.content_matches(
863                            span.id,
864                            &span.label,
865                            attribution_basis,
866                            self.sweep_interval_ms,
867                        )
868                    });
869
870                if content_unchanged {
871                    let dir = self.dir.clone();
872                    let pid = self.pid;
873                    let touch_result = tokio::task::spawn_blocking(move || {
874                        crate::walpin::touch_heartbeat(&dir, pid)
875                    })
876                    .await;
877                    match touch_result {
878                        Ok(Ok(())) => {
879                            self.refresh_beacon().await;
880                            return;
881                        }
882                        Ok(Err(e)) => {
883                            tracing::warn!(
884                                error = %e,
885                                "ADR-091 Amendment 3 Plank F1: walpin heartbeat touch failed; \
886                                 recreating with a full body write"
887                            );
888                        }
889                        Err(join_err) => {
890                            tracing::warn!(
891                                error = %join_err,
892                                "ADR-091 Amendment 3 Plank F1: walpin heartbeat touch task \
893                                 panicked; recreating with a full body write"
894                            );
895                        }
896                    }
897                    // Recovery rule: the touch path must never assume the
898                    // target still exists — enumeration can delete a slow
899                    // writer's heartbeat while its span is still live. Fall
900                    // through to the full write below unconditionally.
901                }
902
903                // The oldest span's registration instant is fixed for as
904                // long as it stays the SAME span: reuse the previously
905                // recorded value rather than re-deriving it from `now -
906                // age`, which would drift by measurement noise across ticks
907                // for no reason. A genuinely new oldest span (or the first
908                // observation) derives it fresh.
909                let oldest_tx_started_at = self
910                    .last_heartbeat
911                    .as_ref()
912                    .filter(|last| last.span_id == span.id)
913                    .map(|last| last.oldest_tx_started_at)
914                    .unwrap_or_else(|| now_epoch_secs().saturating_sub(span.age.as_secs() as i64));
915
916                let heartbeat = crate::walpin::WalpinHeartbeat {
917                    pid: self.pid,
918                    process_role: self.role.to_string(),
919                    started_at: self.started_at,
920                    oldest_tx_age_secs: span.age.as_secs_f64(),
921                    oldest_tx_label: span.label.clone(),
922                    oldest_tx_started_at: Some(oldest_tx_started_at),
923                    updated_at: now_epoch_secs(),
924                    sweep_interval_ms: self.sweep_interval_ms,
925                    attribution_basis: Some(attribution_basis.to_string()),
926                };
927                let dir = self.dir.clone();
928                let result = tokio::task::spawn_blocking(move || {
929                    crate::walpin::write_heartbeat(&dir, &heartbeat)
930                })
931                .await;
932                // The beacon refresh is gated on the heartbeat write
933                // landing: a fresh beacon with no heartbeat file classifies
934                // as `registered-silent` at enumeration, so a failed write
935                // would exonerate a process that currently holds an
936                // over-threshold transaction. Skipping the refresh alone is
937                // not enough — the previous touch stays inside the freshness
938                // window for up to three producer ticks — so the failure
939                // path removes the beacon outright (`drop_beacon_fail_closed`);
940                // the next successful tick re-registers it.
941                match result {
942                    Ok(Ok(())) => {
943                        self.wrote = true;
944                        self.last_heartbeat = Some(LastHeartbeatState {
945                            span_id: span.id,
946                            label: span.label,
947                            attribution_basis,
948                            sweep_interval_ms: self.sweep_interval_ms,
949                            oldest_tx_started_at,
950                        });
951                        self.refresh_beacon().await;
952                    }
953                    Ok(Err(e)) => {
954                        tracing::warn!(
955                            error = %e,
956                            "ADR-091 Amendment 2 Plank B: failed to write walpin heartbeat; \
957                             removing beacon so this process cannot read as \
958                             registered-silent while over threshold"
959                        );
960                        // Unknown what (if anything) is on disk now — the
961                        // next tick must go through a full write, never a
962                        // touch, until a write actually lands.
963                        self.last_heartbeat = None;
964                        self.drop_beacon_fail_closed().await;
965                    }
966                    Err(join_err) => {
967                        tracing::warn!(
968                            error = %join_err,
969                            "ADR-091 Amendment 2 Plank B: walpin heartbeat write task panicked"
970                        );
971                        self.last_heartbeat = None;
972                        self.drop_beacon_fail_closed().await;
973                    }
974                }
975            }
976            _ => {
977                self.refresh_beacon().await;
978                if self.wrote {
979                    let dir = self.dir.clone();
980                    let pid = self.pid;
981                    let result = tokio::task::spawn_blocking(move || {
982                        crate::walpin::remove_heartbeat(&dir, pid)
983                    })
984                    .await;
985                    match result {
986                        Ok(Ok(())) => {}
987                        Ok(Err(e)) => tracing::warn!(
988                            error = %e,
989                            "ADR-091 Amendment 2 Plank B: failed to remove walpin heartbeat"
990                        ),
991                        Err(join_err) => tracing::warn!(
992                            error = %join_err,
993                            "ADR-091 Amendment 2 Plank B: walpin heartbeat removal task panicked"
994                        ),
995                    }
996                    self.wrote = false;
997                    self.last_heartbeat = None;
998                }
999            }
1000        }
1001    }
1002
1003    async fn shutdown(&mut self) {
1004        if self.wrote {
1005            let dir = self.dir.clone();
1006            let pid = self.pid;
1007            let _ = tokio::task::spawn_blocking(move || crate::walpin::remove_heartbeat(&dir, pid))
1008                .await;
1009            self.wrote = false;
1010        }
1011    }
1012}
1013
1014fn now_epoch_secs() -> i64 {
1015    std::time::SystemTime::now()
1016        .duration_since(std::time::UNIX_EPOCH)
1017        .map(|d| d.as_secs() as i64)
1018        .unwrap_or(0)
1019}
1020
1021/// ADR-091 Amendment 2 Plank A: config for the observe-only per-session
1022/// sweep. Sessions never checkpoint — that stays daemon-owned so N session
1023/// processes never compete for the writer mutex — this only watches
1024/// `tx_registry` (and, Plank B, refreshes this process's walpin heartbeat).
1025#[derive(Clone, Debug)]
1026pub struct SessionSweepConfig {
1027    /// How often a session polls the registry. Coarser than the daemon's
1028    /// tick: sessions do not need the daemon's 500ms checkpoint cadence.
1029    ///
1030    /// Overridable via `KHIVE_SESSION_SWEEP_INTERVAL_MS`. Default: 5000 ms.
1031    pub interval: Duration,
1032    /// Same semantics and default as [`CheckpointConfig::tx_warn_secs`].
1033    pub tx_warn_secs: Duration,
1034    /// Same semantics and default as [`CheckpointConfig::tx_max_age_secs`].
1035    pub tx_max_age_secs: Duration,
1036}
1037
1038impl Default for SessionSweepConfig {
1039    fn default() -> Self {
1040        Self {
1041            interval: Duration::from_secs(5),
1042            tx_warn_secs: Duration::from_secs(30),
1043            tx_max_age_secs: Duration::from_secs(120),
1044        }
1045    }
1046}
1047
1048impl SessionSweepConfig {
1049    /// Build from the environment. Reuses `KHIVE_TX_WARN_SECS` /
1050    /// `KHIVE_TX_MAX_AGE_SECS` (the same knobs the daemon's checkpoint task
1051    /// reads) so a session and the daemon agree on the same thresholds.
1052    pub fn from_env() -> Self {
1053        let mut cfg = Self::default();
1054
1055        if let Ok(ms) = std::env::var("KHIVE_SESSION_SWEEP_INTERVAL_MS") {
1056            if let Ok(v) = ms.parse::<u64>() {
1057                if v > 0 {
1058                    cfg.interval = Duration::from_millis(v);
1059                }
1060            }
1061        }
1062        // Shares `tx_age_thresholds_from_env` with `CheckpointConfig::from_env`
1063        // (minor, ADR-091 Amendment 2) so a session and the daemon
1064        // parse and validate `KHIVE_TX_WARN_SECS`/`KHIVE_TX_MAX_AGE_SECS`
1065        // identically from one source, not two hand-copied blocks.
1066        (cfg.tx_warn_secs, cfg.tx_max_age_secs) =
1067            tx_age_thresholds_from_env(cfg.tx_warn_secs, cfg.tx_max_age_secs);
1068
1069        cfg
1070    }
1071}
1072
1073/// One file-backed backend the session sweep observes (ADR-091 Amendment 3
1074/// fan-out). `is_main` selects which [`khive_storage::tx_registry::TxOriginFilter`]
1075/// variant scopes this backend's view of the registry: the main backend's
1076/// `Main` filter additionally observes `Unscoped` spans (the
1077/// never-silently-drop fallback for call sites not yet threaded to an
1078/// origin); a secondary backend's `Secondary` filter is scoped to exactly
1079/// its own identity. A pool whose origin is `Memory` contributes no entry —
1080/// in-memory backends have no sidecar and nothing to attribute
1081/// cross-process.
1082pub struct SweepBackend {
1083    pub pool: Arc<ConnectionPool>,
1084    pub is_main: bool,
1085}
1086
1087/// Per-backend state the session sweep carries across ticks: this backend's
1088/// registry view, its own edge-triggered age-sweep state machine (so a
1089/// sustained stale span on one backend logs independently of the others),
1090/// and its own walpin sidecar (`None` if the sidecar is disabled or this
1091/// backend's origin is `Memory`).
1092struct BackendSweep {
1093    filter: khive_storage::tx_registry::TxOriginFilter,
1094    tx_age_state: TxAgeSweepState,
1095    sidecar: Option<WalpinSidecarState>,
1096}
1097
1098/// ADR-091 Amendment 2 Plank A (Amendment 3: per-backend fan-out): run the
1099/// observe-only per-session sweep.
1100///
1101/// Every non-daemon `kkernel mcp` process runs this instead of the daemon's
1102/// `run_checkpoint_task`: same `tx_registry` age check and Plank B heartbeat
1103/// refresh, but no PASSIVE/TRUNCATE checkpointing — checkpointing stays
1104/// daemon-owned. Stays ONE task for the whole process, but fans out
1105/// internally: each file-backed backend in `backends` gets its own
1106/// registry view, age-sweep state, and sidecar directory, so a long span on
1107/// a secondary backend is attributed (and heartbeats) only in that
1108/// backend's own sidecar — never the main backend's. Loops until
1109/// `shutdown_rx` observes a change (or its sender is dropped), removing
1110/// every written heartbeat on the way out.
1111pub async fn run_session_sweep_task(
1112    backends: Vec<SweepBackend>,
1113    config: SessionSweepConfig,
1114    mut shutdown_rx: tokio::sync::watch::Receiver<()>,
1115) {
1116    let mut interval = tokio::time::interval(config.interval);
1117    interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
1118
1119    let mut sweeps: Vec<BackendSweep> = Vec::with_capacity(backends.len());
1120    for backend in backends {
1121        let identity = match backend.pool.origin() {
1122            khive_storage::tx_registry::TxOrigin::Database(id) => id,
1123            // No on-disk file, so no sidecar and no cross-process
1124            // attribution surface — nothing for this sweep to fan out to.
1125            khive_storage::tx_registry::TxOrigin::Memory
1126            | khive_storage::tx_registry::TxOrigin::Unscoped => continue,
1127        };
1128        let filter = if backend.is_main {
1129            khive_storage::tx_registry::TxOriginFilter::Main(identity)
1130        } else {
1131            khive_storage::tx_registry::TxOriginFilter::Secondary(identity)
1132        };
1133        let sidecar = WalpinSidecarState::new(
1134            backend.pool.canonical_path(),
1135            true,
1136            "session",
1137            config.interval,
1138        );
1139        sweeps.push(BackendSweep {
1140            filter,
1141            tx_age_state: TxAgeSweepState::default(),
1142            sidecar,
1143        });
1144    }
1145    for sweep in sweeps.iter_mut() {
1146        if let Some(sidecar) = sweep.sidecar.as_mut() {
1147            sidecar.register_beacon().await;
1148        }
1149    }
1150
1151    loop {
1152        tokio::select! {
1153            _ = interval.tick() => {}
1154            _ = shutdown_rx.changed() => break,
1155        }
1156
1157        for sweep in sweeps.iter_mut() {
1158            let oldest = khive_storage::tx_registry::oldest_for(&sweep.filter);
1159            for emission in sweep.tx_age_state.observe(
1160                oldest.as_ref().map(|s| (s.id, s.age, s.label.clone())),
1161                config.tx_warn_secs,
1162                config.tx_max_age_secs,
1163            ) {
1164                log_tx_age_emission(&emission);
1165            }
1166            if let Some(sidecar) = sweep.sidecar.as_mut() {
1167                sidecar.observe(oldest, config.tx_warn_secs).await;
1168            }
1169        }
1170    }
1171
1172    for sweep in sweeps.iter_mut() {
1173        if let Some(sidecar) = sweep.sidecar.as_mut() {
1174            sidecar.shutdown().await;
1175        }
1176    }
1177}
1178
1179/// The event sink and namespace owned by one checkpoint task in a fan-out.
1180///
1181/// Backend role and lifecycle ownership are separate: a secondary task may
1182/// own lifecycle emission when the deployment's main backend is in-memory.
1183#[derive(Clone)]
1184pub struct CheckpointLifecycleOwner {
1185    event_store: Arc<dyn khive_storage::EventStore>,
1186    namespace: String,
1187}
1188
1189impl CheckpointLifecycleOwner {
1190    /// Designate `event_store` as the lifecycle sink for one checkpoint task.
1191    pub fn new(
1192        event_store: Arc<dyn khive_storage::EventStore>,
1193        namespace: impl Into<String>,
1194    ) -> Self {
1195        Self {
1196            event_store,
1197            namespace: namespace.into(),
1198        }
1199    }
1200}
1201
1202/// Maximum number of checkpoint lifecycle events waiting behind the append
1203/// currently owned by the worker. One queued row preserves a recent outcome
1204/// without allowing sustained writer contention to grow memory without bound.
1205const CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY: usize = 1;
1206
1207/// Zero-wait handoff from the checkpoint scheduler to its lifecycle sink.
1208///
1209/// The worker serializes appends, preserving the order of every event that is
1210/// accepted. The scheduler only calls [`tokio::sync::mpsc::Sender::try_send`]:
1211/// if the worker and its single queue slot are both occupied, telemetry is
1212/// dropped rather than delaying the next checkpoint cycle. The first drop in
1213/// each uninterrupted full-queue episode warns; a successful enqueue re-arms
1214/// that warning without producing per-tick log spam.
1215struct CheckpointLifecycleEmitter {
1216    namespace: Option<String>,
1217    sender: Option<tokio::sync::mpsc::Sender<khive_storage::Event>>,
1218    worker: Option<tokio::task::JoinHandle<()>>,
1219    busy_warning_emitted: bool,
1220}
1221
1222impl CheckpointLifecycleEmitter {
1223    fn new(owner: Option<CheckpointLifecycleOwner>) -> Self {
1224        let Some(owner) = owner else {
1225            return Self {
1226                namespace: None,
1227                sender: None,
1228                worker: None,
1229                busy_warning_emitted: false,
1230            };
1231        };
1232
1233        let namespace = owner.namespace.clone();
1234        let (sender, mut receiver) =
1235            tokio::sync::mpsc::channel::<khive_storage::Event>(CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY);
1236        let worker = tokio::spawn(async move {
1237            while let Some(event) = receiver.recv().await {
1238                let kind = event.kind;
1239                if let Err(err) = owner.event_store.append_event(event).await {
1240                    tracing::warn!(
1241                        error = %err,
1242                        event_kind = %kind.name(),
1243                        "checkpoint lifecycle event append failed"
1244                    );
1245                }
1246            }
1247        });
1248
1249        Self {
1250            namespace: Some(namespace),
1251            sender: Some(sender),
1252            worker: Some(worker),
1253            busy_warning_emitted: false,
1254        }
1255    }
1256
1257    /// Serialize and enqueue one lifecycle event without awaiting sink I/O.
1258    /// Returns whether the row was accepted for delivery (or no sink exists).
1259    fn try_emit<P: serde::Serialize>(&mut self, kind: khive_types::EventKind, payload: P) -> bool {
1260        let (Some(namespace), Some(sender)) = (&self.namespace, &self.sender) else {
1261            return true;
1262        };
1263        let payload_value = match serde_json::to_value(&payload) {
1264            Ok(value) => value,
1265            Err(err) => {
1266                tracing::warn!(
1267                    error = %err,
1268                    event_kind = %kind.name(),
1269                    "failed to serialize checkpoint lifecycle event payload"
1270                );
1271                return false;
1272            }
1273        };
1274        let event = khive_storage::Event::new(
1275            namespace,
1276            "checkpoint.lifecycle",
1277            kind,
1278            khive_types::SubstrateKind::Event,
1279            "daemon:checkpoint_task",
1280        )
1281        .with_payload(payload_value);
1282
1283        match sender.try_send(event) {
1284            Ok(()) => {
1285                self.busy_warning_emitted = false;
1286                true
1287            }
1288            Err(tokio::sync::mpsc::error::TrySendError::Full(event)) => {
1289                if !self.busy_warning_emitted {
1290                    tracing::warn!(
1291                        event_kind = %event.kind.name(),
1292                        queue_capacity = CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY,
1293                        "checkpoint lifecycle event dropped because the append worker is busy"
1294                    );
1295                    self.busy_warning_emitted = true;
1296                }
1297                false
1298            }
1299            Err(tokio::sync::mpsc::error::TrySendError::Closed(event)) => {
1300                tracing::warn!(
1301                    event_kind = %event.kind.name(),
1302                    "checkpoint lifecycle event dropped because the append worker stopped"
1303                );
1304                false
1305            }
1306        }
1307    }
1308
1309    /// Stop the scheduler-owned async worker without making
1310    /// [`run_checkpoint_task`] wait for its current append future.
1311    ///
1312    /// This bounds checkpoint-task shutdown only. If the event store already
1313    /// admitted the append to `spawn_blocking` or a `WriterTask`, aborting this
1314    /// worker cannot cancel that downstream operation; at most one such sink
1315    /// operation may outlive the checkpoint task.
1316    async fn shutdown(mut self) {
1317        drop(self.sender.take());
1318        let Some(worker) = self.worker.take() else {
1319            return;
1320        };
1321        worker.abort();
1322        match worker.await {
1323            Ok(()) => {}
1324            Err(err) if err.is_cancelled() => {}
1325            Err(err) => tracing::warn!(
1326                error = %err,
1327                "checkpoint lifecycle event append worker terminated unexpectedly"
1328            ),
1329        }
1330    }
1331}
1332
1333impl Drop for CheckpointLifecycleEmitter {
1334    fn drop(&mut self) {
1335        // The normal watch-signal path calls `shutdown` and takes the handle
1336        // first. This fallback covers an externally-aborted or panicking
1337        // checkpoint task so the scheduler-owned async worker itself is never
1338        // detached. One already-admitted downstream sink operation may outlive
1339        // it; see `shutdown`'s contract above.
1340        if let Some(worker) = &self.worker {
1341            worker.abort();
1342        }
1343    }
1344}
1345
1346/// Run the WAL checkpoint background task.
1347///
1348/// Long-running async task — spawn with `tokio::spawn`. Loops until
1349/// `shutdown_rx` observes a change (or its sender is dropped). Callers MUST
1350/// hold the paired `tokio::sync::watch::Sender` for the daemon's run scope
1351/// and send on it to shut down — do NOT rely on `pool`'s `Arc` refcount
1352/// reaching zero; a sibling owner (e.g. `event_store`) holding its own clone
1353/// makes that check unreachable (issue #774).
1354///
1355/// Issues `PRAGMA wal_checkpoint(PASSIVE)` every tick via `try_writer_nowait`
1356/// (zero-wait try-lock): a busy writer skips the tick rather than stalling
1357/// write traffic. A WARNING fires once per below→above threshold crossing,
1358/// not every tick.
1359///
1360/// `lifecycle_owner` (ADR-094): exactly one task in a multi-backend fan-out
1361/// should receive `Some`. That task appends a best-effort
1362/// `CheckpointOutcomeRecorded` event on every at/above-`warn_pages` tick,
1363/// plus one drain row when pressure falls back below `warn_pages`. `None`
1364/// explicitly marks a non-owner. See `crates/khive-db/docs/api/checkpoint.md`
1365/// for the full shutdown-mechanism and event-emission design history.
1366///
1367/// `is_main` (ADR-091 Amendment 3): whether `pool` is the deployment's main
1368/// backend. A daemon owning several file-backed backends spawns one task per
1369/// backend, each with its own pool and shutdown-channel clone (the sender
1370/// broadcasts to every receiver clone alike). Lifecycle ownership is selected
1371/// independently through `lifecycle_owner`; `is_main` only controls registry
1372/// filtering. See the `tx_filter` construction below.
1373pub async fn run_checkpoint_task(
1374    pool: Arc<ConnectionPool>,
1375    config: CheckpointConfig,
1376    lifecycle_owner: Option<CheckpointLifecycleOwner>,
1377    mut shutdown_rx: tokio::sync::watch::Receiver<()>,
1378    is_main: bool,
1379) {
1380    let mut interval = tokio::time::interval(config.interval);
1381    interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
1382    let mut severity_state = CheckpointSeverityState::default();
1383    let mut tx_age_state = TxAgeSweepState::default();
1384    let mut was_above_high_water = false;
1385    let mut truncate_state = TruncateState::default();
1386    let mut lifecycle_emitter = CheckpointLifecycleEmitter::new(lifecycle_owner);
1387    // Independent of `severity_state` (which owns the WARN-episode ladder
1388    // internally): this tracks whether an accepted elevated lifecycle row
1389    // still needs its matching drain row. A full queue leaves it unchanged,
1390    // so a dropped drain is retried on the next healthy tick rather than
1391    // leaving consumers with a permanently open elevation episode.
1392    let mut event_elevation_open = false;
1393    // ADR-091 Amendment 3: this task's own backend-scoped view of the
1394    // registry. `is_main` selects which `TxOriginFilter` variant applies —
1395    // the caller passes `true` for exactly the one checkpoint task covering
1396    // the deployment's main backend, so only that task also observes legacy
1397    // `Unscoped` spans from any call site not yet threaded to an origin, the
1398    // designed never-silently-drop fallback. A secondary backend's task
1399    // never falls back to `Unscoped`: those spans belong to the main view or
1400    // to no view, never to a database they were never registered against.
1401    // `None` only when this pool's own origin isn't `Database` (an in-memory
1402    // checkpoint pool) — degrades to "no open span observed" for the tick
1403    // rather than panicking a long-running daemon loop on an
1404    // assumed-impossible state.
1405    let tx_filter = match pool.origin() {
1406        khive_storage::tx_registry::TxOrigin::Database(id) => Some(if is_main {
1407            khive_storage::tx_registry::TxOriginFilter::Main(id)
1408        } else {
1409            khive_storage::tx_registry::TxOriginFilter::Secondary(id)
1410        }),
1411        khive_storage::tx_registry::TxOrigin::Memory
1412        | khive_storage::tx_registry::TxOrigin::Unscoped => None,
1413    };
1414    // ADR-091 Amendment 2 Plank B: the checkpoint pool is only ever wired for
1415    // file-backed backends (`checkpoint_pool_for`), so `is_file_backed: true`
1416    // is always correct here. `canonical_path()` (not `pool.config().path`)
1417    // so the sidecar directory is keyed off the same minted identity every
1418    // alias of this backend's configured path converges to.
1419    #[cfg(unix)]
1420    let mut walpin_state =
1421        WalpinSidecarState::new(pool.canonical_path(), true, "daemon", config.interval);
1422    #[cfg(unix)]
1423    if let Some(sidecar) = walpin_state.as_mut() {
1424        sidecar.register_beacon().await;
1425    }
1426
1427    loop {
1428        // A closed sender (the daemon returning without an explicit send)
1429        // makes `changed()` resolve with `Err` immediately, which `select!`
1430        // treats as ready — so shutdown is observed either way, not just on
1431        // an explicit send.
1432        tokio::select! {
1433            _ = interval.tick() => {}
1434            _ = shutdown_rx.changed() => break,
1435        }
1436
1437        let tick = checkpoint_once(&pool, &config, &mut truncate_state);
1438
1439        // ADR-091 Plank 1: age-based sweep over the registry's oldest entry
1440        // MUST run on every tick, including a Skipped one — deliberately
1441        // BEFORE the Skipped early-continue below. A registered
1442        // `WriterGuard::transaction` span (`pool.rs`) holds the writer mutex
1443        // for its entire registered lifetime, so exactly the long-running
1444        // transaction this sweep exists to name is the one that makes an
1445        // ordinary checkpoint tick observe `Skipped` — gating the sweep on
1446        // `Observed` would silence it for precisely that scenario, defeating
1447        // the WAL-independent diagnostic the ADR specifies. Independent of
1448        // WAL page pressure by the same design: a registered span can go
1449        // stale (KHIVE_TX_WARN_SECS / KHIVE_TX_MAX_AGE_SECS) while
1450        // wal_pages sits well under warn_pages, or isn't sampled at all this
1451        // tick. Edge-triggered per rung, same debounce idiom as the severity
1452        // ladder below, so a sustained stale span logs once per rung rather
1453        // than once per tick.
1454        let oldest_tx = tx_filter
1455            .as_ref()
1456            .and_then(khive_storage::tx_registry::oldest_for);
1457        for emission in tx_age_state.observe(
1458            oldest_tx.as_ref().map(|s| (s.id, s.age, s.label.clone())),
1459            config.tx_warn_secs,
1460            config.tx_max_age_secs,
1461        ) {
1462            log_tx_age_emission(&emission);
1463        }
1464        // ADR-091 Amendment 2 Plank B: refresh (or clear) this daemon
1465        // process's own walpin heartbeat on the same cadence, so its own
1466        // pin — if any — is attributable the same way a session's is.
1467        #[cfg(unix)]
1468        if let Some(sidecar) = walpin_state.as_mut() {
1469            sidecar
1470                .observe(oldest_tx.clone(), config.tx_warn_secs)
1471                .await;
1472        }
1473
1474        // Skipped ticks leave crossing state unchanged — a busy tick must not
1475        // re-arm the rate limit while WAL pressure is still elevated.
1476        let wal_pages = match tick {
1477            CheckpointTick::Skipped => continue,
1478            CheckpointTick::Observed(n) => n,
1479        };
1480
1481        let above_warn = wal_pages >= config.warn_pages;
1482        let above_high_water = wal_pages >= config.high_water_pages;
1483        let above_truncate_high_water = wal_pages >= config.truncate_high_water_pages;
1484
1485        // Per-tick debug for the oldest open entry always fires (cheap —
1486        // reuses this tick's already-computed `oldest_tx`); the two
1487        // `warn!`-level registry logs below are gated on the SAME crossing
1488        // state as the WAL-threshold WARNs above, so sustained pressure
1489        // logs once per crossing, not once per tick.
1490        log_tx_registry_oldest_debug(wal_pages, oldest_tx.as_ref());
1491
1492        // ADR-091 severity ladder: INFO on the first below→above crossing,
1493        // WARN once `warn_sustained_cycles` consecutive ticks stay elevated.
1494        // The oldest-entry registry WARN rides the same INFO edge the old
1495        // binary crossing_warn used to gate on.
1496        for emission in severity_state.observe_wal_pages(wal_pages, &config) {
1497            match emission.rung {
1498                CheckpointSeverityRung::Info => {
1499                    log_tx_registry_oldest_warn(wal_pages, oldest_tx.as_ref());
1500                    tracing::info!(
1501                        wal_pages = emission.wal_pages,
1502                        warn_threshold = emission.threshold_pages,
1503                        "WAL page count crossed warn threshold"
1504                    );
1505                }
1506                CheckpointSeverityRung::Warn => {
1507                    tracing::warn!(
1508                        wal_pages = emission.wal_pages,
1509                        warn_threshold = emission.threshold_pages,
1510                        consecutive_cycles = emission.consecutive_cycles,
1511                        "WAL page count failed to drain below warn threshold"
1512                    );
1513                }
1514                CheckpointSeverityRung::Alarm => {
1515                    // Never produced by `observe_wal_pages`; see its doc.
1516                }
1517            }
1518        }
1519
1520        let high_water_crossed = crossing_warn(above_high_water, &mut was_above_high_water);
1521        if high_water_crossed {
1522            log_tx_registry_snapshot_warn(wal_pages);
1523            tracing::warn!(
1524                wal_pages,
1525                high_water = config.high_water_pages,
1526                "WAL high-water mark exceeded; sustained WAL pressure — \
1527                 a long-lived reader may be pinning an old snapshot that PASSIVE cannot reclaim"
1528            );
1529        }
1530
1531        // ADR-094: emit every elevated tick, plus exactly one drain row on
1532        // the tick that observes the episode end — never on every ordinary
1533        // below-warn tick.
1534        if checkpoint_outcome_should_emit(above_warn, event_elevation_open) {
1535            let payload = khive_storage::CheckpointOutcomeRecordedPayload {
1536                wal_pages,
1537                warn_pages: config.warn_pages,
1538                high_water_pages: config.high_water_pages,
1539                truncate_high_water_pages: config.truncate_high_water_pages,
1540                above_warn,
1541                above_high_water,
1542                above_truncate_high_water,
1543            };
1544            if lifecycle_emitter
1545                .try_emit(khive_types::EventKind::CheckpointOutcomeRecorded, payload)
1546            {
1547                event_elevation_open = above_warn;
1548            }
1549        }
1550    }
1551
1552    lifecycle_emitter.shutdown().await;
1553
1554    #[cfg(unix)]
1555    if let Some(sidecar) = walpin_state.as_mut() {
1556        sidecar.shutdown().await;
1557    }
1558}
1559
1560/// Whether a `CheckpointOutcomeRecorded` event should be emitted for this
1561/// tick: every elevated (`above_warn`) tick, plus exactly one drain row on
1562/// the first tick that observes a return to below-warn after an elevated
1563/// episode (`was_elevated`). An ordinary below-warn tick following another
1564/// below-warn tick emits nothing.
1565fn checkpoint_outcome_should_emit(above_warn: bool, was_elevated: bool) -> bool {
1566    above_warn || was_elevated
1567}
1568
1569/// ADR-091 Plank 0 (Amendment 3: takes the tick's already-computed,
1570/// backend-scoped oldest span instead of re-querying the process-wide
1571/// aggregate): log the oldest open transaction registry entry alongside the
1572/// WAL frame count at `debug!`, on EVERY tick regardless of threshold
1573/// state. This is the low-volume per-tick trace; the WARN-level escalations
1574/// live in [`log_tx_registry_oldest_warn`] and
1575/// debug-level, unconditional per-tick trace. See
1576/// crates/khive-db/docs/api/checkpoint.md#private-tx-registry-logging-helpers-plank-0
1577fn log_tx_registry_oldest_debug(
1578    wal_pages: u64,
1579    oldest: Option<&khive_storage::tx_registry::OldestSpan>,
1580) {
1581    if let Some(span) = oldest {
1582        tracing::debug!(
1583            wal_pages,
1584            oldest_tx_age_secs = span.age.as_secs_f64(),
1585            oldest_tx_label = span.label.as_deref().unwrap_or("<unlabeled>"),
1586            "WAL checkpoint tick: oldest open transaction registry entry"
1587        );
1588    }
1589}
1590
1591/// Escalates the oldest open registry entry to `warn!`. NOT internally
1592/// rate-limited — caller MUST gate on a below→above `warn_pages` crossing
1593/// (`crossing_warn`) or every tick reproduces the log-spam bug this fixes.
1594fn log_tx_registry_oldest_warn(
1595    wal_pages: u64,
1596    oldest: Option<&khive_storage::tx_registry::OldestSpan>,
1597) {
1598    if let Some(span) = oldest {
1599        tracing::warn!(
1600            wal_pages,
1601            oldest_tx_age_secs = span.age.as_secs_f64(),
1602            oldest_tx_label = span.label.as_deref().unwrap_or("<unlabeled>"),
1603            "WAL checkpoint tick: oldest open transaction registry entry"
1604        );
1605    }
1606}
1607
1608/// Enumerates every open registry entry at `warn!`. NOT internally
1609/// rate-limited — caller MUST gate on a below→above `high_water_pages`
1610/// crossing (`crossing_warn`) or every tick repeats the full enumeration.
1611fn log_tx_registry_snapshot_warn(wal_pages: u64) {
1612    for (age, label) in khive_storage::tx_registry::snapshot() {
1613        tracing::warn!(
1614            wal_pages,
1615            tx_age_secs = age.as_secs_f64(),
1616            tx_label = label.as_deref().unwrap_or("<unlabeled>"),
1617            "WAL high-water: open transaction registry entry"
1618        );
1619    }
1620}
1621
1622/// Issue one checkpoint cycle against the writer connection.
1623///
1624/// Returns [`CheckpointTick::Skipped`] when the writer mutex is already held
1625/// (the tick is a no-op) and [`CheckpointTick::Observed`] with the WAL page
1626/// count otherwise. All checkpoint errors are logged at warn level and treated
1627/// as non-fatal; the next tick retries.
1628///
1629/// Uses `try_writer_nowait` so that a busy active writer causes this tick to
1630/// be skipped immediately rather than stalling for up to `checkout_timeout`.
1631/// The caller (`run_checkpoint_task`) owns all threshold-crossing WARN logging
1632/// so that warnings fire at most once per crossing, not every tick.
1633///
1634/// ADR-091 Plank 2: after the PASSIVE pass, this is also the single point
1635/// that may escalate to TRUNCATE (`maybe_truncate`) — under the SAME writer
1636/// guard acquired above, never a second checkout. A busy writer (`Skipped`)
1637/// short-circuits before either PASSIVE or TRUNCATE run.
1638pub fn checkpoint_once(
1639    pool: &ConnectionPool,
1640    config: &CheckpointConfig,
1641    truncate_state: &mut TruncateState,
1642) -> CheckpointTick {
1643    let writer = match pool.try_writer_nowait() {
1644        Ok(w) => w,
1645        Err(_) => {
1646            note_checkpoint_skipped();
1647            return CheckpointTick::Skipped;
1648        }
1649    };
1650
1651    let wal_pages = query_wal_pages(writer.conn());
1652
1653    if let Err(e) = writer
1654        .conn()
1655        .execute_batch("PRAGMA wal_checkpoint(PASSIVE)")
1656    {
1657        tracing::warn!(error = %e, "WAL checkpoint failed");
1658    } else {
1659        tracing::debug!(wal_pages, "WAL checkpoint issued");
1660    }
1661
1662    maybe_truncate(pool, &writer, config, wal_pages, truncate_state);
1663
1664    CheckpointTick::Observed(wal_pages)
1665}
1666
1667/// Evaluate and, if due, attempt a TRUNCATE escalation under the writer
1668/// guard the caller already holds (never its own checkout). `last_attempt`
1669/// is stamped ONLY on an actual attempt, never on a skip. See
1670/// crates/khive-db/docs/api/checkpoint.md#maybe_truncate--truncate-attempt-gating-plank-2
1671fn maybe_truncate(
1672    pool: &ConnectionPool,
1673    writer: &WriterGuard<'_>,
1674    config: &CheckpointConfig,
1675    wal_pages_before: u64,
1676    truncate_state: &mut TruncateState,
1677) {
1678    if wal_pages_before < config.truncate_high_water_pages {
1679        return;
1680    }
1681
1682    if let Some(last) = truncate_state.last_attempt {
1683        if last.elapsed() < config.truncate_min_interval {
1684            return;
1685        }
1686    }
1687
1688    // Which caller (if any) is pinning the WAL — logged before the attempt so
1689    // it is available even if the attempt itself succeeds.
1690    log_tx_registry_snapshot_warn(wal_pages_before);
1691
1692    let conn = writer.conn();
1693    let original_busy_timeout = pool.config().busy_timeout;
1694
1695    if let Err(e) = conn.busy_timeout(config.truncate_busy_timeout) {
1696        // Setup failed before the TRUNCATE pragma ever ran — this is a skip,
1697        // not an attempt. `last_attempt` must NOT advance here (ADR-091
1698        // §377-382): stamping now would suppress the next eligible attempt
1699        // for the full `truncate_min_interval` on a path that never touched
1700        // the WAL at all.
1701        tracing::warn!(error = %e, "failed to lower busy_timeout for TRUNCATE attempt; skipping");
1702        return;
1703    }
1704
1705    #[cfg(unix)]
1706    let holder_census = capture_wal_holder_census(pool);
1707
1708    // Only now is this a genuine attempt: the writer is held, the threshold
1709    // and interval gates passed, and the busy_timeout override is in effect
1710    // immediately before the TRUNCATE pragma itself.
1711    truncate_state.last_attempt = Some(Instant::now());
1712
1713    let start = Instant::now();
1714    let outcome = conn.execute_batch("PRAGMA wal_checkpoint(TRUNCATE)");
1715    let elapsed = start.elapsed();
1716
1717    // Restore the pool's configured busy_timeout immediately after the
1718    // attempt, win or lose, before any other logging or bookkeeping.
1719    if let Err(e) = conn.busy_timeout(original_busy_timeout) {
1720        tracing::warn!(error = %e, "failed to restore busy_timeout after TRUNCATE attempt");
1721    }
1722
1723    match outcome {
1724        Ok(()) => {
1725            let wal_pages_after = query_wal_pages(conn);
1726            tracing::info!(
1727                wal_pages_before,
1728                wal_pages_after,
1729                elapsed_ms = elapsed.as_millis() as u64,
1730                "WAL TRUNCATE checkpoint attempted"
1731            );
1732
1733            let made_progress = wal_pages_after < wal_pages_before;
1734            if !made_progress {
1735                tracing::warn!(
1736                    wal_pages_before,
1737                    wal_pages_after,
1738                    "WAL TRUNCATE attempt made no progress; \
1739                     a long-lived reader may still be pinning the WAL snapshot"
1740                );
1741                log_tx_registry_snapshot_warn(wal_pages_after);
1742                #[cfg(test)]
1743                if let Some(path) = pool.canonical_path() {
1744                    truncate_report_test_sync::after_no_progress_before_report(path);
1745                }
1746                #[cfg(unix)]
1747                log_walpin_sidecar_report(pool, holder_census);
1748                log_wal_pin_depth(conn);
1749            }
1750
1751            note_truncate_outcome(config, wal_pages_after, truncate_state);
1752        }
1753        Err(e) => {
1754            tracing::warn!(error = %e, wal_pages_before, "WAL TRUNCATE attempt failed");
1755            log_tx_registry_snapshot_warn(wal_pages_before);
1756            note_truncate_outcome(config, wal_pages_before, truncate_state);
1757        }
1758    }
1759}
1760
1761#[cfg(test)]
1762mod truncate_report_test_sync {
1763    use std::path::{Path, PathBuf};
1764    use std::sync::mpsc::{sync_channel, Receiver, SyncSender};
1765    use std::sync::Mutex;
1766
1767    struct Hook {
1768        db_path: PathBuf,
1769        reached_tx: SyncSender<()>,
1770        proceed_rx: Receiver<()>,
1771    }
1772
1773    static HOOK: Mutex<Option<Hook>> = Mutex::new(None);
1774
1775    pub(crate) fn install(db_path: PathBuf) -> (Receiver<()>, SyncSender<()>) {
1776        let (reached_tx, reached_rx) = sync_channel(0);
1777        let (proceed_tx, proceed_rx) = sync_channel(0);
1778        let replaced = HOOK
1779            .lock()
1780            .unwrap_or_else(|poisoned| poisoned.into_inner())
1781            .replace(Hook {
1782                db_path,
1783                reached_tx,
1784                proceed_rx,
1785            });
1786        assert!(replaced.is_none(), "truncate report hook already installed");
1787        (reached_rx, proceed_tx)
1788    }
1789
1790    pub(crate) fn uninstall() {
1791        *HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
1792    }
1793
1794    pub(crate) fn after_no_progress_before_report(db_path: &Path) {
1795        let hook = {
1796            let mut guard = HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
1797            match guard.as_ref() {
1798                Some(hook) if hook.db_path == db_path => guard.take(),
1799                _ => None,
1800            }
1801        };
1802        let Some(hook) = hook else {
1803            return;
1804        };
1805        let _ = hook.reached_tx.send(());
1806        let _ = hook.proceed_rx.recv();
1807    }
1808}
1809
1810/// ADR-091 Plank 2: track consecutive TRUNCATE attempts that fail to bring
1811/// `wal_pages` back below `warn_pages`, firing a one-shot escalated WARN at
1812/// exactly the third consecutive failure (does not repeat every attempt
1813/// thereafter — mirrors the crossing-WARN debounce used elsewhere in this
1814/// module). A single attempt that clears `warn_pages` resets the counter.
1815fn note_truncate_outcome(
1816    config: &CheckpointConfig,
1817    wal_pages_after: u64,
1818    state: &mut TruncateState,
1819) {
1820    // Metrics read-surface (load/perf harness): this function runs exactly
1821    // once per genuine TRUNCATE attempt (both the `Ok` and `Err` outcome
1822    // arms in `maybe_truncate` call it once each), so incrementing here
1823    // counts total attempts without a separate call site.
1824    TRUNCATE_ATTEMPTS.fetch_add(1, Ordering::Relaxed);
1825
1826    if wal_pages_after >= config.warn_pages {
1827        state.consecutive_failures = state.consecutive_failures.saturating_add(1);
1828        if state.consecutive_failures == 3 {
1829            tracing::warn!(
1830                wal_pages_after,
1831                warn_threshold = config.warn_pages,
1832                "WAL TRUNCATE has failed to clear WAL pressure for 3 consecutive attempts"
1833            );
1834        }
1835    } else {
1836        state.consecutive_failures = 0;
1837    }
1838
1839    TRUNCATE_CONSECUTIVE_FAILURES.store(state.consecutive_failures as u64, Ordering::Relaxed);
1840}
1841
1842/// ADR-091 Amendment 2 Plank B: capture the OS holder census immediately
1843/// before a TRUNCATE attempt so a holder that releases during the bounded wait
1844/// remains attributable if the attempt reports no progress. A no-op if the
1845/// sidecar is disabled or this backend has no on-disk path.
1846#[cfg(unix)]
1847fn capture_wal_holder_census(
1848    pool: &ConnectionPool,
1849) -> Option<Result<crate::walpin::CensusResult, String>> {
1850    let path = pool.canonical_path()?;
1851    if !crate::walpin::sidecar_enabled(true) {
1852        return None;
1853    }
1854    Some(crate::walpin::census_holders(path).map_err(|e| e.to_string()))
1855}
1856
1857/// When a TRUNCATE attempt makes no progress, enumerate the walpin sidecar and
1858/// combine it with the holder census captured immediately before that attempt.
1859/// Sidecar enumeration remains deferred until this consumed diagnostic path;
1860/// holder identity cannot be deferred because a transient blocker may have
1861/// released by then.
1862///
1863/// Sidecar-health attribution (ADR-091 Amendment 2):
1864/// the sharper "unregistered/native mechanism" conclusion is licensed only
1865/// when every discovered PID is `reporting` or `registered-silent`
1866/// (`WalpinReport::fully_attributed`); any `unknown` PID — including the
1867/// directory itself failing the trust-boundary check — makes attribution
1868/// inconclusive, and the WARN below names exactly which PIDs are unresolved
1869/// instead of silently exonerating them.
1870#[cfg(unix)]
1871fn log_walpin_sidecar_report(
1872    pool: &ConnectionPool,
1873    census: Option<Result<crate::walpin::CensusResult, String>>,
1874) {
1875    let Some(census) = census else {
1876        return;
1877    };
1878    let Some(path) = pool.canonical_path() else {
1879        return;
1880    };
1881    let dir = crate::walpin::sidecar_dir_for(path);
1882    // Each record carries its producer's own sweep cadence
1883    // (`sweep_interval_ms`), which is what freshness is judged against; the
1884    // interval passed here is only the fallback for records written before
1885    // that field existed.
1886    let sweep_interval = SessionSweepConfig::from_env().interval;
1887    let report = match crate::walpin::enumerate_live(&dir, sweep_interval) {
1888        Ok(report) => report,
1889        Err(e) => {
1890            tracing::warn!(
1891                error = %e,
1892                "ADR-091 Amendment 2 Plank B: sidecar directory failed the trust-boundary \
1893                 check; cross-process WAL-pin attribution is unestablished for this tick"
1894            );
1895            return;
1896        }
1897    };
1898    let now = now_epoch_secs();
1899    for hb in report.reporting() {
1900        // ADR-091 Amendment 3 Plank F2 fail-closed reading rule: the
1901        // logger must never let a fallback-confidence entry read as live
1902        // cross-process ground truth, so the confidence distinction is
1903        // always emitted alongside the raw field — never inferred by the
1904        // reader of this log line.
1905        tracing::warn!(
1906            walpin_pid = hb.pid,
1907            walpin_role = %hb.process_role,
1908            walpin_oldest_tx_age_secs = hb.current_oldest_tx_age_secs(now),
1909            walpin_oldest_tx_label = hb.oldest_tx_label.as_deref().unwrap_or("<unlabeled>"),
1910            walpin_attribution_basis = hb.attribution_basis.as_deref().unwrap_or("<unspecified>"),
1911            walpin_attribution_evidence_backed = hb.attribution_is_evidence_backed(),
1912            walpin_health = "reporting",
1913            "ADR-091 Amendment 2 Plank B: live cross-process WAL-pin attribution report"
1914        );
1915    }
1916    for pid in report.registered_silent_pids() {
1917        tracing::debug!(
1918            walpin_pid = pid,
1919            walpin_health = "registered_silent",
1920            "ADR-091 Amendment 2 Plank B: process affirmatively reports no over-threshold span"
1921        );
1922    }
1923    let mut unknown_pids: Vec<u32> = report.unknown_pids().collect();
1924
1925    // The sidecar directory alone can only speak for PIDs that wrote
1926    // something there. Widen the universe to every PID the OS reports as
1927    // holding the database immediately before the TRUNCATE attempt; any holder
1928    // absent from `report` is unknown.
1929    match census {
1930        Ok(census) => {
1931            let sidecar_known: std::collections::HashSet<u32> = report
1932                .reporting()
1933                .map(|hb| hb.pid)
1934                .chain(report.registered_silent_pids())
1935                .chain(unknown_pids.iter().copied())
1936                .collect();
1937            let mut census_only: Vec<u32> =
1938                census.holders.difference(&sidecar_known).copied().collect();
1939            if !census_only.is_empty() {
1940                census_only.sort_unstable();
1941                tracing::warn!(
1942                    ?census_only,
1943                    "ADR-091 Amendment 2: these PIDs hold the database file open \
1944                     at the OS level but have no sidecar data at all (pre-feature binary, \
1945                     sidecar disabled, or wedged before its first write)"
1946                );
1947                unknown_pids.extend(census_only);
1948            }
1949            if !census.is_complete() {
1950                let mut uninspectable = census.uninspectable_pids.clone();
1951                uninspectable.sort_unstable();
1952                tracing::warn!(
1953                    ?uninspectable,
1954                    truncated = census.truncated,
1955                    "ADR-091 Amendment 2: the OS-derived holder census is \
1956                     INCOMPLETE — either specific PIDs' open file descriptors could not be \
1957                     inspected (permission denied, or a listing race), or the enumeration walk \
1958                     itself has positive evidence it did not see the full live-process universe \
1959                     (namespace/visibility check, directory-iterator error, self-canary, or a \
1960                     libproc buffer that stayed at capacity after bounded retries) — cannot \
1961                     rule out an unregistered holder"
1962                );
1963                if uninspectable.is_empty() {
1964                    // `truncated` fired with no specific PID list (a
1965                    // namespace/visibility or buffer-truncation signal, not
1966                    // a per-PID inspection failure) — still makes
1967                    // attribution inconclusive. Mirror the census-failure
1968                    // arm below with the same non-PID sentinel rather than
1969                    // silently trusting a walk we know was incomplete.
1970                    unknown_pids.push(0);
1971                } else {
1972                    unknown_pids.extend(uninspectable);
1973                }
1974            }
1975        }
1976        Err(e) => {
1977            tracing::warn!(
1978                error = %e,
1979                "ADR-091 Amendment 2: OS-derived holder census failed; \
1980                 attribution cannot rule out an unregistered database holder this tick"
1981            );
1982            // A failed census is itself a health failure for the sharper
1983            // conclusion below — treat it as if at least one PID were
1984            // unresolved, without fabricating a specific PID number.
1985            unknown_pids.push(0);
1986        }
1987    }
1988
1989    if !unknown_pids.is_empty() {
1990        tracing::warn!(
1991            ?unknown_pids,
1992            "ADR-091 Amendment 2 Plank B: sidecar health unestablished for these PIDs; \
1993             attribution is inconclusive and the native/unregistered-mechanism conclusion \
1994             is NOT licensed this tick"
1995        );
1996    } else if report.reporting().next().is_none() {
1997        tracing::info!(
1998            "ADR-091 Amendment 2 Plank B: every live PID is reporting or registered-silent \
1999             with none pinning; the WAL pin is not attributable to any in-process registry \
2000             span this sidecar covers"
2001        );
2002    }
2003}
2004
2005/// ADR-091 Amendment 2 Plank C: on a TRUNCATE no-progress event, run a fresh
2006/// `PRAGMA wal_checkpoint(PASSIVE)` (never blocks readers or writers) and
2007/// report pin depth as `log` minus `checkpointed` from its 3-column return
2008/// row — the number of frames pinned behind the backfill boundary. Zero
2009/// dependence on SQLite's shm WAL-index layout.
2010fn log_wal_pin_depth(conn: &rusqlite::Connection) {
2011    match query_wal_pin_depth(conn) {
2012        Ok((log, checkpointed)) => {
2013            tracing::warn!(
2014                wal_log_frames = log,
2015                wal_checkpointed_frames = checkpointed,
2016                wal_pin_depth = (log - checkpointed).max(0),
2017                "ADR-091 Amendment 2 Plank C: WAL pin depth after TRUNCATE no-progress"
2018            );
2019        }
2020        Err(e) => {
2021            tracing::warn!(
2022                error = %e,
2023                "ADR-091 Amendment 2 Plank C: failed to query WAL pin depth"
2024            );
2025        }
2026    }
2027}
2028
2029/// ADR-091 Amendment 2 Plank C: issue `PRAGMA wal_checkpoint(PASSIVE)` and
2030/// return its `(log, checkpointed)` columns (index 1 and 2 of the 3-column
2031/// return row). PASSIVE never blocks readers or writers. Pin depth is
2032/// `log - checkpointed`; extracted as its own pure query so the arithmetic is
2033/// unit-testable against a real SQLite connection without depending on
2034/// `tracing` capture.
2035fn query_wal_pin_depth(conn: &rusqlite::Connection) -> rusqlite::Result<(i64, i64)> {
2036    conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
2037        Ok((row.get::<_, i64>(1)?, row.get::<_, i64>(2)?))
2038    })
2039}
2040
2041/// Evaluate whether a threshold-crossing WARN should fire and advance the
2042/// crossing-state flag.
2043///
2044/// Returns `true` on a false→true transition in `now_above` (first observed
2045/// above-threshold tick after a below-threshold tick), `false` on any other
2046/// tick. The `was_above` flag is updated in-place to track state across calls.
2047/// Used by `run_checkpoint_task` for both the `warn_pages` band and the
2048/// `high_water_pages` threshold.
2049fn crossing_warn(now_above: bool, was_above: &mut bool) -> bool {
2050    let fire = now_above && !*was_above;
2051    *was_above = now_above;
2052    fire
2053}
2054
2055/// Query the current WAL frame count via `PRAGMA wal_checkpoint`.
2056///
2057/// The pragma returns a 3-column row `(busy, log, checkpointed)`, where `log`
2058/// (column index 1) is the number of frames currently in the WAL file — the
2059/// backlog the high-water threshold keys off. (Column 2 is `checkpointed`, the
2060/// frames moved *by this call*, which is not the WAL size.) The no-arg pragma
2061/// also performs a PASSIVE checkpoint as a side effect; the subsequent explicit
2062/// `PRAGMA wal_checkpoint(PASSIVE)` in `checkpoint_once` is a deliberate second
2063/// pass that can checkpoint any frames written between the two calls.
2064///
2065/// Returns 0 on any error (e.g. in-memory DB where WAL is not active, which
2066/// reports `log = -1`).
2067fn query_wal_pages(conn: &rusqlite::Connection) -> u64 {
2068    let pages = conn
2069        .query_row("PRAGMA wal_checkpoint", [], |row| row.get::<_, i64>(1))
2070        .unwrap_or(0)
2071        .max(0) as u64;
2072    // Metrics read-surface (load/perf harness): mirror every observation into
2073    // the process-wide gauge, regardless of which caller (`checkpoint_once`
2074    // or `maybe_truncate`) triggered it.
2075    LAST_WAL_PAGES.store(pages, Ordering::Relaxed);
2076    note_checkpoint_observed(pages);
2077    pages
2078}
2079
2080#[cfg(test)]
2081mod tests {
2082    use super::*;
2083    use crate::pool::PoolConfig;
2084    use serial_test::serial;
2085    use tracing::field::{Field, Visit};
2086
2087    #[derive(Clone, Debug, Default)]
2088    struct CapturedEvent {
2089        message: Option<String>,
2090        oldest_tx_label: Option<String>,
2091        tx_label: Option<String>,
2092        census_only: Option<String>,
2093    }
2094
2095    #[derive(Default)]
2096    struct CapturedEventVisitor(CapturedEvent);
2097
2098    impl Visit for CapturedEventVisitor {
2099        fn record_str(&mut self, field: &Field, value: &str) {
2100            match field.name() {
2101                "message" => self.0.message = Some(value.to_string()),
2102                "oldest_tx_label" => self.0.oldest_tx_label = Some(value.to_string()),
2103                "tx_label" => self.0.tx_label = Some(value.to_string()),
2104                _ => {}
2105            }
2106        }
2107
2108        fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
2109            let formatted = format!("{value:?}");
2110            let cleaned = formatted
2111                .trim_start_matches('"')
2112                .trim_end_matches('"')
2113                .to_string();
2114            match field.name() {
2115                "message" => self.0.message = Some(cleaned),
2116                "oldest_tx_label" => self.0.oldest_tx_label = Some(cleaned),
2117                "tx_label" => self.0.tx_label = Some(cleaned),
2118                "census_only" => self.0.census_only = Some(cleaned),
2119                _ => {}
2120            }
2121        }
2122    }
2123
2124    /// Minimal `tracing::Subscriber` that captures events into a thread-local
2125    /// vec, installed as the thread-local default for the duration of one
2126    /// test closure via `tracing::subscriber::with_default`. Mirrors the
2127    /// capture subscriber in `khive-runtime/src/pack.rs`'s gate-dispatch tests.
2128    struct CaptureSubscriber {
2129        events: std::sync::Arc<std::sync::Mutex<Vec<CapturedEvent>>>,
2130    }
2131
2132    impl tracing::Subscriber for CaptureSubscriber {
2133        fn enabled(&self, _: &tracing::Metadata<'_>) -> bool {
2134            true
2135        }
2136        fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
2137            tracing::span::Id::from_u64(1)
2138        }
2139        fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
2140        fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
2141        fn event(&self, event: &tracing::Event<'_>) {
2142            let mut visitor = CapturedEventVisitor::default();
2143            event.record(&mut visitor);
2144            self.events.lock().unwrap().push(visitor.0);
2145        }
2146        fn enter(&self, _: &tracing::span::Id) {}
2147        fn exit(&self, _: &tracing::span::Id) {}
2148    }
2149
2150    /// `log_tx_registry_oldest_debug` names the oldest open registry entry.
2151    /// See crates/khive-db/docs/api/checkpoint.md#log_tx_registry_oldest_debug_reports_oldest_open_entry
2152    #[test]
2153    #[serial(tx_registry)]
2154    fn log_tx_registry_oldest_debug_reports_oldest_open_entry() {
2155        let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
2156        let subscriber = CaptureSubscriber {
2157            events: std::sync::Arc::clone(&buffer),
2158        };
2159
2160        let _handle =
2161            khive_storage::tx_registry::register(Some("checkpoint_tick_test".to_string()));
2162
2163        let oldest = khive_storage::tx_registry::oldest().map(|(id, age, label)| {
2164            khive_storage::tx_registry::OldestSpan {
2165                id,
2166                age,
2167                label,
2168                origin: khive_storage::tx_registry::TxOrigin::Unscoped,
2169            }
2170        });
2171        let expected_label = oldest
2172            .as_ref()
2173            .and_then(|s| s.label.clone())
2174            .unwrap_or_else(|| "<unlabeled>".to_string());
2175
2176        tracing::subscriber::with_default(subscriber, || {
2177            log_tx_registry_oldest_debug(100, oldest.as_ref());
2178        });
2179
2180        let events = buffer.lock().unwrap();
2181        assert!(
2182            events.iter().any(|e| {
2183                e.message.as_deref()
2184                    == Some("WAL checkpoint tick: oldest open transaction registry entry")
2185                    && e.oldest_tx_label.as_deref() == Some(expected_label.as_str())
2186            }),
2187            "expected a log line naming the open registry entry's label, got: {events:?}"
2188        );
2189    }
2190
2191    /// ADR-091 Plank 0: the oldest-entry WARN and the
2192    /// high-water snapshot-enumeration WARN are gated by `crossing_warn` at
2193    /// the call site (mirroring the WAL-threshold WARNs), so driving two
2194    /// consecutive above-threshold ticks through that same gate must produce
2195    /// exactly one of each — never a repeat on the second tick.
2196    #[test]
2197    #[serial(tx_registry)]
2198    fn registry_warns_fire_on_crossing_and_do_not_repeat() {
2199        let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
2200        let subscriber = CaptureSubscriber {
2201            events: std::sync::Arc::clone(&buffer),
2202        };
2203
2204        let _handle =
2205            khive_storage::tx_registry::register(Some("registry_warn_crossing_test".to_string()));
2206        let oldest = khive_storage::tx_registry::oldest().map(|(id, age, label)| {
2207            khive_storage::tx_registry::OldestSpan {
2208                id,
2209                age,
2210                label,
2211                origin: khive_storage::tx_registry::TxOrigin::Unscoped,
2212            }
2213        });
2214
2215        let mut was_above_warn = false;
2216        let mut was_above_high_water = false;
2217
2218        tracing::subscriber::with_default(subscriber, || {
2219            // Tick 1: below→above crossing for both bands — both WARNs fire.
2220            if crossing_warn(true, &mut was_above_warn) {
2221                log_tx_registry_oldest_warn(6000, oldest.as_ref());
2222            }
2223            if crossing_warn(true, &mut was_above_high_water) {
2224                log_tx_registry_snapshot_warn(6000);
2225            }
2226
2227            // Tick 2: still above both thresholds — neither must repeat.
2228            if crossing_warn(true, &mut was_above_warn) {
2229                log_tx_registry_oldest_warn(6000, oldest.as_ref());
2230            }
2231            if crossing_warn(true, &mut was_above_high_water) {
2232                log_tx_registry_snapshot_warn(6000);
2233            }
2234        });
2235
2236        let events = buffer.lock().unwrap();
2237
2238        // `tracing::subscriber::with_default` scopes capture to THIS thread for
2239        // the duration of the closure, so `events` contains only the two
2240        // `log_tx_registry_oldest_warn` calls made above — no concurrent test's
2241        // log calls land in this buffer. This lets the crossing/no-repeat
2242        // assertion match on message text alone: unlike the "names MY label"
2243        // assertion in the sibling test above, WHICH label `oldest()` reports
2244        // is irrelevant here (a concurrent write path elsewhere in the binary
2245        // may transiently be the registry's genuine oldest entry) — only the
2246        // fire-once-per-crossing COUNT is under test.
2247        let oldest_warn_count = events
2248            .iter()
2249            .filter(|e| {
2250                e.message.as_deref()
2251                    == Some("WAL checkpoint tick: oldest open transaction registry entry")
2252            })
2253            .count();
2254        assert_eq!(
2255            oldest_warn_count, 1,
2256            "oldest-entry WARN must fire exactly once across two above-threshold ticks, got: {events:?}"
2257        );
2258
2259        let snapshot_warn_count = events
2260            .iter()
2261            .filter(|e| {
2262                e.message.as_deref() == Some("WAL high-water: open transaction registry entry")
2263                    && e.tx_label.as_deref() == Some("registry_warn_crossing_test")
2264            })
2265            .count();
2266        assert_eq!(
2267            snapshot_warn_count, 1,
2268            "high-water snapshot WARN must fire exactly once across two above-threshold ticks, got: {events:?}"
2269        );
2270    }
2271
2272    /// ADR-091 Plank 1: `log_tx_age_emission` emits the correct message text
2273    /// and carries the entry's label, for both the `Warn` and `Stale` rungs.
2274    #[test]
2275    fn log_tx_age_emission_carries_label_for_both_rungs() {
2276        let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
2277        let subscriber = CaptureSubscriber {
2278            events: std::sync::Arc::clone(&buffer),
2279        };
2280
2281        tracing::subscriber::with_default(subscriber, || {
2282            log_tx_age_emission(&TxAgeEmission {
2283                rung: TxAgeRung::Warn,
2284                age: Duration::from_secs(45),
2285                label: Some("plank1_warn_test".to_string()),
2286            });
2287            log_tx_age_emission(&TxAgeEmission {
2288                rung: TxAgeRung::Stale,
2289                age: Duration::from_secs(150),
2290                label: Some("plank1_stale_test".to_string()),
2291            });
2292        });
2293
2294        let events = buffer.lock().unwrap();
2295        assert!(
2296            events.iter().any(|e| {
2297                e.message.as_deref()
2298                    == Some(
2299                        "ADR-091 Plank 1: open transaction registry entry exceeded soft-cap age",
2300                    )
2301                    && e.tx_label.as_deref() == Some("plank1_warn_test")
2302            }),
2303            "expected a Warn-rung log line naming the entry, got: {events:?}"
2304        );
2305        assert!(
2306            events.iter().any(|e| {
2307                e.message.as_deref().is_some_and(|m| {
2308                    m.starts_with(
2309                        "ADR-091 Plank 1: open transaction registry entry exceeded the cooperative",
2310                    )
2311                }) && e.tx_label.as_deref() == Some("plank1_stale_test")
2312            }),
2313            "expected a Stale-rung log line naming the entry, got: {events:?}"
2314        );
2315    }
2316
2317    fn file_pool(path: &std::path::Path) -> Arc<ConnectionPool> {
2318        let cfg = PoolConfig {
2319            path: Some(path.to_path_buf()),
2320            ..PoolConfig::default()
2321        };
2322        Arc::new(ConnectionPool::new(cfg).expect("pool open"))
2323    }
2324
2325    struct TruncateReportHookGuard;
2326
2327    impl Drop for TruncateReportHookGuard {
2328        fn drop(&mut self) {
2329            truncate_report_test_sync::uninstall();
2330        }
2331    }
2332
2333    struct ReaderProcess {
2334        child: std::process::Child,
2335        _stdout: std::io::BufReader<std::process::ChildStdout>,
2336    }
2337
2338    impl ReaderProcess {
2339        fn spawn(db_path: &std::path::Path) -> Self {
2340            use std::io::BufRead;
2341            use std::process::Stdio;
2342
2343            let mut child = std::process::Command::new(
2344                std::env::current_exe().expect("resolve current test executable"),
2345            )
2346            .args([
2347                "--exact",
2348                "checkpoint::tests::walpin_transient_reader_process_helper",
2349                "--nocapture",
2350            ])
2351            .env("KHIVE_CHECKPOINT_READER_HELPER_PATH", db_path)
2352            .stdin(Stdio::piped())
2353            .stdout(Stdio::piped())
2354            .spawn()
2355            .expect("spawn transient WAL reader helper");
2356
2357            let stdout = child.stdout.take().expect("capture helper stdout");
2358            let mut reader = std::io::BufReader::new(stdout);
2359            let mut line = String::new();
2360            loop {
2361                line.clear();
2362                let bytes = reader
2363                    .read_line(&mut line)
2364                    .expect("read transient reader readiness signal");
2365                assert!(bytes > 0, "reader helper exited before readiness signal");
2366                if line.contains("KHIVE_CHECKPOINT_READER_READY") {
2367                    break;
2368                }
2369            }
2370            Self {
2371                child,
2372                _stdout: reader,
2373            }
2374        }
2375
2376        fn pid(&self) -> u32 {
2377            self.child.id()
2378        }
2379
2380        fn release(&mut self) {
2381            use std::io::Write;
2382
2383            let mut stdin = self.child.stdin.take().expect("helper stdin is available");
2384            stdin
2385                .write_all(b"release\n")
2386                .expect("release transient reader");
2387            drop(stdin);
2388            let status = self.child.wait().expect("wait for transient reader helper");
2389            assert!(status.success(), "transient reader helper failed: {status}");
2390        }
2391    }
2392
2393    impl Drop for ReaderProcess {
2394        fn drop(&mut self) {
2395            if self.child.try_wait().ok().flatten().is_none() {
2396                let _ = self.child.kill();
2397                let _ = self.child.wait();
2398            }
2399        }
2400    }
2401
2402    #[test]
2403    fn walpin_transient_reader_process_helper() {
2404        use std::io::Write;
2405
2406        let Some(path) = std::env::var_os("KHIVE_CHECKPOINT_READER_HELPER_PATH") else {
2407            return;
2408        };
2409        let conn = rusqlite::Connection::open(path).expect("helper opens database");
2410        conn.execute_batch("BEGIN DEFERRED; SELECT * FROM t;")
2411            .expect("helper pins a read snapshot");
2412        println!("KHIVE_CHECKPOINT_READER_READY");
2413        std::io::stdout().flush().expect("flush readiness signal");
2414        let mut release = String::new();
2415        std::io::stdin()
2416            .read_line(&mut release)
2417            .expect("wait for release signal");
2418        conn.execute_batch("COMMIT")
2419            .expect("helper releases read snapshot");
2420    }
2421
2422    #[test]
2423    #[cfg(unix)]
2424    #[serial(checkpoint_skip_metrics, walpin_report_seam)]
2425    fn no_progress_report_keeps_holder_released_after_truncate_timeout() {
2426        let dir = tempfile::tempdir().expect("tempdir");
2427        let path = dir.path().join("transient-reader.db");
2428        let pool = file_pool(&path);
2429        {
2430            let writer = pool.try_writer().expect("writer");
2431            writer
2432                .conn()
2433                .execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
2434                .expect("seed WAL before reader snapshot");
2435        }
2436
2437        let mut reader = ReaderProcess::spawn(&path);
2438        let reader_pid = reader.pid();
2439        {
2440            let writer = pool.try_writer().expect("writer");
2441            writer
2442                .conn()
2443                .execute_batch("INSERT INTO t VALUES (2);")
2444                .expect("append WAL behind reader snapshot");
2445        }
2446
2447        let canonical_path = pool
2448            .canonical_path()
2449            .expect("file-backed pool has canonical path")
2450            .to_path_buf();
2451        let (reached_rx, proceed_tx) = truncate_report_test_sync::install(canonical_path.clone());
2452        let _hook_guard = TruncateReportHookGuard;
2453        let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
2454        let thread_buffer = std::sync::Arc::clone(&buffer);
2455        let checkpoint_pool = Arc::clone(&pool);
2456        let checkpoint = std::thread::spawn(move || {
2457            let subscriber = CaptureSubscriber {
2458                events: thread_buffer,
2459            };
2460            tracing::subscriber::with_default(subscriber, || {
2461                checkpoint_once(
2462                    &checkpoint_pool,
2463                    &CheckpointConfig {
2464                        truncate_high_water_pages: 0,
2465                        truncate_min_interval: Duration::ZERO,
2466                        truncate_busy_timeout: Duration::from_millis(50),
2467                        ..CheckpointConfig::default()
2468                    },
2469                    &mut TruncateState::default(),
2470                )
2471            })
2472        });
2473
2474        reached_rx
2475            .recv_timeout(Duration::from_secs(5))
2476            .expect("TRUNCATE must report no progress while the reader is pinned");
2477        reader.release();
2478        let post_attempt_census =
2479            crate::walpin::census_holders(&canonical_path).expect("post-attempt holder census");
2480        assert!(
2481            !post_attempt_census.holders.contains(&reader_pid),
2482            "released reader PID must be absent from a post-attempt census"
2483        );
2484        proceed_tx
2485            .send(())
2486            .expect("allow no-progress reporting to continue");
2487        checkpoint.join().expect("checkpoint thread");
2488
2489        let events = buffer.lock().expect("captured events");
2490        assert!(
2491            events.iter().any(|event| {
2492                event
2493                    .census_only
2494                    .as_deref()
2495                    .is_some_and(|pids| pids.contains(&reader_pid.to_string()))
2496            }),
2497            "the no-progress report must retain PID {reader_pid} from the pre-attempt census: {events:?}"
2498        );
2499    }
2500
2501    // `checkpoint_once` -> `query_wal_pages` writes the process-wide
2502    // `LAST_WAL_PAGES` gauge and resets `CHECKPOINT_CONSECUTIVE_SKIPS`
2503    // (see the reset-discipline comment on `reset_checkpoint_metrics_for_tests`
2504    // above) — this must join the `checkpoint_skip_metrics` group so it can
2505    // never interleave with a test asserting on those same gauges.
2506    #[test]
2507    #[serial(checkpoint_skip_metrics)]
2508    fn checkpoint_once_succeeds_on_file_backed_pool() {
2509        let dir = tempfile::tempdir().unwrap();
2510        let path = dir.path().join("wal_test.db");
2511        let pool = file_pool(&path);
2512
2513        // Create a table so the DB is not completely empty.
2514        {
2515            let writer = pool.try_writer().unwrap();
2516            writer
2517                .conn()
2518                .execute_batch("CREATE TABLE IF NOT EXISTS t (x INTEGER);")
2519                .unwrap();
2520            writer
2521                .conn()
2522                .execute_batch("INSERT INTO t VALUES (1);")
2523                .unwrap();
2524        }
2525
2526        checkpoint_once(
2527            &pool,
2528            &CheckpointConfig::default(),
2529            &mut TruncateState::default(),
2530        );
2531    }
2532
2533    #[test]
2534    #[serial(checkpoint_skip_metrics)]
2535    fn checkpoint_once_is_noop_on_in_memory_pool() {
2536        // In-memory databases do not use WAL; checkpoint_once must not panic.
2537        let cfg = PoolConfig {
2538            path: None,
2539            ..PoolConfig::default()
2540        };
2541        let pool = Arc::new(ConnectionPool::new(cfg).expect("in-memory pool"));
2542        checkpoint_once(
2543            &pool,
2544            &CheckpointConfig::default(),
2545            &mut TruncateState::default(),
2546        );
2547    }
2548
2549    #[tokio::test]
2550    #[serial(checkpoint_skip_metrics)]
2551    async fn checkpoint_task_exits_on_shutdown_signal() {
2552        let dir = tempfile::tempdir().unwrap();
2553        let path = dir.path().join("wal_task_shutdown.db");
2554        let pool = file_pool(&path);
2555
2556        // Use a very short interval so the task ticks quickly in the test.
2557        let cfg = CheckpointConfig {
2558            interval: Duration::from_millis(10),
2559            ..Default::default()
2560        };
2561
2562        let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
2563        let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
2564
2565        shutdown_tx.send(()).expect("send shutdown signal");
2566
2567        tokio::time::timeout(Duration::from_secs(1), handle)
2568            .await
2569            .expect("checkpoint task should exit within 1s")
2570            .expect("checkpoint task panicked");
2571    }
2572
2573    /// Regression #774: exits via watch-signal even with a live event_store
2574    /// pool clone (rules out a strong-count-based exit condition). See
2575    /// crates/khive-db/docs/api/checkpoint.md#checkpoint_task_exits_via_shutdown_signal_with_live_event_store_pool_clone
2576    #[tokio::test]
2577    #[serial(checkpoint_skip_metrics)]
2578    async fn checkpoint_task_exits_via_shutdown_signal_with_live_event_store_pool_clone() {
2579        let dir = tempfile::tempdir().unwrap();
2580        let path = dir.path().join("wal_task_event_store.db");
2581        let pool = file_pool(&path);
2582
2583        let cfg = CheckpointConfig {
2584            interval: Duration::from_millis(10),
2585            ..Default::default()
2586        };
2587
2588        let event_store: Arc<dyn khive_storage::EventStore> =
2589            Arc::new(crate::stores::event::SqlEventStore::new_scoped(
2590                Arc::clone(&pool),
2591                true,
2592                "local".to_string(),
2593            ));
2594        // A second, independent sibling clone of `pool` outlives this test
2595        // function's own binding — mirrors `StorageBackend` retaining
2596        // `self.pool` alongside the `SqlEventStore` it hands to the
2597        // checkpoint task in production.
2598        let sibling_pool_clone = Arc::clone(&pool);
2599
2600        let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
2601        let handle = tokio::spawn(run_checkpoint_task(
2602            pool,
2603            cfg,
2604            Some(CheckpointLifecycleOwner::new(event_store, "local")),
2605            shutdown_rx,
2606            true,
2607        ));
2608
2609        // Confirm strong_count is well above 1 — the old check would spin
2610        // forever here — before proving the new signal-based exit works
2611        // regardless.
2612        assert!(
2613            Arc::strong_count(&sibling_pool_clone) > 1,
2614            "test setup must reproduce the multi-owner shape the bug depends on"
2615        );
2616
2617        shutdown_tx.send(()).expect("send shutdown signal");
2618
2619        tokio::time::timeout(Duration::from_secs(1), handle)
2620            .await
2621            .expect(
2622                "checkpoint task should exit within 1s via the watch signal, \
2623                 even with a live sibling Arc<ConnectionPool> clone held by \
2624                 the event store",
2625            )
2626            .expect("checkpoint task panicked");
2627    }
2628
2629    #[test]
2630    #[serial]
2631    fn checkpoint_config_env_override() {
2632        std::env::set_var("KHIVE_CHECKPOINT_INTERVAL_MS", "250");
2633        std::env::set_var("KHIVE_WAL_WARN_PAGES", "1500");
2634        std::env::set_var("KHIVE_WAL_HIGH_WATER_PAGES", "8000");
2635        std::env::set_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES", "12000");
2636        std::env::set_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS", "60");
2637        std::env::set_var("KHIVE_WAL_TRUNCATE_BUSY_MS", "500");
2638        std::env::set_var("KHIVE_TX_WARN_SECS", "15");
2639        std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "90");
2640
2641        let cfg = CheckpointConfig::from_env();
2642
2643        std::env::remove_var("KHIVE_CHECKPOINT_INTERVAL_MS");
2644        std::env::remove_var("KHIVE_WAL_WARN_PAGES");
2645        std::env::remove_var("KHIVE_WAL_HIGH_WATER_PAGES");
2646        std::env::remove_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES");
2647        std::env::remove_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS");
2648        std::env::remove_var("KHIVE_WAL_TRUNCATE_BUSY_MS");
2649        std::env::remove_var("KHIVE_TX_WARN_SECS");
2650        std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
2651
2652        assert_eq!(cfg.interval, Duration::from_millis(250));
2653        assert_eq!(cfg.warn_pages, 1500);
2654        assert_eq!(cfg.high_water_pages, 8000);
2655        assert_eq!(cfg.truncate_high_water_pages, 12000);
2656        assert_eq!(cfg.truncate_min_interval, Duration::from_secs(60));
2657        assert_eq!(cfg.truncate_busy_timeout, Duration::from_millis(500));
2658        assert_eq!(cfg.tx_warn_secs, Duration::from_secs(15));
2659        assert_eq!(cfg.tx_max_age_secs, Duration::from_secs(90));
2660    }
2661
2662    #[test]
2663    #[serial]
2664    fn checkpoint_config_defaults_on_invalid_env() {
2665        let default = CheckpointConfig::default();
2666
2667        std::env::set_var("KHIVE_CHECKPOINT_INTERVAL_MS", "not_a_number");
2668        std::env::set_var("KHIVE_WAL_WARN_PAGES", "");
2669        std::env::set_var("KHIVE_WAL_HIGH_WATER_PAGES", "0");
2670        std::env::set_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES", "not_a_number");
2671        std::env::set_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS", "");
2672        std::env::set_var("KHIVE_WAL_TRUNCATE_BUSY_MS", "0");
2673        std::env::set_var("KHIVE_TX_WARN_SECS", "not_a_number");
2674        std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "0");
2675
2676        let cfg = CheckpointConfig::from_env();
2677
2678        std::env::remove_var("KHIVE_CHECKPOINT_INTERVAL_MS");
2679        std::env::remove_var("KHIVE_WAL_WARN_PAGES");
2680        std::env::remove_var("KHIVE_WAL_HIGH_WATER_PAGES");
2681        std::env::remove_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES");
2682        std::env::remove_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS");
2683        std::env::remove_var("KHIVE_WAL_TRUNCATE_BUSY_MS");
2684        std::env::remove_var("KHIVE_TX_WARN_SECS");
2685        std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
2686
2687        assert_eq!(cfg.interval, default.interval);
2688        assert_eq!(cfg.warn_pages, default.warn_pages);
2689        assert_eq!(cfg.high_water_pages, default.high_water_pages);
2690        assert_eq!(
2691            cfg.truncate_high_water_pages,
2692            default.truncate_high_water_pages
2693        );
2694        assert_eq!(cfg.truncate_min_interval, default.truncate_min_interval);
2695        assert_eq!(cfg.truncate_busy_timeout, default.truncate_busy_timeout);
2696        assert_eq!(cfg.tx_warn_secs, default.tx_warn_secs);
2697        assert_eq!(cfg.tx_max_age_secs, default.tx_max_age_secs);
2698    }
2699
2700    /// Regression: a high-water tick must NOT block behind an active read
2701    /// transaction (isomorphism guarantee — fails if `checkpoint_once`
2702    /// regresses to TRUNCATE). See
2703    /// crates/khive-db/docs/api/checkpoint.md#checkpoint_high_water_does_not_block_behind_reader
2704    #[test]
2705    #[serial(checkpoint_skip_metrics)]
2706    fn checkpoint_high_water_does_not_block_behind_reader() {
2707        let dir = tempfile::tempdir().unwrap();
2708        let path = dir.path().join("high_water_test.db");
2709
2710        // busy_timeout = 2000ms: a TRUNCATE regression blocks ~2s (clearly caught by
2711        // the <500ms assertion below), but PASSIVE returns well within 500ms even on
2712        // a heavily loaded CI runner. 4x margin on both sides vs. the old 200ms/50ms.
2713        let pool = Arc::new(
2714            ConnectionPool::new(PoolConfig {
2715                path: Some(path.clone()),
2716                busy_timeout: Duration::from_millis(2000),
2717                ..PoolConfig::default()
2718            })
2719            .expect("pool open"),
2720        );
2721
2722        // Write data so the WAL has frames to checkpoint.
2723        {
2724            let writer = pool.try_writer().unwrap();
2725            writer
2726                .conn()
2727                .execute_batch(
2728                    "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
2729                )
2730                .unwrap();
2731        }
2732
2733        // Open a reader and start a real read transaction so it holds a WAL
2734        // snapshot. An idle connection (no BEGIN) does NOT pin frames and would
2735        // not cause TRUNCATE to wait — the transaction is required for isomorphism.
2736        let reader = pool.reader().expect("reader");
2737        reader
2738            .execute_batch("BEGIN DEFERRED; SELECT * FROM t;")
2739            .expect("begin read tx");
2740
2741        // Write another row AFTER the snapshot is established. These new WAL
2742        // frames are now pinned by the open reader snapshot — TRUNCATE cannot
2743        // reclaim them without waiting; PASSIVE skips them and returns immediately.
2744        {
2745            let writer = pool.try_writer().unwrap();
2746            writer
2747                .conn()
2748                .execute_batch("INSERT INTO t VALUES (2);")
2749                .unwrap();
2750        }
2751
2752        let start = std::time::Instant::now();
2753        checkpoint_once(
2754            &pool,
2755            &CheckpointConfig::default(),
2756            &mut TruncateState::default(),
2757        );
2758        let elapsed = start.elapsed();
2759
2760        // Commit and release the read snapshot only after checkpoint_once returns.
2761        reader.execute_batch("COMMIT;").ok();
2762        drop(reader);
2763
2764        // PASSIVE returns in <1ms even with an open reader snapshot.
2765        // A TRUNCATE regression would block ~busy_timeout (2000ms) and fail here.
2766        // 500ms threshold is generous for CI jitter while staying well below 2000ms.
2767        assert!(
2768            elapsed < std::time::Duration::from_millis(500),
2769            "checkpoint_once with active reader snapshot took {:?}; \
2770             expected <500ms (PASSIVE must not block on readers; \
2771             a TRUNCATE regression would block ~2000ms)",
2772            elapsed
2773        );
2774    }
2775
2776    #[test]
2777    #[serial]
2778    fn checkpoint_config_rejects_zero_for_all_fields() {
2779        let default = CheckpointConfig::default();
2780        std::env::set_var("KHIVE_CHECKPOINT_INTERVAL_MS", "0");
2781        std::env::set_var("KHIVE_WAL_WARN_PAGES", "0");
2782        std::env::set_var("KHIVE_WAL_HIGH_WATER_PAGES", "0");
2783        std::env::set_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES", "0");
2784        std::env::set_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS", "0");
2785        std::env::set_var("KHIVE_WAL_TRUNCATE_BUSY_MS", "0");
2786        std::env::set_var("KHIVE_TX_WARN_SECS", "0");
2787        std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "0");
2788
2789        let cfg = CheckpointConfig::from_env();
2790
2791        std::env::remove_var("KHIVE_CHECKPOINT_INTERVAL_MS");
2792        std::env::remove_var("KHIVE_WAL_WARN_PAGES");
2793        std::env::remove_var("KHIVE_WAL_HIGH_WATER_PAGES");
2794        std::env::remove_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES");
2795        std::env::remove_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS");
2796        std::env::remove_var("KHIVE_WAL_TRUNCATE_BUSY_MS");
2797        std::env::remove_var("KHIVE_TX_WARN_SECS");
2798        std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
2799
2800        assert_eq!(
2801            cfg.interval, default.interval,
2802            "zero interval must fall back to default"
2803        );
2804        assert_eq!(
2805            cfg.warn_pages, default.warn_pages,
2806            "zero warn_pages must fall back to default"
2807        );
2808        assert_eq!(
2809            cfg.high_water_pages, default.high_water_pages,
2810            "zero high_water_pages must fall back to default"
2811        );
2812        assert_eq!(
2813            cfg.truncate_high_water_pages, default.truncate_high_water_pages,
2814            "zero truncate_high_water_pages must fall back to default"
2815        );
2816        assert_eq!(
2817            cfg.truncate_min_interval, default.truncate_min_interval,
2818            "zero truncate_min_interval must fall back to default"
2819        );
2820        assert_eq!(
2821            cfg.truncate_busy_timeout, default.truncate_busy_timeout,
2822            "zero truncate_busy_timeout must fall back to default"
2823        );
2824        assert_eq!(
2825            cfg.tx_warn_secs, default.tx_warn_secs,
2826            "zero tx_warn_secs must fall back to default"
2827        );
2828        assert_eq!(
2829            cfg.tx_max_age_secs, default.tx_max_age_secs,
2830            "zero tx_max_age_secs must fall back to default"
2831        );
2832    }
2833
2834    /// Fix: a reversed threshold pair must not be honored independently. See
2835    /// crates/khive-db/docs/api/checkpoint.md#checkpoint_config_rejects_reversed_tx_thresholds
2836    #[test]
2837    #[serial]
2838    fn checkpoint_config_rejects_reversed_tx_thresholds() {
2839        let default = CheckpointConfig::default();
2840        std::env::set_var("KHIVE_TX_WARN_SECS", "120");
2841        std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "30");
2842
2843        let cfg = CheckpointConfig::from_env();
2844
2845        std::env::remove_var("KHIVE_TX_WARN_SECS");
2846        std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
2847
2848        assert_eq!(
2849            cfg.tx_warn_secs, default.tx_warn_secs,
2850            "a reversed pair must fall back tx_warn_secs to its default, got: {:?}",
2851            cfg.tx_warn_secs
2852        );
2853        assert_eq!(
2854            cfg.tx_max_age_secs, default.tx_max_age_secs,
2855            "a reversed pair must fall back tx_max_age_secs to its default, got: {:?}",
2856            cfg.tx_max_age_secs
2857        );
2858    }
2859
2860    /// Degenerate equal-thresholds case; see
2861    /// crates/khive-db/docs/api/checkpoint.md#checkpoint_config_rejects_equal_tx_thresholds
2862    #[test]
2863    #[serial]
2864    fn checkpoint_config_rejects_equal_tx_thresholds() {
2865        let default = CheckpointConfig::default();
2866        std::env::set_var("KHIVE_TX_WARN_SECS", "60");
2867        std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "60");
2868
2869        let cfg = CheckpointConfig::from_env();
2870
2871        std::env::remove_var("KHIVE_TX_WARN_SECS");
2872        std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
2873
2874        assert_eq!(
2875            cfg.tx_warn_secs, default.tx_warn_secs,
2876            "an equal pair must fall back tx_warn_secs to its default, got: {:?}",
2877            cfg.tx_warn_secs
2878        );
2879        assert_eq!(
2880            cfg.tx_max_age_secs, default.tx_max_age_secs,
2881            "an equal pair must fall back tx_max_age_secs to its default, got: {:?}",
2882            cfg.tx_max_age_secs
2883        );
2884    }
2885
2886    /// Regression: a Skipped tick must NOT reset `was_above_high_water`. See
2887    /// crates/khive-db/docs/api/checkpoint.md#skipped_tick_does_not_reset_high_water_crossing_state
2888    #[test]
2889    fn skipped_tick_does_not_reset_high_water_crossing_state() {
2890        let mut was_above = false;
2891
2892        // First observed tick: above threshold — fires WARN, sets was_above=true.
2893        assert!(
2894            crossing_warn(true, &mut was_above),
2895            "should fire on first crossing"
2896        );
2897        assert!(was_above);
2898
2899        // Simulate several skipped ticks: crossing state must remain true.
2900        // (In the task, Skipped causes `continue` so crossing_warn is never called.)
2901        // We verify by calling crossing_warn with the SAME above=true value, which
2902        // is what Observed(high_count) would produce — but a Skipped tick skips
2903        // the call entirely, so was_above stays as-is. Test the invariant directly:
2904        // if we leave was_above unchanged (no call at all), was_above remains true.
2905        assert!(was_above, "was_above must stay true across skipped ticks");
2906
2907        // Another observed tick still above threshold — must NOT re-fire.
2908        let fired = crossing_warn(true, &mut was_above);
2909        assert!(!fired, "WARN must not re-fire while still above threshold");
2910
2911        // Observed tick below threshold — resets was_above.
2912        let fired = crossing_warn(false, &mut was_above);
2913        assert!(!fired);
2914        assert!(!was_above);
2915
2916        // Next observed tick above threshold — fires again (legitimate new crossing).
2917        let fired = crossing_warn(true, &mut was_above);
2918        assert!(fired, "WARN must fire again on a new below→above crossing");
2919    }
2920
2921    /// Regression: warn_pages WARN fires once on crossing, not every tick.
2922    ///
2923    /// Before the fix, the WARN was emitted inside `checkpoint_once` on every tick
2924    /// while WAL sat in the warn band — log spam under sustained moderate pressure.
2925    /// With the fix, `crossing_warn` gates the WARN on the first in-band tick only;
2926    /// subsequent ticks while still in the band return false.
2927    #[test]
2928    fn warn_pages_fires_once_on_crossing_not_every_tick() {
2929        let mut was_above_warn = false;
2930
2931        // Simulate three consecutive ticks with WAL in the warn band.
2932        let fired_1 = crossing_warn(true, &mut was_above_warn);
2933        let fired_2 = crossing_warn(true, &mut was_above_warn);
2934        let fired_3 = crossing_warn(true, &mut was_above_warn);
2935
2936        assert!(fired_1, "WARN must fire on the first in-band tick");
2937        assert!(
2938            !fired_2,
2939            "WARN must not fire on the second consecutive in-band tick"
2940        );
2941        assert!(
2942            !fired_3,
2943            "WARN must not fire on the third consecutive in-band tick"
2944        );
2945
2946        // Drop below warn band — resets state.
2947        crossing_warn(false, &mut was_above_warn);
2948        assert!(!was_above_warn);
2949
2950        // Re-enter warn band — fires again.
2951        let fired_reentry = crossing_warn(true, &mut was_above_warn);
2952        assert!(
2953            fired_reentry,
2954            "WARN must fire again on re-entry into warn band"
2955        );
2956    }
2957
2958    // ADR-091 Plank 2: TRUNCATE escalation state machine tests.
2959
2960    /// Trigger threshold: once `wal_pages` (as observed by `checkpoint_once`) is
2961    /// at/above `truncate_high_water_pages` and no prior attempt has run, the
2962    /// escalation fires and stamps `last_attempt`.
2963    #[test]
2964    #[serial(tx_registry, checkpoint_skip_metrics)]
2965    fn truncate_attempts_when_high_water_crossed_with_no_prior_attempt() {
2966        let dir = tempfile::tempdir().unwrap();
2967        let path = dir.path().join("truncate_trigger.db");
2968        let pool = file_pool(&path);
2969
2970        {
2971            let writer = pool.try_writer().unwrap();
2972            writer
2973                .conn()
2974                .execute_batch(
2975                    "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
2976                )
2977                .unwrap();
2978        }
2979
2980        let config = CheckpointConfig {
2981            // Force the escalation to arm regardless of the tiny WAL this test
2982            // actually produces — isolates the trigger-threshold behavior from
2983            // needing to stuff 20,000 real WAL pages.
2984            truncate_high_water_pages: 0,
2985            truncate_min_interval: Duration::from_secs(300),
2986            ..CheckpointConfig::default()
2987        };
2988        let mut state = TruncateState::default();
2989
2990        assert!(
2991            state.last_attempt.is_none(),
2992            "precondition: no attempt has run yet"
2993        );
2994
2995        let tick = checkpoint_once(&pool, &config, &mut state);
2996        assert!(matches!(tick, CheckpointTick::Observed(_)));
2997        assert!(
2998            state.last_attempt.is_some(),
2999            "an attempt must be stamped once the high-water threshold is crossed"
3000        );
3001    }
3002
3003    /// Below-threshold skip: `wal_pages < truncate_high_water_pages` must never
3004    /// stamp `last_attempt` — only an actual attempt advances it.
3005    #[test]
3006    #[serial(tx_registry, checkpoint_skip_metrics)]
3007    fn truncate_does_not_attempt_below_high_water() {
3008        let dir = tempfile::tempdir().unwrap();
3009        let path = dir.path().join("truncate_below_threshold.db");
3010        let pool = file_pool(&path);
3011
3012        {
3013            let writer = pool.try_writer().unwrap();
3014            writer
3015                .conn()
3016                .execute_batch(
3017                    "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
3018                )
3019                .unwrap();
3020        }
3021
3022        // Effectively unreachable threshold for this test's tiny WAL.
3023        let config = CheckpointConfig {
3024            truncate_high_water_pages: u64::MAX,
3025            ..CheckpointConfig::default()
3026        };
3027        let mut state = TruncateState::default();
3028
3029        checkpoint_once(&pool, &config, &mut state);
3030
3031        assert!(
3032            state.last_attempt.is_none(),
3033            "a below-threshold tick must never stamp last_attempt"
3034        );
3035    }
3036
3037    /// Min-interval skip: once an attempt has run, a subsequent tick that is
3038    /// still above threshold but within `truncate_min_interval` must skip
3039    /// without re-stamping `last_attempt` (the timestamp must not move).
3040    #[test]
3041    #[serial(tx_registry, checkpoint_skip_metrics)]
3042    fn truncate_min_interval_skip_does_not_restamp_last_attempt() {
3043        let dir = tempfile::tempdir().unwrap();
3044        let path = dir.path().join("truncate_min_interval.db");
3045        let pool = file_pool(&path);
3046
3047        {
3048            let writer = pool.try_writer().unwrap();
3049            writer
3050                .conn()
3051                .execute_batch(
3052                    "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
3053                )
3054                .unwrap();
3055        }
3056
3057        let config = CheckpointConfig {
3058            truncate_high_water_pages: 0,
3059            truncate_min_interval: Duration::from_secs(300),
3060            ..CheckpointConfig::default()
3061        };
3062        let mut state = TruncateState::default();
3063
3064        checkpoint_once(&pool, &config, &mut state);
3065        let first_attempt = state.last_attempt.expect("first tick must attempt");
3066
3067        // Second tick, immediately after: still above threshold, but the
3068        // min-interval has clearly not elapsed — must skip and leave
3069        // last_attempt exactly as it was.
3070        checkpoint_once(&pool, &config, &mut state);
3071        let second_attempt = state.last_attempt.expect("attempt timestamp must persist");
3072
3073        assert_eq!(
3074            first_attempt, second_attempt,
3075            "a tick within truncate_min_interval must not re-stamp last_attempt"
3076        );
3077    }
3078
3079    /// Busy fallback: when the writer mutex is already held, `checkpoint_once`
3080    /// must return `Skipped` and never touch the TRUNCATE state at all — both
3081    /// PASSIVE and any due TRUNCATE are skipped together (one writer checkout
3082    /// per tick). Also asserts #646 checkpoint-pressure telemetry: a skipped
3083    /// tick must bump the skipped/consecutive-skip counters and snapshot the
3084    /// last-known WAL pressure.
3085    #[test]
3086    #[serial(tx_registry, checkpoint_skip_metrics)]
3087    fn busy_writer_skips_both_passive_and_truncate() {
3088        reset_checkpoint_metrics_for_tests();
3089
3090        let dir = tempfile::tempdir().unwrap();
3091        let path = dir.path().join("truncate_busy_skip.db");
3092        let pool = file_pool(&path);
3093
3094        {
3095            let writer = pool.try_writer().unwrap();
3096            writer
3097                .conn()
3098                .execute_batch(
3099                    "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
3100                )
3101                .unwrap();
3102        }
3103
3104        // An observed tick first, so the skip below has a last-known WAL
3105        // pressure snapshot to carry forward.
3106        let mut warmup_state = TruncateState::default();
3107        let warmup_tick = checkpoint_once(&pool, &CheckpointConfig::default(), &mut warmup_state);
3108        let observed_pages = match warmup_tick {
3109            CheckpointTick::Observed(n) => n,
3110            CheckpointTick::Skipped => panic!("warmup tick must observe, not skip"),
3111        };
3112        assert_eq!(
3113            checkpoint_consecutive_skips(),
3114            0,
3115            "an observed tick must not itself count as a skip"
3116        );
3117
3118        // Hold the writer mutex for the duration of the checkpoint_once call so
3119        // try_writer_nowait() fails, exactly like a concurrent write in progress.
3120        let _held = pool.try_writer().unwrap();
3121
3122        let config = CheckpointConfig {
3123            truncate_high_water_pages: 0,
3124            ..CheckpointConfig::default()
3125        };
3126        let mut state = TruncateState::default();
3127
3128        let tick = checkpoint_once(&pool, &config, &mut state);
3129
3130        assert_eq!(
3131            tick,
3132            CheckpointTick::Skipped,
3133            "a busy writer must skip the tick entirely"
3134        );
3135        assert!(
3136            state.last_attempt.is_none(),
3137            "a skipped tick (writer busy) must never stamp last_attempt, \
3138             even with a threshold that would otherwise arm immediately"
3139        );
3140
3141        assert_eq!(
3142            checkpoint_skipped_ticks(),
3143            1,
3144            "one skipped tick must bump the lifetime skipped-tick counter"
3145        );
3146        assert_eq!(
3147            checkpoint_consecutive_skips(),
3148            1,
3149            "one skipped tick must bump the consecutive-skip run length"
3150        );
3151        assert_eq!(
3152            checkpoint_last_skip_wal_pages(),
3153            Some(observed_pages),
3154            "the skip must snapshot the last-observed WAL pressure"
3155        );
3156    }
3157
3158    /// Regression guard for #845 (a recurrence of the #828 shared-statics
3159    /// race): every test in this module that calls `checkpoint_once` or
3160    /// `run_checkpoint_task` — both funnel through `query_wal_pages`, which
3161    /// writes the process-wide `LAST_WAL_PAGES` / `CHECKPOINT_*` atomics —
3162    /// must be tagged with a `#[serial(...)]` group that includes
3163    /// `checkpoint_skip_metrics`. Before #828, six such call sites carried no
3164    /// serial tag at all: cargo's default test thread pool ran them
3165    /// concurrently with `busy_writer_skips_both_passive_and_truncate`, and an
3166    /// untagged tick's `query_wal_pages` call clobbered the gauges between
3167    /// this test's warmup tick and its skip assertion (`left: Some(0), right:
3168    /// Some(3)` on CI — the two ticks never actually raced against each
3169    /// other, a third test's tick did). This scans the module's own source so
3170    /// a future test that calls either function without the tag fails this
3171    /// assertion instead of flaking on a loaded CI runner.
3172    #[test]
3173    fn all_checkpoint_metrics_callers_are_serial_tagged() {
3174        const SELF_SRC: &str = include_str!("checkpoint.rs");
3175        let lines: Vec<&str> = SELF_SRC.lines().collect();
3176
3177        let attr_starts: Vec<usize> = lines
3178            .iter()
3179            .enumerate()
3180            .filter(|(_, l)| {
3181                let t = l.trim();
3182                t == "#[test]" || t.starts_with("#[tokio::test")
3183            })
3184            .map(|(i, _)| i)
3185            .collect();
3186
3187        let mut offenders = Vec::new();
3188
3189        for (idx, &start) in attr_starts.iter().enumerate() {
3190            let end = attr_starts.get(idx + 1).copied().unwrap_or(lines.len());
3191            let span = &lines[start..end];
3192
3193            let touches_shared_metrics = span
3194                .iter()
3195                .any(|l| l.contains("checkpoint_once(") || l.contains("run_checkpoint_task("));
3196            if !touches_shared_metrics {
3197                continue;
3198            }
3199
3200            let has_group_tag = span
3201                .iter()
3202                .any(|l| l.contains("#[serial") && l.contains("checkpoint_skip_metrics"));
3203
3204            if !has_group_tag {
3205                let name = span
3206                    .iter()
3207                    .find_map(|l| {
3208                        let t = l.trim_start();
3209                        let t = t.strip_prefix("pub(crate) ").unwrap_or(t);
3210                        let t = t.strip_prefix("pub ").unwrap_or(t);
3211                        let t = t.strip_prefix("async ").unwrap_or(t);
3212                        t.strip_prefix("fn ")
3213                            .map(|rest| rest.split(['(', '<']).next().unwrap_or("").trim())
3214                    })
3215                    .unwrap_or("<unknown test>");
3216                offenders.push(name.to_string());
3217            }
3218        }
3219
3220        assert!(
3221            offenders.is_empty(),
3222            "these tests call checkpoint_once/run_checkpoint_task (which write the \
3223             process-wide LAST_WAL_PAGES/CHECKPOINT_* atomics via query_wal_pages) but \
3224             are not tagged #[serial(checkpoint_skip_metrics)] (or a group including it); \
3225             an untagged caller running concurrently on cargo's default test thread pool \
3226             can clobber those atomics mid-assertion in another test (the #828/#845 race): \
3227             {offenders:?}"
3228        );
3229    }
3230
3231    /// Observation branch: a checkpoint tick that is actually observed (writer
3232    /// free) must close out a prior skip streak, resetting the
3233    /// consecutive-skip counter to 0 without touching the lifetime total.
3234    #[test]
3235    #[serial(tx_registry, checkpoint_skip_metrics)]
3236    fn observed_tick_resets_consecutive_skips_but_not_lifetime_total() {
3237        reset_checkpoint_metrics_for_tests();
3238
3239        let dir = tempfile::tempdir().unwrap();
3240        let path = dir.path().join("skip_then_observe.db");
3241        let pool = file_pool(&path);
3242
3243        {
3244            let writer = pool.try_writer().unwrap();
3245            writer
3246                .conn()
3247                .execute_batch(
3248                    "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
3249                )
3250                .unwrap();
3251        }
3252
3253        // Two consecutive skipped ticks.
3254        {
3255            let _held = pool.try_writer().unwrap();
3256            let mut state = TruncateState::default();
3257            for _ in 0..2 {
3258                let tick = checkpoint_once(&pool, &CheckpointConfig::default(), &mut state);
3259                assert_eq!(tick, CheckpointTick::Skipped);
3260            }
3261        }
3262        assert_eq!(checkpoint_skipped_ticks(), 2);
3263        assert_eq!(checkpoint_consecutive_skips(), 2);
3264
3265        // Now the writer is free: an observed tick must reset the streak.
3266        let mut state = TruncateState::default();
3267        let tick = checkpoint_once(&pool, &CheckpointConfig::default(), &mut state);
3268        assert!(matches!(tick, CheckpointTick::Observed(_)));
3269
3270        assert_eq!(
3271            checkpoint_skipped_ticks(),
3272            2,
3273            "an observed tick must not change the lifetime skipped-tick total"
3274        );
3275        assert_eq!(
3276            checkpoint_consecutive_skips(),
3277            0,
3278            "an observed tick must reset the consecutive-skip run length"
3279        );
3280    }
3281
3282    /// Edge-triggered escalation WARN: `note_truncate_outcome` fires exactly
3283    /// once, on the third consecutive attempt that fails to clear
3284    /// `warn_pages`, and does not repeat on a fourth consecutive failure. A
3285    /// single attempt that clears `warn_pages` resets the counter.
3286    #[test]
3287    fn note_truncate_outcome_warns_once_at_third_consecutive_failure() {
3288        let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
3289        let subscriber = CaptureSubscriber {
3290            events: std::sync::Arc::clone(&buffer),
3291        };
3292
3293        let config = CheckpointConfig {
3294            warn_pages: 2000,
3295            ..CheckpointConfig::default()
3296        };
3297        let mut state = TruncateState::default();
3298
3299        tracing::subscriber::with_default(subscriber, || {
3300            // Three consecutive attempts that fail to clear warn_pages.
3301            note_truncate_outcome(&config, 5000, &mut state);
3302            note_truncate_outcome(&config, 5000, &mut state);
3303            note_truncate_outcome(&config, 5000, &mut state);
3304            // A fourth consecutive failure must not re-fire the escalation.
3305            note_truncate_outcome(&config, 5000, &mut state);
3306        });
3307
3308        assert_eq!(state.consecutive_failures, 4);
3309
3310        let events = buffer.lock().unwrap();
3311        let escalation_count = events
3312            .iter()
3313            .filter(|e| {
3314                e.message.as_deref()
3315                    == Some(
3316                        "WAL TRUNCATE has failed to clear WAL pressure for 3 consecutive attempts",
3317                    )
3318            })
3319            .count();
3320        assert_eq!(
3321            escalation_count, 1,
3322            "escalation WARN must fire exactly once at the 3rd consecutive failure, got: {events:?}"
3323        );
3324
3325        // A clearing attempt resets the counter.
3326        note_truncate_outcome(&config, 100, &mut state);
3327        assert_eq!(
3328            state.consecutive_failures, 0,
3329            "an attempt that clears warn_pages must reset the consecutive-failure counter"
3330        );
3331    }
3332
3333    // ADR-091 #617: graduated severity ladder state-machine tests.
3334
3335    fn severity_test_config() -> CheckpointConfig {
3336        CheckpointConfig {
3337            warn_pages: 100,
3338            warn_sustained_cycles: 3,
3339            ..CheckpointConfig::default()
3340        }
3341    }
3342
3343    /// INFO rung: a below→above crossing emits exactly one INFO and no WARN
3344    /// (default `warn_sustained_cycles = 3`, only one above-warn tick here).
3345    #[test]
3346    fn severity_ladder_info_on_first_crossing_no_warn() {
3347        let config = severity_test_config();
3348        let mut state = CheckpointSeverityState::default();
3349
3350        let below = state.observe_wal_pages(10, &config);
3351        assert!(below.is_empty(), "below-warn tick must emit nothing");
3352
3353        let above = state.observe_wal_pages(150, &config);
3354        assert_eq!(
3355            above,
3356            vec![CheckpointSeverityEmission {
3357                rung: CheckpointSeverityRung::Info,
3358                wal_pages: 150,
3359                threshold_pages: 100,
3360                consecutive_cycles: 1,
3361            }],
3362            "first below->above crossing must emit exactly one INFO and no WARN"
3363        );
3364    }
3365
3366    /// WARN rung: `warn_sustained_cycles` (3) consecutive above-warn ticks
3367    /// emit WARN exactly on the third tick, not before and not repeated after.
3368    #[test]
3369    fn severity_ladder_warn_on_third_consecutive_cycle() {
3370        let config = severity_test_config();
3371        let mut state = CheckpointSeverityState::default();
3372
3373        let tick1 = state.observe_wal_pages(150, &config);
3374        assert_eq!(tick1.len(), 1);
3375        assert_eq!(tick1[0].rung, CheckpointSeverityRung::Info);
3376
3377        let tick2 = state.observe_wal_pages(150, &config);
3378        assert!(
3379            tick2.is_empty(),
3380            "second consecutive above-warn tick must emit nothing yet"
3381        );
3382
3383        let tick3 = state.observe_wal_pages(150, &config);
3384        assert_eq!(
3385            tick3,
3386            vec![CheckpointSeverityEmission {
3387                rung: CheckpointSeverityRung::Warn,
3388                wal_pages: 150,
3389                threshold_pages: 100,
3390                consecutive_cycles: 3,
3391            }],
3392            "WARN must fire exactly on the third consecutive above-warn tick"
3393        );
3394
3395        let tick4 = state.observe_wal_pages(150, &config);
3396        assert!(
3397            tick4.is_empty(),
3398            "WARN must not repeat on a fourth consecutive above-warn tick"
3399        );
3400    }
3401
3402    /// Re-arm: after a WARN episode drains below warn_pages, a fresh episode
3403    /// of `warn_sustained_cycles` above-warn ticks must WARN again.
3404    #[test]
3405    fn severity_ladder_rearms_warn_after_drain() {
3406        let config = severity_test_config();
3407        let mut state = CheckpointSeverityState::default();
3408
3409        // First episode reaches WARN.
3410        for _ in 0..3 {
3411            state.observe_wal_pages(150, &config);
3412        }
3413        assert!(state.warn_emitted_for_episode);
3414
3415        // Drain below warn_pages: resets the episode.
3416        let drain = state.observe_wal_pages(10, &config);
3417        assert!(drain.is_empty(), "a draining tick must emit nothing");
3418
3419        // Second episode: INFO on first tick, no WARN until the third again.
3420        let reentry = state.observe_wal_pages(150, &config);
3421        assert_eq!(reentry.len(), 1);
3422        assert_eq!(reentry[0].rung, CheckpointSeverityRung::Info);
3423
3424        let mid = state.observe_wal_pages(150, &config);
3425        assert!(mid.is_empty());
3426
3427        let second_warn = state.observe_wal_pages(150, &config);
3428        assert_eq!(
3429            second_warn,
3430            vec![CheckpointSeverityEmission {
3431                rung: CheckpointSeverityRung::Warn,
3432                wal_pages: 150,
3433                threshold_pages: 100,
3434                consecutive_cycles: 3,
3435            }],
3436            "a fresh elevation episode after a drain must WARN again"
3437        );
3438    }
3439
3440    /// False-positive guard: three isolated single-tick crossings, each
3441    /// followed by a drain, must never reach WARN — only INFO fires each time.
3442    #[test]
3443    fn severity_ladder_isolated_crossings_never_warn() {
3444        let config = severity_test_config();
3445        let mut state = CheckpointSeverityState::default();
3446
3447        for _ in 0..3 {
3448            let crossing = state.observe_wal_pages(150, &config);
3449            assert_eq!(
3450                crossing.len(),
3451                1,
3452                "each isolated crossing must emit exactly one INFO"
3453            );
3454            assert_eq!(crossing[0].rung, CheckpointSeverityRung::Info);
3455
3456            let drain = state.observe_wal_pages(10, &config);
3457            assert!(drain.is_empty(), "the drain tick must emit nothing");
3458        }
3459
3460        assert!(
3461            !state.warn_emitted_for_episode,
3462            "isolated single-tick crossings must never accumulate into a WARN"
3463        );
3464    }
3465
3466    /// ALARM rung: the existing TRUNCATE-attempt gate is the ADR-091 ALARM
3467    /// tier. `observe_wal_pages` never produces it; this test documents and
3468    /// locks in that boundary so a future change can't silently reroute
3469    /// ALARM through the INFO/WARN ladder.
3470    #[test]
3471    fn severity_ladder_never_emits_alarm() {
3472        let config = CheckpointConfig {
3473            warn_pages: 100,
3474            warn_sustained_cycles: 1,
3475            ..CheckpointConfig::default()
3476        };
3477        let mut state = CheckpointSeverityState::default();
3478
3479        for wal_pages in [150, 200, 250, u64::MAX] {
3480            let emissions = state.observe_wal_pages(wal_pages, &config);
3481            assert!(
3482                emissions
3483                    .iter()
3484                    .all(|e| e.rung != CheckpointSeverityRung::Alarm),
3485                "observe_wal_pages must never emit the ALARM rung, got: {emissions:?}"
3486            );
3487        }
3488    }
3489
3490    // ADR-091 Plank 1: `TxAgeSweepState` background-sweep state-machine tests.
3491    // Pure unit tests mirroring the severity-ladder tests above — no I/O.
3492
3493    fn tx_age_test_config() -> CheckpointConfig {
3494        CheckpointConfig {
3495            tx_warn_secs: Duration::from_secs(30),
3496            tx_max_age_secs: Duration::from_secs(120),
3497            ..CheckpointConfig::default()
3498        }
3499    }
3500
3501    /// Synthetic identity for `TxAgeSweepState::observe`'s pure unit tests
3502    /// below, which exercise identity-change detection without paying for a
3503    /// real `tx_registry::register` call. `TxId`'s wrapped value is public
3504    /// exactly to support this (see its doc comment in `khive-storage`).
3505    fn tx_id(n: u64) -> khive_storage::tx_registry::TxId {
3506        khive_storage::tx_registry::TxId(n)
3507    }
3508
3509    /// No open entry: nothing fires, and any prior latch state clears.
3510    #[test]
3511    fn tx_age_sweep_empty_registry_emits_nothing() {
3512        let config = tx_age_test_config();
3513        let mut state = TxAgeSweepState::default();
3514
3515        let emissions = state.observe(None, config.tx_warn_secs, config.tx_max_age_secs);
3516        assert!(emissions.is_empty(), "no open entry must emit nothing");
3517    }
3518
3519    /// A fresh entry (age below both thresholds) emits nothing.
3520    #[test]
3521    fn tx_age_sweep_fresh_entry_emits_nothing() {
3522        let config = tx_age_test_config();
3523        let mut state = TxAgeSweepState::default();
3524
3525        let emissions = state.observe(
3526            Some((
3527                tx_id(1),
3528                Duration::from_secs(5),
3529                Some("fresh_span".to_string()),
3530            )),
3531            config.tx_warn_secs,
3532            config.tx_max_age_secs,
3533        );
3534        assert!(emissions.is_empty(), "a fresh entry must emit nothing");
3535    }
3536
3537    /// Below→above crossing of `tx_warn_secs` fires exactly one `Warn`
3538    /// emission carrying the entry's label; it must not repeat on a second
3539    /// tick that is still above `tx_warn_secs` but below `tx_max_age_secs`.
3540    #[test]
3541    fn tx_age_sweep_warn_fires_once_on_crossing() {
3542        let config = tx_age_test_config();
3543        let mut state = TxAgeSweepState::default();
3544
3545        let tick1 = state.observe(
3546            Some((
3547                tx_id(1),
3548                Duration::from_secs(45),
3549                Some("stale_span".to_string()),
3550            )),
3551            config.tx_warn_secs,
3552            config.tx_max_age_secs,
3553        );
3554        assert_eq!(
3555            tick1,
3556            vec![TxAgeEmission {
3557                rung: TxAgeRung::Warn,
3558                age: Duration::from_secs(45),
3559                label: Some("stale_span".to_string()),
3560            }],
3561            "crossing tx_warn_secs must emit exactly one Warn"
3562        );
3563
3564        let tick2 = state.observe(
3565            Some((
3566                tx_id(1),
3567                Duration::from_secs(50),
3568                Some("stale_span".to_string()),
3569            )),
3570            config.tx_warn_secs,
3571            config.tx_max_age_secs,
3572        );
3573        assert!(
3574            tick2.is_empty(),
3575            "Warn must not repeat while the entry stays in the warn band"
3576        );
3577    }
3578
3579    /// Crossing `tx_max_age_secs` fires `Stale`; a further tick still above
3580    /// the cap must not repeat it.
3581    #[test]
3582    fn tx_age_sweep_stale_fires_once_on_crossing() {
3583        let config = tx_age_test_config();
3584        let mut state = TxAgeSweepState::default();
3585
3586        // Drive through the warn crossing first, matching real elapsed-time
3587        // progression (an entry ages through the warn band before the max).
3588        state.observe(
3589            Some((
3590                tx_id(1),
3591                Duration::from_secs(45),
3592                Some("stuck_writer_task_tx".to_string()),
3593            )),
3594            config.tx_warn_secs,
3595            config.tx_max_age_secs,
3596        );
3597
3598        let tick = state.observe(
3599            Some((
3600                tx_id(1),
3601                Duration::from_secs(130),
3602                Some("stuck_writer_task_tx".to_string()),
3603            )),
3604            config.tx_warn_secs,
3605            config.tx_max_age_secs,
3606        );
3607        assert_eq!(
3608            tick,
3609            vec![TxAgeEmission {
3610                rung: TxAgeRung::Stale,
3611                age: Duration::from_secs(130),
3612                label: Some("stuck_writer_task_tx".to_string()),
3613            }],
3614            "crossing tx_max_age_secs must emit exactly one Stale"
3615        );
3616
3617        let tick_repeat = state.observe(
3618            Some((
3619                tx_id(1),
3620                Duration::from_secs(200),
3621                Some("stuck_writer_task_tx".to_string()),
3622            )),
3623            config.tx_warn_secs,
3624            config.tx_max_age_secs,
3625        );
3626        assert!(
3627            tick_repeat.is_empty(),
3628            "Stale must not repeat while the entry stays above tx_max_age_secs"
3629        );
3630    }
3631
3632    /// An entry already stale the first time the sweep observes it (e.g.
3633    /// right after process start with a pre-existing registry entry) crosses
3634    /// both rungs on the same tick.
3635    #[test]
3636    fn tx_age_sweep_already_stale_entry_emits_both_rungs_same_tick() {
3637        let config = tx_age_test_config();
3638        let mut state = TxAgeSweepState::default();
3639
3640        let tick = state.observe(
3641            Some((
3642                tx_id(1),
3643                Duration::from_secs(300),
3644                Some("ancient_tx".to_string()),
3645            )),
3646            config.tx_warn_secs,
3647            config.tx_max_age_secs,
3648        );
3649        assert_eq!(
3650            tick,
3651            vec![
3652                TxAgeEmission {
3653                    rung: TxAgeRung::Warn,
3654                    age: Duration::from_secs(300),
3655                    label: Some("ancient_tx".to_string()),
3656                },
3657                TxAgeEmission {
3658                    rung: TxAgeRung::Stale,
3659                    age: Duration::from_secs(300),
3660                    label: Some("ancient_tx".to_string()),
3661                },
3662            ],
3663            "an already-stale entry must cross both rungs on its first observed tick"
3664        );
3665    }
3666
3667    /// Re-arm: once the stale entry closes (registry reports a fresher
3668    /// oldest entry, or none at all), a future stale span must fire again.
3669    #[test]
3670    fn tx_age_sweep_rearms_after_entry_clears() {
3671        let config = tx_age_test_config();
3672        let mut state = TxAgeSweepState::default();
3673
3674        state.observe(
3675            Some((
3676                tx_id(1),
3677                Duration::from_secs(150),
3678                Some("first_span".to_string()),
3679            )),
3680            config.tx_warn_secs,
3681            config.tx_max_age_secs,
3682        );
3683
3684        // The stale span closed; nothing is open now.
3685        let cleared = state.observe(None, config.tx_warn_secs, config.tx_max_age_secs);
3686        assert!(cleared.is_empty(), "a clearing tick must emit nothing");
3687
3688        // A fresh entry (unrelated span) is now oldest — still below threshold.
3689        let fresh = state.observe(
3690            Some((
3691                tx_id(2),
3692                Duration::from_secs(2),
3693                Some("second_span".to_string()),
3694            )),
3695            config.tx_warn_secs,
3696            config.tx_max_age_secs,
3697        );
3698        assert!(fresh.is_empty(), "a fresh oldest entry must emit nothing");
3699
3700        // That second span goes stale in turn — must WARN again (re-armed).
3701        let rewarn = state.observe(
3702            Some((
3703                tx_id(2),
3704                Duration::from_secs(35),
3705                Some("second_span".to_string()),
3706            )),
3707            config.tx_warn_secs,
3708            config.tx_max_age_secs,
3709        );
3710        assert_eq!(
3711            rewarn,
3712            vec![TxAgeEmission {
3713                rung: TxAgeRung::Warn,
3714                age: Duration::from_secs(35),
3715                label: Some("second_span".to_string()),
3716            }],
3717            "a fresh stale episode after a clear must Warn again"
3718        );
3719    }
3720
3721    /// Fix: an already-stale entry replacing a stale one on the next tick,
3722    /// with no intervening clear, must still emit both rungs. See
3723    /// crates/khive-db/docs/api/checkpoint.md#tx_age_sweep_stale_replacement_without_intervening_clear_still_names_new_entry
3724    #[test]
3725    fn tx_age_sweep_stale_replacement_without_intervening_clear_still_names_new_entry() {
3726        let config = tx_age_test_config();
3727        let mut state = TxAgeSweepState::default();
3728
3729        let tick_a = state.observe(
3730            Some((
3731                tx_id(1),
3732                Duration::from_secs(300),
3733                Some("stale_entry_a".to_string()),
3734            )),
3735            config.tx_warn_secs,
3736            config.tx_max_age_secs,
3737        );
3738        assert_eq!(
3739            tick_a.len(),
3740            2,
3741            "entry A must cross both rungs on its first observed tick, got: {tick_a:?}"
3742        );
3743
3744        // B replaces A as the oldest entry on the VERY NEXT tick — already
3745        // stale itself, with no intervening None/below-threshold tick.
3746        let tick_b = state.observe(
3747            Some((
3748                tx_id(2),
3749                Duration::from_secs(400),
3750                Some("stale_entry_b".to_string()),
3751            )),
3752            config.tx_warn_secs,
3753            config.tx_max_age_secs,
3754        );
3755        assert_eq!(
3756            tick_b,
3757            vec![
3758                TxAgeEmission {
3759                    rung: TxAgeRung::Warn,
3760                    age: Duration::from_secs(400),
3761                    label: Some("stale_entry_b".to_string()),
3762                },
3763                TxAgeEmission {
3764                    rung: TxAgeRung::Stale,
3765                    age: Duration::from_secs(400),
3766                    label: Some("stale_entry_b".to_string()),
3767                },
3768            ],
3769            "a same-tick identity change to an already-stale successor must re-emit both \
3770             rungs naming the NEW entry, got: {tick_b:?}"
3771        );
3772    }
3773
3774    /// Closes the loop from env var to actual emitted rung. See
3775    /// crates/khive-db/docs/api/checkpoint.md#tx_age_sweep_uses_configured_thresholds_not_hardcoded_defaults
3776    #[test]
3777    fn tx_age_sweep_uses_configured_thresholds_not_hardcoded_defaults() {
3778        let config = CheckpointConfig {
3779            tx_warn_secs: Duration::from_millis(1),
3780            tx_max_age_secs: Duration::from_millis(2),
3781            ..CheckpointConfig::default()
3782        };
3783        let mut state = TxAgeSweepState::default();
3784
3785        let tick = state.observe(
3786            Some((
3787                tx_id(1),
3788                Duration::from_millis(5),
3789                Some("fast_cap_span".to_string()),
3790            )),
3791            config.tx_warn_secs,
3792            config.tx_max_age_secs,
3793        );
3794        assert_eq!(
3795            tick.len(),
3796            2,
3797            "a millisecond-scale cap must cross both rungs immediately, got: {tick:?}"
3798        );
3799    }
3800
3801    /// Integration-level regression for the incident this ADR fixes. See
3802    /// crates/khive-db/docs/api/checkpoint.md#tx_age_sweep_names_long_lived_reader_pinning_wal_past_high_water
3803    #[test]
3804    #[serial(tx_registry, checkpoint_skip_metrics)]
3805    fn tx_age_sweep_names_long_lived_reader_pinning_wal_past_high_water() {
3806        let dir = tempfile::tempdir().unwrap();
3807        let path = dir.path().join("tx_age_sweep_reader_pin.db");
3808        let pool = file_pool(&path);
3809
3810        {
3811            let writer = pool.try_writer().unwrap();
3812            writer
3813                .conn()
3814                .execute_batch(
3815                    "CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
3816                )
3817                .unwrap();
3818        }
3819
3820        // Open a real read transaction so it holds a WAL snapshot (same
3821        // isomorphism as `checkpoint_high_water_does_not_block_behind_reader`),
3822        // AND register it in tx_registry — the telemetry a real long-lived
3823        // reader call site (e.g. `graph_traverse_read`) is expected to carry.
3824        let reader = pool.reader().expect("reader");
3825        reader
3826            .execute_batch("BEGIN DEFERRED; SELECT * FROM t;")
3827            .expect("begin read tx");
3828        let _tx_handle =
3829            khive_storage::tx_registry::register(Some("tx_age_sweep_reader_pin_test".to_string()));
3830
3831        // Drive writes past high_water_pages while the reader snapshot pins
3832        // the WAL tail — PASSIVE cannot reclaim these frames.
3833        let config = CheckpointConfig {
3834            high_water_pages: 1,
3835            tx_warn_secs: Duration::from_millis(1),
3836            tx_max_age_secs: Duration::from_millis(1),
3837            ..CheckpointConfig::default()
3838        };
3839        {
3840            let writer = pool.try_writer().unwrap();
3841            for i in 0..50 {
3842                writer
3843                    .conn()
3844                    .execute_batch(&format!("INSERT INTO t VALUES ({i});"))
3845                    .unwrap();
3846            }
3847        }
3848
3849        let tick = checkpoint_once(&pool, &config, &mut TruncateState::default());
3850        let wal_pages = match tick {
3851            CheckpointTick::Observed(n) => n,
3852            CheckpointTick::Skipped => panic!("writer must not be busy in this test"),
3853        };
3854        assert!(
3855            wal_pages >= config.high_water_pages,
3856            "test setup must actually drive wal_pages ({wal_pages}) past high_water_pages \
3857             ({}) for this regression to mean anything",
3858            config.high_water_pages
3859        );
3860
3861        // The Plank 1 sweep, given the SAME registry state, must name the
3862        // pinning reader at the Stale rung. The handle's age must exceed the
3863        // 1ms `tx_max_age_secs` cap deterministically: the inserts plus one
3864        // PASSIVE checkpoint above can complete in under a millisecond on a
3865        // warm page cache, so sleep past the cap instead of assuming
3866        // the elapsed work already crossed it.
3867        std::thread::sleep(Duration::from_millis(5));
3868        // `tx_registry` is a process-wide singleton shared by every test in
3869        // this binary (cargo runs `#[test]`s in parallel threads of the same
3870        // process): `#[serial(tx_registry)]` only excludes other tests that
3871        // carry the same key, not every production write path elsewhere in
3872        // the crate (e.g. `graph_upsert_edges`) that also calls `register()`
3873        // as ordinary telemetry. If one of those happens to still be open and
3874        // was registered before this test's own handle, raw `oldest()` would
3875        // return THAT entry instead of the fixture's reader — see #926. Look
3876        // up this test's own entry by its known label instead of trusting
3877        // global `oldest()`, so the assertion is immune to that noise.
3878        let our_entry = khive_storage::tx_registry::snapshot()
3879            .into_iter()
3880            .find(|(_, label)| label.as_deref() == Some("tx_age_sweep_reader_pin_test"))
3881            .expect("this test's own tx_registry entry must still be open");
3882        let mut tx_age_state = TxAgeSweepState::default();
3883        let emissions = tx_age_state.observe(
3884            Some((tx_id(1), our_entry.0, our_entry.1)),
3885            config.tx_warn_secs,
3886            config.tx_max_age_secs,
3887        );
3888        assert!(
3889            emissions.iter().any(|e| e.rung == TxAgeRung::Stale
3890                && e.label.as_deref() == Some("tx_age_sweep_reader_pin_test")),
3891            "expected a Stale emission naming the pinning reader, got: {emissions:?}"
3892        );
3893
3894        reader.execute_batch("COMMIT;").ok();
3895        drop(reader);
3896        drop(_tx_handle);
3897    }
3898
3899    /// Regression #926: reproduces the exact tx_registry race directly. See
3900    /// crates/khive-db/docs/api/checkpoint.md#tx_age_sweep_own_entry_survives_concurrent_older_registration
3901    #[test]
3902    #[serial(tx_registry, checkpoint_skip_metrics)]
3903    fn tx_age_sweep_own_entry_survives_concurrent_older_registration() {
3904        let _decoy = khive_storage::tx_registry::register(Some("decoy_unrelated_span".to_string()));
3905        std::thread::sleep(Duration::from_millis(2));
3906        let _own = khive_storage::tx_registry::register(Some("this_test_own_span".to_string()));
3907        std::thread::sleep(Duration::from_millis(5));
3908
3909        // Confirm the race condition is actually reproduced: an entry older
3910        // than this test's own span must currently lead the process-wide
3911        // registry. Another concurrently running test may have registered an
3912        // entry before the decoy, so do not assume the decoy is globally
3913        // oldest; the required invariant is only that our span is not.
3914        let global_oldest = khive_storage::tx_registry::oldest().expect("registry not empty");
3915        assert_ne!(
3916            global_oldest.2.as_deref(),
3917            Some("this_test_own_span"),
3918            "test setup must reproduce the race: an older, unrelated entry must be \
3919             the current global oldest, got: {global_oldest:?}"
3920        );
3921
3922        let our_entry = khive_storage::tx_registry::snapshot()
3923            .into_iter()
3924            .find(|(_, label)| label.as_deref() == Some("this_test_own_span"))
3925            .expect("this test's own tx_registry entry must still be open");
3926
3927        let config = CheckpointConfig {
3928            tx_warn_secs: Duration::from_millis(1),
3929            tx_max_age_secs: Duration::from_millis(1),
3930            ..CheckpointConfig::default()
3931        };
3932        let mut state = TxAgeSweepState::default();
3933        let emissions = state.observe(
3934            Some((tx_id(2), our_entry.0, our_entry.1)),
3935            config.tx_warn_secs,
3936            config.tx_max_age_secs,
3937        );
3938        assert!(
3939            emissions
3940                .iter()
3941                .any(|e| e.rung == TxAgeRung::Stale
3942                    && e.label.as_deref() == Some("this_test_own_span")),
3943            "expected a Stale emission naming this test's own span despite an older, \
3944             unrelated concurrent registration, got: {emissions:?}"
3945        );
3946    }
3947
3948    /// `KHIVE_WAL_WARN_SUSTAINED_CYCLES` overrides the default and rejects 0.
3949    #[test]
3950    #[serial]
3951    fn checkpoint_config_warn_sustained_cycles_env_override() {
3952        let default = CheckpointConfig::default();
3953        assert_eq!(default.warn_sustained_cycles, DEFAULT_WARN_SUSTAINED_CYCLES);
3954
3955        std::env::set_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES", "5");
3956        let cfg = CheckpointConfig::from_env();
3957        std::env::remove_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES");
3958        assert_eq!(cfg.warn_sustained_cycles, 5);
3959
3960        std::env::set_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES", "0");
3961        let cfg_zero = CheckpointConfig::from_env();
3962        std::env::remove_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES");
3963        assert_eq!(
3964            cfg_zero.warn_sustained_cycles, DEFAULT_WARN_SUSTAINED_CYCLES,
3965            "zero must fall back to the default"
3966        );
3967
3968        std::env::set_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES", "not_a_number");
3969        let cfg_invalid = CheckpointConfig::from_env();
3970        std::env::remove_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES");
3971        assert_eq!(
3972            cfg_invalid.warn_sustained_cycles,
3973            DEFAULT_WARN_SUSTAINED_CYCLES
3974        );
3975    }
3976
3977    // ADR-094: `CheckpointOutcomeRecorded` lifecycle event tests.
3978
3979    #[derive(Clone, Copy)]
3980    enum FakeAppendBehavior {
3981        Record,
3982        Fail,
3983    }
3984
3985    struct FakeEventStore {
3986        events: std::sync::Mutex<Vec<khive_storage::Event>>,
3987        append_attempts: std::sync::atomic::AtomicUsize,
3988        append_behavior: FakeAppendBehavior,
3989    }
3990
3991    impl Default for FakeEventStore {
3992        fn default() -> Self {
3993            Self {
3994                events: std::sync::Mutex::new(Vec::new()),
3995                append_attempts: std::sync::atomic::AtomicUsize::new(0),
3996                append_behavior: FakeAppendBehavior::Record,
3997            }
3998        }
3999    }
4000
4001    impl FakeEventStore {
4002        fn failing() -> Self {
4003            Self {
4004                append_behavior: FakeAppendBehavior::Fail,
4005                ..Self::default()
4006            }
4007        }
4008    }
4009
4010    #[async_trait::async_trait]
4011    impl khive_storage::EventStore for FakeEventStore {
4012        async fn append_event(
4013            &self,
4014            event: khive_storage::Event,
4015        ) -> khive_storage::StorageResult<()> {
4016            self.append_attempts
4017                .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
4018            match self.append_behavior {
4019                FakeAppendBehavior::Record => {
4020                    self.events.lock().unwrap().push(event);
4021                    Ok(())
4022                }
4023                FakeAppendBehavior::Fail => Err(khive_storage::StorageError::Internal(
4024                    "synthetic checkpoint lifecycle append failure".to_string(),
4025                )),
4026            }
4027        }
4028
4029        async fn append_events(
4030            &self,
4031            events: Vec<khive_storage::Event>,
4032        ) -> khive_storage::StorageResult<khive_storage::BatchWriteSummary> {
4033            let count = events.len() as u64;
4034            self.events.lock().unwrap().extend(events);
4035            Ok(khive_storage::BatchWriteSummary {
4036                attempted: count,
4037                affected: count,
4038                failed: 0,
4039                first_error: String::new(),
4040            })
4041        }
4042
4043        async fn get_event(
4044            &self,
4045            id: uuid::Uuid,
4046        ) -> khive_storage::StorageResult<Option<khive_storage::Event>> {
4047            Ok(self
4048                .events
4049                .lock()
4050                .unwrap()
4051                .iter()
4052                .find(|e| e.id == id)
4053                .cloned())
4054        }
4055
4056        async fn query_events(
4057            &self,
4058            _filter: khive_storage::EventFilter,
4059            _page: khive_storage::PageRequest,
4060        ) -> khive_storage::StorageResult<khive_storage::Page<khive_storage::Event>> {
4061            unimplemented!("not exercised by the checkpoint lifecycle-event tests")
4062        }
4063
4064        async fn count_events(
4065            &self,
4066            _filter: khive_storage::EventFilter,
4067        ) -> khive_storage::StorageResult<u64> {
4068            Ok(self.events.lock().unwrap().len() as u64)
4069        }
4070    }
4071
4072    /// Pure decision-table coverage for every input combination
4073    /// `checkpoint_outcome_should_emit` can see: a first elevated tick, a
4074    /// sustained elevated tick, the single drain row, and the ordinary
4075    /// healthy tick that must emit nothing.
4076    #[test]
4077    fn checkpoint_outcome_should_emit_covers_all_transitions() {
4078        assert!(
4079            checkpoint_outcome_should_emit(true, false),
4080            "first elevated tick must emit"
4081        );
4082        assert!(
4083            checkpoint_outcome_should_emit(true, true),
4084            "sustained elevated tick must emit"
4085        );
4086        assert!(
4087            checkpoint_outcome_should_emit(false, true),
4088            "the single drain row (elevated -> healthy) must emit"
4089        );
4090        assert!(
4091            !checkpoint_outcome_should_emit(false, false),
4092            "an ordinary below-warn tick must not emit"
4093        );
4094    }
4095
4096    #[tokio::test]
4097    #[serial(checkpoint_skip_metrics)]
4098    async fn checkpoint_task_emits_outcome_events_while_elevated_and_stops_after_drain() {
4099        let dir = tempfile::tempdir().unwrap();
4100        let path = dir.path().join("outcome_emit.db");
4101        let pool = file_pool(&path);
4102
4103        // warn_pages: 0 means any observed WAL page count (even 0) is
4104        // "elevated" for the duration this config is active.
4105        let cfg = CheckpointConfig {
4106            interval: Duration::from_millis(10),
4107            warn_pages: 0,
4108            ..CheckpointConfig::default()
4109        };
4110        let store = Arc::new(FakeEventStore::default());
4111        let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
4112
4113        let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4114        let handle = tokio::spawn(run_checkpoint_task(
4115            pool,
4116            cfg,
4117            Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
4118            shutdown_rx,
4119            true,
4120        ));
4121
4122        // Poll for the first emitted event instead of a fixed sleep (same
4123        // slowdown-flake class as the stale-sweep test above).
4124        let emitted = wait_for(Duration::from_secs(10), || {
4125            !store.events.lock().unwrap().is_empty()
4126        })
4127        .await;
4128        shutdown_tx.send(()).expect("send shutdown signal");
4129        tokio::time::timeout(Duration::from_secs(1), handle)
4130            .await
4131            .expect("checkpoint task should exit within 1s")
4132            .expect("checkpoint task panicked");
4133
4134        let events = store.events.lock().unwrap();
4135        assert!(
4136            emitted,
4137            "an always-elevated config must append at least one CheckpointOutcomeRecorded event \
4138             within the poll deadline"
4139        );
4140        assert!(
4141            events
4142                .iter()
4143                .all(|e| e.kind == khive_types::EventKind::CheckpointOutcomeRecorded),
4144            "every appended event must be CheckpointOutcomeRecorded, got: {events:?}"
4145        );
4146        assert!(
4147            events.iter().all(|e| e.namespace == "local"),
4148            "events must be stamped with the namespace passed to run_checkpoint_task"
4149        );
4150    }
4151
4152    /// Regression #1434: the event sink's ordinary writer checkout can wait
4153    /// five seconds, but a full lifecycle queue must drop telemetry while the
4154    /// checkpoint scheduler continues through later elevated ticks. The
4155    /// separate in-memory event backend makes that writer contention
4156    /// deterministic without also preventing `checkpoint_once` from
4157    /// observing the file-backed checkpoint pool.
4158    #[tokio::test]
4159    #[serial(checkpoint_skip_metrics)]
4160    async fn checkpoint_cycles_and_task_shutdown_do_not_wait_for_a_contended_lifecycle_writer() {
4161        let dir = tempfile::tempdir().unwrap();
4162        let path = dir.path().join("outcome_contended_sink.db");
4163        let checkpoint_pool = file_pool(&path);
4164
4165        let event_pool = Arc::new(
4166            ConnectionPool::new(PoolConfig {
4167                path: None,
4168                checkout_timeout: Duration::from_secs(5),
4169                write_queue_enabled: false,
4170                ..PoolConfig::default()
4171            })
4172            .expect("event pool"),
4173        );
4174        {
4175            let writer = event_pool.try_writer().expect("initialize event schema");
4176            crate::stores::event::ensure_events_schema(writer.conn())
4177                .expect("initialize event schema");
4178        }
4179        let event_store: Arc<dyn khive_storage::EventStore> =
4180            Arc::new(crate::stores::event::SqlEventStore::new_scoped(
4181                Arc::clone(&event_pool),
4182                false,
4183                "local",
4184            ));
4185        let held_event_writer = event_pool
4186            .try_writer()
4187            .expect("hold the event-store writer");
4188
4189        let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
4190        let subscriber = CaptureSubscriber {
4191            events: std::sync::Arc::clone(&buffer),
4192        };
4193        let _tracing_guard = tracing::subscriber::set_default(subscriber);
4194
4195        let cfg = CheckpointConfig {
4196            interval: Duration::from_millis(10),
4197            warn_pages: 0,
4198            ..CheckpointConfig::default()
4199        };
4200        let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4201        let handle = tokio::spawn(run_checkpoint_task(
4202            checkpoint_pool,
4203            cfg,
4204            Some(CheckpointLifecycleOwner::new(event_store, "local")),
4205            shutdown_rx,
4206            true,
4207        ));
4208
4209        let dropped = wait_for(Duration::from_secs(2), || {
4210            buffer.lock().unwrap().iter().any(|event| {
4211                event.message.as_deref()
4212                    == Some("checkpoint lifecycle event dropped because the append worker is busy")
4213            })
4214        })
4215        .await;
4216        assert!(
4217            dropped,
4218            "later elevated ticks must reach the non-blocking enqueue while the first append is \
4219             still waiting for the held event writer; got: {:?}",
4220            buffer.lock().unwrap()
4221        );
4222
4223        shutdown_tx.send(()).expect("send shutdown signal");
4224        tokio::time::timeout(Duration::from_secs(1), handle)
4225            .await
4226            .expect(
4227                "the run_checkpoint_task handle must not wait for the event store's \
4228                 five-second writer checkout",
4229            )
4230            .expect("checkpoint task panicked");
4231
4232        // The bound above is deliberately checkpoint-task-local. Aborting the
4233        // lifecycle worker cannot cancel the `spawn_blocking` checkout already
4234        // admitted by `SqlEventStore`; release its fixture contention only
4235        // after the `run_checkpoint_task` handle has returned.
4236        drop(held_event_writer);
4237    }
4238
4239    /// A sink error is observable, and the worker remains alive to accept a
4240    /// later checkpoint outcome instead of terminating the scheduler.
4241    #[tokio::test]
4242    #[serial(checkpoint_skip_metrics)]
4243    async fn checkpoint_task_continues_after_lifecycle_append_failure() {
4244        let dir = tempfile::tempdir().unwrap();
4245        let path = dir.path().join("outcome_failing_sink.db");
4246        let pool = file_pool(&path);
4247        let store = Arc::new(FakeEventStore::failing());
4248        let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
4249
4250        let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
4251        let subscriber = CaptureSubscriber {
4252            events: std::sync::Arc::clone(&buffer),
4253        };
4254        let _tracing_guard = tracing::subscriber::set_default(subscriber);
4255
4256        let cfg = CheckpointConfig {
4257            interval: Duration::from_millis(10),
4258            warn_pages: 0,
4259            ..CheckpointConfig::default()
4260        };
4261        let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4262        let handle = tokio::spawn(run_checkpoint_task(
4263            pool,
4264            cfg,
4265            Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
4266            shutdown_rx,
4267            true,
4268        ));
4269
4270        let retried = wait_for(Duration::from_secs(2), || {
4271            store
4272                .append_attempts
4273                .load(std::sync::atomic::Ordering::Relaxed)
4274                >= 2
4275        })
4276        .await;
4277        shutdown_tx.send(()).expect("send shutdown signal");
4278        tokio::time::timeout(Duration::from_secs(1), handle)
4279            .await
4280            .expect("checkpoint task should remain responsive after sink failure")
4281            .expect("checkpoint task panicked");
4282
4283        assert!(
4284            retried,
4285            "a failed append must not terminate the worker or checkpoint task"
4286        );
4287        let captured = buffer.lock().unwrap().clone();
4288        assert!(
4289            captured.iter().any(|event| event.message.as_deref()
4290                == Some("checkpoint lifecycle event append failed")),
4291            "lifecycle append failures must remain observable; got: {:?}",
4292            captured
4293        );
4294    }
4295
4296    #[tokio::test]
4297    #[serial(checkpoint_skip_metrics)]
4298    async fn secondary_checkpoint_task_with_lifecycle_ownership_emits_outcome_events() {
4299        let dir = tempfile::tempdir().unwrap();
4300        let path = dir.path().join("secondary_outcome.db");
4301        let pool = file_pool(&path);
4302        let cfg = CheckpointConfig {
4303            interval: Duration::from_millis(10),
4304            warn_pages: 0,
4305            ..CheckpointConfig::default()
4306        };
4307        let store = Arc::new(FakeEventStore::default());
4308        let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
4309
4310        let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4311        let handle = tokio::spawn(run_checkpoint_task(
4312            pool,
4313            cfg,
4314            Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
4315            shutdown_rx,
4316            false,
4317        ));
4318
4319        // Poll for the first emitted event instead of a fixed sleep (same
4320        // slowdown-flake class as the stale-sweep test above).
4321        let emitted = wait_for(Duration::from_secs(10), || {
4322            !store.events.lock().unwrap().is_empty()
4323        })
4324        .await;
4325        shutdown_tx.send(()).expect("send shutdown signal");
4326        tokio::time::timeout(Duration::from_secs(1), handle)
4327            .await
4328            .expect("checkpoint task should exit within 1s")
4329            .expect("checkpoint task panicked");
4330
4331        assert!(
4332            emitted,
4333            "a designated secondary lifecycle owner must append outcome events within the poll \
4334             deadline"
4335        );
4336    }
4337
4338    #[tokio::test]
4339    #[serial(checkpoint_skip_metrics)]
4340    async fn checkpoint_task_emits_nothing_while_healthy() {
4341        let dir = tempfile::tempdir().unwrap();
4342        let path = dir.path().join("outcome_no_emit.db");
4343        let pool = file_pool(&path);
4344
4345        // An unreachable warn_pages threshold for this test's tiny WAL: every
4346        // tick stays below warn, so no event should ever be appended.
4347        let cfg = CheckpointConfig {
4348            interval: Duration::from_millis(10),
4349            warn_pages: u64::MAX,
4350            ..CheckpointConfig::default()
4351        };
4352        let store = Arc::new(FakeEventStore::default());
4353        let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
4354
4355        let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4356        let handle = tokio::spawn(run_checkpoint_task(
4357            pool,
4358            cfg,
4359            Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
4360            shutdown_rx,
4361            true,
4362        ));
4363
4364        tokio::time::sleep(Duration::from_millis(60)).await;
4365        shutdown_tx.send(()).expect("send shutdown signal");
4366        tokio::time::timeout(Duration::from_secs(1), handle)
4367            .await
4368            .expect("checkpoint task should exit within 1s")
4369            .expect("checkpoint task panicked");
4370
4371        assert!(
4372            store.events.lock().unwrap().is_empty(),
4373            "a config that never crosses warn_pages must never append a lifecycle event"
4374        );
4375    }
4376
4377    #[tokio::test]
4378    #[serial(checkpoint_skip_metrics)]
4379    async fn checkpoint_task_with_no_event_store_does_not_panic() {
4380        let dir = tempfile::tempdir().unwrap();
4381        let path = dir.path().join("outcome_none_store.db");
4382        let pool = file_pool(&path);
4383
4384        let cfg = CheckpointConfig {
4385            interval: Duration::from_millis(10),
4386            warn_pages: 0,
4387            ..CheckpointConfig::default()
4388        };
4389
4390        let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4391        let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
4392
4393        tokio::time::sleep(Duration::from_millis(40)).await;
4394        shutdown_tx.send(()).expect("send shutdown signal");
4395        tokio::time::timeout(Duration::from_secs(1), handle)
4396            .await
4397            .expect("checkpoint task should exit within 1s")
4398            .expect("checkpoint task panicked");
4399    }
4400
4401    // Fix: task-level regressions
4402    // that actually spawn `run_checkpoint_task` and capture its `tracing`
4403    // output, so the wiring at the `tx_age_state.observe(...)` call site
4404    // itself is under test — the pure `TxAgeSweepState` unit tests above
4405    // stay green even if that call site is deleted; these do not. All three
4406    // share `#[serial(tx_registry, checkpoint_skip_metrics)]`: `tx_registry`
4407    // because they read the process-wide registry singleton (see the
4408    // `log_tx_registry_oldest_debug_reports_oldest_open_entry` doc comment
4409    // above for why other tests in this same binary can transiently touch
4410    // it too), `checkpoint_skip_metrics` because they spawn the real task
4411    // that updates the module's skip-tracking atomics.
4412
4413    /// (1) A stale labeled entry with a healthy WAL: the spawned task itself
4414    /// must sweep and escalate it to `Stale`, with WAL-pressure thresholds
4415    /// set unreachably high so only the age sweep — never the WAL-pressure
4416    /// ladder — could be responsible for the captured emission.
4417    #[tokio::test]
4418    #[serial(tx_registry, checkpoint_skip_metrics)]
4419    async fn checkpoint_task_sweeps_stale_registry_entry_while_wal_is_healthy() {
4420        let dir = tempfile::tempdir().unwrap();
4421        let path = dir.path().join("tx_age_sweep_task_healthy_wal.db");
4422        let pool = file_pool(&path);
4423
4424        let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
4425        let subscriber = CaptureSubscriber {
4426            events: std::sync::Arc::clone(&buffer),
4427        };
4428        let _tracing_guard = tracing::subscriber::set_default(subscriber);
4429
4430        let _tx_handle = khive_storage::tx_registry::register(Some(
4431            "checkpoint_task_healthy_wal_sweep_test".to_string(),
4432        ));
4433
4434        let cfg = CheckpointConfig {
4435            interval: Duration::from_millis(10),
4436            warn_pages: u64::MAX,
4437            high_water_pages: u64::MAX,
4438            truncate_high_water_pages: u64::MAX,
4439            tx_warn_secs: Duration::from_millis(1),
4440            tx_max_age_secs: Duration::from_millis(1),
4441            ..CheckpointConfig::default()
4442        };
4443
4444        let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4445        let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
4446
4447        // Poll for the sweep instead of a fixed sleep: a fixed wall-clock
4448        // budget assumes the spawned task completes its tick (registry scan
4449        // + walpin sidecar heartbeat) within that window, which widens and
4450        // flakes under slowdown (coverage instrumentation, contended CI
4451        // runners) — see `wait_for`'s doc comment for the same reasoning
4452        // applied to the sibling walpin tests.
4453        let swept = wait_for(Duration::from_secs(10), || {
4454            buffer.lock().unwrap().iter().any(|e| {
4455                e.tx_label.as_deref() == Some("checkpoint_task_healthy_wal_sweep_test")
4456                    && e.message
4457                        .as_deref()
4458                        .is_some_and(|m| m.contains("stale-op cap"))
4459            })
4460        })
4461        .await;
4462
4463        shutdown_tx.send(()).expect("send shutdown signal");
4464        tokio::time::timeout(Duration::from_secs(1), handle)
4465            .await
4466            .expect("checkpoint task should exit within 1s")
4467            .expect("checkpoint task panicked");
4468
4469        drop(_tx_handle);
4470
4471        let events = buffer.lock().unwrap();
4472        assert!(
4473            swept,
4474            "expected the spawned task to sweep and escalate the stale registry entry \
4475             to Stale on its own within the poll deadline, got: {events:?}"
4476        );
4477    }
4478
4479    /// (2) An empty registry must never produce a Plank 1 age emission from
4480    /// the real spawned task, mirroring the pure
4481    /// `tx_age_sweep_empty_registry_emits_nothing` unit test above.
4482    #[tokio::test]
4483    #[serial(tx_registry, checkpoint_skip_metrics)]
4484    async fn checkpoint_task_emits_no_age_alert_for_an_empty_registry() {
4485        let dir = tempfile::tempdir().unwrap();
4486        let path = dir.path().join("tx_age_sweep_task_empty_registry.db");
4487        let pool = file_pool(&path);
4488
4489        let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
4490        let subscriber = CaptureSubscriber {
4491            events: std::sync::Arc::clone(&buffer),
4492        };
4493        let _tracing_guard = tracing::subscriber::set_default(subscriber);
4494
4495        let cfg = CheckpointConfig {
4496            interval: Duration::from_millis(10),
4497            tx_warn_secs: Duration::from_millis(1),
4498            tx_max_age_secs: Duration::from_millis(1),
4499            ..CheckpointConfig::default()
4500        };
4501
4502        let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4503        let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
4504
4505        tokio::time::sleep(Duration::from_millis(40)).await;
4506        shutdown_tx.send(()).expect("send shutdown signal");
4507        tokio::time::timeout(Duration::from_secs(1), handle)
4508            .await
4509            .expect("checkpoint task should exit within 1s")
4510            .expect("checkpoint task panicked");
4511
4512        let events = buffer.lock().unwrap();
4513        assert!(
4514            events.iter().all(|e| e
4515                .message
4516                .as_deref()
4517                .is_none_or(|m| !m.contains("ADR-091 Plank 1"))),
4518            "an empty registry must never produce a Plank 1 age emission, got: {events:?}"
4519        );
4520    }
4521
4522    /// (3) High-finding regression: a writer-busy tick must NOT silence the
4523    /// age sweep. Holds the pool's writer mutex (via `pool.try_writer()`,
4524    /// never released for the task's entire run) across several checkpoint
4525    /// intervals alongside a stale registered entry, and asserts the age
4526    /// alert still fires even though `checkpoint_once` observes
4527    /// `CheckpointTick::Skipped` on every single tick. Before the fix, the
4528    /// sweep call sat after the `Skipped` early-continue and never ran here.
4529    #[tokio::test]
4530    #[serial(tx_registry, checkpoint_skip_metrics)]
4531    async fn checkpoint_task_sweeps_stale_entry_even_when_writer_is_busy_every_tick() {
4532        reset_checkpoint_metrics_for_tests();
4533
4534        let dir = tempfile::tempdir().unwrap();
4535        let path = dir.path().join("tx_age_sweep_task_writer_busy.db");
4536        let pool = file_pool(&path);
4537        {
4538            let writer = pool.try_writer().unwrap();
4539            writer
4540                .conn()
4541                .execute_batch("CREATE TABLE IF NOT EXISTS t (x INTEGER);")
4542                .unwrap();
4543        }
4544
4545        let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
4546        let subscriber = CaptureSubscriber {
4547            events: std::sync::Arc::clone(&buffer),
4548        };
4549        let _tracing_guard = tracing::subscriber::set_default(subscriber);
4550
4551        let _tx_handle = khive_storage::tx_registry::register(Some(
4552            "checkpoint_task_writer_busy_sweep_test".to_string(),
4553        ));
4554
4555        // Held for the checkpoint task's entire run, acquired BEFORE spawn
4556        // (and with no `.await` in between) so the task cannot possibly
4557        // observe a free writer on any tick.
4558        let _writer_guard = pool.try_writer().expect("acquire writer for busy hold");
4559
4560        let cfg = CheckpointConfig {
4561            interval: Duration::from_millis(10),
4562            tx_warn_secs: Duration::from_millis(1),
4563            tx_max_age_secs: Duration::from_millis(1),
4564            ..CheckpointConfig::default()
4565        };
4566
4567        let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4568        let handle = tokio::spawn(run_checkpoint_task(
4569            Arc::clone(&pool),
4570            cfg,
4571            None,
4572            shutdown_rx,
4573            true,
4574        ));
4575
4576        // Wait until the task has actually recorded a writer-busy Skipped
4577        // tick rather than sleeping a fixed real-time budget: each tick also
4578        // does registry queries and sidecar filesystem writes, so under
4579        // instrumented (coverage) or loaded runners a fixed sleep races the
4580        // first completed tick. Bounded, fail-loud.
4581        let deadline = tokio::time::Instant::now() + Duration::from_secs(10);
4582        while checkpoint_skipped_ticks() == 0 {
4583            assert!(
4584                tokio::time::Instant::now() < deadline,
4585                "test setup must actually drive at least one Skipped tick for this \
4586                 regression to mean anything (none within 10s)"
4587            );
4588            tokio::time::sleep(Duration::from_millis(10)).await;
4589        }
4590
4591        loop {
4592            let events = buffer.lock().unwrap().clone();
4593            if events.iter().any(|e| {
4594                e.tx_label.as_deref() == Some("checkpoint_task_writer_busy_sweep_test")
4595                    && e.message
4596                        .as_deref()
4597                        .is_some_and(|m| m.contains("stale-op cap"))
4598            }) {
4599                break;
4600            }
4601            assert!(
4602                tokio::time::Instant::now() < deadline,
4603                "expected the age sweep to fire even though every tick's writer checkout \
4604                 was skipped within 10s, got: {events:?}"
4605            );
4606            tokio::time::sleep(Duration::from_millis(10)).await;
4607        }
4608
4609        shutdown_tx.send(()).expect("send shutdown signal");
4610        tokio::time::timeout(Duration::from_secs(1), handle)
4611            .await
4612            .expect("checkpoint task should exit within 1s")
4613            .expect("checkpoint task panicked");
4614
4615        drop(_writer_guard);
4616        drop(_tx_handle);
4617    }
4618
4619    // ── ADR-091 Amendment 2: Plank A (session sweep), Plank B (walpin
4620    // sidecar), Plank C (pin-depth probe) ────────────────────────────────
4621
4622    #[tokio::test]
4623    async fn session_sweep_task_exits_on_shutdown_signal() {
4624        let cfg = SessionSweepConfig {
4625            interval: Duration::from_millis(10),
4626            ..SessionSweepConfig::default()
4627        };
4628        let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4629        let handle = tokio::spawn(run_session_sweep_task(Vec::new(), cfg, shutdown_rx));
4630
4631        shutdown_tx.send(()).expect("send shutdown signal");
4632
4633        tokio::time::timeout(Duration::from_secs(1), handle)
4634            .await
4635            .expect("session sweep task should exit within 1s")
4636            .expect("session sweep task panicked");
4637    }
4638
4639    /// Bounded condition poll for filesystem effects of the async sweep
4640    /// task — fixed sleeps flake under parallel test load because sidecar
4641    /// writes fsync.
4642    async fn wait_for(deadline: Duration, mut cond: impl FnMut() -> bool) -> bool {
4643        let start = std::time::Instant::now();
4644        while start.elapsed() < deadline {
4645            if cond() {
4646                return true;
4647            }
4648            tokio::time::sleep(Duration::from_millis(5)).await;
4649        }
4650        cond()
4651    }
4652
4653    #[tokio::test]
4654    #[serial(khive_walpin_sidecar_env)]
4655    async fn walpin_observe_drops_beacon_when_heartbeat_write_fails() {
4656        let dir = tempfile::tempdir().unwrap();
4657        let db_path = dir.path().join("observe_gate.db");
4658        let sidecar_dir = crate::walpin::sidecar_dir_for(&db_path);
4659        let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
4660        std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
4661
4662        let mut state = WalpinSidecarState::new(
4663            Some(db_path.as_path()),
4664            true,
4665            "session",
4666            Duration::from_millis(500),
4667        )
4668        .expect("sidecar enabled for a file-backed path");
4669        state.register_beacon().await;
4670        let pid = std::process::id();
4671        let beacon_path = sidecar_dir.join(format!("{pid}.beacon"));
4672        let before = std::fs::metadata(&beacon_path)
4673            .expect("beacon registered")
4674            .modified()
4675            .unwrap();
4676
4677        // Force the heartbeat write to fail without touching directory
4678        // permissions (which would confound with the dir-mode validation):
4679        // occupy the exclusive-create temp name with a directory, so the
4680        // tolerant unlink and the O_EXCL create both fail.
4681        let obstruction = sidecar_dir.join(format!(".{pid}.json.tmp"));
4682        std::fs::create_dir(&obstruction).unwrap();
4683
4684        tokio::time::sleep(Duration::from_millis(20)).await;
4685        let over_threshold = Some(khive_storage::tx_registry::OldestSpan {
4686            id: khive_storage::tx_registry::TxId(1),
4687            age: Duration::from_secs(60),
4688            label: None,
4689            origin: khive_storage::tx_registry::TxOrigin::Unscoped,
4690        });
4691        state
4692            .observe(over_threshold.clone(), Duration::from_secs(30))
4693            .await;
4694
4695        assert!(
4696            !sidecar_dir.join(format!("{pid}.json")).exists(),
4697            "heartbeat write must have failed"
4698        );
4699        // Skipping the refresh alone would leave `before` fresh inside the
4700        // three-tick window; the fail-closed contract removes the beacon.
4701        assert!(
4702            !beacon_path.exists(),
4703            "a failed heartbeat write must remove the beacon — a still-fresh \
4704             beacon with no heartbeat would classify registered-silent \
4705             (before-mtime {before:?})"
4706        );
4707
4708        // Recovery: clear the obstruction; the next over-threshold tick
4709        // writes the heartbeat and re-registers the beacon.
4710        std::fs::remove_dir(&obstruction).unwrap();
4711        state.observe(over_threshold, Duration::from_secs(30)).await;
4712        assert!(
4713            sidecar_dir.join(format!("{pid}.json")).exists(),
4714            "heartbeat must land once the write path recovers"
4715        );
4716        assert!(
4717            beacon_path.exists(),
4718            "beacon must re-register on the first healthy tick after removal"
4719        );
4720    }
4721
4722    #[tokio::test]
4723    #[serial(khive_walpin_sidecar_env)]
4724    async fn walpin_observe_touches_mtime_without_rewriting_body_when_content_unchanged() {
4725        let dir = tempfile::tempdir().unwrap();
4726        let db_path = dir.path().join("observe_touch.db");
4727        let sidecar_dir = crate::walpin::sidecar_dir_for(&db_path);
4728        let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
4729        std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
4730
4731        let mut state = WalpinSidecarState::new(
4732            Some(db_path.as_path()),
4733            true,
4734            "session",
4735            Duration::from_millis(500),
4736        )
4737        .expect("sidecar enabled for a file-backed path");
4738        let pid = std::process::id();
4739        let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
4740        let span = khive_storage::tx_registry::OldestSpan {
4741            id: khive_storage::tx_registry::TxId(1),
4742            age: Duration::from_secs(60),
4743            label: None,
4744            origin: khive_storage::tx_registry::TxOrigin::Unscoped,
4745        };
4746
4747        state
4748            .observe(Some(span.clone()), Duration::from_secs(30))
4749            .await;
4750        let body_after_create = std::fs::read(&heartbeat_path).expect("heartbeat written");
4751
4752        // Backdate the mtime so the touch is unambiguous: if `observe`
4753        // rewrote the body instead of touching it, the write would also
4754        // reset the mtime, making this assertion pass for the wrong reason —
4755        // the body-byte comparison below is what actually distinguishes
4756        // touch from rewrite.
4757        let backdated = std::time::SystemTime::now() - Duration::from_secs(120);
4758        std::fs::File::open(&heartbeat_path)
4759            .unwrap()
4760            .set_modified(backdated)
4761            .unwrap();
4762
4763        state.observe(Some(span), Duration::from_secs(30)).await;
4764
4765        let body_after_second_observe =
4766            std::fs::read(&heartbeat_path).expect("heartbeat still present");
4767        assert_eq!(
4768            body_after_create, body_after_second_observe,
4769            "unchanged oldest-span identity/label/attribution/cadence must touch mtime, \
4770             not rewrite the body"
4771        );
4772        let mtime_after = std::fs::metadata(&heartbeat_path)
4773            .unwrap()
4774            .modified()
4775            .unwrap();
4776        assert!(
4777            mtime_after > backdated,
4778            "the touch must advance mtime past the backdated value"
4779        );
4780    }
4781
4782    #[tokio::test]
4783    #[serial(khive_walpin_sidecar_env)]
4784    async fn walpin_observe_recreates_heartbeat_after_it_is_deleted_while_span_still_live() {
4785        let dir = tempfile::tempdir().unwrap();
4786        let db_path = dir.path().join("observe_recreate.db");
4787        let sidecar_dir = crate::walpin::sidecar_dir_for(&db_path);
4788        let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
4789        std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
4790
4791        let mut state = WalpinSidecarState::new(
4792            Some(db_path.as_path()),
4793            true,
4794            "session",
4795            Duration::from_millis(500),
4796        )
4797        .expect("sidecar enabled for a file-backed path");
4798        let pid = std::process::id();
4799        let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
4800        let span = khive_storage::tx_registry::OldestSpan {
4801            id: khive_storage::tx_registry::TxId(1),
4802            age: Duration::from_secs(60),
4803            label: None,
4804            origin: khive_storage::tx_registry::TxOrigin::Unscoped,
4805        };
4806
4807        state
4808            .observe(Some(span.clone()), Duration::from_secs(30))
4809            .await;
4810        assert!(heartbeat_path.exists(), "heartbeat written on first tick");
4811
4812        // Simulate enumeration deleting a slow writer's heartbeat while its
4813        // span is still live: the next tick still sees unchanged content
4814        // (same span, same label, same attribution, same cadence) so it
4815        // takes the touch path — which must detect the missing target and
4816        // fall through to a full recreate rather than silently no-op.
4817        std::fs::remove_file(&heartbeat_path).unwrap();
4818        assert!(!heartbeat_path.exists());
4819
4820        state.observe(Some(span), Duration::from_secs(30)).await;
4821
4822        assert!(
4823            heartbeat_path.exists(),
4824            "a touch failure against a deleted heartbeat must recreate it via a full write"
4825        );
4826        let recreated: crate::walpin::WalpinHeartbeat =
4827            serde_json::from_slice(&std::fs::read(&heartbeat_path).unwrap()).unwrap();
4828        assert_eq!(recreated.pid, pid);
4829        assert_eq!(recreated.oldest_tx_age_secs, 60.0);
4830    }
4831
4832    #[tokio::test]
4833    #[serial(tx_registry, khive_walpin_sidecar_env)]
4834    async fn session_sweep_task_writes_and_clears_walpin_heartbeat() {
4835        let dir = tempfile::tempdir().unwrap();
4836        let db_path = dir.path().join("session_sweep.db");
4837        let pool = file_pool(&db_path);
4838        let sidecar_dir =
4839            crate::walpin::sidecar_dir_for(pool.canonical_path().expect("file-backed pool"));
4840        let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
4841        std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
4842
4843        let cfg = SessionSweepConfig {
4844            interval: Duration::from_millis(10),
4845            tx_warn_secs: Duration::from_millis(20),
4846            tx_max_age_secs: Duration::from_millis(500),
4847        };
4848        let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4849        let handle = tokio::spawn(run_session_sweep_task(
4850            vec![SweepBackend {
4851                pool: Arc::clone(&pool),
4852                is_main: true,
4853            }],
4854            cfg,
4855            shutdown_rx,
4856        ));
4857
4858        // No open span yet: a quiet process must write no *heartbeat*, but
4859        // it DOES register its one-time beacon at startup (ADR-091
4860        // Amendment 2 sidecar-health attribution) — the sidecar dir is not
4861        // empty, only heartbeat-free. Poll-wait rather than a fixed sleep:
4862        // the first tick fsyncs the beacon, and under parallel test load
4863        // that write can take longer than any small fixed window.
4864        let pid = std::process::id();
4865        let beacon = crate::walpin::beacon_path(&sidecar_dir, pid);
4866        assert!(
4867            wait_for(Duration::from_secs(2), || beacon.exists()).await,
4868            "a quiet process must still register its one-time beacon"
4869        );
4870        assert!(
4871            !sidecar_dir.join(format!("{pid}.json")).exists(),
4872            "a quiet process must not write a walpin heartbeat"
4873        );
4874
4875        let tx_handle =
4876            khive_storage::tx_registry::register(Some("session_sweep_walpin_test".to_string()));
4877        let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
4878        assert!(
4879            wait_for(Duration::from_secs(2), || heartbeat_path.exists()).await,
4880            "expected a walpin heartbeat once the span crossed tx_warn_secs"
4881        );
4882        let body = std::fs::read_to_string(&heartbeat_path).unwrap();
4883        let hb: crate::walpin::WalpinHeartbeat = serde_json::from_str(&body).unwrap();
4884        assert_eq!(hb.pid, pid);
4885        assert_eq!(hb.process_role, "session");
4886        assert_eq!(
4887            hb.oldest_tx_label.as_deref(),
4888            Some("session_sweep_walpin_test")
4889        );
4890        assert_eq!(
4891            hb.attribution_basis.as_deref(),
4892            Some("fallback"),
4893            "an Unscoped span observed only through the main view's fallback \
4894             must carry attribution_basis=\"fallback\", never \"origin\""
4895        );
4896
4897        drop(tx_handle);
4898        assert!(
4899            wait_for(Duration::from_secs(2), || !heartbeat_path.exists()).await,
4900            "heartbeat must be removed once the stale span clears"
4901        );
4902
4903        shutdown_tx.send(()).expect("send shutdown signal");
4904        tokio::time::timeout(Duration::from_secs(1), handle)
4905            .await
4906            .expect("session sweep task should exit within 1s")
4907            .expect("session sweep task panicked");
4908    }
4909
4910    /// ADR-091 Amendment 3 fan-out: two file-backed pools in one process,
4911    /// each its own `SweepBackend`. A span scoped to the SECONDARY pool's
4912    /// own origin must produce a heartbeat only in the secondary's sidecar
4913    /// — never the main backend's — and, because a `Secondary` filter never
4914    /// falls back to `Unscoped`, its heartbeat carries the evidence-backed
4915    /// `attribution_basis="origin"`. Uses the `graph_traverse_read` label
4916    /// (`stores/graph.rs`'s `traverse`) — the design note's own example of
4917    /// "the most WAL-pin-relevant span in the store" — as the registered
4918    /// span's label, so this doubles as coverage that a traversal read span
4919    /// surfaces correctly in a secondary backend's filtered view.
4920    #[tokio::test]
4921    #[serial(tx_registry, khive_walpin_sidecar_env)]
4922    async fn session_sweep_fan_out_scopes_secondary_span_to_secondary_sidecar_only() {
4923        let main_dir = tempfile::tempdir().unwrap();
4924        let secondary_dir = tempfile::tempdir().unwrap();
4925        let main_pool = file_pool(&main_dir.path().join("main.db"));
4926        let secondary_pool = file_pool(&secondary_dir.path().join("secondary.db"));
4927        let main_sidecar =
4928            crate::walpin::sidecar_dir_for(main_pool.canonical_path().expect("file-backed"));
4929        let secondary_sidecar =
4930            crate::walpin::sidecar_dir_for(secondary_pool.canonical_path().expect("file-backed"));
4931        let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
4932        std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
4933
4934        let cfg = SessionSweepConfig {
4935            interval: Duration::from_millis(10),
4936            tx_warn_secs: Duration::from_millis(20),
4937            tx_max_age_secs: Duration::from_millis(500),
4938        };
4939        let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
4940        let handle = tokio::spawn(run_session_sweep_task(
4941            vec![
4942                SweepBackend {
4943                    pool: Arc::clone(&main_pool),
4944                    is_main: true,
4945                },
4946                SweepBackend {
4947                    pool: Arc::clone(&secondary_pool),
4948                    is_main: false,
4949                },
4950            ],
4951            cfg,
4952            shutdown_rx,
4953        ));
4954
4955        let pid = std::process::id();
4956        let secondary_heartbeat = secondary_sidecar.join(format!("{pid}.json"));
4957        let main_heartbeat = main_sidecar.join(format!("{pid}.json"));
4958
4959        let tx_handle = khive_storage::tx_registry::register_scoped(
4960            Some("graph_traverse_read".to_string()),
4961            secondary_pool.origin(),
4962        );
4963        assert!(
4964            wait_for(Duration::from_secs(2), || secondary_heartbeat.exists()).await,
4965            "expected a walpin heartbeat in the secondary backend's own sidecar"
4966        );
4967        assert!(
4968            !main_heartbeat.exists(),
4969            "a span scoped to the secondary backend's origin must never produce \
4970             a heartbeat in the main backend's sidecar"
4971        );
4972
4973        let body = std::fs::read_to_string(&secondary_heartbeat).unwrap();
4974        let hb: crate::walpin::WalpinHeartbeat = serde_json::from_str(&body).unwrap();
4975        assert_eq!(hb.oldest_tx_label.as_deref(), Some("graph_traverse_read"));
4976        assert_eq!(
4977            hb.attribution_basis.as_deref(),
4978            Some("origin"),
4979            "a Secondary-view winner is always Database-origin-backed — never fallback"
4980        );
4981
4982        drop(tx_handle);
4983        assert!(
4984            wait_for(Duration::from_secs(2), || !secondary_heartbeat.exists()).await,
4985            "secondary heartbeat must be removed once its span clears"
4986        );
4987        assert!(
4988            !main_heartbeat.exists(),
4989            "the main sidecar must have stayed untouched for the whole tick sequence"
4990        );
4991
4992        shutdown_tx.send(()).expect("send shutdown signal");
4993        tokio::time::timeout(Duration::from_secs(1), handle)
4994            .await
4995            .expect("session sweep task should exit within 1s")
4996            .expect("session sweep task panicked");
4997    }
4998
4999    /// ADR-091 Amendment 3: a `run_checkpoint_task` instance for backend A
5000    /// (`is_main: false`, a `Secondary` filter scoped to A's own identity)
5001    /// must never observe a span registered against a DIFFERENT backend's
5002    /// `Database` origin, nor an `Unscoped` span — a `Secondary` filter never
5003    /// falls back to `Unscoped` (that fallback is the main view's alone).
5004    /// Drives the real task for several ticks and asserts neither the
5005    /// captured `tracing` emissions nor backend A's own sidecar ever name
5006    /// either span.
5007    #[tokio::test]
5008    #[serial(tx_registry, checkpoint_skip_metrics, khive_walpin_sidecar_env)]
5009    async fn checkpoint_task_ignores_span_registered_against_other_backend_origin_and_unscoped() {
5010        let dir_a = tempfile::tempdir().unwrap();
5011        let dir_b = tempfile::tempdir().unwrap();
5012        let pool_a = file_pool(&dir_a.path().join("backend_a.db"));
5013        // Only used to mint a real, distinct `DbIdentity` for backend B — no
5014        // checkpoint task is spawned for it.
5015        let pool_b = file_pool(&dir_b.path().join("backend_b.db"));
5016        let sidecar_a =
5017            crate::walpin::sidecar_dir_for(pool_a.canonical_path().expect("file-backed"));
5018        let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
5019        std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
5020
5021        let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5022        let subscriber = CaptureSubscriber {
5023            events: std::sync::Arc::clone(&buffer),
5024        };
5025        let _tracing_guard = tracing::subscriber::set_default(subscriber);
5026
5027        let _b_origin_handle = khive_storage::tx_registry::register_scoped(
5028            Some("b_origin_span_ignored_by_a".to_string()),
5029            pool_b.origin(),
5030        );
5031        let _unscoped_handle = khive_storage::tx_registry::register(Some(
5032            "unscoped_span_ignored_by_secondary".to_string(),
5033        ));
5034
5035        let cfg = CheckpointConfig {
5036            interval: Duration::from_millis(10),
5037            tx_warn_secs: Duration::from_millis(1),
5038            tx_max_age_secs: Duration::from_millis(1),
5039            ..CheckpointConfig::default()
5040        };
5041        let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
5042        let handle = tokio::spawn(run_checkpoint_task(
5043            pool_a,
5044            cfg,
5045            None,
5046            shutdown_rx,
5047            false, // is_main: backend A is a secondary backend here
5048        ));
5049
5050        // No positive condition to poll for — this asserts an absence over a
5051        // bounded run of several ticks, mirroring
5052        // `checkpoint_task_emits_no_age_alert_for_an_empty_registry`'s same
5053        // fixed-window shape (there is nothing to wait-until for a negative).
5054        tokio::time::sleep(Duration::from_millis(60)).await;
5055        shutdown_tx.send(()).expect("send shutdown signal");
5056        tokio::time::timeout(Duration::from_secs(1), handle)
5057            .await
5058            .expect("checkpoint task should exit within 1s")
5059            .expect("checkpoint task panicked");
5060
5061        let events = buffer.lock().unwrap();
5062        assert!(
5063            events.iter().all(|e| {
5064                e.tx_label.as_deref() != Some("b_origin_span_ignored_by_a")
5065                    && e.tx_label.as_deref() != Some("unscoped_span_ignored_by_secondary")
5066            }),
5067            "backend A's Secondary filter must never emit an age alert naming a span \
5068             registered against a different backend's origin or an Unscoped span, got: \
5069             {events:?}"
5070        );
5071        assert!(
5072            !sidecar_a
5073                .join(format!("{}.json", std::process::id()))
5074                .exists(),
5075            "backend A's own sidecar must never gain a heartbeat from a span it does not own"
5076        );
5077    }
5078
5079    /// ADR-091 Amendment 3: a secondary backend's own `run_checkpoint_task`
5080    /// must detect a stall on its OWN backend (never main-only ownership) —
5081    /// both the Plank 1 age-sweep emission and the sidecar heartbeat, with
5082    /// `attribution_basis="origin"` (a `Secondary` filter winner is always
5083    /// `Database`-origin-backed, never the `Unscoped` fallback) and a
5084    /// nonzero reflected age.
5085    #[tokio::test]
5086    #[serial(tx_registry, checkpoint_skip_metrics, khive_walpin_sidecar_env)]
5087    async fn checkpoint_task_detects_and_enumerates_secondary_backend_stall() {
5088        let dir = tempfile::tempdir().unwrap();
5089        let pool = file_pool(&dir.path().join("secondary_stall.db"));
5090        let sidecar_dir =
5091            crate::walpin::sidecar_dir_for(pool.canonical_path().expect("file-backed"));
5092        let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
5093        std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
5094
5095        let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5096        let subscriber = CaptureSubscriber {
5097            events: std::sync::Arc::clone(&buffer),
5098        };
5099        let _tracing_guard = tracing::subscriber::set_default(subscriber);
5100
5101        let tx_handle = khive_storage::tx_registry::register_scoped(
5102            Some("secondary_stall_test".to_string()),
5103            pool.origin(),
5104        );
5105
5106        let cfg = CheckpointConfig {
5107            interval: Duration::from_millis(10),
5108            tx_warn_secs: Duration::from_millis(5),
5109            tx_max_age_secs: Duration::from_millis(500),
5110            ..CheckpointConfig::default()
5111        };
5112        let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
5113        let pid = std::process::id();
5114        let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
5115        let handle = tokio::spawn(run_checkpoint_task(
5116            pool,
5117            cfg,
5118            None,
5119            shutdown_rx,
5120            false, // is_main: this is a secondary backend's own checkpoint task
5121        ));
5122
5123        assert!(
5124            wait_for(Duration::from_secs(2), || heartbeat_path.exists()).await,
5125            "expected a walpin heartbeat once the secondary backend's own span crossed \
5126             tx_warn_secs"
5127        );
5128        let body = std::fs::read_to_string(&heartbeat_path).unwrap();
5129        let hb: crate::walpin::WalpinHeartbeat = serde_json::from_str(&body).unwrap();
5130        assert_eq!(hb.oldest_tx_label.as_deref(), Some("secondary_stall_test"));
5131        assert_eq!(
5132            hb.attribution_basis.as_deref(),
5133            Some("origin"),
5134            "a Secondary-view winner is always Database-origin-backed — never fallback"
5135        );
5136        assert!(
5137            hb.oldest_tx_age_secs > 0.0,
5138            "the heartbeat must reflect a nonzero stale age for the secondary backend's own \
5139             span, got {hb:?}"
5140        );
5141
5142        shutdown_tx.send(()).expect("send shutdown signal");
5143        tokio::time::timeout(Duration::from_secs(1), handle)
5144            .await
5145            .expect("checkpoint task should exit within 1s")
5146            .expect("checkpoint task panicked");
5147
5148        drop(tx_handle);
5149
5150        let events = buffer.lock().unwrap();
5151        assert!(
5152            events.iter().any(|e| {
5153                e.tx_label.as_deref() == Some("secondary_stall_test")
5154                    && e.message
5155                        .as_deref()
5156                        .is_some_and(|m| m.contains("ADR-091 Plank 1"))
5157            }),
5158            "expected the secondary backend's own checkpoint task to emit a Plank 1 age alert \
5159             for its own stalled span, got: {events:?}"
5160        );
5161    }
5162
5163    #[test]
5164    fn wal_pin_depth_arithmetic_against_real_connection() {
5165        let dir = tempfile::tempdir().unwrap();
5166        let path = dir.path().join("pin_depth.db");
5167        let pool = file_pool(&path);
5168        let writer = pool.try_writer().expect("acquire writer");
5169        let conn = writer.conn();
5170
5171        conn.execute_batch("CREATE TABLE t (v INTEGER)").unwrap();
5172        conn.execute_batch("INSERT INTO t (v) VALUES (1)").unwrap();
5173
5174        let (log, checkpointed) =
5175            query_wal_pin_depth(conn).expect("PRAGMA wal_checkpoint(PASSIVE) must succeed");
5176        // Nothing pins the WAL open in this test (no concurrent reader), so a
5177        // PASSIVE checkpoint must fully drain what it just wrote: pin depth
5178        // (log - checkpointed) is zero.
5179        assert!(
5180            log >= checkpointed,
5181            "checkpointed frames cannot exceed log frames"
5182        );
5183        assert_eq!(
5184            log - checkpointed,
5185            0,
5186            "an unpinned WAL must fully checkpoint under PASSIVE"
5187        );
5188    }
5189
5190    #[test]
5191    fn wal_pin_depth_arithmetic_on_in_memory_pool_errors_cleanly() {
5192        // In-memory databases report `log = -1` (no WAL); the pragma read
5193        // itself does not panic and the caller (`log_wal_pin_depth`) treats
5194        // any error as a logged warning, never a crash.
5195        let cfg = PoolConfig {
5196            path: None,
5197            ..PoolConfig::default()
5198        };
5199        let pool = ConnectionPool::new(cfg).expect("in-memory pool");
5200        let writer = pool.try_writer().expect("acquire writer");
5201        // Either an explicit error or a nonsensical negative `log` value is
5202        // acceptable here — the requirement is just "does not panic".
5203        let _ = query_wal_pin_depth(writer.conn());
5204    }
5205}