use std::collections::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;
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);
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RoutineWalObservation {
pub busy: i64,
pub log_frames: u64,
pub checkpointed_frames: u64,
pub pending_frames: u64,
pub physical_wal_bytes: Option<u64>,
pub observed_at_unix_ms: u64,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(default)]
pub struct CheckpointTiming {
pub ticks: u64,
pub elapsed_us_sum: u64,
pub elapsed_us_max: u64,
pub busy_ticks: u64,
pub error_ticks: u64,
}
static CHECKPOINT_TIMINGS: OnceLock<Mutex<HashMap<Option<PathBuf>, CheckpointTiming>>> =
OnceLock::new();
fn checkpoint_timings() -> &'static Mutex<HashMap<Option<PathBuf>, CheckpointTiming>> {
CHECKPOINT_TIMINGS.get_or_init(|| Mutex::new(HashMap::new()))
}
fn record_checkpoint_timing(pool: &ConnectionPool, elapsed_us: u64, busy: Option<i64>) {
let mut timings = checkpoint_timings()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let timing = timings.entry(checkpoint_db_key(pool)).or_default();
timing.ticks = timing.ticks.saturating_add(1);
timing.elapsed_us_sum = timing.elapsed_us_sum.saturating_add(elapsed_us);
timing.elapsed_us_max = timing.elapsed_us_max.max(elapsed_us);
timing.busy_ticks = timing
.busy_ticks
.saturating_add(u64::from(busy.is_some_and(|value| value != 0)));
timing.error_ticks = timing.error_ticks.saturating_add(u64::from(busy.is_none()));
}
pub fn checkpoint_timing(pool: &ConnectionPool) -> CheckpointTiming {
checkpoint_timings()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(&checkpoint_db_key(pool))
.copied()
.unwrap_or_default()
}
static ROUTINE_WAL_OBSERVATIONS: OnceLock<Mutex<HashMap<Option<PathBuf>, RoutineWalObservation>>> =
OnceLock::new();
fn routine_wal_observations() -> &'static Mutex<HashMap<Option<PathBuf>, RoutineWalObservation>> {
ROUTINE_WAL_OBSERVATIONS.get_or_init(|| Mutex::new(HashMap::new()))
}
fn checkpoint_db_key_from_path(path: Option<&Path>) -> Option<PathBuf> {
path.map(Path::to_path_buf)
}
fn checkpoint_db_key(pool: &ConnectionPool) -> Option<PathBuf> {
checkpoint_db_key_from_path(pool.canonical_path())
}
fn observed_at_unix_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|duration| duration.as_millis() as u64)
.unwrap_or(0)
}
fn physical_wal_bytes(pool: &ConnectionPool) -> Option<u64> {
let path = pool.canonical_path()?;
let mut sidecar = path.as_os_str().to_os_string();
sidecar.push("-wal");
std::fs::metadata(PathBuf::from(sidecar))
.ok()
.map(|metadata| metadata.len())
}
fn record_routine_wal_observation(
pool: &ConnectionPool,
raw: RawCheckpointObservation,
) -> RoutineWalObservation {
let log_frames = raw.log_frames.max(0) as u64;
let checkpointed_frames = raw.checkpointed_frames.max(0) as u64;
let observation = RoutineWalObservation {
busy: raw.busy,
log_frames,
checkpointed_frames,
pending_frames: log_frames.saturating_sub(checkpointed_frames),
physical_wal_bytes: physical_wal_bytes(pool),
observed_at_unix_ms: observed_at_unix_ms(),
};
routine_wal_observations()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(checkpoint_db_key(pool), observation.clone());
observation
}
pub fn routine_wal_observation(pool: &ConnectionPool) -> Option<RoutineWalObservation> {
routine_wal_observations()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(&checkpoint_db_key(pool))
.cloned()
}
pub fn last_observed_wal_pages() -> Option<u64> {
match LAST_WAL_PAGES.load(Ordering::Relaxed) {
u64::MAX => None,
pages => Some(pages),
}
}
pub fn truncate_attempts() -> u64 {
TRUNCATE_ATTEMPTS.load(Ordering::Relaxed)
}
pub fn truncate_consecutive_failures() -> u64 {
TRUNCATE_CONSECUTIVE_FAILURES.load(Ordering::Relaxed)
}
pub fn checkpoint_skipped_ticks() -> u64 {
CHECKPOINT_SKIPPED_TICKS.load(Ordering::Relaxed)
}
pub fn checkpoint_consecutive_skips() -> u64 {
CHECKPOINT_CONSECUTIVE_SKIPS.load(Ordering::Relaxed)
}
pub fn checkpoint_last_skip_wal_pages() -> Option<u64> {
match CHECKPOINT_LAST_SKIP_WAL_PAGES.load(Ordering::Relaxed) {
u64::MAX => None,
pages => Some(pages),
}
}
pub fn checkpoint_pressure_elevated_ticks() -> u64 {
CHECKPOINT_PRESSURE_ELEVATED_TICKS.load(Ordering::Relaxed)
}
pub fn checkpoint_pressure_episodes_started() -> u64 {
CHECKPOINT_PRESSURE_EPISODES_STARTED.load(Ordering::Relaxed)
}
pub fn checkpoint_pressure_episodes_recovered() -> u64 {
CHECKPOINT_PRESSURE_EPISODES_RECOVERED.load(Ordering::Relaxed)
}
pub fn checkpoint_lifecycle_append_attempts() -> u64 {
CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.load(Ordering::Relaxed)
}
pub fn checkpoint_lifecycle_append_failures() -> u64 {
CHECKPOINT_LIFECYCLE_APPEND_FAILURES.load(Ordering::Relaxed)
}
pub fn read_tx_max_age_evictions() -> u64 {
READ_TX_MAX_AGE_EVICTIONS.load(Ordering::Relaxed)
}
pub(crate) fn note_read_tx_max_age_eviction() {
READ_TX_MAX_AGE_EVICTIONS.fetch_add(1, Ordering::Relaxed);
}
pub fn checkpoint_lifecycle_enqueue_drops() -> u64 {
CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.load(Ordering::Relaxed)
}
fn note_checkpoint_skipped() {
CHECKPOINT_SKIPPED_TICKS.fetch_add(1, Ordering::Relaxed);
CHECKPOINT_CONSECUTIVE_SKIPS.fetch_add(1, Ordering::Relaxed);
if let Some(pages) = last_observed_wal_pages() {
CHECKPOINT_LAST_SKIP_WAL_PAGES.store(pages, Ordering::Relaxed);
}
}
fn note_checkpoint_observed(_wal_pages: u64) {
CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
}
fn note_checkpoint_pressure_observation(above_warn: bool, was_above_warn: bool) {
if above_warn {
CHECKPOINT_PRESSURE_ELEVATED_TICKS.fetch_add(1, Ordering::Relaxed);
if !was_above_warn {
CHECKPOINT_PRESSURE_EPISODES_STARTED.fetch_add(1, Ordering::Relaxed);
}
} else if was_above_warn {
CHECKPOINT_PRESSURE_EPISODES_RECOVERED.fetch_add(1, Ordering::Relaxed);
}
}
#[cfg(test)]
pub(crate) fn reset_checkpoint_metrics_for_tests() {
CHECKPOINT_SKIPPED_TICKS.store(0, Ordering::Relaxed);
CHECKPOINT_CONSECUTIVE_SKIPS.store(0, Ordering::Relaxed);
CHECKPOINT_LAST_SKIP_WAL_PAGES.store(u64::MAX, Ordering::Relaxed);
CHECKPOINT_PRESSURE_ELEVATED_TICKS.store(0, Ordering::Relaxed);
CHECKPOINT_PRESSURE_EPISODES_STARTED.store(0, Ordering::Relaxed);
CHECKPOINT_PRESSURE_EPISODES_RECOVERED.store(0, Ordering::Relaxed);
CHECKPOINT_LIFECYCLE_APPEND_ATTEMPTS.store(0, Ordering::Relaxed);
CHECKPOINT_LIFECYCLE_APPEND_FAILURES.store(0, Ordering::Relaxed);
CHECKPOINT_LIFECYCLE_ENQUEUE_DROPS.store(0, Ordering::Relaxed);
READ_TX_MAX_AGE_EVICTIONS.store(0, Ordering::Relaxed);
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CheckpointTick {
Skipped,
Observed(u64),
}
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,
) {
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 checkpoint_once_core(&pool, &conn, &config, &mut truncate_state) {
Ok(outcome) => {
#[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"
);
}
}
CheckpointTick::Observed(outcome.wal_pages)
}
Err(e) => {
tracing::warn!(
error = %e,
"dedicated checkpoint connection failed a pragma; \
dropping it 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 open transaction older than the age \
threshold is pinning a snapshot PASSIVE cannot reclaim"
),
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 with no open transaction old enough to pin \
a snapshot; the WAL is growing faster than PASSIVE checkpoints reclaim it"
),
}
}
#[derive(Debug)]
#[must_use]
struct CheckpointCoreOutcome {
wal_pages: u64,
sidecar_attribution: Option<WalpinAttributionRequest>,
}
pub fn checkpoint_once(
pool: &ConnectionPool,
conn: &rusqlite::Connection,
config: &CheckpointConfig,
truncate_state: &mut TruncateState,
) -> Result<u64, rusqlite::Error> {
checkpoint_once_core(pool, conn, config, truncate_state).map(|outcome| outcome.wal_pages)
}
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_checkpoint_observation(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) => {
tracing::warn!(error = %e, elapsed_us, "WAL checkpoint failed");
return Err(e);
}
};
let observation = record_routine_wal_observation(pool, raw_observation);
let wal_pages = observation.log_frames;
LAST_WAL_PAGES.store(wal_pages, Ordering::Relaxed);
note_checkpoint_observed(wal_pages);
if raw_observation.busy != 0 {
tracing::debug!(
busy = raw_observation.busy,
wal_log_frames = raw_observation.log_frames,
wal_checkpointed_frames = raw_observation.checkpointed_frames,
wal_pending_frames = observation.pending_frames,
wal_physical_bytes = ?observation.physical_wal_bytes,
"WAL PASSIVE checkpoint reported incomplete progress"
);
}
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,
sidecar_attribution,
})
}
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());
let start = Instant::now();
let outcome = conn.execute_batch("PRAGMA wal_checkpoint(TRUNCATE)");
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(conn);
tracing::info!(
wal_pages_before,
wal_pages_after,
elapsed_ms = elapsed.as_millis() as u64,
"WAL TRUNCATE checkpoint attempted"
);
let made_progress = wal_pages_after < wal_pages_before;
if !made_progress {
let snapshot = khive_storage::tx_registry::snapshot();
log_truncate_no_progress_warn(wal_pages_before, wal_pages_after, &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_wal_pin_depth(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, 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: u64,
state: &mut TruncateState,
) {
TRUNCATE_ATTEMPTS.fetch_add(1, Ordering::Relaxed);
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_wal_pin_depth(conn: &rusqlite::Connection) {
match query_wal_pin_depth(conn) {
Ok((log, checkpointed)) => {
tracing::warn!(
wal_log_frames = log,
wal_checkpointed_frames = checkpointed,
wal_pin_depth = (log - checkpointed).max(0),
"ADR-091 Amendment 2 Plank C: WAL pin depth after TRUNCATE no-progress"
);
}
Err(e) => {
tracing::warn!(
error = %e,
"ADR-091 Amendment 2 Plank C: failed to query WAL pin depth"
);
}
}
}
fn query_wal_pin_depth(conn: &rusqlite::Connection) -> rusqlite::Result<(i64, i64)> {
conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
Ok((row.get::<_, i64>(1)?, row.get::<_, i64>(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,
}
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_wal_pages(conn: &rusqlite::Connection) -> u64 {
let pages = query_checkpoint_observation(conn)
.map(|observation| observation.log_frames)
.unwrap_or(0)
.max(0) as u64;
LAST_WAL_PAGES.store(pages, Ordering::Relaxed);
note_checkpoint_observed(pages);
pages
}
#[cfg(test)]
mod tests {
use super::*;
use crate::pool::PoolConfig;
use crate::writer_task::WriterTaskHandle;
use rusqlite::hooks::{AuthAction, Authorization};
use serial_test::serial;
use tracing::field::{Field, Visit};
#[derive(Clone, Debug, Default)]
struct CapturedEvent {
message: Option<String>,
open_tx_count: Option<u64>,
oldest_tx_age_secs: Option<String>,
elapsed_us: Option<u64>,
busy: Option<i64>,
oldest_tx_label: Option<String>,
tx_label: Option<String>,
census_only: Option<String>,
}
#[derive(Default)]
struct CapturedEventVisitor(CapturedEvent);
impl Visit for CapturedEventVisitor {
fn record_u64(&mut self, field: &Field, value: u64) {
match field.name() {
"open_tx_count" => self.0.open_tx_count = Some(value),
"elapsed_us" => self.0.elapsed_us = Some(value),
_ => {}
}
}
fn record_i64(&mut self, field: &Field, value: i64) {
if field.name() == "busy" {
self.0.busy = Some(value);
}
}
fn record_str(&mut self, field: &Field, value: &str) {
match field.name() {
"message" => self.0.message = Some(value.to_string()),
"oldest_tx_label" => self.0.oldest_tx_label = Some(value.to_string()),
"tx_label" => self.0.tx_label = Some(value.to_string()),
_ => {}
}
}
fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
let formatted = format!("{value:?}");
let cleaned = formatted
.trim_start_matches('"')
.trim_end_matches('"')
.to_string();
match field.name() {
"message" => self.0.message = Some(cleaned),
"oldest_tx_label" => self.0.oldest_tx_label = Some(cleaned),
"tx_label" => self.0.tx_label = Some(cleaned),
"census_only" => self.0.census_only = Some(cleaned),
"oldest_tx_age_secs" => self.0.oldest_tx_age_secs = Some(cleaned),
_ => {}
}
}
}
struct CaptureSubscriber {
events: std::sync::Arc<std::sync::Mutex<Vec<CapturedEvent>>>,
}
impl tracing::Subscriber for CaptureSubscriber {
fn enabled(&self, _: &tracing::Metadata<'_>) -> bool {
true
}
fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
tracing::span::Id::from_u64(1)
}
fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
fn event(&self, event: &tracing::Event<'_>) {
let mut visitor = CapturedEventVisitor::default();
event.record(&mut visitor);
self.events.lock().unwrap().push(visitor.0);
}
fn enter(&self, _: &tracing::span::Id) {}
fn exit(&self, _: &tracing::span::Id) {}
}
fn spans_straddling(
threshold: Duration,
) -> (
khive_storage::tx_registry::OldestSpan,
khive_storage::tx_registry::OldestSpan,
) {
let _handle = khive_storage::tx_registry::register(Some("writer_task_tx".to_string()));
let (id, _age, _label) =
khive_storage::tx_registry::oldest().expect("a registration is open");
let base = khive_storage::tx_registry::OldestSpan {
id,
age: Duration::ZERO,
label: Some("writer_task_tx".to_string()),
origin: khive_storage::tx_registry::TxOrigin::Unscoped,
};
let aged = khive_storage::tx_registry::OldestSpan {
age: threshold + Duration::from_secs(1),
..base.clone()
};
let young = khive_storage::tx_registry::OldestSpan {
age: Duration::from_micros(5_849),
..base
};
(aged, young)
}
fn capture<F: FnOnce()>(f: F) -> Vec<CapturedEvent> {
let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let subscriber = CaptureSubscriber {
events: std::sync::Arc::clone(&buffer),
};
tracing::subscriber::with_default(subscriber, f);
let events = buffer.lock().unwrap();
events.clone()
}
#[test]
fn truncate_no_progress_warn_reports_nonempty_registry_snapshot_facts() {
let snapshot = [
(Duration::from_secs(2), Some("younger-entry".to_string())),
(Duration::from_secs(7), Some("older-entry".to_string())),
];
let events = capture(|| log_truncate_no_progress_warn(6003, 6003, &snapshot));
assert_eq!(events.len(), 3, "one summary and both captured entries");
let summary = &events[0];
assert_eq!(summary.open_tx_count, Some(2));
assert_eq!(summary.oldest_tx_age_secs.as_deref(), Some("Some(7.0)"));
let message = summary.message.as_deref().expect("summary message");
assert_eq!(
message,
"WAL TRUNCATE attempt made no progress; open transactions observed in this process's registry"
);
assert!(!message.contains("pinning"));
assert!(!message.contains("long-lived reader"));
assert_eq!(events[1].tx_label.as_deref(), Some("younger-entry"));
assert_eq!(events[2].tx_label.as_deref(), Some("older-entry"));
assert!(events[1..].iter().all(|event| {
event.message.as_deref() == Some("WAL high-water: open transaction registry entry")
}));
}
#[test]
fn truncate_no_progress_warn_reports_empty_process_registry_without_pin_claim() {
let events = capture(|| log_truncate_no_progress_warn(6003, 6003, &[]));
assert_eq!(
events.len(),
1,
"an empty snapshot has no entries to enumerate"
);
let summary = &events[0];
assert_eq!(summary.open_tx_count, Some(0));
assert_eq!(summary.oldest_tx_age_secs.as_deref(), Some("None"));
let message = summary.message.as_deref().expect("summary message");
assert_eq!(
message,
"WAL TRUNCATE attempt made no progress; no open transaction in this process's registry"
);
assert!(!message.contains("pinning"));
assert!(!message.contains("long-lived reader"));
let nonempty = capture(|| {
log_truncate_no_progress_warn(6003, 6003, &[(Duration::ZERO, None)]);
});
assert_ne!(summary.message, nonempty[0].message);
assert_eq!(nonempty[0].open_tx_count, Some(1));
assert_eq!(nonempty[0].oldest_tx_age_secs.as_deref(), Some("Some(0.0)"));
}
#[test]
#[serial(tx_registry)]
fn high_water_warn_names_the_pin_when_an_aged_entry_exists() {
let threshold = Duration::from_secs(30);
let (aged, _young) = spans_straddling(threshold);
let events = capture(|| log_wal_high_water_warn(6003, 6000, Some(&aged), threshold));
let message = events
.iter()
.find_map(|e| e.message.clone())
.expect("one WARN is emitted");
assert!(
message.contains("is pinning a snapshot"),
"an aged entry must produce the pin wording, got {message:?}"
);
assert_eq!(
events.iter().find_map(|e| e.oldest_tx_label.clone()),
Some("writer_task_tx".to_string()),
"the named entry is the one handed in"
);
}
#[test]
#[serial(tx_registry)]
fn high_water_warn_rules_the_pin_out_when_the_oldest_entry_is_young() {
let threshold = Duration::from_secs(30);
let (aged, young) = spans_straddling(threshold);
let young_message =
capture(|| log_wal_high_water_warn(6003, 6000, Some(&young), threshold))
.iter()
.find_map(|e| e.message.clone())
.expect("one WARN is emitted");
let aged_message = capture(|| log_wal_high_water_warn(6003, 6000, Some(&aged), threshold))
.iter()
.find_map(|e| e.message.clone())
.expect("one WARN is emitted");
assert_ne!(
young_message, aged_message,
"the two registry states must produce different text"
);
assert!(
young_message.contains("no open transaction old enough to pin"),
"got {young_message:?}"
);
assert!(
!young_message.contains("pinning"),
"the young branch must not assert a pin, got {young_message:?}"
);
assert!(
!young_message.contains("long-lived reader"),
"the young branch must not name a reader, got {young_message:?}"
);
}
#[test]
#[serial(tx_registry)]
fn high_water_warn_with_an_empty_registry_takes_the_no_pin_branch() {
let events = capture(|| log_wal_high_water_warn(6003, 6000, None, Duration::from_secs(30)));
let message = events
.iter()
.find_map(|e| e.message.clone())
.expect("one WARN is emitted");
assert!(
message.contains("no open transaction old enough to pin"),
"got {message:?}"
);
assert_eq!(
events.iter().find_map(|e| e.oldest_tx_label.clone()),
Some("<none>".to_string()),
"an absent entry is labelled as absent, never as unlabeled"
);
}
#[test]
#[serial(tx_registry)]
fn log_tx_registry_oldest_debug_reports_oldest_open_entry() {
let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let subscriber = CaptureSubscriber {
events: std::sync::Arc::clone(&buffer),
};
let _handle =
khive_storage::tx_registry::register(Some("checkpoint_tick_test".to_string()));
let oldest = khive_storage::tx_registry::oldest().map(|(id, age, label)| {
khive_storage::tx_registry::OldestSpan {
id,
age,
label,
origin: khive_storage::tx_registry::TxOrigin::Unscoped,
}
});
let expected_label = oldest
.as_ref()
.and_then(|s| s.label.clone())
.unwrap_or_else(|| "<unlabeled>".to_string());
tracing::subscriber::with_default(subscriber, || {
log_tx_registry_oldest_debug(100, oldest.as_ref());
});
let events = buffer.lock().unwrap();
assert!(
events.iter().any(|e| {
e.message.as_deref()
== Some("WAL checkpoint tick: oldest open transaction registry entry")
&& e.oldest_tx_label.as_deref() == Some(expected_label.as_str())
}),
"expected a log line naming the open registry entry's label, got: {events:?}"
);
}
#[test]
#[serial(tx_registry)]
fn registry_warns_fire_on_crossing_and_do_not_repeat() {
let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let subscriber = CaptureSubscriber {
events: std::sync::Arc::clone(&buffer),
};
let _handle =
khive_storage::tx_registry::register(Some("registry_warn_crossing_test".to_string()));
let oldest = khive_storage::tx_registry::oldest().map(|(id, age, label)| {
khive_storage::tx_registry::OldestSpan {
id,
age,
label,
origin: khive_storage::tx_registry::TxOrigin::Unscoped,
}
});
let mut was_above_warn = false;
let mut was_above_high_water = false;
tracing::subscriber::with_default(subscriber, || {
if crossing_warn(true, &mut was_above_warn) {
log_tx_registry_oldest_warn(6000, oldest.as_ref());
}
if crossing_warn(true, &mut was_above_high_water) {
log_tx_registry_snapshot_warn(6000);
}
if crossing_warn(true, &mut was_above_warn) {
log_tx_registry_oldest_warn(6000, oldest.as_ref());
}
if crossing_warn(true, &mut was_above_high_water) {
log_tx_registry_snapshot_warn(6000);
}
});
let events = buffer.lock().unwrap();
let oldest_warn_count = events
.iter()
.filter(|e| {
e.message.as_deref()
== Some("WAL checkpoint tick: oldest open transaction registry entry")
})
.count();
assert_eq!(
oldest_warn_count, 1,
"oldest-entry WARN must fire exactly once across two above-threshold ticks, got: {events:?}"
);
let snapshot_warn_count = events
.iter()
.filter(|e| {
e.message.as_deref() == Some("WAL high-water: open transaction registry entry")
&& e.tx_label.as_deref() == Some("registry_warn_crossing_test")
})
.count();
assert_eq!(
snapshot_warn_count, 1,
"high-water snapshot WARN must fire exactly once across two above-threshold ticks, got: {events:?}"
);
}
#[test]
fn log_tx_age_emission_carries_label_for_both_rungs() {
let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let subscriber = CaptureSubscriber {
events: std::sync::Arc::clone(&buffer),
};
tracing::subscriber::with_default(subscriber, || {
log_tx_age_emission(&TxAgeEmission {
rung: TxAgeRung::Warn,
age: Duration::from_secs(45),
label: Some("plank1_warn_test".to_string()),
});
log_tx_age_emission(&TxAgeEmission {
rung: TxAgeRung::Stale,
age: Duration::from_secs(150),
label: Some("plank1_stale_test".to_string()),
});
});
let events = buffer.lock().unwrap();
assert!(
events.iter().any(|e| {
e.message.as_deref()
== Some(
"ADR-091 Plank 1: open transaction registry entry exceeded soft-cap age",
)
&& e.tx_label.as_deref() == Some("plank1_warn_test")
}),
"expected a Warn-rung log line naming the entry, got: {events:?}"
);
assert!(
events.iter().any(|e| {
e.message.as_deref().is_some_and(|m| {
m.starts_with(
"ADR-091 Plank 1: open transaction registry entry exceeded the cooperative",
)
}) && e.tx_label.as_deref() == Some("plank1_stale_test")
}),
"expected a Stale-rung log line naming the entry, got: {events:?}"
);
}
fn file_pool(path: &std::path::Path) -> Arc<ConnectionPool> {
let cfg = PoolConfig {
path: Some(path.to_path_buf()),
..PoolConfig::for_test()
};
Arc::new(ConnectionPool::new(cfg).expect("pool open"))
}
async fn writer_task_wal_autocheckpoint_pages(handle: &WriterTaskHandle) -> u32 {
handle
.send_top_level(|conn| {
conn.pragma_query_value(None, "wal_autocheckpoint", |row| row.get::<_, u32>(0))
.map_err(|error| khive_storage::error::StorageError::Pool {
operation: "test_wal_autocheckpoint".into(),
message: error.to_string(),
})
})
.await
.expect("query writer-task connection pragma")
}
fn checkpoint_conn(pool: &ConnectionPool) -> rusqlite::Connection {
pool.open_standalone_writer()
.expect("open dedicated checkpoint connection")
}
struct TruncateReportHookGuard;
impl Drop for TruncateReportHookGuard {
fn drop(&mut self) {
truncate_report_test_sync::uninstall();
}
}
#[cfg(unix)]
struct WalpinAttributionHookGuard;
#[cfg(unix)]
impl Drop for WalpinAttributionHookGuard {
fn drop(&mut self) {
walpin_attribution_test_sync::uninstall();
}
}
#[tokio::test(flavor = "current_thread")]
#[cfg(unix)]
#[serial(khive_walpin_sidecar_env)]
async fn diagnostic_legacy_forecast_matches_housekeeping_with_distinct_cadences() {
let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
let root = tempfile::tempdir().unwrap();
let pool = file_pool(&root.path().join("forecast.db"));
let path = pool.canonical_path().unwrap();
let checkpoint_interval = Duration::from_millis(500);
let session_interval = SessionSweepConfig::default().interval;
assert_eq!(session_interval, Duration::from_secs(5));
let state = TruncateState::default();
assert_eq!(state.legacy_walpin_fallback_interval, session_interval);
let sidecar = WalpinSidecarState::new(Some(path), true, "daemon", checkpoint_interval)
.expect("enabled fixture sidecar");
crate::walpin::ensure_sidecar_dir(&sidecar.dir).unwrap();
let mut paths = Vec::new();
for (pid, age) in [(2_000_000_001, 5), (2_000_000_002, 40)] {
assert!(!crate::walpin::is_process_alive(pid));
let temp = sidecar.dir.join(format!(".{pid}.beacon.tmp"));
std::fs::write(
&temp,
serde_json::to_vec(&serde_json::json!({
"pid": pid, "process_role": "session", "started_at": 1
}))
.unwrap(),
)
.unwrap();
std::fs::File::options()
.write(true)
.open(&temp)
.unwrap()
.set_modified(std::time::SystemTime::now() - Duration::from_secs(age))
.unwrap();
paths.push(temp);
}
let fast = crate::walpin::inspect_live(&sidecar.dir, checkpoint_interval).unwrap();
assert_eq!(
fast.cleanup_would_reap, 2,
"control must distinguish the cadences"
);
let forecast = crate::diagnostics::wal_pin_attribution(path, session_interval);
assert_eq!(forecast.sidecar_listing_truncated, Some(false));
assert_eq!(forecast.sidecar_entries_cleanup_would_reap, Some(1));
assert!(
paths.iter().all(|path| path.exists()),
"inspection retains evidence"
);
let cleanup = sidecar
.reap_dead_entries_bounded(state.legacy_walpin_fallback_interval)
.await
.expect("housekeeping report");
assert_eq!(
Some(cleanup.orphan_temps_reaped),
forecast.sidecar_entries_cleanup_would_reap
);
assert!(paths[0].exists(), "the temp inside the 15s window remains");
assert!(!paths[1].exists(), "the trusted older temp is reaped");
}
#[test]
#[cfg(unix)]
fn walpin_full_scan_cadence_refreshes_first_then_reuses_until_boundary() {
let cadence = Duration::from_secs(30);
let started_at = Instant::now();
let mut state = TruncateState::with_walpin_full_scan_cadence(cadence);
assert!(matches!(
state.plan_walpin_attribution_at(started_at),
WalpinFullScanPlan::Refresh { .. }
));
state.cache_walpin_attribution(
crate::walpin::WalpinReport::default(),
Ok(crate::walpin::CensusResult::default()),
started_at,
);
assert!(matches!(
state.plan_walpin_attribution_at(started_at + cadence - Duration::from_nanos(1)),
WalpinFullScanPlan::Cached(_)
));
assert!(matches!(
state.plan_walpin_attribution_at(started_at + cadence),
WalpinFullScanPlan::Refresh { .. }
));
}
#[test]
#[cfg(unix)]
fn walpin_full_scan_failure_retries_only_after_cadence() {
let cadence = Duration::from_secs(30);
let started_at = Instant::now();
let mut state = TruncateState::with_walpin_full_scan_cadence(cadence);
assert!(matches!(
state.plan_walpin_attribution_at(started_at),
WalpinFullScanPlan::Refresh { .. }
));
assert!(matches!(
state.plan_walpin_attribution_at(started_at + cadence - Duration::from_nanos(1)),
WalpinFullScanPlan::Suppressed
));
assert!(matches!(
state.plan_walpin_attribution_at(started_at + cadence),
WalpinFullScanPlan::Refresh { .. }
));
}
#[test]
#[cfg(unix)]
fn cached_walpin_report_is_diagnostic_only_even_when_fully_attributed() {
let report = crate::walpin::WalpinReport::default();
assert!(
report.fully_attributed(),
"the fixture must otherwise license the sharp conclusion"
);
let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let subscriber = CaptureSubscriber {
events: std::sync::Arc::clone(&buffer),
};
tracing::subscriber::with_default(subscriber, || {
log_walpin_sidecar_report(
&report,
Ok(crate::walpin::CensusResult::default()),
WalpinReportFreshness::Cached {
age: Duration::from_secs(1),
},
);
});
let events = buffer.lock().unwrap();
assert!(
events.iter().any(|event| {
event.message.as_deref()
== Some(
"cached WAL-pin attribution is diagnostic-only; fully-attributed \
conclusion is not licensed",
)
}),
"cached attribution must declare its fail-closed status: {events:?}"
);
assert!(
!events.iter().any(|event| {
event.message.as_deref().is_some_and(|message| {
message.starts_with("ADR-091 Amendment 2 Plank B: every live PID is reporting")
})
}),
"cached attribution must never authorize the fully-attributed conclusion: {events:?}"
);
}
#[tokio::test(flavor = "current_thread")]
#[cfg(unix)]
#[serial(checkpoint_skip_metrics, khive_walpin_sidecar_env)]
async fn progressing_truncate_releases_full_scan_reservation_to_housekeeping() {
let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("walpin-progress-reservation.db");
let pool = file_pool(&path);
{
let writer = pool.try_writer().unwrap();
writer
.conn()
.execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
.unwrap();
}
let conn = checkpoint_conn(&pool);
let mut state = TruncateState::default();
let config = CheckpointConfig {
truncate_high_water_pages: 0,
truncate_min_interval: Duration::ZERO,
..CheckpointConfig::default()
};
assert!(
maybe_truncate(&pool, &conn, &config, u64::MAX, &mut state).is_none(),
"a progressing TRUNCATE must not schedule no-progress attribution"
);
assert!(
state.housekeeping_due(),
"unused pre-TRUNCATE reservation must be restored before housekeeping"
);
let sidecar =
WalpinSidecarState::new(pool.canonical_path(), true, "daemon", config.interval)
.expect("file-backed test sidecar");
assert!(
run_walpin_housekeeping_if_due(&sidecar, &mut state, DEFAULT_SESSION_SWEEP_INTERVAL,)
.await,
"the production housekeeping arm must consume one full scan"
);
assert!(state.walpin_cached_attribution.is_some());
}
#[tokio::test(flavor = "current_thread")]
#[cfg(unix)]
#[serial(checkpoint_skip_metrics, khive_walpin_sidecar_env)]
async fn erroring_truncate_releases_full_scan_reservation_to_housekeeping() {
let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("walpin-error-reservation.db");
let pool = file_pool(&path);
{
let writer = pool.try_writer().unwrap();
writer
.conn()
.execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
.unwrap();
}
let conn = checkpoint_conn(&pool);
conn.authorizer(Some(
|context: rusqlite::hooks::AuthContext<'_>| match context.action {
AuthAction::Pragma { pragma_name, .. }
if pragma_name.eq_ignore_ascii_case("wal_checkpoint") =>
{
Authorization::Deny
}
_ => Authorization::Allow,
},
))
.unwrap();
let mut state = TruncateState::default();
let config = CheckpointConfig {
truncate_high_water_pages: 0,
truncate_min_interval: Duration::ZERO,
..CheckpointConfig::default()
};
assert!(
maybe_truncate(&pool, &conn, &config, u64::MAX, &mut state).is_none(),
"an erroring TRUNCATE must not schedule no-progress attribution"
);
conn.authorizer(None::<fn(rusqlite::hooks::AuthContext<'_>) -> Authorization>)
.unwrap();
assert!(
state.housekeeping_due(),
"failed TRUNCATE must restore its unused full-scan reservation"
);
let sidecar =
WalpinSidecarState::new(pool.canonical_path(), true, "daemon", config.interval)
.expect("file-backed test sidecar");
assert!(
run_walpin_housekeeping_if_due(&sidecar, &mut state, DEFAULT_SESSION_SWEEP_INTERVAL,)
.await,
"the production housekeeping arm must consume one full scan"
);
assert!(state.walpin_cached_attribution.is_some());
}
struct ReaderProcess {
child: std::process::Child,
_stdout: std::io::BufReader<std::process::ChildStdout>,
}
impl ReaderProcess {
fn spawn(db_path: &std::path::Path) -> Self {
use std::io::BufRead;
use std::process::Stdio;
let mut child = std::process::Command::new(
std::env::current_exe().expect("resolve current test executable"),
)
.args([
"--exact",
"checkpoint::tests::walpin_transient_reader_process_helper",
"--nocapture",
])
.env("KHIVE_CHECKPOINT_READER_HELPER_PATH", db_path)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.spawn()
.expect("spawn transient WAL reader helper");
let stdout = child.stdout.take().expect("capture helper stdout");
let mut reader = std::io::BufReader::new(stdout);
let mut line = String::new();
loop {
line.clear();
let bytes = reader
.read_line(&mut line)
.expect("read transient reader readiness signal");
assert!(bytes > 0, "reader helper exited before readiness signal");
if line.contains("KHIVE_CHECKPOINT_READER_READY") {
break;
}
}
Self {
child,
_stdout: reader,
}
}
fn pid(&self) -> u32 {
self.child.id()
}
fn release(&mut self) {
use std::io::Write;
let mut stdin = self.child.stdin.take().expect("helper stdin is available");
stdin
.write_all(b"release\n")
.expect("release transient reader");
drop(stdin);
let status = self.child.wait().expect("wait for transient reader helper");
assert!(status.success(), "transient reader helper failed: {status}");
}
}
impl Drop for ReaderProcess {
fn drop(&mut self) {
if self.child.try_wait().ok().flatten().is_none() {
let _ = self.child.kill();
let _ = self.child.wait();
}
}
}
#[test]
fn walpin_transient_reader_process_helper() {
use std::io::Write;
let Some(path) = std::env::var_os("KHIVE_CHECKPOINT_READER_HELPER_PATH") else {
return;
};
let conn = rusqlite::Connection::open(path).expect("helper opens database");
conn.execute_batch("BEGIN DEFERRED; SELECT * FROM t;")
.expect("helper pins a read snapshot");
println!("KHIVE_CHECKPOINT_READER_READY");
std::io::stdout().flush().expect("flush readiness signal");
let mut release = String::new();
std::io::stdin()
.read_line(&mut release)
.expect("wait for release signal");
conn.execute_batch("COMMIT")
.expect("helper releases read snapshot");
}
#[tokio::test(flavor = "current_thread")]
#[cfg(unix)]
#[serial(
checkpoint_skip_metrics,
khive_walpin_sidecar_env,
walpin_attribution_async,
walpin_report_seam
)]
async fn no_progress_report_keeps_holder_released_after_truncate_timeout() {
let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("transient-reader.db");
let pool = file_pool(&path);
{
let writer = pool.try_writer().expect("writer");
writer
.conn()
.execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
.expect("seed WAL before reader snapshot");
}
let mut reader = ReaderProcess::spawn(&path);
let reader_pid = reader.pid();
{
let writer = pool.try_writer().expect("writer");
writer
.conn()
.execute_batch("INSERT INTO t VALUES (2);")
.expect("append WAL behind reader snapshot");
}
let canonical_path = pool
.canonical_path()
.expect("file-backed pool has canonical path")
.to_path_buf();
let (reached_rx, proceed_tx) = truncate_report_test_sync::install(canonical_path.clone());
let _hook_guard = TruncateReportHookGuard;
let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let checkpoint_pool = Arc::clone(&pool);
let dedicated_conn = checkpoint_conn(&checkpoint_pool);
let checkpoint_events = Arc::clone(&buffer);
let checkpoint = std::thread::spawn(move || {
let subscriber = CaptureSubscriber {
events: checkpoint_events,
};
let _subscriber_guard = tracing::subscriber::set_default(subscriber);
let mut state = TruncateState::default();
let result = checkpoint_once_core(
&checkpoint_pool,
&dedicated_conn,
&CheckpointConfig {
truncate_high_water_pages: 0,
truncate_min_interval: Duration::ZERO,
truncate_busy_timeout: Duration::from_millis(50),
..CheckpointConfig::default()
},
&mut state,
);
(result, state)
});
reached_rx
.recv_timeout(Duration::from_secs(5))
.expect("TRUNCATE must report no progress while the reader is pinned");
reader.release();
let post_attempt_census =
crate::walpin::census_holders(&canonical_path).expect("post-attempt holder census");
assert!(
!post_attempt_census.holders.contains(&reader_pid),
"released reader PID must be absent from a post-attempt census"
);
proceed_tx
.send(())
.expect("allow no-progress reporting to continue");
let (checkpoint_result, mut state) = checkpoint.join().expect("checkpoint thread");
let outcome = checkpoint_result.expect("checkpoint succeeds");
assert!(
outcome.sidecar_attribution.is_some(),
"the synchronous checkpoint result must carry a separate attribution request"
);
assert!(
!state.sidecar_attribution_attempted_this_tick,
"capturing a request is not the same as attempting its directory enumeration"
);
let subscriber = CaptureSubscriber {
events: std::sync::Arc::clone(&buffer),
};
let _subscriber_guard = tracing::subscriber::set_default(subscriber);
let attribution_attempted =
complete_walpin_attribution(outcome.sidecar_attribution, &mut state)
.await
.expect("deferred attribution succeeds");
assert!(
attribution_attempted,
"a no-progress attribution pass must suppress the redundant healthy-housekeeping \
pass for the same tick"
);
let events = buffer.lock().expect("captured events");
let summaries: Vec<_> = events
.iter()
.filter(|event| {
event.message.as_deref().is_some_and(|message| {
message.starts_with("WAL TRUNCATE attempt made no progress;")
})
})
.collect();
assert_eq!(
summaries.len(),
1,
"the real no-progress path emits one summary"
);
let summary = summaries[0];
assert!(summary.open_tx_count.is_some());
assert!(summary.oldest_tx_age_secs.is_some());
let message = summary.message.as_deref().expect("summary message");
assert!(message.contains("in this process's registry"));
assert!(!message.contains("pinning"));
assert!(!message.contains("long-lived reader"));
assert!(
events.iter().any(|event| {
event
.census_only
.as_deref()
.is_some_and(|pids| pids.contains(&reader_pid.to_string()))
}),
"the no-progress report must retain PID {reader_pid} from the pre-attempt census: {events:?}"
);
}
#[tokio::test(flavor = "current_thread")]
#[cfg(unix)]
#[serial(walpin_attribution_async)]
async fn no_progress_attribution_is_off_runtime_and_awaited_before_report_use() {
use std::sync::atomic::Ordering;
let dir = tempfile::tempdir().expect("tempdir");
let sidecar_dir = dir.path().join("checkpoint.db.walpin");
let (reached_rx, proceed_tx, report_counter) =
walpin_attribution_test_sync::install_pause(sidecar_dir.clone());
let _hook_guard = WalpinAttributionHookGuard;
let runtime_thread = std::thread::current().id();
let mut state = TruncateState::default();
let request = Some(WalpinAttributionRequest::Fresh {
dir: sidecar_dir,
census: Ok(crate::walpin::CensusResult::default()),
legacy_fallback_interval: DEFAULT_SESSION_SWEEP_INTERVAL,
previous_last_attempt: None,
});
let completion = tokio::spawn(async move {
let result = complete_walpin_attribution(request, &mut state).await;
(result, state)
});
let blocking_thread = reached_rx
.await
.expect("spawn_blocking attribution reached test seam");
assert_ne!(
blocking_thread, runtime_thread,
"sidecar enumeration must not execute on the current-thread Tokio runtime worker"
);
assert!(
!completion.is_finished(),
"the async attribution owner must await the still-paused blocking enumeration"
);
assert_eq!(
report_counter.load(Ordering::SeqCst),
0,
"the attribution report must not be consumed before enumeration completes"
);
proceed_tx
.send(())
.expect("release blocking attribution enumeration");
let (result, state) = completion.await.expect("attribution task joins");
assert_eq!(result, Ok(true));
assert_eq!(
report_counter.load(Ordering::SeqCst),
1,
"the completed enumeration must feed exactly one report use"
);
assert!(state.sidecar_attribution_attempted_this_tick);
assert!(
!state.housekeeping_due(),
"completed attribution must suppress same-tick housekeeping"
);
}
#[tokio::test(flavor = "current_thread")]
#[cfg(unix)]
#[serial(walpin_attribution_async)]
async fn no_progress_attribution_join_failure_is_honest_and_suppresses_retry() {
let dir = tempfile::tempdir().expect("tempdir");
let sidecar_dir = dir.path().join("checkpoint.db.walpin");
walpin_attribution_test_sync::install_panic(sidecar_dir.clone());
let _hook_guard = WalpinAttributionHookGuard;
let mut state = TruncateState::default();
let request = Some(WalpinAttributionRequest::Fresh {
dir: sidecar_dir,
census: Ok(crate::walpin::CensusResult::default()),
legacy_fallback_interval: DEFAULT_SESSION_SWEEP_INTERVAL,
previous_last_attempt: None,
});
let error = complete_walpin_attribution(request, &mut state)
.await
.expect_err("injected worker panic must surface as failure");
assert!(
matches!(error, WalpinAttributionFailure::Worker(_)),
"join failure must retain its worker classification: {error:?}"
);
assert!(state.sidecar_attribution_attempted_this_tick);
assert!(
!state.housekeeping_due(),
"an indeterminate partial pass must not authorize a second scan"
);
}
#[test]
#[cfg(unix)]
#[serial(checkpoint_skip_metrics)]
fn async_checkpoint_source_keeps_enumeration_behind_awaited_spawn_blocking() {
fn section<'a>(source: &'a str, start: &str, end: &str) -> &'a str {
source
.split_once(start)
.unwrap_or_else(|| panic!("missing source marker {start:?}"))
.1
.split_once(end)
.unwrap_or_else(|| panic!("missing source marker {end:?}"))
.0
}
let source = include_str!("checkpoint.rs");
let checkpoint_core = section(
source,
"fn checkpoint_once_core(",
"/// Evaluate and, if due, attempt a TRUNCATE escalation",
);
let truncate_core = section(
source,
"fn maybe_truncate(",
"#[cfg(test)]\nmod truncate_report_test_sync",
);
let report_logger = section(
source,
"fn log_walpin_sidecar_report(",
"/// ADR-091 Amendment 2 Plank C",
);
for (name, body) in [
("checkpoint_once_core", checkpoint_core),
("maybe_truncate", truncate_core),
("log_walpin_sidecar_report", report_logger),
] {
assert!(
!body.contains("enumerate_live("),
"{name} must not perform direct sidecar enumeration"
);
}
let async_completion = section(
source,
"async fn complete_walpin_attribution(",
"/// When a TRUNCATE attempt makes no progress",
);
let spawn = async_completion
.find("tokio::task::spawn_blocking")
.expect("completion must spawn blocking work");
let enumerate = async_completion
.find("crate::walpin::enumerate_live")
.expect("blocking closure must perform the attribution enumeration");
let awaited = async_completion[enumerate..]
.find(".await")
.map(|offset| enumerate + offset)
.expect("blocking worker must be awaited");
assert!(spawn < enumerate && enumerate < awaited);
let task = section(
source,
"pub async fn run_checkpoint_task(",
"/// Whether a `CheckpointOutcomeRecorded` transition should be enqueued",
);
let checkpoint = task
.find("checkpoint_once_core(")
.expect("checkpoint core call");
let completion = task
.find("complete_walpin_attribution(")
.expect("awaited attribution completion");
let housekeeping = task
.find("run_walpin_housekeeping_if_due(")
.expect("fallback housekeeping");
let outcome = task
.find("observe_checkpoint_pressure_tick(")
.expect("lifecycle outcome use");
assert!(
checkpoint < completion && completion < housekeeping && housekeeping < outcome,
"tick ordering must be checkpoint -> awaited attribution -> housekeeping decision -> outcome"
);
let housekeeping_helper = section(
source,
"async fn run_walpin_housekeeping_if_due(",
"fn now_epoch_secs()",
);
assert!(
housekeeping_helper.contains("reap_dead_entries_bounded(legacy_fallback_interval)"),
"the ordered housekeeping arm must retain the bounded full scan"
);
let pressure_tick = section(
source,
"fn observe_checkpoint_pressure_tick(",
"/// ADR-091 Plank 0",
);
assert!(
pressure_tick.contains("checkpoint_outcome_should_emit"),
"extracted pressure tick helper must gate on the lifecycle emit decision"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn fts_maintenance_off_worker_matches_the_direct_call() {
fn fragmented_fixture() -> (tempfile::TempDir, rusqlite::Connection) {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("fts-off-worker.db");
let conn = rusqlite::Connection::open(&path).expect("open sqlite");
conn.execute_batch(
"CREATE VIRTUAL TABLE fts_entities USING fts5(namespace UNINDEXED, subject_id UNINDEXED, title, body, tokenize='trigram');
CREATE VIRTUAL TABLE fts_notes USING fts5(namespace UNINDEXED, subject_id UNINDEXED, title, body, tokenize='trigram');
INSERT INTO fts_entities(fts_entities, rank) VALUES('automerge', 0);
INSERT INTO fts_notes(fts_notes, rank) VALUES('automerge', 0);",
)
.expect("create FTS fixtures");
for index in 0..80 {
let body = format!(
"segment fixture {index} keeps enough repeated production recall text to span pages {}",
"memory query corpus ".repeat(40)
);
conn.execute(
"INSERT INTO fts_entities(namespace, subject_id, title, body) VALUES(?1, ?2, ?3, ?4)",
rusqlite::params![
"local",
format!("id-{index}"),
format!("title {index}"),
body
],
)
.expect("one autocommit FTS write");
}
(dir, conn)
}
let config = crate::fts_maintenance::FtsMaintenanceConfig {
enabled: true,
interval: Duration::ZERO,
merge_pages: 8,
minimum_segments: 2,
};
let (_direct_dir, direct_conn) = fragmented_fixture();
let mut direct_state = crate::fts_maintenance::FtsMaintenanceState::new(Instant::now());
let direct_step = crate::fts_maintenance::run_if_due(
&direct_conn,
&config,
&mut direct_state,
Instant::now(),
)
.expect("direct maintenance step")
.expect("fragmented fixture has a due step");
let (_wrapped_dir, wrapped_conn) = fragmented_fixture();
let wrapped_state = crate::fts_maintenance::FtsMaintenanceState::new(Instant::now());
let (_conn, _state, wrapped_result) =
run_fts_maintenance_off_worker(wrapped_conn, config, wrapped_state, Instant::now())
.await
.expect("maintenance step did not panic");
let wrapped_step = wrapped_result
.expect("wrapped maintenance step")
.expect("fragmented fixture has a due step");
assert_eq!(
direct_step, wrapped_step,
"the spawn_blocking wrapper must produce the same step outcome as calling \
run_if_due directly"
);
assert_eq!(direct_step.table, "fts_entities");
}
#[test]
#[serial(checkpoint_skip_metrics)]
fn checkpoint_once_succeeds_on_file_backed_pool() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("wal_test.db");
let pool = file_pool(&path);
{
let writer = pool.try_writer().unwrap();
writer
.conn()
.execute_batch("CREATE TABLE IF NOT EXISTS t (x INTEGER);")
.unwrap();
writer
.conn()
.execute_batch("INSERT INTO t VALUES (1);")
.unwrap();
}
let conn = checkpoint_conn(&pool);
checkpoint_once(
&pool,
&conn,
&CheckpointConfig::default(),
&mut TruncateState::default(),
)
.expect("checkpoint_once must succeed against a healthy dedicated connection");
}
#[test]
fn open_standalone_writer_fails_on_in_memory_pool() {
let cfg = PoolConfig {
path: None,
..PoolConfig::default()
};
let pool = Arc::new(ConnectionPool::new(cfg).expect("in-memory pool"));
assert!(
pool.open_standalone_writer().is_err(),
"an in-memory pool must not be able to open a dedicated checkpoint connection"
);
}
#[test]
fn ensure_open_does_not_move_writer_acquisition_counters() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("checkpoint_ensure_open.db");
let pool = file_pool(&path);
let before_first_open = pool.writer_acquisition_snapshot();
let mut checkpoint_conn = CheckpointConnection::new();
checkpoint_conn
.ensure_open(&pool)
.expect("dedicated checkpoint connection must open against a file-backed pool");
assert_eq!(
pool.writer_acquisition_snapshot(),
before_first_open,
"the checkpoint connection's initial open must not count as a writer acquisition"
);
checkpoint_conn.conn = None;
let before_reopen = pool.writer_acquisition_snapshot();
checkpoint_conn
.ensure_open(&pool)
.expect("dedicated checkpoint connection must reopen after invalidation");
assert_eq!(
pool.writer_acquisition_snapshot(),
before_reopen,
"reopening the checkpoint connection must not count as a writer acquisition either"
);
}
#[test]
fn checkpoint_connection_disables_wal_autocheckpoint_on_open_and_reopen() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("checkpoint_autocheckpoint.db");
let pool = file_pool(&path);
let mut checkpoint_conn = CheckpointConnection::new();
let initial: u32 = checkpoint_conn
.ensure_open(&pool)
.expect("dedicated checkpoint connection must open")
.pragma_query_value(None, "wal_autocheckpoint", |row| row.get(0))
.expect("read initial autocheckpoint setting");
assert_eq!(initial, 0);
checkpoint_conn.conn = None;
let reopened: u32 = checkpoint_conn
.ensure_open(&pool)
.expect("dedicated checkpoint connection must reopen")
.pragma_query_value(None, "wal_autocheckpoint", |row| row.get(0))
.expect("read reopened autocheckpoint setting");
assert_eq!(reopened, 0);
}
#[tokio::test(flavor = "current_thread")]
#[serial(checkpoint_skip_metrics)]
async fn failed_checkpoint_claim_keeps_existing_writer_task_on_fallback() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("failed_claim_writer_task.db");
let pool = Arc::new(
ConnectionPool::new(PoolConfig {
path: Some(path),
checkout_timeout: Duration::from_millis(1),
write_queue_enabled: Some(true),
..PoolConfig::for_test()
})
.expect("pool open"),
);
let writer_task = pool
.writer_task_handle()
.expect("writer-task resolution")
.expect("writer task enabled");
assert_eq!(
writer_task_wal_autocheckpoint_pages(&writer_task).await,
crate::pool::FALLBACK_WAL_AUTOCHECKPOINT_PAGES
);
let legacy_conn = pool.legacy_conn();
let (held_tx, held_rx) = tokio::sync::oneshot::channel();
let (release_tx, release_rx) = std::sync::mpsc::channel();
let holder = tokio::task::spawn_blocking(move || {
let _held_writer = legacy_conn.lock();
held_tx.send(()).expect("signal held pooled writer");
release_rx.recv().expect("release held pooled writer");
});
held_rx.await.expect("pooled writer holder started");
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
drop(shutdown_tx);
run_checkpoint_task(
Arc::clone(&pool),
CheckpointConfig {
interval: Duration::from_secs(60),
..CheckpointConfig::default()
},
None,
shutdown_rx,
true,
)
.await;
assert_eq!(pool.writer_acquisition_snapshot().timeouts, 1);
assert_eq!(
writer_task_wal_autocheckpoint_pages(&writer_task).await,
crate::pool::FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
"failed pooled-writer claim must not partially propagate ownership"
);
release_tx.send(()).expect("release pooled writer");
holder.await.expect("pooled writer holder joined");
}
#[tokio::test]
#[serial(checkpoint_skip_metrics)]
async fn checkpoint_task_exits_on_shutdown_signal() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("wal_task_shutdown.db");
let pool = file_pool(&path);
let cfg = CheckpointConfig {
interval: Duration::from_millis(10),
..Default::default()
};
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("checkpoint task should exit within 1s")
.expect("checkpoint task panicked");
}
#[cfg(unix)]
#[tokio::test]
#[serial(checkpoint_skip_metrics, khive_walpin_sidecar_env)]
async fn healthy_checkpoint_tick_reaps_a_dead_walpin_beacon_without_truncate() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("healthy_sidecar_reap.db");
let pool = file_pool(&path);
let sidecar_dir =
crate::walpin::sidecar_dir_for(pool.canonical_path().expect("file-backed pool"));
let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
let dead_pid = 2_000_000_000;
let dead_beacon = crate::walpin::WalpinBeacon {
pid: dead_pid,
process_role: "session".to_string(),
started_at: 1,
sweep_interval_ms: 5_000,
};
crate::walpin::write_beacon(&sidecar_dir, &dead_beacon)
.expect("seed a crashed process's orphan beacon");
let dead_beacon_path = crate::walpin::beacon_path(&sidecar_dir, dead_pid);
assert!(dead_beacon_path.exists(), "orphan fixture must exist");
let cfg = CheckpointConfig {
interval: Duration::from_millis(10),
warn_pages: u64::MAX,
high_water_pages: u64::MAX,
truncate_high_water_pages: u64::MAX,
..CheckpointConfig::default()
};
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
let reaped = wait_for(Duration::from_secs(2), || !dead_beacon_path.exists()).await;
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("checkpoint task should exit within 1s")
.expect("checkpoint task panicked");
assert!(
reaped,
"the ordinary healthy tick must reap positively dead sidecar residue independently of \
TRUNCATE diagnostics"
);
}
#[tokio::test]
#[serial(checkpoint_skip_metrics)]
async fn checkpoint_task_exits_via_shutdown_signal_with_live_event_store_pool_clone() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("wal_task_event_store.db");
let pool = file_pool(&path);
let cfg = CheckpointConfig {
interval: Duration::from_millis(10),
..Default::default()
};
let event_store: Arc<dyn khive_storage::EventStore> =
Arc::new(crate::stores::event::SqlEventStore::new_scoped(
Arc::clone(&pool),
true,
"local".to_string(),
));
let sibling_pool_clone = Arc::clone(&pool);
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let handle = tokio::spawn(run_checkpoint_task(
pool,
cfg,
Some(CheckpointLifecycleOwner::new(event_store, "local")),
shutdown_rx,
true,
));
assert!(
Arc::strong_count(&sibling_pool_clone) > 1,
"test setup must reproduce the multi-owner shape the bug depends on"
);
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect(
"checkpoint task should exit within 1s via the watch signal, \
even with a live sibling Arc<ConnectionPool> clone held by \
the event store",
)
.expect("checkpoint task panicked");
}
#[test]
#[serial]
fn checkpoint_config_env_override() {
std::env::set_var("KHIVE_CHECKPOINT_INTERVAL_MS", "250");
std::env::set_var("KHIVE_WAL_WARN_PAGES", "1500");
std::env::set_var("KHIVE_WAL_HIGH_WATER_PAGES", "8000");
std::env::set_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES", "12000");
std::env::set_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS", "60");
std::env::set_var("KHIVE_WAL_TRUNCATE_BUSY_MS", "500");
std::env::set_var("KHIVE_TX_WARN_SECS", "15");
std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "90");
let cfg = CheckpointConfig::from_env();
std::env::remove_var("KHIVE_CHECKPOINT_INTERVAL_MS");
std::env::remove_var("KHIVE_WAL_WARN_PAGES");
std::env::remove_var("KHIVE_WAL_HIGH_WATER_PAGES");
std::env::remove_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES");
std::env::remove_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS");
std::env::remove_var("KHIVE_WAL_TRUNCATE_BUSY_MS");
std::env::remove_var("KHIVE_TX_WARN_SECS");
std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
assert_eq!(cfg.interval, Duration::from_millis(250));
assert_eq!(cfg.warn_pages, 1500);
assert_eq!(cfg.high_water_pages, 8000);
assert_eq!(cfg.truncate_high_water_pages, 12000);
assert_eq!(cfg.truncate_min_interval, Duration::from_secs(60));
assert_eq!(cfg.truncate_busy_timeout, Duration::from_millis(500));
assert_eq!(cfg.tx_warn_secs, Duration::from_secs(15));
assert_eq!(cfg.tx_max_age_secs, Duration::from_secs(90));
}
#[test]
#[serial]
fn checkpoint_config_defaults_on_invalid_env() {
let default = CheckpointConfig::default();
std::env::set_var("KHIVE_CHECKPOINT_INTERVAL_MS", "not_a_number");
std::env::set_var("KHIVE_WAL_WARN_PAGES", "");
std::env::set_var("KHIVE_WAL_HIGH_WATER_PAGES", "0");
std::env::set_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES", "not_a_number");
std::env::set_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS", "");
std::env::set_var("KHIVE_WAL_TRUNCATE_BUSY_MS", "0");
std::env::set_var("KHIVE_TX_WARN_SECS", "not_a_number");
std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "0");
let cfg = CheckpointConfig::from_env();
std::env::remove_var("KHIVE_CHECKPOINT_INTERVAL_MS");
std::env::remove_var("KHIVE_WAL_WARN_PAGES");
std::env::remove_var("KHIVE_WAL_HIGH_WATER_PAGES");
std::env::remove_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES");
std::env::remove_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS");
std::env::remove_var("KHIVE_WAL_TRUNCATE_BUSY_MS");
std::env::remove_var("KHIVE_TX_WARN_SECS");
std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
assert_eq!(cfg.interval, default.interval);
assert_eq!(cfg.warn_pages, default.warn_pages);
assert_eq!(cfg.high_water_pages, default.high_water_pages);
assert_eq!(
cfg.truncate_high_water_pages,
default.truncate_high_water_pages
);
assert_eq!(cfg.truncate_min_interval, default.truncate_min_interval);
assert_eq!(cfg.truncate_busy_timeout, default.truncate_busy_timeout);
assert_eq!(cfg.tx_warn_secs, default.tx_warn_secs);
assert_eq!(cfg.tx_max_age_secs, default.tx_max_age_secs);
}
#[test]
#[serial(checkpoint_skip_metrics)]
fn checkpoint_high_water_does_not_block_behind_reader() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("high_water_test.db");
let pool = Arc::new(
ConnectionPool::new(PoolConfig {
path: Some(path.clone()),
busy_timeout: Duration::from_millis(2000),
..PoolConfig::for_test()
})
.expect("pool open"),
);
{
let writer = pool.try_writer().unwrap();
writer
.conn()
.execute_batch(
"CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
)
.unwrap();
}
let reader = pool.reader().expect("reader");
reader
.conn()
.execute_batch("BEGIN DEFERRED; SELECT * FROM t;")
.expect("begin read tx");
{
let writer = pool.try_writer().unwrap();
writer
.conn()
.execute_batch("INSERT INTO t VALUES (2);")
.unwrap();
}
let checkpoint_config = CheckpointConfig::default();
let conn = checkpoint_conn(&pool);
let start = std::time::Instant::now();
checkpoint_once(
&pool,
&conn,
&checkpoint_config,
&mut TruncateState::default(),
)
.expect("checkpoint_once must succeed against a healthy dedicated connection");
let elapsed = start.elapsed();
reader.conn().execute_batch("COMMIT;").ok();
drop(reader);
let max_elapsed = checkpoint_config.truncate_busy_timeout / 2;
assert!(
elapsed < max_elapsed,
"checkpoint_once with active reader snapshot took {:?}; expected <{:?} \
(PASSIVE must not block on readers; a TRUNCATE regression would block \
for the configured {:?})",
elapsed,
max_elapsed,
checkpoint_config.truncate_busy_timeout
);
}
#[test]
#[serial]
fn checkpoint_config_rejects_zero_for_all_fields() {
let default = CheckpointConfig::default();
std::env::set_var("KHIVE_CHECKPOINT_INTERVAL_MS", "0");
std::env::set_var("KHIVE_WAL_WARN_PAGES", "0");
std::env::set_var("KHIVE_WAL_HIGH_WATER_PAGES", "0");
std::env::set_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES", "0");
std::env::set_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS", "0");
std::env::set_var("KHIVE_WAL_TRUNCATE_BUSY_MS", "0");
std::env::set_var("KHIVE_TX_WARN_SECS", "0");
std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "0");
let cfg = CheckpointConfig::from_env();
std::env::remove_var("KHIVE_CHECKPOINT_INTERVAL_MS");
std::env::remove_var("KHIVE_WAL_WARN_PAGES");
std::env::remove_var("KHIVE_WAL_HIGH_WATER_PAGES");
std::env::remove_var("KHIVE_WAL_TRUNCATE_HIGH_WATER_PAGES");
std::env::remove_var("KHIVE_WAL_TRUNCATE_MIN_INTERVAL_SECS");
std::env::remove_var("KHIVE_WAL_TRUNCATE_BUSY_MS");
std::env::remove_var("KHIVE_TX_WARN_SECS");
std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
assert_eq!(
cfg.interval, default.interval,
"zero interval must fall back to default"
);
assert_eq!(
cfg.warn_pages, default.warn_pages,
"zero warn_pages must fall back to default"
);
assert_eq!(
cfg.high_water_pages, default.high_water_pages,
"zero high_water_pages must fall back to default"
);
assert_eq!(
cfg.truncate_high_water_pages, default.truncate_high_water_pages,
"zero truncate_high_water_pages must fall back to default"
);
assert_eq!(
cfg.truncate_min_interval, default.truncate_min_interval,
"zero truncate_min_interval must fall back to default"
);
assert_eq!(
cfg.truncate_busy_timeout, default.truncate_busy_timeout,
"zero truncate_busy_timeout must fall back to default"
);
assert_eq!(
cfg.tx_warn_secs, default.tx_warn_secs,
"zero tx_warn_secs must fall back to default"
);
assert_eq!(
cfg.tx_max_age_secs, default.tx_max_age_secs,
"zero tx_max_age_secs must fall back to default"
);
}
#[test]
#[serial]
fn checkpoint_config_rejects_reversed_tx_thresholds() {
let default = CheckpointConfig::default();
std::env::set_var("KHIVE_TX_WARN_SECS", "120");
std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "30");
let cfg = CheckpointConfig::from_env();
std::env::remove_var("KHIVE_TX_WARN_SECS");
std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
assert_eq!(
cfg.tx_warn_secs, default.tx_warn_secs,
"a reversed pair must fall back tx_warn_secs to its default, got: {:?}",
cfg.tx_warn_secs
);
assert_eq!(
cfg.tx_max_age_secs, default.tx_max_age_secs,
"a reversed pair must fall back tx_max_age_secs to its default, got: {:?}",
cfg.tx_max_age_secs
);
}
#[test]
#[serial]
fn checkpoint_config_rejects_equal_tx_thresholds() {
let default = CheckpointConfig::default();
std::env::set_var("KHIVE_TX_WARN_SECS", "60");
std::env::set_var("KHIVE_TX_MAX_AGE_SECS", "60");
let cfg = CheckpointConfig::from_env();
std::env::remove_var("KHIVE_TX_WARN_SECS");
std::env::remove_var("KHIVE_TX_MAX_AGE_SECS");
assert_eq!(
cfg.tx_warn_secs, default.tx_warn_secs,
"an equal pair must fall back tx_warn_secs to its default, got: {:?}",
cfg.tx_warn_secs
);
assert_eq!(
cfg.tx_max_age_secs, default.tx_max_age_secs,
"an equal pair must fall back tx_max_age_secs to its default, got: {:?}",
cfg.tx_max_age_secs
);
}
#[test]
fn skipped_tick_does_not_reset_high_water_crossing_state() {
let mut was_above = false;
assert!(
crossing_warn(true, &mut was_above),
"should fire on first crossing"
);
assert!(was_above);
assert!(was_above, "was_above must stay true across skipped ticks");
let fired = crossing_warn(true, &mut was_above);
assert!(!fired, "WARN must not re-fire while still above threshold");
let fired = crossing_warn(false, &mut was_above);
assert!(!fired);
assert!(!was_above);
let fired = crossing_warn(true, &mut was_above);
assert!(fired, "WARN must fire again on a new below→above crossing");
}
#[test]
fn warn_pages_fires_once_on_crossing_not_every_tick() {
let mut was_above_warn = false;
let fired_1 = crossing_warn(true, &mut was_above_warn);
let fired_2 = crossing_warn(true, &mut was_above_warn);
let fired_3 = crossing_warn(true, &mut was_above_warn);
assert!(fired_1, "WARN must fire on the first in-band tick");
assert!(
!fired_2,
"WARN must not fire on the second consecutive in-band tick"
);
assert!(
!fired_3,
"WARN must not fire on the third consecutive in-band tick"
);
crossing_warn(false, &mut was_above_warn);
assert!(!was_above_warn);
let fired_reentry = crossing_warn(true, &mut was_above_warn);
assert!(
fired_reentry,
"WARN must fire again on re-entry into warn band"
);
}
#[test]
#[serial(tx_registry, checkpoint_skip_metrics)]
fn truncate_attempts_when_high_water_crossed_with_no_prior_attempt() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("truncate_trigger.db");
let pool = file_pool(&path);
{
let writer = pool.try_writer().unwrap();
writer
.conn()
.execute_batch(
"CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
)
.unwrap();
}
let config = CheckpointConfig {
truncate_high_water_pages: 0,
truncate_min_interval: Duration::from_secs(300),
..CheckpointConfig::default()
};
let mut state = TruncateState::default();
assert!(
state.last_attempt.is_none(),
"precondition: no attempt has run yet"
);
let conn = checkpoint_conn(&pool);
checkpoint_once(&pool, &conn, &config, &mut state)
.expect("checkpoint_once must succeed against a healthy dedicated connection");
assert!(
state.last_attempt.is_some(),
"an attempt must be stamped once the high-water threshold is crossed"
);
}
#[test]
#[serial(tx_registry, checkpoint_skip_metrics)]
fn truncate_does_not_attempt_below_high_water() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("truncate_below_threshold.db");
let pool = file_pool(&path);
{
let writer = pool.try_writer().unwrap();
writer
.conn()
.execute_batch(
"CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
)
.unwrap();
}
let config = CheckpointConfig {
truncate_high_water_pages: u64::MAX,
..CheckpointConfig::default()
};
let mut state = TruncateState::default();
let conn = checkpoint_conn(&pool);
checkpoint_once(&pool, &conn, &config, &mut state)
.expect("checkpoint_once must succeed against a healthy dedicated connection");
assert!(
state.last_attempt.is_none(),
"a below-threshold tick must never stamp last_attempt"
);
}
#[test]
#[serial(tx_registry, checkpoint_skip_metrics)]
fn truncate_min_interval_skip_does_not_restamp_last_attempt() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("truncate_min_interval.db");
let pool = file_pool(&path);
{
let writer = pool.try_writer().unwrap();
writer
.conn()
.execute_batch(
"CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
)
.unwrap();
}
let config = CheckpointConfig {
truncate_high_water_pages: 0,
truncate_min_interval: Duration::from_secs(300),
..CheckpointConfig::default()
};
let mut state = TruncateState::default();
let conn = checkpoint_conn(&pool);
checkpoint_once(&pool, &conn, &config, &mut state)
.expect("checkpoint_once must succeed against a healthy dedicated connection");
let first_attempt = state.last_attempt.expect("first tick must attempt");
checkpoint_once(&pool, &conn, &config, &mut state)
.expect("checkpoint_once must succeed against a healthy dedicated connection");
let second_attempt = state.last_attempt.expect("attempt timestamp must persist");
assert_eq!(
first_attempt, second_attempt,
"a tick within truncate_min_interval must not re-stamp last_attempt"
);
}
#[test]
#[serial(tx_registry, checkpoint_skip_metrics)]
fn checkpoint_once_proceeds_and_can_attempt_truncate_while_pool_writer_held() {
reset_checkpoint_metrics_for_tests();
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("truncate_busy_skip.db");
let pool = file_pool(&path);
{
let writer = pool.try_writer().unwrap();
writer
.conn()
.execute_batch(
"CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
)
.unwrap();
}
let conn = checkpoint_conn(&pool);
let _held = pool.try_writer().unwrap();
let config = CheckpointConfig {
truncate_high_water_pages: 0,
..CheckpointConfig::default()
};
let mut state = TruncateState::default();
checkpoint_once(&pool, &conn, &config, &mut state).expect(
"checkpoint_once must observe normally on its own dedicated connection even \
while a concurrent caller holds the pool's writer mutex",
);
assert!(
state.last_attempt.is_some(),
"a threshold-armed tick must still evaluate (and attempt) TRUNCATE even while \
the pool writer is held — the dedicated connection is unaffected by it"
);
assert_eq!(
checkpoint_skipped_ticks(),
0,
"a busy pool writer must no longer count as a skipped checkpoint tick"
);
assert_eq!(
checkpoint_consecutive_skips(),
0,
"a busy pool writer must not bump the consecutive-skip run length"
);
}
#[test]
#[serial(checkpoint_skip_metrics)]
fn all_checkpoint_metrics_callers_are_serial_tagged() {
const SELF_SRC: &str = include_str!("checkpoint.rs");
let lines: Vec<&str> = SELF_SRC.lines().collect();
let attr_starts: Vec<usize> = lines
.iter()
.enumerate()
.filter(|(_, l)| {
let t = l.trim();
t == "#[test]" || t.starts_with("#[tokio::test")
})
.map(|(i, _)| i)
.collect();
let mut offenders = Vec::new();
for (idx, &start) in attr_starts.iter().enumerate() {
let end = attr_starts.get(idx + 1).copied().unwrap_or(lines.len());
let span = &lines[start..end];
let touches_shared_metrics = span.iter().any(|l| {
l.contains("checkpoint_once(")
|| l.contains("checkpoint_once_core(")
|| l.contains("run_checkpoint_task(")
});
if !touches_shared_metrics {
continue;
}
let mut in_serial_attr = false;
let has_group_tag = span.iter().any(|line| {
let trimmed = line.trim();
if !in_serial_attr {
in_serial_attr = trimmed.starts_with("#[serial(");
}
if !in_serial_attr {
return false;
}
let has_group = trimmed
.split(|ch: char| !ch.is_ascii_alphanumeric() && ch != '_')
.any(|token| token == "checkpoint_skip_metrics");
if trimmed.ends_with(")]") {
in_serial_attr = false;
}
has_group
});
if !has_group_tag {
let name = span
.iter()
.find_map(|l| {
let t = l.trim_start();
let t = t.strip_prefix("pub(crate) ").unwrap_or(t);
let t = t.strip_prefix("pub ").unwrap_or(t);
let t = t.strip_prefix("async ").unwrap_or(t);
t.strip_prefix("fn ")
.map(|rest| rest.split(['(', '<']).next().unwrap_or("").trim())
})
.unwrap_or("<unknown test>");
offenders.push(name.to_string());
}
}
assert!(
offenders.is_empty(),
"these tests call checkpoint_once/checkpoint_once_core/run_checkpoint_task (which write the \
process-wide LAST_WAL_PAGES/CHECKPOINT_* atomics via query_wal_pages) but \
are not tagged #[serial(checkpoint_skip_metrics)] (or a group including it); \
an untagged caller running concurrently on cargo's default test thread pool \
can clobber those atomics mid-assertion in another test (the #828/#845 race): \
{offenders:?}"
);
}
#[test]
#[serial(tx_registry, checkpoint_skip_metrics)]
fn observed_tick_resets_consecutive_skips_but_not_lifetime_total() {
reset_checkpoint_metrics_for_tests();
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("skip_then_observe.db");
let pool = file_pool(&path);
{
let writer = pool.try_writer().unwrap();
writer
.conn()
.execute_batch(
"CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
)
.unwrap();
}
note_checkpoint_skipped();
note_checkpoint_skipped();
assert_eq!(checkpoint_skipped_ticks(), 2);
assert_eq!(checkpoint_consecutive_skips(), 2);
let conn = checkpoint_conn(&pool);
let mut state = TruncateState::default();
checkpoint_once(&pool, &conn, &CheckpointConfig::default(), &mut state)
.expect("checkpoint_once must succeed against a healthy dedicated connection");
assert_eq!(
checkpoint_skipped_ticks(),
2,
"an observed tick must not change the lifetime skipped-tick total"
);
assert_eq!(
checkpoint_consecutive_skips(),
0,
"an observed tick must reset the consecutive-skip run length"
);
}
#[test]
fn note_truncate_outcome_warns_once_at_third_consecutive_failure() {
let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let subscriber = CaptureSubscriber {
events: std::sync::Arc::clone(&buffer),
};
let config = CheckpointConfig {
warn_pages: 2000,
..CheckpointConfig::default()
};
let mut state = TruncateState::default();
tracing::subscriber::with_default(subscriber, || {
note_truncate_outcome(&config, 5000, &mut state);
note_truncate_outcome(&config, 5000, &mut state);
note_truncate_outcome(&config, 5000, &mut state);
note_truncate_outcome(&config, 5000, &mut state);
});
assert_eq!(state.consecutive_failures, 4);
let events = buffer.lock().unwrap();
let escalation_count = events
.iter()
.filter(|e| {
e.message.as_deref()
== Some(
"WAL TRUNCATE has failed to clear WAL pressure for 3 consecutive attempts",
)
})
.count();
assert_eq!(
escalation_count, 1,
"escalation WARN must fire exactly once at the 3rd consecutive failure, got: {events:?}"
);
note_truncate_outcome(&config, 100, &mut state);
assert_eq!(
state.consecutive_failures, 0,
"an attempt that clears warn_pages must reset the consecutive-failure counter"
);
}
fn severity_test_config() -> CheckpointConfig {
CheckpointConfig {
warn_pages: 100,
warn_sustained_cycles: 3,
..CheckpointConfig::default()
}
}
#[test]
fn severity_ladder_info_on_first_crossing_no_warn() {
let config = severity_test_config();
let mut state = CheckpointSeverityState::default();
let below = state.observe_wal_pages(10, &config);
assert!(below.is_empty(), "below-warn tick must emit nothing");
let above = state.observe_wal_pages(150, &config);
assert_eq!(
above,
vec![CheckpointSeverityEmission {
rung: CheckpointSeverityRung::Info,
wal_pages: 150,
threshold_pages: 100,
consecutive_cycles: 1,
}],
"first below->above crossing must emit exactly one INFO and no WARN"
);
}
#[test]
fn severity_ladder_warn_on_third_consecutive_cycle() {
let config = severity_test_config();
let mut state = CheckpointSeverityState::default();
let tick1 = state.observe_wal_pages(150, &config);
assert_eq!(tick1.len(), 1);
assert_eq!(tick1[0].rung, CheckpointSeverityRung::Info);
let tick2 = state.observe_wal_pages(150, &config);
assert!(
tick2.is_empty(),
"second consecutive above-warn tick must emit nothing yet"
);
let tick3 = state.observe_wal_pages(150, &config);
assert_eq!(
tick3,
vec![CheckpointSeverityEmission {
rung: CheckpointSeverityRung::Warn,
wal_pages: 150,
threshold_pages: 100,
consecutive_cycles: 3,
}],
"WARN must fire exactly on the third consecutive above-warn tick"
);
let tick4 = state.observe_wal_pages(150, &config);
assert!(
tick4.is_empty(),
"WARN must not repeat on a fourth consecutive above-warn tick"
);
}
#[test]
fn severity_ladder_rearms_warn_after_drain() {
let config = severity_test_config();
let mut state = CheckpointSeverityState::default();
for _ in 0..3 {
state.observe_wal_pages(150, &config);
}
assert!(state.warn_emitted_for_episode);
let drain = state.observe_wal_pages(10, &config);
assert!(drain.is_empty(), "a draining tick must emit nothing");
let reentry = state.observe_wal_pages(150, &config);
assert_eq!(reentry.len(), 1);
assert_eq!(reentry[0].rung, CheckpointSeverityRung::Info);
let mid = state.observe_wal_pages(150, &config);
assert!(mid.is_empty());
let second_warn = state.observe_wal_pages(150, &config);
assert_eq!(
second_warn,
vec![CheckpointSeverityEmission {
rung: CheckpointSeverityRung::Warn,
wal_pages: 150,
threshold_pages: 100,
consecutive_cycles: 3,
}],
"a fresh elevation episode after a drain must WARN again"
);
}
#[test]
fn severity_ladder_isolated_crossings_never_warn() {
let config = severity_test_config();
let mut state = CheckpointSeverityState::default();
for _ in 0..3 {
let crossing = state.observe_wal_pages(150, &config);
assert_eq!(
crossing.len(),
1,
"each isolated crossing must emit exactly one INFO"
);
assert_eq!(crossing[0].rung, CheckpointSeverityRung::Info);
let drain = state.observe_wal_pages(10, &config);
assert!(drain.is_empty(), "the drain tick must emit nothing");
}
assert!(
!state.warn_emitted_for_episode,
"isolated single-tick crossings must never accumulate into a WARN"
);
}
#[test]
fn severity_ladder_never_emits_alarm() {
let config = CheckpointConfig {
warn_pages: 100,
warn_sustained_cycles: 1,
..CheckpointConfig::default()
};
let mut state = CheckpointSeverityState::default();
for wal_pages in [150, 200, 250, u64::MAX] {
let emissions = state.observe_wal_pages(wal_pages, &config);
assert!(
emissions
.iter()
.all(|e| e.rung != CheckpointSeverityRung::Alarm),
"observe_wal_pages must never emit the ALARM rung, got: {emissions:?}"
);
}
}
fn tx_age_test_config() -> CheckpointConfig {
CheckpointConfig {
tx_warn_secs: Duration::from_secs(30),
tx_max_age_secs: Duration::from_secs(120),
..CheckpointConfig::default()
}
}
fn tx_id(n: u64) -> khive_storage::tx_registry::TxId {
khive_storage::tx_registry::TxId(n)
}
#[test]
fn tx_age_sweep_empty_registry_emits_nothing() {
let config = tx_age_test_config();
let mut state = TxAgeSweepState::default();
let emissions = state.observe(None, config.tx_warn_secs, config.tx_max_age_secs);
assert!(emissions.is_empty(), "no open entry must emit nothing");
}
#[test]
fn tx_age_sweep_fresh_entry_emits_nothing() {
let config = tx_age_test_config();
let mut state = TxAgeSweepState::default();
let emissions = state.observe(
Some((
tx_id(1),
Duration::from_secs(5),
Some("fresh_span".to_string()),
)),
config.tx_warn_secs,
config.tx_max_age_secs,
);
assert!(emissions.is_empty(), "a fresh entry must emit nothing");
}
#[test]
fn tx_age_sweep_warn_fires_once_on_crossing() {
let config = tx_age_test_config();
let mut state = TxAgeSweepState::default();
let tick1 = state.observe(
Some((
tx_id(1),
Duration::from_secs(45),
Some("stale_span".to_string()),
)),
config.tx_warn_secs,
config.tx_max_age_secs,
);
assert_eq!(
tick1,
vec![TxAgeEmission {
rung: TxAgeRung::Warn,
age: Duration::from_secs(45),
label: Some("stale_span".to_string()),
}],
"crossing tx_warn_secs must emit exactly one Warn"
);
let tick2 = state.observe(
Some((
tx_id(1),
Duration::from_secs(50),
Some("stale_span".to_string()),
)),
config.tx_warn_secs,
config.tx_max_age_secs,
);
assert!(
tick2.is_empty(),
"Warn must not repeat while the entry stays in the warn band"
);
}
#[test]
fn tx_age_sweep_stale_fires_once_on_crossing() {
let config = tx_age_test_config();
let mut state = TxAgeSweepState::default();
state.observe(
Some((
tx_id(1),
Duration::from_secs(45),
Some("stuck_writer_task_tx".to_string()),
)),
config.tx_warn_secs,
config.tx_max_age_secs,
);
let tick = state.observe(
Some((
tx_id(1),
Duration::from_secs(130),
Some("stuck_writer_task_tx".to_string()),
)),
config.tx_warn_secs,
config.tx_max_age_secs,
);
assert_eq!(
tick,
vec![TxAgeEmission {
rung: TxAgeRung::Stale,
age: Duration::from_secs(130),
label: Some("stuck_writer_task_tx".to_string()),
}],
"crossing tx_max_age_secs must emit exactly one Stale"
);
let tick_repeat = state.observe(
Some((
tx_id(1),
Duration::from_secs(200),
Some("stuck_writer_task_tx".to_string()),
)),
config.tx_warn_secs,
config.tx_max_age_secs,
);
assert!(
tick_repeat.is_empty(),
"Stale must not repeat while the entry stays above tx_max_age_secs"
);
}
#[test]
fn tx_age_sweep_already_stale_entry_emits_both_rungs_same_tick() {
let config = tx_age_test_config();
let mut state = TxAgeSweepState::default();
let tick = state.observe(
Some((
tx_id(1),
Duration::from_secs(300),
Some("ancient_tx".to_string()),
)),
config.tx_warn_secs,
config.tx_max_age_secs,
);
assert_eq!(
tick,
vec![
TxAgeEmission {
rung: TxAgeRung::Warn,
age: Duration::from_secs(300),
label: Some("ancient_tx".to_string()),
},
TxAgeEmission {
rung: TxAgeRung::Stale,
age: Duration::from_secs(300),
label: Some("ancient_tx".to_string()),
},
],
"an already-stale entry must cross both rungs on its first observed tick"
);
}
#[test]
fn tx_age_sweep_rearms_after_entry_clears() {
let config = tx_age_test_config();
let mut state = TxAgeSweepState::default();
state.observe(
Some((
tx_id(1),
Duration::from_secs(150),
Some("first_span".to_string()),
)),
config.tx_warn_secs,
config.tx_max_age_secs,
);
let cleared = state.observe(None, config.tx_warn_secs, config.tx_max_age_secs);
assert!(cleared.is_empty(), "a clearing tick must emit nothing");
let fresh = state.observe(
Some((
tx_id(2),
Duration::from_secs(2),
Some("second_span".to_string()),
)),
config.tx_warn_secs,
config.tx_max_age_secs,
);
assert!(fresh.is_empty(), "a fresh oldest entry must emit nothing");
let rewarn = state.observe(
Some((
tx_id(2),
Duration::from_secs(35),
Some("second_span".to_string()),
)),
config.tx_warn_secs,
config.tx_max_age_secs,
);
assert_eq!(
rewarn,
vec![TxAgeEmission {
rung: TxAgeRung::Warn,
age: Duration::from_secs(35),
label: Some("second_span".to_string()),
}],
"a fresh stale episode after a clear must Warn again"
);
}
#[test]
fn tx_age_sweep_stale_replacement_without_intervening_clear_still_names_new_entry() {
let config = tx_age_test_config();
let mut state = TxAgeSweepState::default();
let tick_a = state.observe(
Some((
tx_id(1),
Duration::from_secs(300),
Some("stale_entry_a".to_string()),
)),
config.tx_warn_secs,
config.tx_max_age_secs,
);
assert_eq!(
tick_a.len(),
2,
"entry A must cross both rungs on its first observed tick, got: {tick_a:?}"
);
let tick_b = state.observe(
Some((
tx_id(2),
Duration::from_secs(400),
Some("stale_entry_b".to_string()),
)),
config.tx_warn_secs,
config.tx_max_age_secs,
);
assert_eq!(
tick_b,
vec![
TxAgeEmission {
rung: TxAgeRung::Warn,
age: Duration::from_secs(400),
label: Some("stale_entry_b".to_string()),
},
TxAgeEmission {
rung: TxAgeRung::Stale,
age: Duration::from_secs(400),
label: Some("stale_entry_b".to_string()),
},
],
"a same-tick identity change to an already-stale successor must re-emit both \
rungs naming the NEW entry, got: {tick_b:?}"
);
}
#[test]
fn tx_age_sweep_uses_configured_thresholds_not_hardcoded_defaults() {
let config = CheckpointConfig {
tx_warn_secs: Duration::from_millis(1),
tx_max_age_secs: Duration::from_millis(2),
..CheckpointConfig::default()
};
let mut state = TxAgeSweepState::default();
let tick = state.observe(
Some((
tx_id(1),
Duration::from_millis(5),
Some("fast_cap_span".to_string()),
)),
config.tx_warn_secs,
config.tx_max_age_secs,
);
assert_eq!(
tick.len(),
2,
"a millisecond-scale cap must cross both rungs immediately, got: {tick:?}"
);
}
#[test]
#[serial(tx_registry, checkpoint_skip_metrics)]
fn tx_age_sweep_names_long_lived_reader_pinning_wal_past_high_water() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("tx_age_sweep_reader_pin.db");
let pool = file_pool(&path);
{
let writer = pool.try_writer().unwrap();
writer
.conn()
.execute_batch(
"CREATE TABLE IF NOT EXISTS t (x INTEGER); INSERT INTO t VALUES (1);",
)
.unwrap();
}
let reader = pool.reader().expect("reader");
reader
.conn()
.execute_batch("BEGIN DEFERRED; SELECT * FROM t;")
.expect("begin read tx");
let _tx_handle =
khive_storage::tx_registry::register(Some("tx_age_sweep_reader_pin_test".to_string()));
let config = CheckpointConfig {
high_water_pages: 1,
tx_warn_secs: Duration::from_millis(1),
tx_max_age_secs: Duration::from_millis(1),
..CheckpointConfig::default()
};
{
let writer = pool.try_writer().unwrap();
for i in 0..50 {
writer
.conn()
.execute_batch(&format!("INSERT INTO t VALUES ({i});"))
.unwrap();
}
}
let conn = checkpoint_conn(&pool);
let wal_pages = checkpoint_once(&pool, &conn, &config, &mut TruncateState::default())
.expect("checkpoint_once must succeed against a healthy dedicated connection");
assert!(
wal_pages >= config.high_water_pages,
"test setup must actually drive wal_pages ({wal_pages}) past high_water_pages \
({}) for this regression to mean anything",
config.high_water_pages
);
std::thread::sleep(Duration::from_millis(5));
let our_entry = khive_storage::tx_registry::snapshot()
.into_iter()
.find(|(_, label)| label.as_deref() == Some("tx_age_sweep_reader_pin_test"))
.expect("this test's own tx_registry entry must still be open");
let mut tx_age_state = TxAgeSweepState::default();
let emissions = tx_age_state.observe(
Some((tx_id(1), our_entry.0, our_entry.1)),
config.tx_warn_secs,
config.tx_max_age_secs,
);
assert!(
emissions.iter().any(|e| e.rung == TxAgeRung::Stale
&& e.label.as_deref() == Some("tx_age_sweep_reader_pin_test")),
"expected a Stale emission naming the pinning reader, got: {emissions:?}"
);
reader.conn().execute_batch("COMMIT;").ok();
drop(reader);
drop(_tx_handle);
}
#[test]
#[serial(tx_registry, checkpoint_skip_metrics)]
fn tx_age_sweep_own_entry_survives_concurrent_older_registration() {
let _decoy = khive_storage::tx_registry::register(Some("decoy_unrelated_span".to_string()));
std::thread::sleep(Duration::from_millis(2));
let _own = khive_storage::tx_registry::register(Some("this_test_own_span".to_string()));
std::thread::sleep(Duration::from_millis(5));
let global_oldest = khive_storage::tx_registry::oldest().expect("registry not empty");
assert_ne!(
global_oldest.2.as_deref(),
Some("this_test_own_span"),
"test setup must reproduce the race: an older, unrelated entry must be \
the current global oldest, got: {global_oldest:?}"
);
let our_entry = khive_storage::tx_registry::snapshot()
.into_iter()
.find(|(_, label)| label.as_deref() == Some("this_test_own_span"))
.expect("this test's own tx_registry entry must still be open");
let config = CheckpointConfig {
tx_warn_secs: Duration::from_millis(1),
tx_max_age_secs: Duration::from_millis(1),
..CheckpointConfig::default()
};
let mut state = TxAgeSweepState::default();
let emissions = state.observe(
Some((tx_id(2), our_entry.0, our_entry.1)),
config.tx_warn_secs,
config.tx_max_age_secs,
);
assert!(
emissions
.iter()
.any(|e| e.rung == TxAgeRung::Stale
&& e.label.as_deref() == Some("this_test_own_span")),
"expected a Stale emission naming this test's own span despite an older, \
unrelated concurrent registration, got: {emissions:?}"
);
}
#[test]
#[serial]
fn checkpoint_config_warn_sustained_cycles_env_override() {
let default = CheckpointConfig::default();
assert_eq!(default.warn_sustained_cycles, DEFAULT_WARN_SUSTAINED_CYCLES);
std::env::set_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES", "5");
let cfg = CheckpointConfig::from_env();
std::env::remove_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES");
assert_eq!(cfg.warn_sustained_cycles, 5);
std::env::set_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES", "0");
let cfg_zero = CheckpointConfig::from_env();
std::env::remove_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES");
assert_eq!(
cfg_zero.warn_sustained_cycles, DEFAULT_WARN_SUSTAINED_CYCLES,
"zero must fall back to the default"
);
std::env::set_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES", "not_a_number");
let cfg_invalid = CheckpointConfig::from_env();
std::env::remove_var("KHIVE_WAL_WARN_SUSTAINED_CYCLES");
assert_eq!(
cfg_invalid.warn_sustained_cycles,
DEFAULT_WARN_SUSTAINED_CYCLES
);
}
#[derive(Clone, Copy)]
enum FakeAppendBehavior {
Record,
Fail,
}
struct FakeEventStore {
events: std::sync::Mutex<Vec<khive_storage::Event>>,
append_attempts: std::sync::atomic::AtomicUsize,
append_behavior: FakeAppendBehavior,
}
impl Default for FakeEventStore {
fn default() -> Self {
Self {
events: std::sync::Mutex::new(Vec::new()),
append_attempts: std::sync::atomic::AtomicUsize::new(0),
append_behavior: FakeAppendBehavior::Record,
}
}
}
impl FakeEventStore {
fn failing() -> Self {
Self {
append_behavior: FakeAppendBehavior::Fail,
..Self::default()
}
}
}
#[async_trait::async_trait]
impl khive_storage::EventStore for FakeEventStore {
async fn append_event(
&self,
event: khive_storage::Event,
) -> khive_storage::StorageResult<()> {
self.append_attempts
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
match self.append_behavior {
FakeAppendBehavior::Record => {
self.events.lock().unwrap().push(event);
Ok(())
}
FakeAppendBehavior::Fail => Err(khive_storage::StorageError::Internal(
"synthetic checkpoint lifecycle append failure".to_string(),
)),
}
}
async fn append_events(
&self,
events: Vec<khive_storage::Event>,
) -> khive_storage::StorageResult<khive_storage::BatchWriteSummary> {
let count = events.len() as u64;
self.events.lock().unwrap().extend(events);
Ok(khive_storage::BatchWriteSummary {
attempted: count,
affected: count,
..khive_storage::BatchWriteSummary::default()
})
}
async fn get_event(
&self,
id: uuid::Uuid,
) -> khive_storage::StorageResult<Option<khive_storage::Event>> {
Ok(self
.events
.lock()
.unwrap()
.iter()
.find(|e| e.id == id)
.cloned())
}
async fn query_events(
&self,
_filter: khive_storage::EventFilter,
_page: khive_storage::PageRequest,
) -> khive_storage::StorageResult<khive_storage::Page<khive_storage::Event>> {
unimplemented!("not exercised by the checkpoint lifecycle-event tests")
}
async fn count_events(
&self,
_filter: khive_storage::EventFilter,
) -> khive_storage::StorageResult<u64> {
Ok(self.events.lock().unwrap().len() as u64)
}
}
#[test]
fn checkpoint_outcome_should_emit_covers_all_transitions() {
assert!(
checkpoint_outcome_should_emit(true, false),
"first elevated tick must emit"
);
assert!(
!checkpoint_outcome_should_emit(true, true),
"sustained elevated ticks must aggregate in memory instead of writing the WAL"
);
assert!(
checkpoint_outcome_should_emit(false, true),
"the single drain row (elevated -> healthy) must emit"
);
assert!(
!checkpoint_outcome_should_emit(false, false),
"an ordinary below-warn tick must not emit"
);
}
#[test]
fn persistent_pressure_lifecycle_rows_are_o_state_transitions() {
let observations = [true; 128].into_iter().chain([false]).chain([false; 128]);
let mut was_elevated = false;
let writes = observations
.filter(|above_warn| {
let emit = checkpoint_outcome_should_emit(*above_warn, was_elevated);
if emit {
was_elevated = *above_warn;
}
emit
})
.count();
assert_eq!(
writes, 2,
"one elevation row plus one recovery summary must cover any number of attempts"
);
}
#[test]
fn checkpoint_pressure_episode_retains_recovery_summary() {
let mut episode = CheckpointPressureEpisode::start(2_500);
episode.observe(2_300);
episode.observe(8_100);
episode.observe(4_000);
assert_eq!(episode.elevated_ticks, 4);
assert_eq!(episode.peak_wal_pages, 8_100);
}
fn drive_pressure_ticks(
config: &CheckpointConfig,
ticks: &[(bool, u64)],
mut fail_on: impl FnMut(usize) -> bool,
) -> Vec<khive_storage::CheckpointOutcomeRecordedPayload> {
let mut event_elevation_open = false;
let mut pressure_episode: Option<CheckpointPressureEpisode> = None;
let mut pending_recovery: Option<khive_storage::CheckpointOutcomeRecordedPayload> = None;
let mut delivered = Vec::new();
let mut call_index = 0usize;
for &(above_warn, wal_pages) in ticks {
observe_checkpoint_pressure_tick(
above_warn,
wal_pages,
false,
false,
config,
&mut event_elevation_open,
&mut pressure_episode,
&mut pending_recovery,
|payload| {
let idx = call_index;
call_index += 1;
if fail_on(idx) {
false
} else {
delivered.push(payload);
true
}
},
);
}
delivered
}
#[test]
fn dropped_recovery_handoff_does_not_merge_pressure_episodes() {
let config = CheckpointConfig {
warn_pages: 1_000,
..CheckpointConfig::default()
};
let ticks = [
(true, 1_500), (true, 1_800), (true, 2_000), (false, 500), (false, 400), (true, 3_000), (true, 3_500), (false, 300), ];
let delivered = drive_pressure_ticks(&config, &ticks, |idx| matches!(idx, 1..=3));
assert_eq!(
delivered.len(),
4,
"expected episode-1 open, episode-1 delayed recovery, episode-2 open, \
episode-2 recovery: {delivered:?}"
);
let ep1_open = &delivered[0];
assert!(ep1_open.above_warn);
assert_eq!(ep1_open.episode_elevated_ticks, Some(1));
assert_eq!(ep1_open.episode_peak_wal_pages, Some(1_500));
let ep1_recovery = &delivered[1];
assert!(
!ep1_recovery.above_warn,
"episode 1's recovery must be delivered BEFORE episode 2's opening; \
an opening in this slot means the barrier failed: {delivered:?}"
);
assert_eq!(
ep1_recovery.episode_elevated_ticks,
Some(3),
"episode 1's delayed recovery must report only its own 3 elevated ticks, \
not ticks absorbed from episode 2"
);
assert_eq!(ep1_recovery.episode_peak_wal_pages, Some(2_000));
let ep2_open = &delivered[2];
assert!(ep2_open.above_warn);
assert_eq!(
ep2_open.episode_elevated_ticks,
Some(2),
"episode 2 opens fresh (never continuing episode 1's count), deferred one \
tick by the barrier, so its opening reports 2 elevated ticks"
);
assert_eq!(ep2_open.episode_peak_wal_pages, Some(3_500));
let ep2_recovery = &delivered[3];
assert!(!ep2_recovery.above_warn);
assert_eq!(
ep2_recovery.episode_elevated_ticks,
Some(2),
"episode 2's recovery must report only its own 2 elevated ticks"
);
assert_eq!(ep2_recovery.episode_peak_wal_pages, Some(3_500));
}
#[test]
fn episode_elapsed_entirely_behind_barrier_is_discarded_not_reordered() {
let config = CheckpointConfig {
warn_pages: 1_000,
..CheckpointConfig::default()
};
let ticks = [
(true, 1_500), (false, 500), (true, 9_000), (false, 400), (false, 300), (true, 2_500), (false, 200), ];
let delivered = drive_pressure_ticks(&config, &ticks, |idx| matches!(idx, 1..=3));
let peaks: Vec<_> = delivered
.iter()
.map(|payload| (payload.above_warn, payload.episode_peak_wal_pages))
.collect();
assert_eq!(
peaks,
vec![
(true, Some(1_500)), (false, Some(1_500)), (true, Some(2_500)), (false, Some(2_500)), ],
"an episode elapsed entirely behind the barrier must not surface late or \
out of order: {delivered:?}"
);
}
#[test]
fn no_dropped_handoff_reports_two_separate_episodes() {
let config = CheckpointConfig {
warn_pages: 1_000,
..CheckpointConfig::default()
};
let ticks = [
(true, 1_500),
(true, 1_800),
(true, 2_000),
(false, 500),
(false, 400),
(true, 3_000),
(true, 3_500),
(false, 300),
];
let delivered = drive_pressure_ticks(&config, &ticks, |_idx| false);
assert_eq!(delivered.len(), 4, "{delivered:?}");
assert_eq!(
(
delivered[0].above_warn,
delivered[0].episode_elevated_ticks,
delivered[0].episode_peak_wal_pages
),
(true, Some(1), Some(1_500)),
"episode 1 open"
);
assert_eq!(
(
delivered[1].above_warn,
delivered[1].episode_elevated_ticks,
delivered[1].episode_peak_wal_pages
),
(false, Some(3), Some(2_000)),
"episode 1 recovery"
);
assert_eq!(
(
delivered[2].above_warn,
delivered[2].episode_elevated_ticks,
delivered[2].episode_peak_wal_pages
),
(true, Some(1), Some(3_000)),
"episode 2 open"
);
assert_eq!(
(
delivered[3].above_warn,
delivered[3].episode_elevated_ticks,
delivered[3].episode_peak_wal_pages
),
(false, Some(2), Some(3_500)),
"episode 2 recovery"
);
}
#[test]
#[serial(checkpoint_skip_metrics)]
fn pressure_diagnostics_count_observations_and_transitions_separately() {
reset_checkpoint_metrics_for_tests();
note_checkpoint_pressure_observation(true, false);
note_checkpoint_pressure_observation(true, true);
note_checkpoint_pressure_observation(true, true);
note_checkpoint_pressure_observation(false, true);
note_checkpoint_pressure_observation(false, false);
assert_eq!(checkpoint_pressure_elevated_ticks(), 3);
assert_eq!(checkpoint_pressure_episodes_started(), 1);
assert_eq!(checkpoint_pressure_episodes_recovered(), 1);
assert_eq!(checkpoint_lifecycle_append_attempts(), 0);
}
#[tokio::test]
#[serial(checkpoint_skip_metrics)]
async fn checkpoint_task_emits_one_opening_for_persistent_pressure() {
reset_checkpoint_metrics_for_tests();
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("outcome_emit.db");
let pool = file_pool(&path);
let cfg = CheckpointConfig {
interval: Duration::from_millis(10),
warn_pages: 0,
..CheckpointConfig::default()
};
let store = Arc::new(FakeEventStore::default());
let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let handle = tokio::spawn(run_checkpoint_task(
pool,
cfg,
Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
shutdown_rx,
true,
));
let progressed = wait_for(Duration::from_secs(10), || {
checkpoint_pressure_elevated_ticks() >= 10
})
.await;
let emitted = wait_for(Duration::from_secs(10), || {
!store.events.lock().unwrap().is_empty()
})
.await;
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("checkpoint task should exit within 1s")
.expect("checkpoint task panicked");
let events = store.events.lock().unwrap();
assert!(
progressed,
"the simulated persistent-pressure episode must span at least ten checkpoint ticks"
);
assert!(
emitted,
"an always-elevated config must append one CheckpointOutcomeRecorded event \
within the poll deadline"
);
assert_eq!(
checkpoint_lifecycle_append_attempts(),
1,
"primary-store lifecycle writes must stay O(state transitions), not O(attempts)"
);
assert_eq!(events.len(), 1);
assert_eq!(events[0].payload_schema_version, 2);
assert_eq!(events[0].payload["episode_elevated_ticks"], 1);
assert_eq!(
events[0].payload["episode_peak_wal_pages"],
events[0].payload["wal_pages"]
);
assert!(
events
.iter()
.all(|e| e.kind == khive_types::EventKind::CheckpointOutcomeRecorded),
"every appended event must be CheckpointOutcomeRecorded, got: {events:?}"
);
assert!(
events.iter().all(|e| e.namespace == "local"),
"events must be stamped with the namespace passed to run_checkpoint_task"
);
}
#[tokio::test]
#[serial(checkpoint_skip_metrics)]
async fn checkpoint_cycles_and_task_shutdown_do_not_wait_for_a_contended_lifecycle_writer() {
reset_checkpoint_metrics_for_tests();
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("outcome_contended_sink.db");
let checkpoint_pool = file_pool(&path);
let event_pool = Arc::new(
ConnectionPool::new(PoolConfig {
path: None,
checkout_timeout: Duration::from_secs(5),
write_queue_enabled: Some(false),
..PoolConfig::default()
})
.expect("event pool"),
);
{
let writer = event_pool.try_writer().expect("initialize event schema");
crate::stores::event::ensure_events_schema(writer.conn())
.expect("initialize event schema");
}
let event_store: Arc<dyn khive_storage::EventStore> =
Arc::new(crate::stores::event::SqlEventStore::new_scoped(
Arc::clone(&event_pool),
false,
"local",
));
let held_event_writer = event_pool
.try_writer()
.expect("hold the event-store writer");
let cfg = CheckpointConfig {
interval: Duration::from_millis(10),
warn_pages: 0,
..CheckpointConfig::default()
};
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let handle = tokio::spawn(run_checkpoint_task(
checkpoint_pool,
cfg,
Some(CheckpointLifecycleOwner::new(event_store, "local")),
shutdown_rx,
true,
));
let progressed = wait_for(Duration::from_secs(2), || {
checkpoint_pressure_elevated_ticks() >= 10
})
.await;
assert!(
progressed,
"checkpoint observations must continue while the lifecycle append is contended"
);
assert_eq!(checkpoint_lifecycle_append_attempts(), 1);
assert_eq!(checkpoint_lifecycle_enqueue_drops(), 0);
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect(
"the run_checkpoint_task handle must not wait for the event store's \
five-second writer checkout",
)
.expect("checkpoint task panicked");
drop(held_event_writer);
}
#[tokio::test]
#[serial(checkpoint_skip_metrics)]
async fn checkpoint_task_continues_after_lifecycle_append_failure() {
reset_checkpoint_metrics_for_tests();
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("outcome_failing_sink.db");
let pool = file_pool(&path);
let store = Arc::new(FakeEventStore::failing());
let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let subscriber = CaptureSubscriber {
events: std::sync::Arc::clone(&buffer),
};
let _tracing_guard = tracing::subscriber::set_default(subscriber);
let cfg = CheckpointConfig {
interval: Duration::from_millis(10),
warn_pages: 0,
..CheckpointConfig::default()
};
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let handle = tokio::spawn(run_checkpoint_task(
pool,
cfg,
Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
shutdown_rx,
true,
));
let progressed = wait_for(Duration::from_secs(2), || {
checkpoint_pressure_elevated_ticks() >= 10
})
.await;
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("checkpoint task should remain responsive after sink failure")
.expect("checkpoint task panicked");
assert!(
progressed,
"a failed append must not terminate or stall the checkpoint task"
);
assert_eq!(
store
.append_attempts
.load(std::sync::atomic::Ordering::Relaxed),
1,
"a persistent pressure state must not retry one primary-store append per tick"
);
assert_eq!(checkpoint_lifecycle_append_attempts(), 1);
assert_eq!(checkpoint_lifecycle_append_failures(), 1);
let captured = buffer.lock().unwrap().clone();
assert!(
captured.iter().any(|event| event.message.as_deref()
== Some("checkpoint lifecycle event append failed")),
"lifecycle append failures must remain observable; got: {:?}",
captured
);
}
#[tokio::test]
#[serial(checkpoint_skip_metrics)]
async fn secondary_checkpoint_task_with_lifecycle_ownership_emits_outcome_events() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("secondary_outcome.db");
let pool = file_pool(&path);
let cfg = CheckpointConfig {
interval: Duration::from_millis(10),
warn_pages: 0,
..CheckpointConfig::default()
};
let store = Arc::new(FakeEventStore::default());
let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let handle = tokio::spawn(run_checkpoint_task(
pool,
cfg,
Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
shutdown_rx,
false,
));
let emitted = wait_for(Duration::from_secs(10), || {
!store.events.lock().unwrap().is_empty()
})
.await;
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("checkpoint task should exit within 1s")
.expect("checkpoint task panicked");
assert!(
emitted,
"a designated secondary lifecycle owner must append outcome events within the poll \
deadline"
);
}
#[tokio::test]
#[serial(checkpoint_skip_metrics)]
async fn checkpoint_task_emits_nothing_while_healthy() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("outcome_no_emit.db");
let pool = file_pool(&path);
let cfg = CheckpointConfig {
interval: Duration::from_millis(10),
warn_pages: u64::MAX,
..CheckpointConfig::default()
};
let store = Arc::new(FakeEventStore::default());
let store_dyn: Arc<dyn khive_storage::EventStore> = store.clone();
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let handle = tokio::spawn(run_checkpoint_task(
pool,
cfg,
Some(CheckpointLifecycleOwner::new(store_dyn, "local")),
shutdown_rx,
true,
));
tokio::time::sleep(Duration::from_millis(60)).await;
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("checkpoint task should exit within 1s")
.expect("checkpoint task panicked");
assert!(
store.events.lock().unwrap().is_empty(),
"a config that never crosses warn_pages must never append a lifecycle event"
);
}
#[tokio::test]
#[serial(checkpoint_skip_metrics)]
async fn checkpoint_task_with_no_event_store_does_not_panic() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("outcome_none_store.db");
let pool = file_pool(&path);
let cfg = CheckpointConfig {
interval: Duration::from_millis(10),
warn_pages: 0,
..CheckpointConfig::default()
};
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
tokio::time::sleep(Duration::from_millis(40)).await;
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("checkpoint task should exit within 1s")
.expect("checkpoint task panicked");
}
#[tokio::test]
#[serial(tx_registry, checkpoint_skip_metrics)]
async fn checkpoint_task_sweeps_stale_registry_entry_while_wal_is_healthy() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("tx_age_sweep_task_healthy_wal.db");
let pool = file_pool(&path);
let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let subscriber = CaptureSubscriber {
events: std::sync::Arc::clone(&buffer),
};
let _tracing_guard = tracing::subscriber::set_default(subscriber);
let _tx_handle = khive_storage::tx_registry::register(Some(
"checkpoint_task_healthy_wal_sweep_test".to_string(),
));
let cfg = CheckpointConfig {
interval: Duration::from_millis(10),
warn_pages: u64::MAX,
high_water_pages: u64::MAX,
truncate_high_water_pages: u64::MAX,
tx_warn_secs: Duration::from_millis(1),
tx_max_age_secs: Duration::from_millis(1),
..CheckpointConfig::default()
};
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
let swept = wait_for(Duration::from_secs(10), || {
buffer.lock().unwrap().iter().any(|e| {
e.tx_label.as_deref() == Some("checkpoint_task_healthy_wal_sweep_test")
&& e.message
.as_deref()
.is_some_and(|m| m.contains("stale-op cap"))
})
})
.await;
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("checkpoint task should exit within 1s")
.expect("checkpoint task panicked");
drop(_tx_handle);
let events = buffer.lock().unwrap();
assert!(
swept,
"expected the spawned task to sweep and escalate the stale registry entry \
to Stale on its own within the poll deadline, got: {events:?}"
);
}
#[tokio::test]
#[serial(tx_registry, checkpoint_skip_metrics)]
async fn checkpoint_task_emits_no_age_alert_for_an_empty_registry() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("tx_age_sweep_task_empty_registry.db");
let pool = file_pool(&path);
let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let subscriber = CaptureSubscriber {
events: std::sync::Arc::clone(&buffer),
};
let _tracing_guard = tracing::subscriber::set_default(subscriber);
let cfg = CheckpointConfig {
interval: Duration::from_millis(10),
tx_warn_secs: Duration::from_millis(1),
tx_max_age_secs: Duration::from_millis(1),
..CheckpointConfig::default()
};
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let handle = tokio::spawn(run_checkpoint_task(pool, cfg, None, shutdown_rx, true));
tokio::time::sleep(Duration::from_millis(40)).await;
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("checkpoint task should exit within 1s")
.expect("checkpoint task panicked");
let events = buffer.lock().unwrap();
assert!(
events.iter().all(|e| e
.message
.as_deref()
.is_none_or(|m| !m.contains("ADR-091 Plank 1"))),
"an empty registry must never produce a Plank 1 age emission, got: {events:?}"
);
}
#[tokio::test]
#[serial(tx_registry, checkpoint_skip_metrics)]
async fn checkpoint_task_sweeps_stale_entry_even_when_dedicated_connection_is_unavailable_every_tick(
) {
reset_checkpoint_metrics_for_tests();
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("tx_age_sweep_task_conn_unavailable.db");
{
let seed_pool = file_pool(&path);
let writer = seed_pool.try_writer().unwrap();
writer
.conn()
.execute_batch("CREATE TABLE IF NOT EXISTS t (x INTEGER);")
.unwrap();
}
#[cfg(unix)]
{
khive_storage::test_support::freeze_snapshot_sidecars(&path);
}
let pool = Arc::new(
ConnectionPool::new(PoolConfig {
path: Some(path.clone()),
read_only: true,
..PoolConfig::for_test()
})
.expect("read-only pool open"),
);
assert!(
pool.open_standalone_writer().is_err(),
"test precondition: a read-only pool must never be able to open a dedicated \
checkpoint connection"
);
let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let subscriber = CaptureSubscriber {
events: std::sync::Arc::clone(&buffer),
};
let _tracing_guard = tracing::subscriber::set_default(subscriber);
let _tx_handle = khive_storage::tx_registry::register(Some(
"checkpoint_task_conn_unavailable_sweep_test".to_string(),
));
let cfg = CheckpointConfig {
interval: Duration::from_millis(10),
tx_warn_secs: Duration::from_millis(1),
tx_max_age_secs: Duration::from_millis(1),
..CheckpointConfig::default()
};
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let handle = tokio::spawn(run_checkpoint_task(
Arc::clone(&pool),
cfg,
None,
shutdown_rx,
true,
));
let deadline = tokio::time::Instant::now() + Duration::from_secs(10);
while checkpoint_skipped_ticks() == 0 {
assert!(
tokio::time::Instant::now() < deadline,
"test setup must actually drive at least one Skipped tick for this \
regression to mean anything (none within 10s)"
);
tokio::time::sleep(Duration::from_millis(10)).await;
}
loop {
let events = buffer.lock().unwrap().clone();
if events.iter().any(|e| {
e.tx_label.as_deref() == Some("checkpoint_task_conn_unavailable_sweep_test")
&& e.message
.as_deref()
.is_some_and(|m| m.contains("stale-op cap"))
}) {
break;
}
assert!(
tokio::time::Instant::now() < deadline,
"expected the age sweep to fire even though every tick's dedicated \
connection was unavailable within 10s, got: {events:?}"
);
tokio::time::sleep(Duration::from_millis(10)).await;
}
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("checkpoint task should exit within 1s")
.expect("checkpoint task panicked");
drop(_tx_handle);
}
#[tokio::test]
async fn session_sweep_task_exits_on_shutdown_signal() {
let cfg = SessionSweepConfig {
interval: Duration::from_millis(10),
..SessionSweepConfig::default()
};
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let handle = tokio::spawn(run_session_sweep_task(Vec::new(), cfg, shutdown_rx));
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("session sweep task should exit within 1s")
.expect("session sweep task panicked");
}
async fn wait_for(deadline: Duration, mut cond: impl FnMut() -> bool) -> bool {
let start = std::time::Instant::now();
while start.elapsed() < deadline {
if cond() {
return true;
}
tokio::time::sleep(Duration::from_millis(5)).await;
}
cond()
}
#[tokio::test]
#[serial(khive_walpin_sidecar_env)]
async fn walpin_observe_drops_beacon_when_heartbeat_write_fails() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("observe_gate.db");
let sidecar_dir = crate::walpin::sidecar_dir_for(&db_path);
let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
let mut state = WalpinSidecarState::new(
Some(db_path.as_path()),
true,
"session",
Duration::from_millis(500),
)
.expect("sidecar enabled for a file-backed path");
let pid = std::process::id();
state.register_beacon().await;
let beacon_path = sidecar_dir.join(format!("{pid}.beacon"));
let before = std::fs::metadata(&beacon_path)
.expect("register_beacon must create the beacon file")
.modified()
.unwrap();
let obstruction = sidecar_dir.join(format!(".{pid}.json.tmp"));
std::fs::create_dir(&obstruction).unwrap();
tokio::time::sleep(Duration::from_millis(20)).await;
let over_threshold = Some(khive_storage::tx_registry::OldestSpan {
id: khive_storage::tx_registry::TxId(1),
age: Duration::from_secs(60),
label: None,
origin: khive_storage::tx_registry::TxOrigin::Unscoped,
});
state
.observe(over_threshold.clone(), Duration::from_secs(30))
.await;
assert!(
!sidecar_dir.join(format!("{pid}.json")).exists(),
"heartbeat write must have failed"
);
assert!(
!beacon_path.exists(),
"a failed heartbeat write must remove the beacon — a still-fresh \
beacon with no heartbeat would classify registered-silent \
(before-mtime {before:?})"
);
std::fs::remove_dir(&obstruction).unwrap();
state.observe(over_threshold, Duration::from_secs(30)).await;
assert!(
sidecar_dir.join(format!("{pid}.json")).exists(),
"heartbeat must land once the write path recovers"
);
assert!(
beacon_path.exists(),
"beacon must re-register on the first healthy tick after removal"
);
}
#[tokio::test]
#[serial(khive_walpin_sidecar_env)]
async fn walpin_observe_touches_mtime_without_rewriting_body_when_content_unchanged() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("observe_touch.db");
let sidecar_dir = crate::walpin::sidecar_dir_for(&db_path);
let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
let mut state = WalpinSidecarState::new(
Some(db_path.as_path()),
true,
"session",
Duration::from_millis(500),
)
.expect("sidecar enabled for a file-backed path");
let pid = std::process::id();
let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
let span = khive_storage::tx_registry::OldestSpan {
id: khive_storage::tx_registry::TxId(1),
age: Duration::from_secs(60),
label: None,
origin: khive_storage::tx_registry::TxOrigin::Unscoped,
};
state
.observe(Some(span.clone()), Duration::from_secs(30))
.await;
let body_after_create = std::fs::read(&heartbeat_path).expect("heartbeat written");
let backdated = std::time::SystemTime::now() - Duration::from_secs(120);
std::fs::OpenOptions::new()
.write(true)
.open(&heartbeat_path)
.unwrap()
.set_modified(backdated)
.unwrap();
state.observe(Some(span), Duration::from_secs(30)).await;
let body_after_second_observe =
std::fs::read(&heartbeat_path).expect("heartbeat still present");
assert_eq!(
body_after_create, body_after_second_observe,
"unchanged oldest-span identity/label/attribution/cadence must touch mtime, \
not rewrite the body"
);
let mtime_after = std::fs::metadata(&heartbeat_path)
.unwrap()
.modified()
.unwrap();
assert!(
mtime_after > backdated,
"the touch must advance mtime past the backdated value"
);
}
#[tokio::test]
#[serial(khive_walpin_sidecar_env)]
async fn walpin_observe_recreates_heartbeat_after_it_is_deleted_while_span_still_live() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("observe_recreate.db");
let sidecar_dir = crate::walpin::sidecar_dir_for(&db_path);
let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
let mut state = WalpinSidecarState::new(
Some(db_path.as_path()),
true,
"session",
Duration::from_millis(500),
)
.expect("sidecar enabled for a file-backed path");
let pid = std::process::id();
let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
let span = khive_storage::tx_registry::OldestSpan {
id: khive_storage::tx_registry::TxId(1),
age: Duration::from_secs(60),
label: None,
origin: khive_storage::tx_registry::TxOrigin::Unscoped,
};
state
.observe(Some(span.clone()), Duration::from_secs(30))
.await;
assert!(heartbeat_path.exists(), "heartbeat written on first tick");
std::fs::remove_file(&heartbeat_path).unwrap();
assert!(!heartbeat_path.exists());
state.observe(Some(span), Duration::from_secs(30)).await;
assert!(
heartbeat_path.exists(),
"a touch failure against a deleted heartbeat must recreate it via a full write"
);
let recreated: crate::walpin::WalpinHeartbeat =
serde_json::from_slice(&std::fs::read(&heartbeat_path).unwrap()).unwrap();
assert_eq!(recreated.pid, pid);
assert_eq!(recreated.oldest_tx_age_secs, 60.0);
}
#[tokio::test]
#[serial(tx_registry, khive_walpin_sidecar_env)]
async fn session_sweep_task_writes_and_clears_walpin_heartbeat() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("session_sweep.db");
let pool = file_pool(&db_path);
let sidecar_dir =
crate::walpin::sidecar_dir_for(pool.canonical_path().expect("file-backed pool"));
let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
let cfg = SessionSweepConfig {
interval: Duration::from_millis(10),
tx_warn_secs: Duration::from_millis(20),
tx_max_age_secs: Duration::from_millis(500),
};
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let handle = tokio::spawn(run_session_sweep_task(
vec![SweepBackend {
pool: Arc::clone(&pool),
is_main: true,
}],
cfg,
shutdown_rx,
));
let pid = std::process::id();
let beacon = crate::walpin::beacon_path(&sidecar_dir, pid);
let beacon_registered = wait_for(Duration::from_secs(2), || beacon.exists()).await;
assert!(
beacon_registered,
"a quiet process must still register its one-time beacon"
);
assert!(
!sidecar_dir.join(format!("{pid}.json")).exists(),
"a quiet process must not write a walpin heartbeat"
);
let tx_handle =
khive_storage::tx_registry::register(Some("session_sweep_walpin_test".to_string()));
let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
assert!(
wait_for(Duration::from_secs(2), || heartbeat_path.exists()).await,
"expected a walpin heartbeat once the span crossed tx_warn_secs"
);
let body = std::fs::read_to_string(&heartbeat_path).unwrap();
let hb: crate::walpin::WalpinHeartbeat = serde_json::from_str(&body).unwrap();
assert_eq!(hb.pid, pid);
assert_eq!(hb.process_role, "session");
assert_eq!(
hb.oldest_tx_label.as_deref(),
Some("session_sweep_walpin_test")
);
assert_eq!(
hb.attribution_basis.as_deref(),
Some("fallback"),
"an Unscoped span observed only through the main view's fallback \
must carry attribution_basis=\"fallback\", never \"origin\""
);
drop(tx_handle);
assert!(
wait_for(Duration::from_secs(2), || !heartbeat_path.exists()).await,
"heartbeat must be removed once the stale span clears"
);
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("session sweep task should exit within 1s")
.expect("session sweep task panicked");
}
#[tokio::test]
#[serial(tx_registry, khive_walpin_sidecar_env)]
async fn session_sweep_fan_out_scopes_secondary_span_to_secondary_sidecar_only() {
let main_dir = tempfile::tempdir().unwrap();
let secondary_dir = tempfile::tempdir().unwrap();
let main_pool = file_pool(&main_dir.path().join("main.db"));
let secondary_pool = file_pool(&secondary_dir.path().join("secondary.db"));
let main_sidecar =
crate::walpin::sidecar_dir_for(main_pool.canonical_path().expect("file-backed"));
let secondary_sidecar =
crate::walpin::sidecar_dir_for(secondary_pool.canonical_path().expect("file-backed"));
let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
let cfg = SessionSweepConfig {
interval: Duration::from_millis(10),
tx_warn_secs: Duration::from_millis(20),
tx_max_age_secs: Duration::from_millis(500),
};
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let handle = tokio::spawn(run_session_sweep_task(
vec![
SweepBackend {
pool: Arc::clone(&main_pool),
is_main: true,
},
SweepBackend {
pool: Arc::clone(&secondary_pool),
is_main: false,
},
],
cfg,
shutdown_rx,
));
let pid = std::process::id();
let secondary_heartbeat = secondary_sidecar.join(format!("{pid}.json"));
let main_heartbeat = main_sidecar.join(format!("{pid}.json"));
let tx_handle = khive_storage::tx_registry::register_scoped(
Some("graph_traverse_read".to_string()),
secondary_pool.origin(),
);
assert!(
wait_for(Duration::from_secs(2), || secondary_heartbeat.exists()).await,
"expected a walpin heartbeat in the secondary backend's own sidecar"
);
assert!(
!main_heartbeat.exists(),
"a span scoped to the secondary backend's origin must never produce \
a heartbeat in the main backend's sidecar"
);
let body = std::fs::read_to_string(&secondary_heartbeat).unwrap();
let hb: crate::walpin::WalpinHeartbeat = serde_json::from_str(&body).unwrap();
assert_eq!(hb.oldest_tx_label.as_deref(), Some("graph_traverse_read"));
assert_eq!(
hb.attribution_basis.as_deref(),
Some("origin"),
"a Secondary-view winner is always Database-origin-backed — never fallback"
);
drop(tx_handle);
assert!(
wait_for(Duration::from_secs(2), || !secondary_heartbeat.exists()).await,
"secondary heartbeat must be removed once its span clears"
);
assert!(
!main_heartbeat.exists(),
"the main sidecar must have stayed untouched for the whole tick sequence"
);
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("session sweep task should exit within 1s")
.expect("session sweep task panicked");
}
#[tokio::test]
#[serial(tx_registry, checkpoint_skip_metrics, khive_walpin_sidecar_env)]
async fn checkpoint_task_ignores_span_registered_against_other_backend_origin_and_unscoped() {
let dir_a = tempfile::tempdir().unwrap();
let dir_b = tempfile::tempdir().unwrap();
let pool_a = file_pool(&dir_a.path().join("backend_a.db"));
let pool_b = file_pool(&dir_b.path().join("backend_b.db"));
let sidecar_a =
crate::walpin::sidecar_dir_for(pool_a.canonical_path().expect("file-backed"));
let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let subscriber = CaptureSubscriber {
events: std::sync::Arc::clone(&buffer),
};
let _tracing_guard = tracing::subscriber::set_default(subscriber);
let _b_origin_handle = khive_storage::tx_registry::register_scoped(
Some("b_origin_span_ignored_by_a".to_string()),
pool_b.origin(),
);
let _unscoped_handle = khive_storage::tx_registry::register(Some(
"unscoped_span_ignored_by_secondary".to_string(),
));
let cfg = CheckpointConfig {
interval: Duration::from_millis(10),
tx_warn_secs: Duration::from_millis(1),
tx_max_age_secs: Duration::from_millis(1),
..CheckpointConfig::default()
};
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let handle = tokio::spawn(run_checkpoint_task(
pool_a,
cfg,
None,
shutdown_rx,
false, ));
tokio::time::sleep(Duration::from_millis(60)).await;
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("checkpoint task should exit within 1s")
.expect("checkpoint task panicked");
let events = buffer.lock().unwrap();
assert!(
events.iter().all(|e| {
e.tx_label.as_deref() != Some("b_origin_span_ignored_by_a")
&& e.tx_label.as_deref() != Some("unscoped_span_ignored_by_secondary")
}),
"backend A's Secondary filter must never emit an age alert naming a span \
registered against a different backend's origin or an Unscoped span, got: \
{events:?}"
);
assert!(
!sidecar_a
.join(format!("{}.json", std::process::id()))
.exists(),
"backend A's own sidecar must never gain a heartbeat from a span it does not own"
);
}
#[tokio::test]
#[serial(tx_registry, checkpoint_skip_metrics, khive_walpin_sidecar_env)]
async fn checkpoint_task_detects_and_enumerates_secondary_backend_stall() {
let dir = tempfile::tempdir().unwrap();
let pool = file_pool(&dir.path().join("secondary_stall.db"));
let sidecar_dir =
crate::walpin::sidecar_dir_for(pool.canonical_path().expect("file-backed"));
let _env_guard = crate::walpin::EnvVarGuard::capture("KHIVE_WALPIN_SIDECAR");
std::env::set_var("KHIVE_WALPIN_SIDECAR", "1");
let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let subscriber = CaptureSubscriber {
events: std::sync::Arc::clone(&buffer),
};
let _tracing_guard = tracing::subscriber::set_default(subscriber);
let tx_handle = khive_storage::tx_registry::register_scoped(
Some("secondary_stall_test".to_string()),
pool.origin(),
);
let cfg = CheckpointConfig {
interval: Duration::from_millis(10),
tx_warn_secs: Duration::from_millis(5),
tx_max_age_secs: Duration::from_millis(500),
..CheckpointConfig::default()
};
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(());
let pid = std::process::id();
let heartbeat_path = sidecar_dir.join(format!("{pid}.json"));
let handle = tokio::spawn(run_checkpoint_task(
pool,
cfg,
None,
shutdown_rx,
false, ));
assert!(
wait_for(Duration::from_secs(2), || heartbeat_path.exists()).await,
"expected a walpin heartbeat once the secondary backend's own span crossed \
tx_warn_secs"
);
let body = std::fs::read_to_string(&heartbeat_path).unwrap();
let hb: crate::walpin::WalpinHeartbeat = serde_json::from_str(&body).unwrap();
assert_eq!(hb.oldest_tx_label.as_deref(), Some("secondary_stall_test"));
assert_eq!(
hb.attribution_basis.as_deref(),
Some("origin"),
"a Secondary-view winner is always Database-origin-backed — never fallback"
);
assert!(
hb.oldest_tx_age_secs > 0.0,
"the heartbeat must reflect a nonzero stale age for the secondary backend's own \
span, got {hb:?}"
);
shutdown_tx.send(()).expect("send shutdown signal");
tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("checkpoint task should exit within 1s")
.expect("checkpoint task panicked");
drop(tx_handle);
let events = buffer.lock().unwrap();
assert!(
events.iter().any(|e| {
e.tx_label.as_deref() == Some("secondary_stall_test")
&& e.message
.as_deref()
.is_some_and(|m| m.contains("ADR-091 Plank 1"))
}),
"expected the secondary backend's own checkpoint task to emit a Plank 1 age alert \
for its own stalled span, got: {events:?}"
);
}
#[test]
fn wal_pin_depth_arithmetic_against_real_connection() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("pin_depth.db");
let pool = file_pool(&path);
let writer = pool.try_writer().expect("acquire writer");
let conn = writer.conn();
conn.execute_batch("CREATE TABLE t (v INTEGER)").unwrap();
conn.execute_batch("INSERT INTO t (v) VALUES (1)").unwrap();
let (log, checkpointed) =
query_wal_pin_depth(conn).expect("PRAGMA wal_checkpoint(PASSIVE) must succeed");
assert!(
log >= checkpointed,
"checkpointed frames cannot exceed log frames"
);
assert_eq!(
log - checkpointed,
0,
"an unpinned WAL must fully checkpoint under PASSIVE"
);
}
#[test]
fn wal_pin_depth_arithmetic_on_in_memory_pool_errors_cleanly() {
let cfg = PoolConfig {
path: None,
..PoolConfig::default()
};
let pool = ConnectionPool::new(cfg).expect("in-memory pool");
let writer = pool.try_writer().expect("acquire writer");
let _ = query_wal_pin_depth(writer.conn());
}
#[cfg(unix)]
#[test]
fn routine_wal_backend_key_preserves_non_utf8_path_bytes() {
use std::ffi::OsString;
use std::os::unix::ffi::OsStringExt;
let path_a = PathBuf::from(OsString::from_vec(b"/tmp/khive-wal-\x80.db".to_vec()));
let path_b = PathBuf::from(OsString::from_vec(b"/tmp/khive-wal-\x81.db".to_vec()));
assert_eq!(
path_a.display().to_string(),
path_b.display().to_string(),
"fixture must reproduce the lossy display-label collision"
);
assert_ne!(
checkpoint_db_key_from_path(Some(&path_a)),
checkpoint_db_key_from_path(Some(&path_b)),
"backend keys must retain the canonical path's exact OS bytes"
);
}
#[test]
#[serial(checkpoint_skip_metrics)]
fn routine_checkpoint_records_one_pass_logical_and_physical_wal_sample() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("routine_wal_sample.db");
let pool = file_pool(&path);
{
let writer = pool.try_writer().expect("writer");
writer
.conn()
.execute_batch(
"PRAGMA wal_autocheckpoint=0; \
CREATE TABLE t (id INTEGER PRIMARY KEY, payload TEXT); \
INSERT INTO t VALUES (0, 'seed');",
)
.unwrap();
}
let reader = rusqlite::Connection::open_with_flags(
&path,
rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY,
)
.unwrap();
reader.execute_batch("BEGIN").unwrap();
let _: i64 = reader
.query_row("SELECT COUNT(*) FROM t", [], |row| row.get(0))
.unwrap();
{
let writer = pool.try_writer().expect("writer");
writer.conn().execute_batch("BEGIN IMMEDIATE").unwrap();
for id in 1..=256_i64 {
writer
.conn()
.execute("INSERT INTO t VALUES (?1, printf('%.*c', 2048, 'x'))", [id])
.unwrap();
}
writer.conn().execute_batch("COMMIT").unwrap();
}
let checkpoint_conn = pool.open_standalone_writer().unwrap();
let pragma_calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let pragma_calls_from_hook = Arc::clone(&pragma_calls);
checkpoint_conn
.authorizer(Some(move |context: rusqlite::hooks::AuthContext<'_>| {
if matches!(
context.action,
AuthAction::Pragma { pragma_name, .. }
if pragma_name.eq_ignore_ascii_case("wal_checkpoint")
) {
pragma_calls_from_hook.fetch_add(1, Ordering::SeqCst);
}
Authorization::Allow
}))
.unwrap();
checkpoint_once(
&pool,
&checkpoint_conn,
&CheckpointConfig::default(),
&mut TruncateState::default(),
)
.unwrap();
checkpoint_conn
.authorizer(None::<fn(rusqlite::hooks::AuthContext<'_>) -> Authorization>)
.unwrap();
assert_eq!(
pragma_calls.load(Ordering::SeqCst),
1,
"one routine tick must issue exactly one PASSIVE checkpoint"
);
let pinned = routine_wal_observation(&pool).expect("routine sample");
assert_eq!(
pinned.busy, 0,
"a pinned reader is not checkpoint-lock contention"
);
let first_timing = checkpoint_timing(&pool);
assert_eq!(first_timing.ticks, 1);
assert_eq!(
first_timing.busy_ticks, 0,
"pending frames must not count as busy"
);
assert!(pinned.log_frames > 0, "the test must create WAL frames");
assert!(
pinned.pending_frames > 0,
"the old reader must leave a logical backlog: {pinned:?}"
);
assert_eq!(
pinned.pending_frames,
pinned.log_frames.saturating_sub(pinned.checkpointed_frames)
);
assert!(
pinned.physical_wal_bytes.is_some_and(|bytes| bytes > 0),
"the physical sidecar high-water must be reported separately: {pinned:?}"
);
reader.execute_batch("COMMIT").unwrap();
checkpoint_once(
&pool,
&checkpoint_conn,
&CheckpointConfig::default(),
&mut TruncateState::default(),
)
.unwrap();
let drained = routine_wal_observation(&pool).expect("drained routine sample");
let drained_timing = checkpoint_timing(&pool);
assert_eq!(drained_timing.ticks, first_timing.ticks + 1);
assert!(drained_timing.elapsed_us_sum >= first_timing.elapsed_us_sum);
assert!(drained_timing.elapsed_us_max >= first_timing.elapsed_us_max);
assert_eq!(drained_timing.busy_ticks, 0);
assert_eq!(drained.pending_frames, 0, "unpinned PASSIVE must drain");
assert!(
drained.physical_wal_bytes.is_some_and(|bytes| bytes > 0),
"PASSIVE may reuse rather than shrink the physical WAL; the two gauges must remain \
independently visible: {drained:?}"
);
}
#[test]
#[serial(checkpoint_skip_metrics)]
fn routine_checkpoint_timing_records_real_call_and_debug_fields() {
let dir = tempfile::tempdir().unwrap();
let pool = file_pool(&dir.path().join("timed_tick.db"));
let conn = checkpoint_conn(&pool);
conn.execute_batch("CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);")
.unwrap();
let (entered_tx, entered_rx) = std::sync::mpsc::sync_channel(0);
let (release_tx, release_rx) = std::sync::mpsc::sync_channel(0);
let release_rx = Mutex::new(release_rx);
conn.authorizer(Some(move |context: rusqlite::hooks::AuthContext<'_>| {
if matches!(context.action, AuthAction::Pragma { pragma_name, .. }
if pragma_name.eq_ignore_ascii_case("wal_checkpoint"))
{
entered_tx.send(()).unwrap();
release_rx
.lock()
.unwrap()
.recv_timeout(Duration::from_secs(5))
.unwrap();
}
Authorization::Allow
}))
.unwrap();
let tick_pool = Arc::clone(&pool);
let tick = std::thread::spawn(move || {
let mut result = None;
let events = capture(|| {
result = Some(checkpoint_once(
&tick_pool,
&conn,
&CheckpointConfig {
truncate_high_water_pages: u64::MAX,
..CheckpointConfig::default()
},
&mut TruncateState::default(),
));
});
(result.unwrap(), events)
});
entered_rx.recv_timeout(Duration::from_secs(5)).unwrap();
let during = checkpoint_timing(&pool);
release_tx.send(()).unwrap();
let (result, events) = tick.join().unwrap();
result.unwrap();
assert_eq!(
during.ticks, 0,
"in-flight call must not publish partial counters"
);
let timing = checkpoint_timing(&pool);
assert_eq!(
timing.ticks, 1,
"one actual PASSIVE call must advance the count"
);
assert!(
timing.elapsed_us_sum > 0,
"channel-held checkpoint call must record elapsed time"
);
assert_eq!(timing.elapsed_us_max, timing.elapsed_us_sum);
assert_eq!(timing.busy_ticks, 0);
assert_eq!(timing.error_ticks, 0);
let issued: Vec<_> = events
.iter()
.filter(|event| event.message.as_deref() == Some("WAL checkpoint issued"))
.collect();
assert_eq!(issued.len(), 1);
assert_eq!(issued[0].elapsed_us, Some(timing.elapsed_us_sum));
assert_eq!(issued[0].busy, Some(0));
}
#[test]
#[serial(checkpoint_skip_metrics)]
fn routine_checkpoint_timing_counts_sqlite_busy_from_competing_checkpoint() {
struct BusyGate {
entered: std::sync::mpsc::SyncSender<()>,
release: std::sync::mpsc::Receiver<()>,
}
static BUSY_GATE: Mutex<Option<BusyGate>> = Mutex::new(None);
fn hold_checkpoint_lock(_attempt: i32) -> bool {
let gate = BUSY_GATE
.lock()
.unwrap()
.take()
.expect("armed busy handler");
gate.entered.send(()).unwrap();
gate.release.recv_timeout(Duration::from_secs(5)).unwrap();
false
}
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("busy_checkpoint.db");
let pool = file_pool(&path);
let conn = checkpoint_conn(&pool);
conn.execute_batch(
"PRAGMA wal_autocheckpoint=0; CREATE TABLE t (x INTEGER); INSERT INTO t VALUES (1);",
)
.unwrap();
let reader = rusqlite::Connection::open_with_flags(
&path,
rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY,
)
.unwrap();
reader.execute_batch("BEGIN").unwrap();
let _: i64 = reader
.query_row("SELECT COUNT(*) FROM t", [], |row| row.get(0))
.unwrap();
conn.execute_batch("INSERT INTO t VALUES (2);").unwrap();
let (entered_tx, entered_rx) = std::sync::mpsc::sync_channel(0);
let (release_tx, release_rx) = std::sync::mpsc::sync_channel(0);
*BUSY_GATE.lock().unwrap() = Some(BusyGate {
entered: entered_tx,
release: release_rx,
});
let competing = std::thread::spawn(move || {
let checkpoint = rusqlite::Connection::open(path).unwrap();
checkpoint.busy_handler(Some(hold_checkpoint_lock)).unwrap();
checkpoint.query_row("PRAGMA wal_checkpoint(FULL)", [], |row| {
row.get::<_, i64>(0)
})
});
entered_rx
.recv_timeout(Duration::from_secs(5))
.expect("FULL checkpoint holds CKPT lock while waiting on reader");
let mut result = None;
let events = capture(|| {
result = Some(checkpoint_once(
&pool,
&conn,
&CheckpointConfig {
truncate_high_water_pages: u64::MAX,
..CheckpointConfig::default()
},
&mut TruncateState::default(),
));
});
release_tx.send(()).unwrap();
let competing_busy = competing.join().unwrap().unwrap();
reader.execute_batch("COMMIT").unwrap();
result.unwrap().unwrap();
assert_eq!(competing_busy, 1);
assert_eq!(
routine_wal_observation(&pool).unwrap().busy,
1,
"fixture must reach SQLite busy"
);
let timing = checkpoint_timing(&pool);
assert_eq!(timing.ticks, 1);
assert_eq!(
timing.busy_ticks, 1,
"SQLite busy result must increment busy ticks"
);
assert_eq!(timing.error_ticks, 0);
let issued = events
.iter()
.find(|event| event.message.as_deref() == Some("WAL checkpoint issued"))
.expect("existing tick record");
assert_eq!(issued.busy, Some(1), "tick record must retain SQLite busy");
assert_eq!(issued.elapsed_us, Some(timing.elapsed_us_sum));
}
#[test]
#[serial(checkpoint_skip_metrics)]
fn routine_checkpoint_timing_counts_errors_and_excludes_post_truncate_probes() {
let dir = tempfile::tempdir().unwrap();
let pool = file_pool(&dir.path().join("failed_tick.db"));
let conn = checkpoint_conn(&pool);
conn.authorizer(Some(|context: rusqlite::hooks::AuthContext<'_>| {
if matches!(context.action, AuthAction::Pragma { pragma_name, .. }
if pragma_name.eq_ignore_ascii_case("wal_checkpoint"))
{
Authorization::Deny
} else {
Authorization::Allow
}
}))
.unwrap();
let result = checkpoint_once(
&pool,
&conn,
&CheckpointConfig::default(),
&mut TruncateState::default(),
);
conn.authorizer(None::<fn(rusqlite::hooks::AuthContext<'_>) -> Authorization>)
.unwrap();
assert!(result.is_err());
let timing = checkpoint_timing(&pool);
assert_eq!(timing.ticks, 1);
assert_eq!(timing.error_ticks, 1);
assert_eq!(timing.busy_ticks, 0);
query_wal_pages(&conn);
assert_eq!(
checkpoint_timing(&pool),
timing,
"post-TRUNCATE observation is not a routine tick"
);
}
#[test]
fn checkpoint_timing_accumulates_per_store_and_saturates() {
let dir = tempfile::tempdir().unwrap();
let a = file_pool(&dir.path().join("a.db"));
let b = file_pool(&dir.path().join("b.db"));
record_checkpoint_timing(&a, 17, Some(0));
record_checkpoint_timing(&a, 31, Some(1));
record_checkpoint_timing(&a, 7, None);
record_checkpoint_timing(&b, 3, Some(0));
assert_eq!(
checkpoint_timing(&a),
CheckpointTiming {
ticks: 3,
elapsed_us_sum: 55,
elapsed_us_max: 31,
busy_ticks: 1,
error_ticks: 1,
}
);
assert_eq!(
checkpoint_timing(&b),
CheckpointTiming {
ticks: 1,
elapsed_us_sum: 3,
elapsed_us_max: 3,
busy_ticks: 0,
error_ticks: 0,
}
);
checkpoint_timings().lock().unwrap().insert(
checkpoint_db_key(&a),
CheckpointTiming {
ticks: u64::MAX,
elapsed_us_sum: u64::MAX,
elapsed_us_max: 31,
busy_ticks: u64::MAX,
error_ticks: u64::MAX,
},
);
record_checkpoint_timing(&a, 1, Some(1));
assert_eq!(
checkpoint_timing(&a).elapsed_us_sum,
u64::MAX,
"elapsed sum must saturate on the first overflowing addition"
);
record_checkpoint_timing(&a, u64::MAX, None);
assert_eq!(
checkpoint_timing(&a),
CheckpointTiming {
ticks: u64::MAX,
elapsed_us_sum: u64::MAX,
elapsed_us_max: u64::MAX,
busy_ticks: u64::MAX,
error_ticks: u64::MAX,
}
);
}
}