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;