Skip to main content

khive_db/
checkpoint.rs

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