Skip to main content

khive_db/checkpoint/
run_state.rs

1//! Checkpoint run state, routine observations, timing, and metric accessors.
2
3use super::{
4    BTreeMap, ConnectionPool, Duration, HashMap, Instant, Mutex, OnceLock, Ordering, Path, PathBuf,
5    RawCheckpointObservation, CHECKPOINT_CONSECUTIVE_SKIPS, CHECKPOINT_LAST_SKIP_WAL_PAGES,
6    CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS, CHECKPOINT_LIFECYCLE_APPEND_FAILURES,
7    CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS, CHECKPOINT_PRESSURE_ELEVATED_TICKS,
8    CHECKPOINT_PRESSURE_EPISODES_RECOVERED, CHECKPOINT_PRESSURE_EPISODES_STARTED,
9    CHECKPOINT_SKIPPED_TICKS, LAST_WAL_PAGES, READ_TX_MAX_AGE_EVICTIONS, TRUNCATE_ATTEMPTS,
10    TRUNCATE_CONSECUTIVE_FAILURES,
11};
12
13/// One backend-scoped observation produced by the periodic checkpoint task's
14/// own PASSIVE pass. Logical frame counts and the physical `-wal` allocation
15/// are intentionally separate: SQLite may retain/reuse the sidecar after the
16/// logical backlog drains (#1849).
17#[derive(Debug, Clone, PartialEq, Eq)]
18pub struct RoutineWalObservation {
19    pub busy: i64,
20    pub log_frames: u64,
21    pub checkpointed_frames: u64,
22    pub pending_frames: u64,
23    pub physical_wal_bytes: Option<u64>,
24    pub observed_at_unix_ms: u64,
25}
26
27/// One backend's current checkpoint run: informative checkpoint results that
28/// stopped at the same frame while frames remained to backfill, allowing
29/// bounded neutral busy results between them. Reported by `db_diagnostics` as
30/// `oldest_pinned_frame_run`.
31#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
32pub struct CheckpointRun {
33    /// The checkpointed frame every result in the run stopped at.
34    pub frame: i64,
35    /// Unix time in milliseconds of the run's first checkpoint result.
36    pub first_observed_at_unix_ms: u64,
37}
38
39#[derive(Debug, Clone, Copy, PartialEq, Eq)]
40pub(crate) enum CheckpointRunStatus {
41    NoTask,
42    NoObservation,
43    Observed(CheckpointRun),
44}
45
46#[derive(Debug, Clone, Copy, PartialEq, Eq)]
47pub(super) struct CheckpointRunEntry {
48    pub(super) run: CheckpointRun,
49    first_observed_at: Instant,
50    pub(super) last_log_frames: i64,
51    last_informative_at: Instant,
52    pub(super) busy_since_last_informative: bool,
53}
54
55#[derive(Debug, Default)]
56pub(super) struct CheckpointRunState {
57    active_tasks: usize,
58    pub(super) checkpoint_interval_ms: u64,
59    owner_intervals_ms: BTreeMap<u64, usize>,
60    pub(super) entry: Option<CheckpointRunEntry>,
61}
62
63static CHECKPOINT_RUNS: OnceLock<Mutex<HashMap<Option<PathBuf>, CheckpointRunState>>> =
64    OnceLock::new();
65
66pub(super) fn checkpoint_runs() -> &'static Mutex<HashMap<Option<PathBuf>, CheckpointRunState>> {
67    CHECKPOINT_RUNS.get_or_init(|| Mutex::new(HashMap::new()))
68}
69
70pub(crate) struct CheckpointRunTaskGuard {
71    key: Option<PathBuf>,
72    interval_ms: u64,
73}
74
75impl CheckpointRunTaskGuard {
76    pub(crate) fn start(pool: &ConnectionPool, interval: Duration) -> Self {
77        let key = checkpoint_db_key(pool);
78        let interval_ms = interval.as_millis().min(u128::from(u64::MAX)) as u64;
79        let interval_ms = interval_ms.max(1);
80        let mut runs = checkpoint_runs()
81            .lock()
82            .unwrap_or_else(std::sync::PoisonError::into_inner);
83        let state = runs.entry(key.clone()).or_default();
84        if state.active_tasks == 0 {
85            state.entry = None;
86        }
87        *state.owner_intervals_ms.entry(interval_ms).or_default() += 1;
88        state.active_tasks = state.active_tasks.saturating_add(1);
89        state.checkpoint_interval_ms = *state
90            .owner_intervals_ms
91            .first_key_value()
92            .expect("active owner has an interval")
93            .0;
94        Self { key, interval_ms }
95    }
96}
97
98impl Drop for CheckpointRunTaskGuard {
99    fn drop(&mut self) {
100        let mut runs = checkpoint_runs()
101            .lock()
102            .unwrap_or_else(std::sync::PoisonError::into_inner);
103        let Some(state) = runs.get_mut(&self.key) else {
104            return;
105        };
106        let Some(count) = state.owner_intervals_ms.get_mut(&self.interval_ms) else {
107            return;
108        };
109        *count -= 1;
110        if *count == 0 {
111            state.owner_intervals_ms.remove(&self.interval_ms);
112        }
113        state.active_tasks = state.active_tasks.saturating_sub(1);
114        if state.active_tasks == 0 {
115            runs.remove(&self.key);
116        } else {
117            state.checkpoint_interval_ms = *state
118                .owner_intervals_ms
119                .first_key_value()
120                .expect("surviving owner has an interval")
121                .0;
122        }
123    }
124}
125
126pub(super) fn advance_checkpoint_run_at(
127    entry: &mut Option<CheckpointRunEntry>,
128    observation: Option<(i64, i64, i64)>,
129    observed_at_unix_ms: u64,
130    observed_at: Instant,
131    checkpoint_interval_ms: u64,
132) {
133    let Some((busy, log_frames, checkpointed_frames)) = observation else {
134        *entry = None;
135        return;
136    };
137    // A busy row has no reliable reading of the ceiling, even if SQLite fills
138    // in the other columns. It cannot advance or end the current run.
139    if busy != 0 {
140        if let Some(current) = entry {
141            current.busy_since_last_informative = true;
142        }
143        return;
144    }
145    if log_frames < 0 || checkpointed_frames < 0 || checkpointed_frames >= log_frames {
146        *entry = None;
147        return;
148    }
149
150    match entry {
151        Some(current)
152            if current.run.frame == checkpointed_frames
153                && log_frames >= current.last_log_frames
154                && (!current.busy_since_last_informative
155                    || observed_at
156                        .checked_duration_since(current.last_informative_at)
157                        .is_some_and(|elapsed| {
158                            elapsed
159                                <= Duration::from_millis(checkpoint_interval_ms.saturating_mul(2))
160                        })) =>
161        {
162            current.last_log_frames = log_frames;
163            current.last_informative_at = observed_at;
164            current.busy_since_last_informative = false;
165        }
166        _ => {
167            *entry = Some(CheckpointRunEntry {
168                run: CheckpointRun {
169                    frame: checkpointed_frames,
170                    first_observed_at_unix_ms: observed_at_unix_ms,
171                },
172                first_observed_at: observed_at,
173                last_log_frames: log_frames,
174                last_informative_at: observed_at,
175                busy_since_last_informative: false,
176            });
177        }
178    }
179}
180
181#[cfg(test)]
182pub(super) fn advance_checkpoint_run(
183    entry: &mut Option<CheckpointRunEntry>,
184    observation: Option<(i64, i64, i64)>,
185    observed_at_unix_ms: u64,
186    checkpoint_interval_ms: u64,
187) {
188    static TEST_ORIGIN: OnceLock<Instant> = OnceLock::new();
189    let observed_at = TEST_ORIGIN
190        .get_or_init(Instant::now)
191        .checked_add(Duration::from_millis(observed_at_unix_ms))
192        .expect("test monotonic timestamp");
193    advance_checkpoint_run_at(
194        entry,
195        observation,
196        observed_at_unix_ms,
197        observed_at,
198        checkpoint_interval_ms,
199    );
200}
201
202pub(crate) fn record_checkpoint_run_result(
203    pool: &ConnectionPool,
204    observation: Option<(i64, i64, i64)>,
205) -> CheckpointRunStatus {
206    let mut runs = checkpoint_runs()
207        .lock()
208        .unwrap_or_else(std::sync::PoisonError::into_inner);
209    let Some(state) = runs.get_mut(&checkpoint_db_key(pool)) else {
210        return CheckpointRunStatus::NoTask;
211    };
212    if state.active_tasks == 0 {
213        return CheckpointRunStatus::NoTask;
214    }
215    let checkpoint_interval_ms = state.checkpoint_interval_ms;
216    advance_checkpoint_run_at(
217        &mut state.entry,
218        observation,
219        observed_at_unix_ms(),
220        Instant::now(),
221        checkpoint_interval_ms,
222    );
223    state
224        .entry
225        .map_or(CheckpointRunStatus::NoObservation, |entry| {
226            CheckpointRunStatus::Observed(entry.run)
227        })
228}
229
230#[cfg(test)]
231pub(crate) fn checkpoint_run_status(pool: &ConnectionPool) -> CheckpointRunStatus {
232    checkpoint_run_snapshot(pool).0
233}
234
235/// A run and its monotonic age from the same state snapshot. The Unix timestamp
236/// on `CheckpointRun` is display metadata and must not gate pin detection.
237pub(crate) fn checkpoint_run_snapshot(
238    pool: &ConnectionPool,
239) -> (CheckpointRunStatus, Option<Duration>) {
240    let runs = checkpoint_runs()
241        .lock()
242        .unwrap_or_else(std::sync::PoisonError::into_inner);
243    let Some(state) = runs.get(&checkpoint_db_key(pool)) else {
244        return (CheckpointRunStatus::NoTask, None);
245    };
246    if state.active_tasks == 0 {
247        return (CheckpointRunStatus::NoTask, None);
248    }
249    state
250        .entry
251        .map_or((CheckpointRunStatus::NoObservation, None), |entry| {
252            (
253                CheckpointRunStatus::Observed(entry.run),
254                Some(entry.first_observed_at.elapsed()),
255            )
256        })
257}
258
259/// Process-lifetime totals for actual routine PASSIVE calls on one store.
260/// Ticks with no PASSIVE call and post-TRUNCATE probes are excluded. Busy counts only
261/// SQLite's returned busy flag, never an incomplete checkpoint's pending frames.
262#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
263#[serde(default)]
264pub struct CheckpointTiming {
265    pub ticks: u64,
266    pub elapsed_us_sum: u64,
267    pub elapsed_us_max: u64,
268    pub busy_ticks: u64,
269    pub error_ticks: u64,
270}
271
272static CHECKPOINT_TIMINGS: OnceLock<Mutex<HashMap<Option<PathBuf>, CheckpointTiming>>> =
273    OnceLock::new();
274
275pub(super) fn checkpoint_timings() -> &'static Mutex<HashMap<Option<PathBuf>, CheckpointTiming>> {
276    CHECKPOINT_TIMINGS.get_or_init(|| Mutex::new(HashMap::new()))
277}
278
279pub(super) fn record_checkpoint_timing(pool: &ConnectionPool, elapsed_us: u64, busy: Option<i64>) {
280    let mut timings = checkpoint_timings()
281        .lock()
282        .unwrap_or_else(std::sync::PoisonError::into_inner);
283    let timing = timings.entry(checkpoint_db_key(pool)).or_default();
284    timing.ticks = timing.ticks.saturating_add(1);
285    timing.elapsed_us_sum = timing.elapsed_us_sum.saturating_add(elapsed_us);
286    timing.elapsed_us_max = timing.elapsed_us_max.max(elapsed_us);
287    timing.busy_ticks = timing
288        .busy_ticks
289        .saturating_add(u64::from(busy.is_some_and(|value| value != 0)));
290    timing.error_ticks = timing.error_ticks.saturating_add(u64::from(busy.is_none()));
291}
292
293/// Pure in-memory read; zero means no routine call has been recorded for this key.
294pub fn checkpoint_timing(pool: &ConnectionPool) -> CheckpointTiming {
295    checkpoint_timings()
296        .lock()
297        .unwrap_or_else(std::sync::PoisonError::into_inner)
298        .get(&checkpoint_db_key(pool))
299        .copied()
300        .unwrap_or_default()
301}
302
303/// Latest routine observation by canonical database identity. Checkpoint
304/// tasks fan out per backend, so a single process-global "last task wins"
305/// gauge would misattribute a secondary backend to the main metrics frame.
306static ROUTINE_WAL_OBSERVATIONS: OnceLock<Mutex<HashMap<Option<PathBuf>, RoutineWalObservation>>> =
307    OnceLock::new();
308
309fn routine_wal_observations() -> &'static Mutex<HashMap<Option<PathBuf>, RoutineWalObservation>> {
310    ROUTINE_WAL_OBSERVATIONS.get_or_init(|| Mutex::new(HashMap::new()))
311}
312
313pub(super) fn checkpoint_db_key_from_path(path: Option<&Path>) -> Option<PathBuf> {
314    path.map(Path::to_path_buf)
315}
316
317pub(super) fn checkpoint_db_key(pool: &ConnectionPool) -> Option<PathBuf> {
318    checkpoint_db_key_from_path(pool.canonical_path())
319}
320
321fn observed_at_unix_ms() -> u64 {
322    std::time::SystemTime::now()
323        .duration_since(std::time::UNIX_EPOCH)
324        .map(|duration| duration.as_millis() as u64)
325        .unwrap_or(0)
326}
327
328fn physical_wal_bytes(pool: &ConnectionPool) -> Option<u64> {
329    let path = pool.canonical_path()?;
330    let mut sidecar = path.as_os_str().to_os_string();
331    sidecar.push("-wal");
332    std::fs::metadata(PathBuf::from(sidecar))
333        .ok()
334        .map(|metadata| metadata.len())
335}
336
337pub(super) fn record_routine_wal_observation(
338    pool: &ConnectionPool,
339    raw: RawCheckpointObservation,
340) -> RoutineWalObservation {
341    let log_frames = raw.log_frames.max(0) as u64;
342    let checkpointed_frames = raw.checkpointed_frames.max(0) as u64;
343    let observation = RoutineWalObservation {
344        busy: raw.busy,
345        log_frames,
346        checkpointed_frames,
347        pending_frames: log_frames.saturating_sub(checkpointed_frames),
348        physical_wal_bytes: physical_wal_bytes(pool),
349        observed_at_unix_ms: observed_at_unix_ms(),
350    };
351    routine_wal_observations()
352        .lock()
353        .unwrap_or_else(std::sync::PoisonError::into_inner)
354        .insert(checkpoint_db_key(pool), observation.clone());
355    observation
356}
357
358/// Latest periodic checkpoint sample for this exact backend. This is a pure
359/// in-memory read: it never issues `wal_checkpoint` or stats the filesystem.
360pub fn routine_wal_observation(pool: &ConnectionPool) -> Option<RoutineWalObservation> {
361    routine_wal_observations()
362        .lock()
363        .unwrap_or_else(std::sync::PoisonError::into_inner)
364        .get(&checkpoint_db_key(pool))
365        .cloned()
366}
367
368/// Last-observed WAL page count, if any checkpoint tick has run yet in this
369/// process. Read surface for the daemon-frame metrics snapshot.
370pub fn last_observed_wal_pages() -> Option<u64> {
371    match LAST_WAL_PAGES.load(Ordering::Relaxed) {
372        u64::MAX => None,
373        pages => Some(pages),
374    }
375}
376
377/// Total WAL TRUNCATE attempts made in this process's lifetime.
378pub fn truncate_attempts() -> u64 {
379    TRUNCATE_ATTEMPTS.load(Ordering::Relaxed)
380}
381
382/// Current measured TRUNCATE-failure streak; unmeasured attempts leave it unchanged.
383pub fn truncate_consecutive_failures() -> u64 {
384    TRUNCATE_CONSECUTIVE_FAILURES.load(Ordering::Relaxed)
385}
386
387/// Total checkpoint ticks without an available WAL frame observation in this
388/// process's lifetime (dedicated connection unavailable or SQLite busy).
389pub fn checkpoint_skipped_ticks() -> u64 {
390    CHECKPOINT_SKIPPED_TICKS.load(Ordering::Relaxed)
391}
392
393/// Current consecutive-skip run length; 0 once the next tick is observed.
394pub fn checkpoint_consecutive_skips() -> u64 {
395    CHECKPOINT_CONSECUTIVE_SKIPS.load(Ordering::Relaxed)
396}
397
398/// WAL page count last known at the time of the most recent skip, if any
399/// skip has occurred yet in this process.
400pub fn checkpoint_last_skip_wal_pages() -> Option<u64> {
401    match CHECKPOINT_LAST_SKIP_WAL_PAGES.load(Ordering::Relaxed) {
402        u64::MAX => None,
403        pages => Some(pages),
404    }
405}
406
407/// Total at/above-`warn_pages` observations aggregated in memory.
408pub fn checkpoint_pressure_elevated_ticks() -> u64 {
409    CHECKPOINT_PRESSURE_ELEVATED_TICKS.load(Ordering::Relaxed)
410}
411
412/// Total pressure episodes observed to start in this process.
413pub fn checkpoint_pressure_episodes_started() -> u64 {
414    CHECKPOINT_PRESSURE_EPISODES_STARTED.load(Ordering::Relaxed)
415}
416
417/// Total pressure episodes observed to recover in this process.
418pub fn checkpoint_pressure_episodes_recovered() -> u64 {
419    CHECKPOINT_PRESSURE_EPISODES_RECOVERED.load(Ordering::Relaxed)
420}
421
422/// Total primary-store append calls made by checkpoint lifecycle workers.
423pub fn checkpoint_lifecycle_append_attempts() -> u64 {
424    CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.load(Ordering::Relaxed)
425}
426
427/// Total checkpoint lifecycle append calls that returned a storage error.
428pub fn checkpoint_lifecycle_append_failures() -> u64 {
429    CHECKPOINT_LIFECYCLE_APPEND_FAILURES.load(Ordering::Relaxed)
430}
431
432/// Total cached-reader read transactions rolled back on reuse for exceeding
433/// `read_tx_max_age` (#1846), across this process's lifetime.
434pub fn read_tx_max_age_evictions() -> u64 {
435    READ_TX_MAX_AGE_EVICTIONS.load(Ordering::Relaxed)
436}
437
438/// Records one cached-reader read transaction rolled back on reuse for
439/// exceeding `read_tx_max_age`. Called from `sql_bridge.rs` at the point the
440/// rollback is issued, regardless of whether the rollback itself succeeds —
441/// this counts the eviction *attempt*, matching the `truncate_attempts`
442/// naming convention above.
443pub(crate) fn note_read_tx_max_age_eviction() {
444    READ_TX_MAX_AGE_EVICTIONS.fetch_add(1, Ordering::Relaxed);
445}
446
447/// Total checkpoint lifecycle transitions rejected before append.
448pub fn checkpoint_lifecycle_enqueue_drops() -> u64 {
449    CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.load(Ordering::Relaxed)
450}
451
452/// A tick had no usable WAL frame observation: bump the
453/// lifetime and consecutive-skip counters and snapshot the last-known WAL
454/// pressure so an operator can see how bad the WAL was heading into the skip
455/// streak.
456pub(super) fn note_checkpoint_skipped() {
457    CHECKPOINT_SKIPPED_TICKS.fetch_add(1, Ordering::Relaxed);
458    CHECKPOINT_CONSECUTIVE_SKIPS.fetch_add(1, Ordering::Relaxed);
459    if let Some(pages) = last_observed_wal_pages() {
460        CHECKPOINT_LAST_SKIP_WAL_PAGES.store(pages, Ordering::Relaxed);
461    }
462}
463
464/// A tick had a valid WAL frame observation: close out any prior skip
465/// streak. `_wal_pages` is accepted for call-site symmetry with
466/// `note_checkpoint_skipped` and to leave room for a future observed-side
467/// gauge without changing this function's signature again.
468pub(super) fn note_checkpoint_observed(_wal_pages: u64) {
469    CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
470}
471
472pub(super) fn note_checkpoint_pressure_observation(above_warn: bool, was_above_warn: bool) {
473    if above_warn {
474        CHECKPOINT_PRESSURE_ELEVATED_TICKS.fetch_add(1, Ordering::Relaxed);
475        if !was_above_warn {
476            CHECKPOINT_PRESSURE_EPISODES_STARTED.fetch_add(1, Ordering::Relaxed);
477        }
478    } else if was_above_warn {
479        CHECKPOINT_PRESSURE_EPISODES_RECOVERED.fetch_add(1, Ordering::Relaxed);
480    }
481}
482
483/// Reset the checkpoint-pressure atomics between tests. Process-wide gauges
484/// are otherwise shared across every test in this binary; tests that assert
485/// on them must reset first and run under a shared `#[serial(...)]` group.
486#[cfg(test)]
487pub(crate) fn reset_checkpoint_metrics_for_tests() {
488    CHECKPOINT_SKIPPED_TICKS.store(0, Ordering::Relaxed);
489    CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
490    CHECKPOINT_LAST_SKIP_WAL_PAGES.store(u64::MAX, Ordering::Relaxed);
491    CHECKPOINT_PRESSURE_ELEVATED_TICKS.store(0, Ordering::Relaxed);
492    CHECKPOINT_PRESSURE_EPISODES_STARTED.store(0, Ordering::Relaxed);
493    CHECKPOINT_PRESSURE_EPISODES_RECOVERED.store(0, Ordering::Relaxed);
494    CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.store(0, Ordering::Relaxed);
495    CHECKPOINT_LIFECYCLE_APPEND_FAILURES.store(0, Ordering::Relaxed);
496    CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.store(0, Ordering::Relaxed);
497    READ_TX_MAX_AGE_EVICTIONS.store(0, Ordering::Relaxed);
498}
499
500/// Outcome of a single checkpoint attempt.
501///
502/// `Skipped` is returned when the dedicated connection is unavailable or its
503/// PASSIVE result is busy or inconsistent and has no usable WAL frame
504/// observation. A concurrent pool writer does not cause a skip: the task never
505/// checks out its writer mutex. `Observed` carries the WAL page count read
506/// during the tick. A skipped tick leaves threshold-crossing state unchanged.
507#[derive(Debug, Clone, Copy, PartialEq, Eq)]
508pub enum CheckpointTick {
509    /// No usable WAL frame observation was available this tick.
510    Skipped,
511    /// A checkpoint was issued; the value is the observed WAL page count.
512    Observed(u64),
513}