use std::collections::{BTreeMap, HashMap};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex, OnceLock};
use std::time::{Duration, Instant};
use crate::pool::ConnectionPool;
mod off_worker;
#[cfg(test)]
mod off_worker_tests;
static LAST_WAL_PAGES: AtomicU64 = AtomicU64::new(u64::MAX);
static TRUNCATE_ATTEMPTS: AtomicU64 = AtomicU64::new(0);
static TRUNCATE_CONSECUTIVE_FAILURES: AtomicU64 = AtomicU64::new(0);
static CHECKPOINT_SKIPPED_TICKS: AtomicU64 = AtomicU64::new(0);
static CHECKPOINT_CONSECUTIVE_SKIPS: AtomicU64 = AtomicU64::new(0);
static CHECKPOINT_LAST_SKIP_WAL_PAGES: AtomicU64 = AtomicU64::new(u64::MAX);
static CHECKPOINT_PRESSURE_ELEVATED_TICKS: AtomicU64 = AtomicU64::new(0);
static CHECKPOINT_PRESSURE_EPISODES_STARTED: AtomicU64 = AtomicU64::new(0);
static CHECKPOINT_PRESSURE_EPISODES_RECOVERED: AtomicU64 = AtomicU64::new(0);
static CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS: AtomicU64 = AtomicU64::new(0);
static CHECKPOINT_LIFECYCLE_APPEND_FAILURES: AtomicU64 = AtomicU64::new(0);
static CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS: AtomicU64 = AtomicU64::new(0);
static READ_TX_MAX_AGE_EVICTIONS: AtomicU64 = AtomicU64::new(0);
mod run_state;
pub use run_state::{
checkpoint_consecutive_skips, checkpoint_last_skip_wal_pages,
checkpoint_lifecycle_append_attempts, checkpoint_lifecycle_append_failures,
checkpoint_lifecycle_enqueue_drops, checkpoint_pressure_elevated_ticks,
checkpoint_pressure_episodes_recovered, checkpoint_pressure_episodes_started,
checkpoint_skipped_ticks, checkpoint_timing, last_observed_wal_pages,
read_tx_max_age_evictions, routine_wal_observation, truncate_attempts,
truncate_consecutive_failures, CheckpointRun, CheckpointTick, CheckpointTiming,
RoutineWalObservation,
};
pub(crate) use run_state::{
checkpoint_run_snapshot, note_read_tx_max_age_eviction, record_checkpoint_run_result,
CheckpointRunStatus, CheckpointRunTaskGuard,
};
use run_state::{
note_checkpoint_observed, note_checkpoint_pressure_observation, note_checkpoint_skipped,
record_checkpoint_timing, record_routine_wal_observation,
};
#[cfg(test)]
pub(crate) use run_state::{checkpoint_run_status, reset_checkpoint_metrics_for_tests};
#[cfg(test)]
use run_state::{
advance_checkpoint_run, advance_checkpoint_run_at, checkpoint_db_key,
checkpoint_db_key_from_path, checkpoint_runs, checkpoint_timings,
};
pub const DEFAULT_WARN_SUSTAINED_CYCLES: u8 = 3;
#[derive(Clone, Debug)]
pub struct CheckpointConfig {
pub interval: Duration,
pub warn_pages: u64,
pub warn_sustained_cycles: u8,
pub high_water_pages: u64,
pub truncate_high_water_pages: u64,
pub truncate_min_interval: Duration,
pub truncate_busy_timeout: Duration,
pub tx_warn_secs: Duration,
pub tx_max_age_secs: Duration,
}
impl Default for CheckpointConfig {
fn default() -> Self {
Self {
interval: Duration::from_millis(500),
warn_pages: 2000,
warn_sustained_cycles: DEFAULT_WARN_SUSTAINED_CYCLES,
high_water_pages: 6000,
truncate_high_water_pages: 20_000,
truncate_min_interval: Duration::from_secs(300),
truncate_busy_timeout: Duration::from_millis(2000),
tx_warn_secs: Duration::from_secs(30),
tx_max_age_secs: Duration::from_secs(120),
}
}
}
impl CheckpointConfig {
pub fn from_env() -> Self {
let mut cfg = Self::default();
if let Ok(ms) = std::env::var("KHIVE_CHECKPOINT_INTERVAL_MS") {
if let Ok(v) = ms.parse::<u64>() {
if v > 0 {
cfg.interval = Duration::from_millis(v);
}
}
}
if let Ok(v) = std::env::var("KHIVE_WAL_WARN_PAGES") {
if let Ok(n) = v.parse::<u64>() {
if n > 0 {
cfg.warn_pages = n;
}
}
}
if let Ok(v) = std::env::var("KHIVE_WAL_WARN_SUSTAINED_CYCLES") {
if let Ok(n) = v.parse::<u8>() {
if n > 0 {
cfg.warn_sustained_cycles = n;
}
}
}
if let Ok(v) = std::env::var("KHIVE_WAL_HIGH_WATER_PAGES") {
if let Ok(n) = v.parse::<u64>() {
if n > 0 {
cfg.high_water_pages = n;
}
}
}
if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES") {
if let Ok(n) = v.parse::<u64>() {
if n > 0 {
cfg.truncate_high_water_pages = n;
}
}
}
if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS") {
if let Ok(n) = v.parse::<u64>() {
if n > 0 {
cfg.truncate_min_interval = Duration::from_secs(n);
}
}
}
if let Ok(v) = std::env::var("KHIVE_WAL_TRUNCATE_BUSY_MS") {
if let Ok(n) = v.parse::<u64>() {
if n > 0 {
cfg.truncate_busy_timeout = Duration::from_millis(n);
}
}
}
(cfg.tx_warn_secs, cfg.tx_max_age_secs) =
tx_age_thresholds_from_env(cfg.tx_warn_secs, cfg.tx_max_age_secs);
cfg
}
}
pub(crate) fn tx_age_thresholds_from_env(
default_warn: Duration,
default_max: Duration,
) -> (Duration, Duration) {
let mut warn_secs = default_warn;
let mut max_age_secs = default_max;
if let Ok(v) = std::env::var("KHIVE_TX_WARN_SECS") {
if let Ok(n) = v.parse::<u64>() {
if n > 0 {
warn_secs = Duration::from_secs(n);
}
}
}
if let Ok(v) = std::env::var("KHIVE_TX_MAX_AGE_SECS") {
if let Ok(n) = v.parse::<u64>() {
if n > 0 {
max_age_secs = Duration::from_secs(n);
}
}
}
if warn_secs >= max_age_secs {
tracing::warn!(
configured_tx_warn_secs = warn_secs.as_secs_f64(),
configured_tx_max_age_secs = max_age_secs.as_secs_f64(),
fallback_tx_warn_secs = default_warn.as_secs_f64(),
fallback_tx_max_age_secs = default_max.as_secs_f64(),
"KHIVE_TX_WARN_SECS must be strictly less than KHIVE_TX_MAX_AGE_SECS; \
both transaction-age thresholds were rejected and reset to their defaults"
);
return (default_warn, default_max);
}
(warn_secs, max_age_secs)
}
#[cfg(unix)]
const DEFAULT_WALPIN_FULL_SCAN_INTERVAL: Duration = Duration::from_secs(30);
#[cfg(unix)]
#[derive(Debug, Clone)]
struct CachedWalpinAttribution {
report: crate::walpin::WalpinReport,
census: Result<crate::walpin::CensusResult, String>,
captured_at: Instant,
}
#[cfg(unix)]
#[derive(Debug)]
enum WalpinFullScanPlan {
Refresh {
previous_last_attempt: Option<Instant>,
},
Cached(CachedWalpinAttribution),
Suppressed,
}
#[derive(Debug)]
pub struct TruncateState {
last_attempt: Option<Instant>,
consecutive_failures: u32,
#[cfg(unix)]
legacy_walpin_fallback_interval: Duration,
#[cfg(unix)]
walpin_full_scan_interval: Duration,
#[cfg(unix)]
walpin_full_scan_last_attempt: Option<Instant>,
#[cfg(unix)]
walpin_cached_attribution: Option<CachedWalpinAttribution>,
#[cfg(unix)]
sidecar_attribution_attempted_this_tick: bool,
}
impl Default for TruncateState {
fn default() -> Self {
Self {
last_attempt: None,
consecutive_failures: 0,
#[cfg(unix)]
legacy_walpin_fallback_interval: DEFAULT_SESSION_SWEEP_INTERVAL,
#[cfg(unix)]
walpin_full_scan_interval: DEFAULT_WALPIN_FULL_SCAN_INTERVAL,
#[cfg(unix)]
walpin_full_scan_last_attempt: None,
#[cfg(unix)]
walpin_cached_attribution: None,
#[cfg(unix)]
sidecar_attribution_attempted_this_tick: false,
}
}
}
impl TruncateState {
#[cfg(unix)]
fn with_legacy_walpin_fallback(interval: Duration) -> Self {
Self {
legacy_walpin_fallback_interval: interval,
..Self::default()
}
}
#[cfg(all(test, unix))]
fn with_walpin_full_scan_cadence(interval: Duration) -> Self {
Self {
walpin_full_scan_interval: interval,
..Self::default()
}
}
#[cfg(unix)]
fn begin_tick(&mut self) {
self.sidecar_attribution_attempted_this_tick = false;
}
#[cfg(unix)]
fn housekeeping_due(&self) -> bool {
!self.sidecar_attribution_attempted_this_tick
&& self.walpin_full_scan_due_at(Instant::now())
}
#[cfg(unix)]
fn walpin_full_scan_due_at(&self, now: Instant) -> bool {
self.walpin_full_scan_last_attempt.is_none_or(|last| {
now.saturating_duration_since(last) >= self.walpin_full_scan_interval
})
}
#[cfg(unix)]
fn claim_walpin_full_scan_at(&mut self, now: Instant) -> bool {
if !self.walpin_full_scan_due_at(now) {
return false;
}
self.walpin_full_scan_last_attempt = Some(now);
true
}
#[cfg(unix)]
fn plan_walpin_attribution_at(&mut self, now: Instant) -> WalpinFullScanPlan {
if self.walpin_full_scan_due_at(now) {
let previous_last_attempt = self.walpin_full_scan_last_attempt.replace(now);
WalpinFullScanPlan::Refresh {
previous_last_attempt,
}
} else if let Some(cached) = self.walpin_cached_attribution.clone() {
WalpinFullScanPlan::Cached(cached)
} else {
WalpinFullScanPlan::Suppressed
}
}
#[cfg(unix)]
fn restore_walpin_full_scan_reservation(&mut self, previous_last_attempt: Option<Instant>) {
self.walpin_full_scan_last_attempt = previous_last_attempt;
}
#[cfg(unix)]
fn cache_walpin_attribution(
&mut self,
report: crate::walpin::WalpinReport,
census: Result<crate::walpin::CensusResult, String>,
captured_at: Instant,
) {
self.walpin_cached_attribution = Some(CachedWalpinAttribution {
report,
census,
captured_at,
});
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CheckpointSeverityRung {
Info,
Warn,
Alarm,
}
#[derive(Debug, Default, Clone)]
pub struct CheckpointSeverityState {
was_above_warn: bool,
consecutive_above_warn: u8,
warn_emitted_for_episode: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct CheckpointSeverityEmission {
pub rung: CheckpointSeverityRung,
pub wal_pages: u64,
pub threshold_pages: u64,
pub consecutive_cycles: u8,
}
impl CheckpointSeverityState {
pub fn observe_wal_pages(
&mut self,
wal_pages: u64,
config: &CheckpointConfig,
) -> Vec<CheckpointSeverityEmission> {
let mut emissions = Vec::new();
let above_warn = wal_pages >= config.warn_pages;
if above_warn {
self.consecutive_above_warn = self.consecutive_above_warn.saturating_add(1);
if !self.was_above_warn {
emissions.push(CheckpointSeverityEmission {
rung: CheckpointSeverityRung::Info,
wal_pages,
threshold_pages: config.warn_pages,
consecutive_cycles: self.consecutive_above_warn,
});
}
if !self.warn_emitted_for_episode
&& self.consecutive_above_warn >= config.warn_sustained_cycles
{
emissions.push(CheckpointSeverityEmission {
rung: CheckpointSeverityRung::Warn,
wal_pages,
threshold_pages: config.warn_pages,
consecutive_cycles: self.consecutive_above_warn,
});
self.warn_emitted_for_episode = true;
}
} else {
self.consecutive_above_warn = 0;
self.warn_emitted_for_episode = false;
}
self.was_above_warn = above_warn;
emissions
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TxAgeRung {
Warn,
Stale,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TxAgeEmission {
pub rung: TxAgeRung,
pub age: Duration,
pub label: Option<String>,
}
#[derive(Debug, Default, Clone)]
pub struct TxAgeSweepState {
was_above_warn: bool,
was_above_max_age: bool,
tracked_id: Option<khive_storage::tx_registry::TxId>,
}
impl TxAgeSweepState {
pub fn observe(
&mut self,
oldest: Option<(khive_storage::tx_registry::TxId, Duration, Option<String>)>,
tx_warn_secs: Duration,
tx_max_age_secs: Duration,
) -> Vec<TxAgeEmission> {
let mut emissions = Vec::new();
let Some((id, age, label)) = oldest else {
self.was_above_warn = false;
self.was_above_max_age = false;
self.tracked_id = None;
return emissions;
};
if self.tracked_id != Some(id) {
self.was_above_warn = false;
self.was_above_max_age = false;
}
self.tracked_id = Some(id);
let above_warn = age >= tx_warn_secs;
let above_max_age = age >= tx_max_age_secs;
if above_warn && !self.was_above_warn {
emissions.push(TxAgeEmission {
rung: TxAgeRung::Warn,
age,
label: label.clone(),
});
}
if above_max_age && !self.was_above_max_age {
emissions.push(TxAgeEmission {
rung: TxAgeRung::Stale,
age,
label,
});
}
self.was_above_warn = above_warn;
self.was_above_max_age = above_max_age;
emissions
}
}
fn log_tx_age_emission(emission: &TxAgeEmission) {
let label = emission.label.as_deref().unwrap_or("<unlabeled>");
match emission.rung {
TxAgeRung::Warn => {
tracing::warn!(
tx_age_secs = emission.age.as_secs_f64(),
tx_label = label,
"ADR-091 Plank 1: open transaction registry entry exceeded soft-cap age"
);
}
TxAgeRung::Stale => {
tracing::error!(
tx_age_secs = emission.age.as_secs_f64(),
tx_label = label,
"ADR-091 Plank 1: open transaction registry entry exceeded the cooperative \
stale-op cap; no in-process mechanism can force-close it — investigate the \
labeled caller directly"
);
}
}
}
struct WalpinSidecarState {
dir: PathBuf,
pid: u32,
role: &'static str,
started_at: i64,
sweep_interval_ms: u64,
wrote: bool,
beacon_registered: bool,
last_heartbeat: Option<LastHeartbeatState>,
}
struct LastHeartbeatState {
span_id: khive_storage::tx_registry::TxId,
label: Option<String>,
attribution_basis: &'static str,
sweep_interval_ms: u64,
oldest_tx_started_at: i64,
}
impl LastHeartbeatState {
fn content_matches(
&self,
span_id: khive_storage::tx_registry::TxId,
label: &Option<String>,
attribution_basis: &str,
sweep_interval_ms: u64,
) -> bool {
self.span_id == span_id
&& self.label == *label
&& self.attribution_basis == attribution_basis
&& self.sweep_interval_ms == sweep_interval_ms
}
}
impl WalpinSidecarState {
fn new(
db_path: Option<&Path>,
is_file_backed: bool,
role: &'static str,
interval: Duration,
) -> Option<Self> {
let path = db_path?;
if !crate::walpin::sidecar_enabled(is_file_backed) {
return None;
}
let pid = std::process::id();
Some(Self {
dir: crate::walpin::sidecar_dir_for(path),
pid,
role,
started_at: crate::walpin::process_start_time_secs(pid).unwrap_or(0),
sweep_interval_ms: interval.as_millis().min(u64::MAX as u128) as u64,
wrote: false,
last_heartbeat: None,
beacon_registered: false,
})
}
async fn register_beacon(&mut self) {
let dir = self.dir.clone();
let beacon = crate::walpin::WalpinBeacon {
pid: self.pid,
process_role: self.role.to_string(),
started_at: self.started_at,
sweep_interval_ms: self.sweep_interval_ms,
};
let result =
tokio::task::spawn_blocking(move || crate::walpin::write_beacon(&dir, &beacon)).await;
match result {
Ok(Ok(())) => {
self.beacon_registered = true;
}
Ok(Err(e)) => {
tracing::warn!(
error = %e,
"ADR-091 Amendment 2: failed to write walpin registration beacon; \
this process's sidecar health will read as unknown, not registered-silent"
);
}
Err(join_err) => {
tracing::warn!(
error = %join_err,
"ADR-091 Amendment 2: walpin beacon write task panicked"
);
}
}
}
#[cfg(unix)]
async fn reap_dead_entries_bounded(
&self,
legacy_fallback_interval: Duration,
) -> Option<crate::walpin::WalpinReport> {
let dir = self.dir.clone();
let result = tokio::task::spawn_blocking(move || {
crate::walpin::housekeep_live(&dir, legacy_fallback_interval)
})
.await;
match result {
Ok(Ok(report)) => Some(report),
Ok(Err(e)) => {
tracing::warn!(
error = %e,
"ADR-091 Amendment 6: bounded walpin sidecar cleanup failed"
);
None
}
Err(join_err) => {
tracing::warn!(
error = %join_err,
"ADR-091 Amendment 6: walpin sidecar cleanup task panicked"
);
None
}
}
}
async fn refresh_beacon(&mut self) {
if !self.beacon_registered {
self.register_beacon().await;
return;
}
let dir = self.dir.clone();
let pid = self.pid;
let result =
tokio::task::spawn_blocking(move || crate::walpin::touch_beacon(&dir, pid)).await;
match result {
Ok(Ok(())) => {}
Ok(Err(e)) => {
self.beacon_registered = false;
tracing::warn!(
error = %e,
"ADR-091 Amendment 2: failed to refresh walpin registration beacon; \
this process's sidecar health will read as unknown, not registered-silent"
);
}
Err(join_err) => {
self.beacon_registered = false;
tracing::warn!(
error = %join_err,
"ADR-091 Amendment 2: walpin beacon refresh task panicked"
);
}
}
}
async fn drop_beacon_fail_closed(&mut self) {
let dir = self.dir.clone();
let pid = self.pid;
self.beacon_registered = false;
let result =
tokio::task::spawn_blocking(move || crate::walpin::remove_beacon(&dir, pid)).await;
match result {
Ok(Ok(())) => {}
Ok(Err(e)) => {
tracing::warn!(
error = %e,
"ADR-091 Amendment 2: failed to remove walpin beacon after a failed \
heartbeat write; beacon will age out of the freshness window instead"
);
}
Err(join_err) => {
tracing::warn!(
error = %join_err,
"ADR-091 Amendment 2: walpin beacon removal task panicked"
);
}
}
}
async fn observe(
&mut self,
oldest: Option<khive_storage::tx_registry::OldestSpan>,
tx_warn_secs: Duration,
) {
match oldest {
Some(span) if span.age >= tx_warn_secs => {
let attribution_basis = match span.origin {
khive_storage::tx_registry::TxOrigin::Database(_) => "origin",
khive_storage::tx_registry::TxOrigin::Unscoped
| khive_storage::tx_registry::TxOrigin::Memory => "fallback",
};
let content_unchanged = self.wrote
&& self.last_heartbeat.as_ref().is_some_and(|last| {
last.content_matches(
span.id,
&span.label,
attribution_basis,
self.sweep_interval_ms,
)
});
if content_unchanged {
let dir = self.dir.clone();
let pid = self.pid;
let touch_result = tokio::task::spawn_blocking(move || {
crate::walpin::touch_heartbeat(&dir, pid)
})
.await;
match touch_result {
Ok(Ok(())) => {
self.refresh_beacon().await;
return;
}
Ok(Err(e)) => {
tracing::warn!(
error = %e,
"ADR-091 Amendment 3 Plank F1: walpin heartbeat touch failed; \
recreating with a full body write"
);
}
Err(join_err) => {
tracing::warn!(
error = %join_err,
"ADR-091 Amendment 3 Plank F1: walpin heartbeat touch task \
panicked; recreating with a full body write"
);
}
}
}
let oldest_tx_started_at = self
.last_heartbeat
.as_ref()
.filter(|last| last.span_id == span.id)
.map(|last| last.oldest_tx_started_at)
.unwrap_or_else(|| now_epoch_secs().saturating_sub(span.age.as_secs() as i64));
let heartbeat = crate::walpin::WalpinHeartbeat {
pid: self.pid,
process_role: self.role.to_string(),
started_at: self.started_at,
oldest_tx_age_secs: span.age.as_secs_f64(),
oldest_tx_label: span.label.clone(),
oldest_tx_started_at: Some(oldest_tx_started_at),
updated_at: now_epoch_secs(),
sweep_interval_ms: self.sweep_interval_ms,
attribution_basis: Some(attribution_basis.to_string()),
};
let dir = self.dir.clone();
let result = tokio::task::spawn_blocking(move || {
crate::walpin::write_heartbeat(&dir, &heartbeat)
})
.await;
match result {
Ok(Ok(())) => {
self.wrote = true;
self.last_heartbeat = Some(LastHeartbeatState {
span_id: span.id,
label: span.label,
attribution_basis,
sweep_interval_ms: self.sweep_interval_ms,
oldest_tx_started_at,
});
self.refresh_beacon().await;
}
Ok(Err(e)) => {
tracing::warn!(
error = %e,
"ADR-091 Amendment 2 Plank B: failed to write walpin heartbeat; \
removing beacon so this process cannot read as \
registered-silent while over threshold"
);
self.last_heartbeat = None;
self.drop_beacon_fail_closed().await;
}
Err(join_err) => {
tracing::warn!(
error = %join_err,
"ADR-091 Amendment 2 Plank B: walpin heartbeat write task panicked"
);
self.last_heartbeat = None;
self.drop_beacon_fail_closed().await;
}
}
}
_ => {
self.refresh_beacon().await;
if self.wrote {
let dir = self.dir.clone();
let pid = self.pid;
let result = tokio::task::spawn_blocking(move || {
crate::walpin::remove_heartbeat(&dir, pid)
})
.await;
match result {
Ok(Ok(())) => {}
Ok(Err(e)) => tracing::warn!(
error = %e,
"ADR-091 Amendment 2 Plank B: failed to remove walpin heartbeat"
),
Err(join_err) => tracing::warn!(
error = %join_err,
"ADR-091 Amendment 2 Plank B: walpin heartbeat removal task panicked"
),
}
self.wrote = false;
self.last_heartbeat = None;
}
}
}
}
async fn shutdown(&mut self) {
if self.wrote {
let dir = self.dir.clone();
let pid = self.pid;
let _ = tokio::task::spawn_blocking(move || crate::walpin::remove_heartbeat(&dir, pid))
.await;
self.wrote = false;
}
}
}
#[cfg(unix)]
async fn run_walpin_housekeeping_if_due(
sidecar: &WalpinSidecarState,
state: &mut TruncateState,
legacy_fallback_interval: Duration,
) -> bool {
if !state.housekeeping_due() || !state.claim_walpin_full_scan_at(Instant::now()) {
return false;
}
if let Some(report) = sidecar
.reap_dead_entries_bounded(legacy_fallback_interval)
.await
{
state.cache_walpin_attribution(
report,
Err("OS holder census is unavailable for a housekeeping-only scan".to_string()),
Instant::now(),
);
}
true
}
fn now_epoch_secs() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0)
}
const DEFAULT_SESSION_SWEEP_INTERVAL: Duration = Duration::from_secs(5);
#[derive(Clone, Debug)]
pub struct SessionSweepConfig {
pub interval: Duration,
pub tx_warn_secs: Duration,
pub tx_max_age_secs: Duration,
}
impl Default for SessionSweepConfig {
fn default() -> Self {
Self {
interval: DEFAULT_SESSION_SWEEP_INTERVAL,
tx_warn_secs: Duration::from_secs(30),
tx_max_age_secs: Duration::from_secs(120),
}
}
}
impl SessionSweepConfig {
pub fn from_env() -> Self {
let mut cfg = Self {
interval: session_sweep_interval_from_env(),
..Self::default()
};
(cfg.tx_warn_secs, cfg.tx_max_age_secs) =
tx_age_thresholds_from_env(cfg.tx_warn_secs, cfg.tx_max_age_secs);
cfg
}
}
fn session_sweep_interval_from_env() -> Duration {
std::env::var("KHIVE_SESSION_SWEEP_INTERVAL_MS")
.ok()
.and_then(|ms| ms.parse::<u64>().ok())
.filter(|ms| *ms > 0)
.map(Duration::from_millis)
.unwrap_or(DEFAULT_SESSION_SWEEP_INTERVAL)
}
pub struct SweepBackend {
pub pool: Arc<ConnectionPool>,
pub is_main: bool,
}
struct BackendSweep {
filter: khive_storage::tx_registry::TxOriginFilter,
tx_age_state: TxAgeSweepState,
sidecar: Option<WalpinSidecarState>,
}
pub async fn run_session_sweep_task(
backends: Vec<SweepBackend>,
config: SessionSweepConfig,
mut shutdown_rx: tokio::sync::watch::Receiver<()>,
) {
let mut interval = tokio::time::interval(config.interval);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
let mut sweeps: Vec<BackendSweep> = Vec::with_capacity(backends.len());
for backend in backends {
let identity = match backend.pool.origin() {
khive_storage::tx_registry::TxOrigin::Database(id) => id,
khive_storage::tx_registry::TxOrigin::Memory
| khive_storage::tx_registry::TxOrigin::Unscoped => continue,
};
let filter = if backend.is_main {
khive_storage::tx_registry::TxOriginFilter::Main(identity)
} else {
khive_storage::tx_registry::TxOriginFilter::Secondary(identity)
};
let sidecar = WalpinSidecarState::new(
backend.pool.canonical_path(),
true,
"session",
config.interval,
);
sweeps.push(BackendSweep {
filter,
tx_age_state: TxAgeSweepState::default(),
sidecar,
});
}
for sweep in sweeps.iter_mut() {
if let Some(sidecar) = sweep.sidecar.as_mut() {
sidecar.register_beacon().await;
}
}
loop {
tokio::select! {
_ = interval.tick() => {}
_ = shutdown_rx.changed() => break,
}
for sweep in sweeps.iter_mut() {
let oldest = khive_storage::tx_registry::oldest_for(&sweep.filter);
for emission in sweep.tx_age_state.observe(
oldest.as_ref().map(|s| (s.id, s.age, s.label.clone())),
config.tx_warn_secs,
config.tx_max_age_secs,
) {
log_tx_age_emission(&emission);
}
if let Some(sidecar) = sweep.sidecar.as_mut() {
sidecar.observe(oldest, config.tx_warn_secs).await;
}
}
}
for sweep in sweeps.iter_mut() {
if let Some(sidecar) = sweep.sidecar.as_mut() {
sidecar.shutdown().await;
}
}
}
#[derive(Clone)]
pub struct CheckpointLifecycleOwner {
event_store: Arc<dyn khive_storage::EventStore>,
namespace: String,
}
impl CheckpointLifecycleOwner {
pub fn new(
event_store: Arc<dyn khive_storage::EventStore>,
namespace: impl Into<String>,
) -> Self {
Self {
event_store,
namespace: namespace.into(),
}
}
}
const CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY: usize = 1;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct CheckpointPressureEpisode {
elevated_ticks: u64,
peak_wal_pages: u64,
}
impl CheckpointPressureEpisode {
fn start(wal_pages: u64) -> Self {
Self {
elevated_ticks: 1,
peak_wal_pages: wal_pages,
}
}
fn observe(&mut self, wal_pages: u64) {
self.elevated_ticks = self.elevated_ticks.saturating_add(1);
self.peak_wal_pages = self.peak_wal_pages.max(wal_pages);
}
}
struct CheckpointLifecycleEmitter {
namespace: Option<String>,
sender: Option<tokio::sync::mpsc::Sender<khive_storage::Event>>,
worker: Option<tokio::task::JoinHandle<()>>,
busy_warning_emitted: bool,
}
impl CheckpointLifecycleEmitter {
fn new(owner: Option<CheckpointLifecycleOwner>) -> Self {
let Some(owner) = owner else {
return Self {
namespace: None,
sender: None,
worker: None,
busy_warning_emitted: false,
};
};
let namespace = owner.namespace.clone();
let (sender, mut receiver) =
tokio::sync::mpsc::channel::<khive_storage::Event>(CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY);
let worker = tokio::spawn(async move {
while let Some(event) = receiver.recv().await {
let kind = event.kind;
CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.fetch_add(1, Ordering::Relaxed);
if let Err(err) = owner.event_store.append_event(event).await {
CHECKPOINT_LIFECYCLE_APPEND_FAILURES.fetch_add(1, Ordering::Relaxed);
tracing::warn!(
error = %err,
event_kind = %kind.name(),
"checkpoint lifecycle event append failed"
);
}
}
});
Self {
namespace: Some(namespace),
sender: Some(sender),
worker: Some(worker),
busy_warning_emitted: false,
}
}
fn try_emit<P: serde::Serialize>(&mut self, kind: khive_types::EventKind, payload: P) -> bool {
let (Some(namespace), Some(sender)) = (&self.namespace, &self.sender) else {
return true;
};
let payload_value = match serde_json::to_value(&payload) {
Ok(value) => value,
Err(err) => {
tracing::warn!(
error = %err,
event_kind = %kind.name(),
"failed to serialize checkpoint lifecycle event payload"
);
CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.fetch_add(1, Ordering::Relaxed);
return false;
}
};
let payload_schema_version = match kind {
khive_types::EventKind::CheckpointOutcomeRecorded => 2,
_ => 1,
};
let event = khive_storage::Event::new(
namespace,
"checkpoint.lifecycle",
kind,
khive_types::SubstrateKind::Event,
"daemon:checkpoint_task",
)
.with_payload(payload_value)
.with_payload_schema_version(payload_schema_version);
match sender.try_send(event) {
Ok(()) => {
self.busy_warning_emitted = false;
true
}
Err(tokio::sync::mpsc::error::TrySendError::Full(event)) => {
CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.fetch_add(1, Ordering::Relaxed);
if !self.busy_warning_emitted {
tracing::warn!(
event_kind = %event.kind.name(),
queue_capacity = CHECKPOINT_LIFECYCLE_QUEUE_CAPACITY,
"checkpoint lifecycle event dropped because the append worker is busy"
);
self.busy_warning_emitted = true;
}
false
}
Err(tokio::sync::mpsc::error::TrySendError::Closed(event)) => {
CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.fetch_add(1, Ordering::Relaxed);
tracing::warn!(
event_kind = %event.kind.name(),
"checkpoint lifecycle event dropped because the append worker stopped"
);
false
}
}
}
async fn shutdown(mut self) {
drop(self.sender.take());
let Some(worker) = self.worker.take() else {
return;
};
worker.abort();
match worker.await {
Ok(()) => {}
Err(err) if err.is_cancelled() => {}
Err(err) => tracing::warn!(
error = %err,
"checkpoint lifecycle event append worker terminated unexpectedly"
),
}
}
}
impl Drop for CheckpointLifecycleEmitter {
fn drop(&mut self) {
if let Some(worker) = &self.worker {
worker.abort();
}
}
}
struct CheckpointConnection {
conn: Option<rusqlite::Connection>,
consecutive_open_failures: u32,
}
impl CheckpointConnection {
fn new() -> Self {
Self {
conn: None,
consecutive_open_failures: 0,
}
}
fn ensure_open(&mut self, pool: &ConnectionPool) -> Option<&rusqlite::Connection> {
if self.conn.is_none() {
match pool.open_standalone_writer_untracked() {
Ok(conn) => {
if let Err(e) = conn.pragma_update(None, "wal_autocheckpoint", 0) {
tracing::warn!(
error = %e,
"could not disable autocheckpoint on the dedicated checkpoint \
connection"
);
}
if self.consecutive_open_failures > 0 {
tracing::info!(
prior_consecutive_failures = self.consecutive_open_failures,
"dedicated checkpoint connection opened successfully, ending a \
failure streak"
);
}
self.consecutive_open_failures = 0;
self.conn = Some(conn);
}
Err(e) => {
if self.consecutive_open_failures == 0 {
tracing::warn!(
error = %e,
"failed to open the dedicated checkpoint connection; \
this tick is skipped and the open retried next tick"
);
} else {
tracing::debug!(
error = %e,
consecutive_failures = self.consecutive_open_failures,
"dedicated checkpoint connection still unavailable; \
this tick is skipped and the open retried next tick"
);
}
self.consecutive_open_failures =
self.consecutive_open_failures.saturating_add(1);
return None;
}
}
}
self.conn.as_ref()
}
}
async fn run_fts_maintenance_off_worker(
conn: rusqlite::Connection,
config: crate::fts_maintenance::FtsMaintenanceConfig,
mut state: crate::fts_maintenance::FtsMaintenanceState,
now: Instant,
) -> Result<
(
rusqlite::Connection,
crate::fts_maintenance::FtsMaintenanceState,
Result<Option<crate::fts_maintenance::FtsMaintenanceStep>, String>,
),
tokio::task::JoinError,
> {
tokio::task::spawn_blocking(move || {
let result = crate::fts_maintenance::run_if_due(&conn, &config, &mut state, now);
(conn, state, result)
})
.await
}
pub async fn run_checkpoint_task(
pool: Arc<ConnectionPool>,
config: CheckpointConfig,
lifecycle_owner: Option<CheckpointLifecycleOwner>,
mut shutdown_rx: tokio::sync::watch::Receiver<()>,
is_main: bool,
) {
let _checkpoint_run_guard = CheckpointRunTaskGuard::start(&pool, config.interval);
match pool.claim_checkpoint_ownership() {
Ok(()) => {
if let Err(e) = pool.propagate_checkpoint_claim_to_writer_task().await {
tracing::warn!(
error = %e,
"checkpoint task could not reach the writer task's connection; it keeps the \
bounded autocheckpoint fallback"
);
}
}
Err(e) => {
tracing::warn!(
error = %e,
"checkpoint task could not re-apply the ownership pragma on the pooled writer; \
writer connections keep the bounded autocheckpoint fallback unless ownership is \
claimed later"
);
}
}
let mut interval = tokio::time::interval(config.interval);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
let mut severity_state = CheckpointSeverityState::default();
let mut tx_age_state = TxAgeSweepState::default();
let mut was_above_high_water = false;
#[cfg(unix)]
let legacy_walpin_fallback_interval = DEFAULT_SESSION_SWEEP_INTERVAL;
#[cfg(unix)]
let mut truncate_state =
TruncateState::with_legacy_walpin_fallback(legacy_walpin_fallback_interval);
#[cfg(not(unix))]
let mut truncate_state = TruncateState::default();
let mut lifecycle_emitter = CheckpointLifecycleEmitter::new(lifecycle_owner);
let mut event_elevation_open = false;
let mut pressure_episode: Option<CheckpointPressureEpisode> = None;
let mut pending_recovery: Option<khive_storage::CheckpointOutcomeRecordedPayload> = None;
let mut was_observed_above_warn = false;
let tx_filter = match pool.origin() {
khive_storage::tx_registry::TxOrigin::Database(id) => Some(if is_main {
khive_storage::tx_registry::TxOriginFilter::Main(id)
} else {
khive_storage::tx_registry::TxOriginFilter::Secondary(id)
}),
khive_storage::tx_registry::TxOrigin::Memory
| khive_storage::tx_registry::TxOrigin::Unscoped => None,
};
#[cfg(unix)]
let mut walpin_state =
WalpinSidecarState::new(pool.canonical_path(), true, "daemon", config.interval);
#[cfg(unix)]
if let Some(sidecar) = walpin_state.as_mut() {
sidecar.register_beacon().await;
}
let mut checkpoint_conn = CheckpointConnection::new();
checkpoint_conn.ensure_open(&pool);
let mut fts_maintenance_config = crate::fts_maintenance::FtsMaintenanceConfig::from_env();
fts_maintenance_config.enabled &= is_main;
let mut fts_maintenance_state =
crate::fts_maintenance::FtsMaintenanceState::new(Instant::now());
loop {
tokio::select! {
_ = interval.tick() => {}
_ = shutdown_rx.changed() => break,
}
#[cfg(unix)]
truncate_state.begin_tick();
#[cfg(unix)]
let mut pending_sidecar_attribution = None;
let tick = if checkpoint_conn.ensure_open(&pool).is_none() {
note_checkpoint_skipped();
CheckpointTick::Skipped
} else {
let conn = checkpoint_conn
.conn
.take()
.expect("ensure_open just confirmed a connection is open");
match off_worker::run_checkpoint_core_off_worker(
Arc::clone(&pool),
conn,
config.clone(),
truncate_state,
)
.await
{
Ok((conn, state, Ok(outcome))) => {
truncate_state = state;
#[cfg(unix)]
{
pending_sidecar_attribution = outcome.sidecar_attribution;
}
#[cfg(not(unix))]
let _ = outcome.sidecar_attribution;
let fts_result = if !fts_maintenance_state
.is_due(&fts_maintenance_config, Instant::now())
{
checkpoint_conn.conn = Some(conn);
Ok(None)
} else {
match run_fts_maintenance_off_worker(
conn,
fts_maintenance_config.clone(),
fts_maintenance_state,
Instant::now(),
)
.await
{
Ok((conn, state, fts_result)) => {
fts_maintenance_state = state;
checkpoint_conn.conn = Some(conn);
fts_result
}
Err(join_err) => {
fts_maintenance_state =
crate::fts_maintenance::FtsMaintenanceState::new(Instant::now());
Err(format!(
"bounded FTS5 segment maintenance task panicked: {join_err}"
))
}
}
};
match fts_result {
Ok(Some(step)) => match step.outcome {
crate::fts_maintenance::FtsMaintenanceOutcome::Worked => {
tracing::info!(
table = step.table,
requested_pages = step.requested_pages,
segments_before = step.segments_before,
segments_after = step.segments_after,
"bounded FTS5 segment maintenance made progress"
);
}
crate::fts_maintenance::FtsMaintenanceOutcome::Busy => {
tracing::debug!(
table = step.table,
requested_pages = step.requested_pages,
segments = step.segments_before,
"bounded FTS5 segment maintenance skipped a busy writer"
);
}
crate::fts_maintenance::FtsMaintenanceOutcome::Noop
| crate::fts_maintenance::FtsMaintenanceOutcome::BelowThreshold => {
tracing::debug!(
table = step.table,
outcome = ?step.outcome,
segments = step.segments_before,
"bounded FTS5 segment maintenance had no work"
);
}
},
Ok(None) => {}
Err(error) => {
tracing::warn!(
error = %error,
"bounded FTS5 segment maintenance failed"
);
}
}
match outcome.wal_pages {
Some(wal_pages) => CheckpointTick::Observed(wal_pages),
None => {
note_checkpoint_skipped();
CheckpointTick::Skipped
}
}
}
Ok((_conn, state, Err(e))) => {
truncate_state = state;
tracing::warn!(
error = %e,
"dedicated checkpoint connection failed a pragma; \
dropping it for a fresh reopen next tick"
);
note_checkpoint_skipped();
CheckpointTick::Skipped
}
Err(panicked) => {
truncate_state = panicked.truncate_state;
tracing::warn!(
error = %panicked.join_error,
"WAL checkpoint cycle panicked on its blocking thread; \
dropping the connection for a fresh reopen next tick"
);
note_checkpoint_skipped();
CheckpointTick::Skipped
}
}
};
#[cfg(unix)]
if let Err(error) =
complete_walpin_attribution(pending_sidecar_attribution, &mut truncate_state).await
{
tracing::warn!(
error = %error,
failure_kind = error.kind(),
"ADR-091 Amendment 2 Plank B: no-progress sidecar attribution failed"
);
}
let oldest_tx = tx_filter
.as_ref()
.and_then(khive_storage::tx_registry::oldest_for);
for emission in tx_age_state.observe(
oldest_tx.as_ref().map(|s| (s.id, s.age, s.label.clone())),
config.tx_warn_secs,
config.tx_max_age_secs,
) {
log_tx_age_emission(&emission);
}
#[cfg(unix)]
if let Some(sidecar) = walpin_state.as_mut() {
sidecar
.observe(oldest_tx.clone(), config.tx_warn_secs)
.await;
let _ = run_walpin_housekeeping_if_due(
sidecar,
&mut truncate_state,
legacy_walpin_fallback_interval,
)
.await;
}
let wal_pages = match tick {
CheckpointTick::Skipped => continue,
CheckpointTick::Observed(n) => n,
};
let above_warn = wal_pages >= config.warn_pages;
let above_high_water = wal_pages >= config.high_water_pages;
let above_truncate_high_water = wal_pages >= config.truncate_high_water_pages;
note_checkpoint_pressure_observation(above_warn, was_observed_above_warn);
was_observed_above_warn = above_warn;
log_tx_registry_oldest_debug(wal_pages, oldest_tx.as_ref());
for emission in severity_state.observe_wal_pages(wal_pages, &config) {
match emission.rung {
CheckpointSeverityRung::Info => {
log_tx_registry_oldest_warn(wal_pages, oldest_tx.as_ref());
tracing::info!(
wal_pages = emission.wal_pages,
warn_threshold = emission.threshold_pages,
"WAL page count crossed warn threshold"
);
}
CheckpointSeverityRung::Warn => {
tracing::warn!(
wal_pages = emission.wal_pages,
warn_threshold = emission.threshold_pages,
consecutive_cycles = emission.consecutive_cycles,
"WAL page count failed to drain below warn threshold"
);
}
CheckpointSeverityRung::Alarm => {
}
}
}
let high_water_crossed = crossing_warn(above_high_water, &mut was_above_high_water);
if high_water_crossed {
log_tx_registry_snapshot_warn(wal_pages);
log_wal_high_water_warn(
wal_pages,
config.high_water_pages,
oldest_tx.as_ref(),
config.tx_warn_secs,
);
}
observe_checkpoint_pressure_tick(
above_warn,
wal_pages,
above_high_water,
above_truncate_high_water,
&config,
&mut event_elevation_open,
&mut pressure_episode,
&mut pending_recovery,
|payload| {
lifecycle_emitter
.try_emit(khive_types::EventKind::CheckpointOutcomeRecorded, payload)
},
);
}
lifecycle_emitter.shutdown().await;
#[cfg(unix)]
if let Some(sidecar) = walpin_state.as_mut() {
sidecar.shutdown().await;
}
}
fn checkpoint_outcome_should_emit(above_warn: bool, was_elevated: bool) -> bool {
above_warn != was_elevated
}
#[allow(clippy::too_many_arguments)]
fn observe_checkpoint_pressure_tick(
above_warn: bool,
wal_pages: u64,
above_high_water: bool,
above_truncate_high_water: bool,
config: &CheckpointConfig,
event_elevation_open: &mut bool,
pressure_episode: &mut Option<CheckpointPressureEpisode>,
pending_recovery: &mut Option<khive_storage::CheckpointOutcomeRecordedPayload>,
mut try_emit: impl FnMut(khive_storage::CheckpointOutcomeRecordedPayload) -> bool,
) {
let pending_blocks_emission = if let Some(payload) = pending_recovery.clone() {
if try_emit(payload) {
*pending_recovery = None;
false
} else {
true
}
} else {
false
};
if above_warn {
match pressure_episode.as_mut() {
Some(episode) => episode.observe(wal_pages),
None => *pressure_episode = Some(CheckpointPressureEpisode::start(wal_pages)),
}
} else if !*event_elevation_open {
if pending_blocks_emission && pressure_episode.is_some() {
CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.fetch_add(1, Ordering::Relaxed);
tracing::warn!(
wal_pages,
"checkpoint pressure episode elapsed unreported behind an undelivered recovery summary"
);
}
*pressure_episode = None;
}
if pending_blocks_emission || !checkpoint_outcome_should_emit(above_warn, *event_elevation_open)
{
return;
}
let Some(episode) = *pressure_episode else {
tracing::warn!(
above_warn,
event_elevation_open = *event_elevation_open,
"checkpoint pressure transition has no episode aggregate"
);
return;
};
let payload = khive_storage::CheckpointOutcomeRecordedPayload {
wal_pages,
warn_pages: config.warn_pages,
high_water_pages: config.high_water_pages,
truncate_high_water_pages: config.truncate_high_water_pages,
above_warn,
above_high_water,
above_truncate_high_water,
episode_elevated_ticks: Some(episode.elevated_ticks),
episode_peak_wal_pages: Some(episode.peak_wal_pages),
};
if try_emit(payload.clone()) {
*event_elevation_open = above_warn;
if !above_warn {
*pressure_episode = None;
}
} else if !above_warn {
debug_assert!(
pending_recovery.is_none(),
"recovery emission attempted while an earlier summary was still pending"
);
*event_elevation_open = false;
*pressure_episode = None;
*pending_recovery = Some(payload);
}
}
fn log_tx_registry_oldest_debug(
wal_pages: u64,
oldest: Option<&khive_storage::tx_registry::OldestSpan>,
) {
if let Some(span) = oldest {
tracing::debug!(
wal_pages,
oldest_tx_age_secs = span.age.as_secs_f64(),
oldest_tx_label = span.label.as_deref().unwrap_or("<unlabeled>"),
"WAL checkpoint tick: oldest open transaction registry entry"
);
}
}
fn log_tx_registry_oldest_warn(
wal_pages: u64,
oldest: Option<&khive_storage::tx_registry::OldestSpan>,
) {
if let Some(span) = oldest {
tracing::warn!(
wal_pages,
oldest_tx_age_secs = span.age.as_secs_f64(),
oldest_tx_label = span.label.as_deref().unwrap_or("<unlabeled>"),
"WAL checkpoint tick: oldest open transaction registry entry"
);
}
}
fn log_tx_registry_snapshot_warn(wal_pages: u64) {
log_tx_registry_entries_warn(wal_pages, &khive_storage::tx_registry::snapshot());
}
fn log_tx_registry_entries_warn(wal_pages: u64, snapshot: &[(Duration, Option<String>)]) {
for (age, label) in snapshot {
tracing::warn!(
wal_pages,
tx_age_secs = age.as_secs_f64(),
tx_label = label.as_deref().unwrap_or("<unlabeled>"),
"WAL high-water: open transaction registry entry"
);
}
}
fn log_truncate_no_progress_warn(
wal_pages_before: u64,
wal_pages_after: u64,
snapshot: &[(Duration, Option<String>)],
) {
let open_tx_count = snapshot.len();
let oldest_tx_age_secs = snapshot
.iter()
.map(|(age, _)| *age)
.max()
.map(|age| age.as_secs_f64());
if snapshot.is_empty() {
tracing::warn!(
wal_pages_before,
wal_pages_after,
open_tx_count,
oldest_tx_age_secs = ?oldest_tx_age_secs,
"WAL TRUNCATE attempt made no progress; no open transaction in this process's registry"
);
} else {
tracing::warn!(
wal_pages_before,
wal_pages_after,
open_tx_count,
oldest_tx_age_secs = ?oldest_tx_age_secs,
"WAL TRUNCATE attempt made no progress; open transactions observed in this process's registry"
);
}
log_tx_registry_entries_warn(wal_pages_after, snapshot);
}
fn log_wal_high_water_warn(
wal_pages: u64,
high_water: u64,
oldest: Option<&khive_storage::tx_registry::OldestSpan>,
warn_after: Duration,
) {
match oldest.filter(|span| span.age >= warn_after) {
Some(span) => tracing::warn!(
wal_pages,
high_water,
oldest_tx_age_secs = span.age.as_secs_f64(),
oldest_tx_label = span.label.as_deref().unwrap_or("<unlabeled>"),
"WAL high-water mark exceeded; an in-process registered transaction is older \
than the age threshold and may hold a snapshot"
),
None => tracing::warn!(
wal_pages,
high_water,
oldest_tx_age_secs = ?oldest.map(|span| span.age.as_secs_f64()),
oldest_tx_label = oldest
.and_then(|span| span.label.as_deref())
.unwrap_or("<none>"),
"WAL high-water mark exceeded; no in-process transaction older than the age \
threshold is visible; a reader in another process may hold the snapshot, \
or writes may outpace PASSIVE checkpoints"
),
}
}
#[derive(Debug)]
#[must_use]
struct CheckpointCoreOutcome {
wal_pages: Option<u64>,
unavailable_reason: Option<CheckpointUnavailableReason>,
sidecar_attribution: Option<WalpinAttributionRequest>,
}
pub fn checkpoint_once(
pool: &ConnectionPool,
conn: &rusqlite::Connection,
config: &CheckpointConfig,
truncate_state: &mut TruncateState,
) -> Result<u64, rusqlite::Error> {
let outcome = checkpoint_once_core(pool, conn, config, truncate_state)?;
outcome.wal_pages.ok_or_else(|| match outcome
.unavailable_reason
.expect("unavailable core outcome has a reason")
{
CheckpointUnavailableReason::Busy => rusqlite::Error::SqliteFailure(
rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_BUSY),
Some("PASSIVE checkpoint returned a busy row without a WAL frame observation".into()),
),
CheckpointUnavailableReason::InconsistentFrames => rusqlite::Error::SqliteFailure(
rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_ERROR),
Some("PASSIVE checkpoint returned an inconsistent frame pair without a WAL frame observation".into()),
),
})
}
fn checkpoint_once_core(
pool: &ConnectionPool,
conn: &rusqlite::Connection,
config: &CheckpointConfig,
truncate_state: &mut TruncateState,
) -> Result<CheckpointCoreOutcome, rusqlite::Error> {
#[cfg(unix)]
truncate_state.begin_tick();
let started = Instant::now();
let checkpoint_result = query_routine_checkpoint_observation(pool, conn);
let elapsed_us = started.elapsed().as_micros().min(u128::from(u64::MAX)) as u64;
record_checkpoint_timing(
pool,
elapsed_us,
checkpoint_result
.as_ref()
.ok()
.map(|observation| observation.busy),
);
let raw_observation = match checkpoint_result {
Ok(observation) => observation,
Err(e) => {
record_checkpoint_run_result(pool, None);
tracing::warn!(error = %e, elapsed_us, "WAL checkpoint failed");
return Err(e);
}
};
record_checkpoint_run_result(
pool,
Some((
raw_observation.busy,
raw_observation.log_frames,
raw_observation.checkpointed_frames,
)),
);
let wal_pages = match observed_wal_pages(raw_observation) {
Ok(wal_pages) => wal_pages,
Err(reason) => {
match reason {
CheckpointUnavailableReason::Busy => tracing::debug!(
busy = raw_observation.busy,
wal_log_frames = raw_observation.log_frames,
wal_checkpointed_frames = raw_observation.checkpointed_frames,
elapsed_us,
"WAL PASSIVE checkpoint returned a busy row; frame observation unavailable"
),
CheckpointUnavailableReason::InconsistentFrames => tracing::warn!(
busy = raw_observation.busy,
wal_log_frames = raw_observation.log_frames,
wal_checkpointed_frames = raw_observation.checkpointed_frames,
elapsed_us,
"WAL PASSIVE checkpoint returned an inconsistent frame pair; frame observation unavailable"
),
}
return Ok(CheckpointCoreOutcome {
wal_pages: None,
unavailable_reason: Some(reason),
sidecar_attribution: None,
});
}
};
let observation = record_routine_wal_observation(pool, raw_observation);
LAST_WAL_PAGES.store(wal_pages, Ordering::Relaxed);
note_checkpoint_observed(wal_pages);
tracing::debug!(
wal_pages,
elapsed_us,
busy = raw_observation.busy,
wal_checkpointed_frames = observation.checkpointed_frames,
wal_pending_frames = observation.pending_frames,
wal_physical_bytes = ?observation.physical_wal_bytes,
"WAL checkpoint issued"
);
let sidecar_attribution = maybe_truncate(pool, conn, config, wal_pages, truncate_state);
Ok(CheckpointCoreOutcome {
wal_pages: Some(wal_pages),
unavailable_reason: None,
sidecar_attribution,
})
}
fn truncate_needs_attribution(wal_pages_before: u64, wal_pages_after: Option<u64>) -> bool {
wal_pages_after.is_none_or(|pages| pages >= wal_pages_before)
}
fn maybe_truncate(
pool: &ConnectionPool,
conn: &rusqlite::Connection,
config: &CheckpointConfig,
wal_pages_before: u64,
truncate_state: &mut TruncateState,
) -> Option<WalpinAttributionRequest> {
if wal_pages_before < config.truncate_high_water_pages {
return None;
}
if let Some(last) = truncate_state.last_attempt {
if last.elapsed() < config.truncate_min_interval {
return None;
}
}
log_tx_registry_snapshot_warn(wal_pages_before);
let original_busy_timeout = pool.config().busy_timeout;
if let Err(e) = conn.busy_timeout(config.truncate_busy_timeout) {
tracing::warn!(error = %e, "failed to lower busy_timeout for TRUNCATE attempt; skipping");
return None;
}
#[cfg(unix)]
let mut holder_attribution = capture_walpin_attribution_request(pool, truncate_state);
#[cfg(unix)]
let mut sidecar_attribution = None;
#[cfg(not(unix))]
let sidecar_attribution = None;
truncate_state.last_attempt = Some(Instant::now());
#[cfg(test)]
off_worker::cycle_panic_seam::after_attempt_decided(pool.canonical_path());
let start = Instant::now();
let outcome = query_truncate_observation(conn);
record_checkpoint_run_result(
pool,
outcome.as_ref().ok().map(|observation| {
(
observation.busy,
observation.log_frames,
observation.checkpointed_frames,
)
}),
);
let elapsed = start.elapsed();
if let Err(e) = conn.busy_timeout(original_busy_timeout) {
tracing::warn!(error = %e, "failed to restore busy_timeout after TRUNCATE attempt");
}
match outcome {
Ok(_) => {
let wal_pages_after = query_wal_pages(pool, conn);
if let Some(pages) = wal_pages_after {
tracing::info!(
wal_pages_before,
wal_pages_after = pages,
elapsed_ms = elapsed.as_millis() as u64,
"WAL TRUNCATE checkpoint attempted"
);
} else {
tracing::info!(
wal_pages_before,
wal_pages_after_unavailable = true,
elapsed_ms = elapsed.as_millis() as u64,
"WAL TRUNCATE checkpoint attempted"
);
}
if truncate_needs_attribution(wal_pages_before, wal_pages_after) {
let snapshot = khive_storage::tx_registry::snapshot();
if let Some(pages) = wal_pages_after {
log_truncate_no_progress_warn(wal_pages_before, pages, &snapshot);
} else {
tracing::warn!(
wal_pages_before,
"WAL TRUNCATE progress unmeasured; checking possible holders in this process and others"
);
log_tx_registry_entries_warn(wal_pages_before, &snapshot);
}
#[cfg(test)]
if let Some(path) = pool.canonical_path() {
truncate_report_test_sync::after_no_progress_before_report(path);
}
#[cfg(unix)]
{
sidecar_attribution = holder_attribution.take();
}
log_backfill_gap(pool, conn);
}
note_truncate_outcome(config, wal_pages_after, truncate_state);
}
Err(e) => {
tracing::warn!(error = %e, wal_pages_before, "WAL TRUNCATE attempt failed");
log_tx_registry_snapshot_warn(wal_pages_before);
note_truncate_outcome(config, Some(wal_pages_before), truncate_state);
}
}
#[cfg(unix)]
if let Some(WalpinAttributionRequest::Fresh {
previous_last_attempt,
..
}) = holder_attribution.as_ref()
{
truncate_state.restore_walpin_full_scan_reservation(*previous_last_attempt);
}
sidecar_attribution
}
#[cfg(test)]
mod truncate_report_test_sync {
use std::path::{Path, PathBuf};
use std::sync::mpsc::{sync_channel, Receiver, SyncSender};
use std::sync::Mutex;
struct Hook {
db_path: PathBuf,
reached_tx: SyncSender<()>,
proceed_rx: Receiver<()>,
}
static HOOK: Mutex<Option<Hook>> = Mutex::new(None);
pub(crate) fn install(db_path: PathBuf) -> (Receiver<()>, SyncSender<()>) {
let (reached_tx, reached_rx) = sync_channel(0);
let (proceed_tx, proceed_rx) = sync_channel(0);
let replaced = HOOK
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.replace(Hook {
db_path,
reached_tx,
proceed_rx,
});
assert!(replaced.is_none(), "truncate report hook already installed");
(reached_rx, proceed_tx)
}
pub(crate) fn uninstall() {
*HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
}
pub(crate) fn after_no_progress_before_report(db_path: &Path) {
let hook = {
let mut guard = HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
match guard.as_ref() {
Some(hook) if hook.db_path == db_path => guard.take(),
_ => None,
}
};
let Some(hook) = hook else {
return;
};
let _ = hook.reached_tx.send(());
let _ = hook.proceed_rx.recv();
}
}
#[cfg(all(test, unix))]
mod walpin_attribution_test_sync {
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::mpsc::{sync_channel, Receiver, SyncSender};
use std::sync::{Arc, Mutex};
enum Behavior {
Pause {
reached_tx: tokio::sync::oneshot::Sender<std::thread::ThreadId>,
proceed_rx: Receiver<()>,
},
Panic,
}
struct Hook {
dir: PathBuf,
behavior: Behavior,
}
static HOOK: Mutex<Option<Hook>> = Mutex::new(None);
static REPORT_COUNTER: Mutex<Option<Arc<AtomicUsize>>> = Mutex::new(None);
pub(crate) fn install_pause(
dir: PathBuf,
) -> (
tokio::sync::oneshot::Receiver<std::thread::ThreadId>,
SyncSender<()>,
Arc<AtomicUsize>,
) {
let (reached_tx, reached_rx) = tokio::sync::oneshot::channel();
let (proceed_tx, proceed_rx) = sync_channel(0);
let report_counter = Arc::new(AtomicUsize::new(0));
let replaced = HOOK
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.replace(Hook {
dir,
behavior: Behavior::Pause {
reached_tx,
proceed_rx,
},
});
assert!(
replaced.is_none(),
"walpin attribution hook already installed"
);
*REPORT_COUNTER
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner()) = Some(Arc::clone(&report_counter));
(reached_rx, proceed_tx, report_counter)
}
pub(crate) fn install_panic(dir: PathBuf) {
let replaced = HOOK
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.replace(Hook {
dir,
behavior: Behavior::Panic,
});
assert!(
replaced.is_none(),
"walpin attribution hook already installed"
);
}
pub(crate) fn uninstall() {
*HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
*REPORT_COUNTER
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner()) = None;
}
pub(crate) fn before_enumeration(dir: &Path) {
let hook = {
let mut guard = HOOK.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
match guard.as_ref() {
Some(hook) if hook.dir == dir => guard.take(),
_ => None,
}
};
let Some(hook) = hook else {
return;
};
match hook.behavior {
Behavior::Pause {
reached_tx,
proceed_rx,
} => {
if reached_tx.send(std::thread::current().id()).is_ok() {
let _ = proceed_rx.recv();
}
}
Behavior::Panic => panic!("injected walpin attribution worker panic"),
}
}
pub(crate) fn report_used() {
if let Some(counter) = REPORT_COUNTER
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.as_ref()
{
counter.fetch_add(1, Ordering::SeqCst);
}
}
}
fn note_truncate_outcome(
config: &CheckpointConfig,
wal_pages_after: Option<u64>,
state: &mut TruncateState,
) {
TRUNCATE_ATTEMPTS.fetch_add(1, Ordering::Relaxed);
if let Some(wal_pages_after) = wal_pages_after {
if wal_pages_after >= config.warn_pages {
state.consecutive_failures = state.consecutive_failures.saturating_add(1);
if state.consecutive_failures == 3 {
tracing::warn!(
wal_pages_after,
warn_threshold = config.warn_pages,
"WAL TRUNCATE has failed to clear WAL pressure for 3 consecutive attempts"
);
}
} else {
state.consecutive_failures = 0;
}
}
TRUNCATE_CONSECUTIVE_FAILURES.store(state.consecutive_failures as u64, Ordering::Relaxed);
}
#[cfg(unix)]
#[derive(Debug)]
enum WalpinAttributionRequest {
Fresh {
dir: PathBuf,
census: Result<crate::walpin::CensusResult, String>,
legacy_fallback_interval: Duration,
previous_last_attempt: Option<Instant>,
},
Cached(CachedWalpinAttribution),
Suppressed,
}
#[cfg(unix)]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum WalpinReportFreshness {
Fresh,
Cached { age: Duration },
}
#[cfg(unix)]
impl WalpinReportFreshness {
fn is_fresh(self) -> bool {
self == Self::Fresh
}
}
#[cfg(not(unix))]
type WalpinAttributionRequest = ();
#[cfg(unix)]
#[derive(Debug, Clone, PartialEq, Eq)]
enum WalpinAttributionFailure {
Enumeration(String),
Worker(String),
}
#[cfg(unix)]
impl WalpinAttributionFailure {
fn kind(&self) -> &'static str {
match self {
Self::Enumeration(_) => "enumeration",
Self::Worker(_) => "blocking_worker",
}
}
}
#[cfg(unix)]
impl std::fmt::Display for WalpinAttributionFailure {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Enumeration(error) => write!(
formatter,
"sidecar directory failed the trust-boundary enumeration; cross-process \
WAL-pin attribution is unestablished for this tick: {error}"
),
Self::Worker(error) => write!(
formatter,
"sidecar attribution blocking worker failed; cross-process WAL-pin \
attribution is unestablished for this tick: {error}"
),
}
}
}
#[cfg(unix)]
fn capture_walpin_attribution_request(
pool: &ConnectionPool,
state: &mut TruncateState,
) -> Option<WalpinAttributionRequest> {
let path = pool.canonical_path()?;
if !crate::walpin::sidecar_enabled(true) {
return None;
}
let legacy_fallback_interval = state.legacy_walpin_fallback_interval;
Some(match state.plan_walpin_attribution_at(Instant::now()) {
WalpinFullScanPlan::Refresh {
previous_last_attempt,
} => WalpinAttributionRequest::Fresh {
dir: crate::walpin::sidecar_dir_for(path),
census: crate::walpin::census_holders(path).map_err(|error| error.to_string()),
legacy_fallback_interval,
previous_last_attempt,
},
WalpinFullScanPlan::Cached(cached) => WalpinAttributionRequest::Cached(cached),
WalpinFullScanPlan::Suppressed => WalpinAttributionRequest::Suppressed,
})
}
#[cfg(unix)]
async fn complete_walpin_attribution(
request: Option<WalpinAttributionRequest>,
state: &mut TruncateState,
) -> Result<bool, WalpinAttributionFailure> {
let Some(request) = request else {
return Ok(false);
};
match request {
WalpinAttributionRequest::Suppressed => Ok(false),
WalpinAttributionRequest::Cached(cached) => {
log_walpin_sidecar_report(
&cached.report,
cached.census,
WalpinReportFreshness::Cached {
age: Instant::now().saturating_duration_since(cached.captured_at),
},
);
Ok(true)
}
WalpinAttributionRequest::Fresh {
dir,
census,
legacy_fallback_interval,
previous_last_attempt: _,
} => {
state.sidecar_attribution_attempted_this_tick = true;
if state.walpin_full_scan_last_attempt.is_none() {
state.walpin_full_scan_last_attempt = Some(Instant::now());
}
let fallback = state.walpin_cached_attribution.clone();
let result = tokio::task::spawn_blocking(move || {
#[cfg(test)]
walpin_attribution_test_sync::before_enumeration(&dir);
crate::walpin::enumerate_live(&dir, legacy_fallback_interval)
})
.await
.map_err(|error| WalpinAttributionFailure::Worker(error.to_string()))
.and_then(|result| {
result.map_err(|error| WalpinAttributionFailure::Enumeration(error.to_string()))
});
match result {
Ok(report) => {
let captured_at = Instant::now();
log_walpin_sidecar_report(
&report,
census.clone(),
WalpinReportFreshness::Fresh,
);
state.cache_walpin_attribution(report, census, captured_at);
Ok(true)
}
Err(error) => {
if let Some(cached) = fallback {
log_walpin_sidecar_report(
&cached.report,
cached.census,
WalpinReportFreshness::Cached {
age: Instant::now().saturating_duration_since(cached.captured_at),
},
);
}
Err(error)
}
}
}
}
}
#[cfg(unix)]
fn log_walpin_sidecar_report(
report: &crate::walpin::WalpinReport,
census: Result<crate::walpin::CensusResult, String>,
freshness: WalpinReportFreshness,
) {
#[cfg(test)]
walpin_attribution_test_sync::report_used();
let now = now_epoch_secs();
for hb in report.reporting() {
tracing::warn!(
walpin_pid = hb.pid,
walpin_role = %hb.process_role,
walpin_oldest_tx_age_secs = hb.current_oldest_tx_age_secs(now),
walpin_oldest_tx_label = hb.oldest_tx_label.as_deref().unwrap_or("<unlabeled>"),
walpin_attribution_basis = hb.attribution_basis.as_deref().unwrap_or("<unspecified>"),
walpin_attribution_evidence_backed = hb.attribution_is_evidence_backed(),
walpin_attribution_fresh = freshness.is_fresh(),
walpin_health = "reporting",
"ADR-091 Amendment 2 Plank B: live cross-process WAL-pin attribution report"
);
}
for pid in report.registered_silent_pids() {
tracing::debug!(
walpin_pid = pid,
walpin_health = "registered_silent",
walpin_attribution_fresh = freshness.is_fresh(),
"ADR-091 Amendment 2 Plank B: process affirmatively reports no over-threshold span"
);
}
let mut unknown_pids: Vec<u32> = report.unknown_pids().collect();
if let WalpinReportFreshness::Cached { age } = freshness {
tracing::warn!(
walpin_cache_age_ms = age.as_millis() as u64,
"cached WAL-pin attribution is diagnostic-only; fully-attributed \
conclusion is not licensed"
);
unknown_pids.push(0);
}
match census {
Ok(census) => {
let sidecar_known: std::collections::HashSet<u32> = report
.reporting()
.map(|hb| hb.pid)
.chain(report.registered_silent_pids())
.chain(unknown_pids.iter().copied())
.collect();
let mut census_only: Vec<u32> =
census.holders.difference(&sidecar_known).copied().collect();
if !census_only.is_empty() {
census_only.sort_unstable();
tracing::warn!(
?census_only,
"ADR-091 Amendment 2: these PIDs hold the database file open \
at the OS level but have no sidecar data at all (pre-feature binary, \
sidecar disabled, or wedged before its first write)"
);
unknown_pids.extend(census_only);
}
if !census.is_complete() {
let mut uninspectable = census.uninspectable_pids.clone();
uninspectable.sort_unstable();
tracing::warn!(
?uninspectable,
truncated = census.truncated,
"ADR-091 Amendment 2: the OS-derived holder census is \
INCOMPLETE — either specific PIDs' open file descriptors could not be \
inspected (permission denied, or a listing race), or the enumeration walk \
itself has positive evidence it did not see the full live-process universe \
(namespace/visibility check, directory-iterator error, self-canary, or a \
libproc buffer that stayed at capacity after bounded retries) — cannot \
rule out an unregistered holder"
);
if uninspectable.is_empty() {
unknown_pids.push(0);
} else {
unknown_pids.extend(uninspectable);
}
}
}
Err(e) => {
tracing::warn!(
error = %e,
"ADR-091 Amendment 2: OS-derived holder census failed; \
attribution cannot rule out an unregistered database holder this tick"
);
unknown_pids.push(0);
}
}
unknown_pids.sort_unstable();
unknown_pids.dedup();
if !unknown_pids.is_empty() {
tracing::warn!(
?unknown_pids,
"ADR-091 Amendment 2 Plank B: sidecar health unestablished for these PIDs; \
attribution is inconclusive and the native/unregistered-mechanism conclusion \
is NOT licensed this tick"
);
} else if report.reporting().next().is_none() {
tracing::info!(
"ADR-091 Amendment 2 Plank B: every live PID is reporting or registered-silent \
with none pinning; the WAL pin is not attributable to any in-process registry \
span this sidecar covers"
);
}
}
fn log_backfill_gap(pool: &ConnectionPool, conn: &rusqlite::Connection) {
match query_backfill_gap(conn) {
Ok(observation) => {
record_checkpoint_run_result(
pool,
Some((
observation.busy,
observation.log_frames,
observation.checkpointed_frames,
)),
);
if observed_wal_pages(observation).is_ok() {
tracing::warn!(
busy = observation.busy,
wal_log_frames = observation.log_frames,
wal_checkpointed_frames = observation.checkpointed_frames,
backfill_gap_frames = observation
.log_frames
.saturating_sub(observation.checkpointed_frames)
.max(0),
"ADR-091 Plank C: WAL backfill gap after TRUNCATE no-progress"
);
} else {
tracing::warn!(
busy = observation.busy,
wal_log_frames = observation.log_frames,
wal_checkpointed_frames = observation.checkpointed_frames,
"ADR-091 Plank C: WAL backfill gap unavailable after TRUNCATE"
);
}
}
Err(e) => {
record_checkpoint_run_result(pool, None);
tracing::warn!(
error = %e,
"ADR-091 Plank C: failed to query WAL backfill gap"
);
}
}
}
fn query_backfill_gap(conn: &rusqlite::Connection) -> rusqlite::Result<RawCheckpointObservation> {
conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
Ok(RawCheckpointObservation {
busy: row.get(0)?,
log_frames: row.get(1)?,
checkpointed_frames: row.get(2)?,
})
})
}
fn crossing_warn(now_above: bool, was_above: &mut bool) -> bool {
let fire = now_above && !*was_above;
*was_above = now_above;
fire
}
#[derive(Debug, Clone, Copy)]
struct RawCheckpointObservation {
busy: i64,
log_frames: i64,
checkpointed_frames: i64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum CheckpointUnavailableReason {
Busy,
InconsistentFrames,
}
fn observed_wal_pages(
observation: RawCheckpointObservation,
) -> Result<u64, CheckpointUnavailableReason> {
if observation.busy != 0 {
return Err(CheckpointUnavailableReason::Busy);
}
if observation.log_frames == -1 && observation.checkpointed_frames == -1 {
return Ok(0);
}
if observation.log_frames >= 0
&& observation.checkpointed_frames >= 0
&& observation.checkpointed_frames <= observation.log_frames
{
Ok(observation.log_frames as u64)
} else {
Err(CheckpointUnavailableReason::InconsistentFrames)
}
}
fn query_checkpoint_observation(
conn: &rusqlite::Connection,
) -> rusqlite::Result<RawCheckpointObservation> {
conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
Ok(RawCheckpointObservation {
busy: row.get(0)?,
log_frames: row.get(1)?,
checkpointed_frames: row.get(2)?,
})
})
}
fn query_routine_checkpoint_observation(
pool: &ConnectionPool,
conn: &rusqlite::Connection,
) -> rusqlite::Result<RawCheckpointObservation> {
#[cfg(test)]
if let Some(row) = test_take_passive_row(pool) {
return Ok(row);
}
#[cfg(not(test))]
let _ = pool;
query_checkpoint_observation(conn)
}
#[cfg(test)]
static TEST_PASSIVE_ROWS: OnceLock<Mutex<HashMap<Option<PathBuf>, RawCheckpointObservation>>> =
OnceLock::new();
#[cfg(test)]
fn test_passive_rows() -> &'static Mutex<HashMap<Option<PathBuf>, RawCheckpointObservation>> {
TEST_PASSIVE_ROWS.get_or_init(|| Mutex::new(HashMap::new()))
}
#[cfg(test)]
fn test_arm_passive_row(pool: &ConnectionPool, row: RawCheckpointObservation) {
test_passive_rows()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(checkpoint_db_key(pool), row);
}
#[cfg(test)]
fn test_take_passive_row(pool: &ConnectionPool) -> Option<RawCheckpointObservation> {
test_passive_rows()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&checkpoint_db_key(pool))
}
fn query_truncate_observation(
conn: &rusqlite::Connection,
) -> rusqlite::Result<RawCheckpointObservation> {
conn.query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |row| {
Ok(RawCheckpointObservation {
busy: row.get(0)?,
log_frames: row.get(1)?,
checkpointed_frames: row.get(2)?,
})
})
}
fn query_wal_pages(pool: &ConnectionPool, conn: &rusqlite::Connection) -> Option<u64> {
let observation = query_checkpoint_observation(conn);
record_checkpoint_run_result(
pool,
observation.as_ref().ok().map(|observation| {
(
observation.busy,
observation.log_frames,
observation.checkpointed_frames,
)
}),
);
let pages = observation
.ok()
.and_then(|row| observed_wal_pages(row).ok());
if let Some(pages) = pages {
LAST_WAL_PAGES.store(pages, Ordering::Relaxed);
note_checkpoint_observed(pages);
}
pages
}
#[cfg(test)]
#[path = "checkpoint_tests.rs"]
mod tests;