use futures_util::future::{FutureExt, join_all};
use std::collections::HashMap;
use std::future::Future;
use std::panic::AssertUnwindSafe;
use std::path::{Path, PathBuf};
use std::sync::{LazyLock, Mutex};
use std::time::{Duration, Instant};
use tracing::{debug, error, info, warn};
use crate::db::failure_record::{self, FailureKind, FailureReport, RoundCounter};
use crate::db::{
CheckpointOutcome, Connection, TicketTitleFtsRuntimeRepair, checkpoint_cause, shrink_gate,
};
use crate::util::UnwrapPoison;
const WAL_CHECKPOINT_CAP_BYTES: u64 = 32 * 1024 * 1024;
const DEFAULT_CHECKPOINT_MIN_FREE_BYTES: u64 = 64 * 1024 * 1024;
const CHECKPOINT_ATTEMPTS: usize = 6;
const CHECKPOINT_RETRY_PAUSES: [Duration; CHECKPOINT_ATTEMPTS - 1] = [
Duration::from_secs(1),
Duration::from_secs(5),
Duration::from_secs(15),
Duration::from_secs(30),
Duration::from_secs(60),
];
const CHECKPOINT_FAILURE_MIN_ROUNDS: u64 = 2;
const CHECKPOINT_FAILURE_WINDOW: Duration = Duration::from_secs(120);
const BLOCK_ON_STDERR: &str = "(could not be filed — block on stderr)";
fn checkpoint_min_free_bytes() -> u64 {
std::env::var("MAHBOT_CHECKPOINT_MIN_FREE_BYTES")
.ok()
.and_then(|v| v.parse::<u64>().ok())
.unwrap_or(DEFAULT_CHECKPOINT_MIN_FREE_BYTES)
}
fn available_free_bytes(path: &Path) -> u64 {
crate::util::disk::free_and_capacity(path).map_or(0, |(free, _)| free)
}
fn truncate_allowed(root: &Path) -> bool {
let min = checkpoint_min_free_bytes();
if min == 0 {
return true;
}
let free = available_free_bytes(root);
let allowed = free >= min;
if !allowed {
warn!(
free_bytes = free,
min_free_bytes = min,
"Free disk space below TRUNCATE checkpoint threshold — running PASSIVE only",
);
}
allowed
}
#[derive(Debug, Clone, Copy)]
enum CheckpointPolicy {
Truncate,
PassiveCapped(u64),
}
impl CheckpointPolicy {
fn periodic() -> Self {
Self::PassiveCapped(WAL_CHECKPOINT_CAP_BYTES)
}
}
async fn for_each_store<F, Fut>(op: F)
where
F: Fn(&'static str, &'static crate::db::Connection) -> Fut,
Fut: Future<Output = ()>,
{
let futs: Vec<_> = crate::db::iter_checkpoint_stores()
.filter_map(|(name, conn_opt)| {
let conn = conn_opt?;
let fut = AssertUnwindSafe(op(name, conn)).catch_unwind();
Some(async move {
if let Err(payload) = fut.await {
error!(
panic = %crate::util::panic_message(&*payload),
db = name,
"Store operation panicked — isolated to this store",
);
}
})
})
.collect();
join_all(futs).await;
}
pub async fn checkpoint_all_databases() {
checkpoint_stores(CheckpointRound::Exit).await;
}
pub async fn periodic_checkpoint_and_verify() {
checkpoint_stores(CheckpointRound::Periodic).await;
}
#[derive(Clone, Copy)]
enum CheckpointRound {
Exit,
Periodic,
}
async fn checkpoint_stores(round: CheckpointRound) {
let root = crate::config::CONFIG.try_storage_root();
let truncate_gate =
crate::util::with_block_in_place(|| root.as_deref().is_none_or(truncate_allowed));
let policy = match round {
CheckpointRound::Exit => CheckpointPolicy::Truncate,
CheckpointRound::Periodic => CheckpointPolicy::periodic(),
};
for_each_store(|name, conn| {
let root = root.clone();
async move {
let truncate = match policy {
CheckpointPolicy::Truncate => truncate_gate,
CheckpointPolicy::PassiveCapped(cap) => {
root.as_deref()
.map(|r| {
let db_path = crate::db::store_db_path(r, name);
crate::db::wal_guard::wal_size(&db_path)
})
.is_some_and(|wal_size| wal_size > cap)
&& truncate_gate
}
};
match round {
CheckpointRound::Exit => {
exit_checkpoint(name, conn, truncate, root.as_deref()).await;
}
CheckpointRound::Periodic => {
periodic_checkpoint(name, conn, truncate, root.as_deref()).await;
verify_integrity(name, conn, root.as_deref()).await;
}
}
}
})
.await;
}
async fn verify_integrity(name: &'static str, conn: &Connection, root: Option<&Path>) {
match conn.quick_check_problems().await {
Ok(problems) if problems.is_empty() => {
debug!(db = %name, "Database integrity check passed");
}
Ok(problems) if crate::db::is_tolerated_integrity_report(&problems) => {
warn!(
db = %name,
problems = %problems.join("; "),
"Database integrity check reported only tolerated findings — not recorded, not fatal",
);
}
Ok(problems) => {
record_integrity_failure(
name,
conn,
&anyhow::anyhow!("{}", problems.join("; ")),
root,
);
}
Err(e) => record_integrity_failure(name, conn, &e, root),
}
}
enum Attempted {
Ran(anyhow::Result<CheckpointOutcome>),
ShrinkRefused,
}
async fn guarded_checkpoint(
attempt: impl Future<Output = anyhow::Result<CheckpointOutcome>>,
) -> anyhow::Result<CheckpointOutcome> {
match AssertUnwindSafe(attempt).catch_unwind().await {
Ok(result) => result,
Err(panic) => Err(anyhow::anyhow!(
"checkpoint attempt panicked: {}",
crate::util::panic_message(&*panic)
)),
}
}
async fn checkpoint_attempt(
name: &'static str,
conn: &Connection,
truncate: bool,
root: Option<&Path>,
) -> Attempted {
if truncate && !shrink_gate::shrink_allowed(conn, name, root).await {
return Attempted::ShrinkRefused;
}
Attempted::Ran(guarded_checkpoint(conn.checkpoint_mode(truncate)).await)
}
async fn exit_checkpoint(
name: &'static str,
conn: &Connection,
truncate: bool,
root: Option<&Path>,
) {
match checkpoint_attempt(name, conn, truncate, root).await {
Attempted::ShrinkRefused => {}
Attempted::Ran(Ok(o)) if o.is_complete() => {
debug!(
db = %name,
log = o.log_frames,
checkpointed = o.checkpointed_frames,
"Database WAL checkpointed",
);
}
Attempted::Ran(Ok(o)) => {
warn!(
db = %name,
busy = o.busy,
log = o.log_frames,
checkpointed = o.checkpointed_frames,
"Checkpoint busy or partial — WAL frames left uncheckpointed",
);
}
Attempted::Ran(Err(e)) => {
warn!(error = %e, db = %name, "Failed to checkpoint database WAL");
record_store_failure(
FailureKind::ExitCheckpointFailure,
"exit checkpoint failure",
name,
conn,
&e,
root,
);
}
}
}
async fn periodic_checkpoint(
name: &'static str,
conn: &Connection,
truncate: bool,
root: Option<&Path>,
) {
periodic_checkpoint_inner(
name,
conn,
root,
CHECKPOINT_RETRY_PAUSES,
CHECKPOINT_FAILURE_WINDOW,
|| checkpoint_attempt(name, conn, truncate, root),
)
.await;
}
#[derive(Clone)]
struct FailureWindow {
rounds: u64,
attempts: u64,
since: Instant,
first_failure: String,
}
impl FailureWindow {
fn extend(&mut self, attempts: u64) {
self.rounds += 1;
self.attempts += attempts;
}
fn elapsed(&self) -> Duration {
self.since.elapsed()
}
fn exhausted(&self, span: Duration) -> bool {
self.rounds >= CHECKPOINT_FAILURE_MIN_ROUNDS && self.elapsed() >= span
}
fn summary(&self) -> String {
format!(
"failure window: {} attempts (blocked and partial ones included) over {} consecutive \
failing rounds, {:.1}s since the first failure",
self.attempts,
self.rounds,
self.elapsed().as_secs_f64(),
)
}
}
static FAILURE_WINDOWS: LazyLock<Mutex<HashMap<PathBuf, FailureWindow>>> =
LazyLock::new(|| Mutex::new(HashMap::new()));
fn failure_window_key(conn: &Connection) -> PathBuf {
conn.db_path().to_path_buf()
}
fn extend_failure_window(
key: PathBuf,
attempts: u64,
since: Instant,
first_failure: &str,
span: Duration,
) -> (FailureWindow, bool) {
let mut windows = FAILURE_WINDOWS.lock().unwrap_poison();
let mut window = windows.remove(&key).unwrap_or_else(|| FailureWindow {
rounds: 0,
attempts: 0,
since,
first_failure: first_failure.to_string(),
});
window.extend(attempts);
let exhausted = window.exhausted(span);
if !exhausted {
windows.insert(key, window.clone());
}
(window, exhausted)
}
fn close_failure_window(name: &'static str, key: &Path, reason: &'static str) {
let Some(window) = FAILURE_WINDOWS.lock().unwrap_poison().remove(key) else {
return;
};
warn!(
db = %name,
window_attempts = window.attempts,
window_rounds = window.rounds,
elapsed_ms = window.elapsed().as_millis(),
reason,
"Checkpoint failure window closed — continuing",
);
}
fn cure_failure_window(name: &'static str, key: &Path, round_failures: usize) {
let Some(window) = FAILURE_WINDOWS.lock().unwrap_poison().remove(key) else {
if round_failures > 0 {
debug!(db = %name, round_failures, "Checkpoint completed on a retry");
}
return;
};
debug!(
db = %name,
round_failures,
window_attempts = window.attempts,
window_rounds = window.rounds,
elapsed_ms = window.elapsed().as_millis(),
"Checkpoint completed — failure window closed",
);
}
#[expect(clippy::too_many_lines)] async fn periodic_checkpoint_inner<'a, Fut>(
name: &'static str,
conn: &'a Connection,
root: Option<&Path>,
retry_pauses: [Duration; CHECKPOINT_ATTEMPTS - 1],
span: Duration,
mut attempt: impl FnMut() -> Fut,
) where
Fut: Future<Output = Attempted> + Send + 'a,
{
let key = failure_window_key(conn);
let mut failures: Vec<(usize, anyhow::Error)> = Vec::new();
let mut repair: Option<TicketTitleFtsRuntimeRepair> = None;
let mut attempts_made = 0u64;
let mut first_failure_at: Option<Instant> = None;
for attempt_no in 1..=CHECKPOINT_ATTEMPTS {
match attempt().await {
Attempted::ShrinkRefused => {
if failures.is_empty() {
close_failure_window(
name,
&key,
"round ended on a refused shrink with no failure",
);
return;
}
break;
}
Attempted::Ran(Ok(o)) if o.is_complete() => {
debug!(
db = %name,
log = o.log_frames,
checkpointed = o.checkpointed_frames,
"Database WAL checkpointed",
);
cure_failure_window(name, &key, failures.len());
return;
}
Attempted::Ran(Ok(o)) => {
info!(
db = %name,
log = o.log_frames,
checkpointed = o.checkpointed_frames,
"Checkpoint folded the journal partially",
);
if failures.is_empty() {
close_failure_window(name, &key, "round folded partially with no failure");
return;
}
}
Attempted::Ran(Err(e)) if checkpoint_cause::is_blocked_checkpoint(&e) => {
info!(attempt = attempt_no, error = %e, db = %name, "Checkpoint blocked — retrying");
}
Attempted::Ran(Err(e)) => {
debug!(
attempt = attempt_no,
error = %e,
db = %name,
"Failed to checkpoint database WAL"
);
first_failure_at.get_or_insert_with(Instant::now);
failures.push((attempt_no, e));
}
}
attempts_made += 1;
if attempt_no == 1 {
repair = Some(crate::db::repair_ticket_title_fts_runtime(conn).await);
}
if let Some(pause) = retry_pauses.get(attempt_no - 1).copied()
&& !crate::shutdown::sleep_or_shutdown_or_drain(pause).await
{
let cause = failures.first().map(|(_, e)| format!("{e:#}"));
warn!(
db = %name,
round_attempts = attempts_made,
error = cause.as_deref(),
"Checkpoint round cut short by the daemon's own shutdown — failure window not decided",
);
return;
}
}
if failures.is_empty() {
info!(db = %name, "Checkpoint round ended without a completion — continuing");
close_failure_window(name, &key, "round ended without a failure or a completion");
return;
}
let (_, first) = failures
.first()
.expect("the window is only extended after a failed attempt");
let first_failure = format!("{first:#}");
let (window, exhausted) = extend_failure_window(
key,
attempts_made,
first_failure_at.expect("a failed attempt set the round's first failure instant"),
&first_failure,
span,
);
if !exhausted {
warn!(
db = %name,
round_attempts = attempts_made,
window_attempts = window.attempts,
window_rounds = window.rounds,
elapsed_ms = window.elapsed().as_millis(),
error = %first_failure,
"Genuine checkpoint failure with no completion — failure window not exhausted, continuing",
);
return;
}
let report = build_failure_report(name, &failures, repair.as_ref(), &window, conn).await;
let pointer = failure_record::recorded_pointer(
"checkpoint failure",
failure_record::record(root, &report.render()),
);
error!(
db = %name,
record = pointer.as_deref().unwrap_or(BLOCK_ON_STDERR),
window_attempts = window.attempts,
window_rounds = window.rounds,
elapsed_ms = window.elapsed().as_millis(),
"Genuine checkpoint failure with no completion across the failure window — recorded, \
initiating graceful shutdown",
);
crate::shutdown::drain_begin();
}
fn store_failure_report(
kind: FailureKind,
name: &'static str,
e: &anyhow::Error,
db_path: &Path,
) -> FailureReport {
let report = FailureReport::new(kind)
.store(name)
.reason(format!("{e:#}"))
.environment(crate::db::is_actionable_signal(e));
with_store_file_state(report, db_path)
}
fn with_store_file_state(report: FailureReport, db_path: &Path) -> FailureReport {
let wal_bytes = crate::db::wal_guard::stat_size(&crate::db::wal_path(db_path)).ok();
report
.db_path(db_path.to_path_buf())
.extra(failure_record::artifact_state_line(wal_bytes))
}
fn record_store_failure(
kind: FailureKind,
what: &str,
name: &'static str,
conn: &Connection,
e: &anyhow::Error,
root: Option<&Path>,
) {
let report = store_failure_report(kind, name, e, conn.db_path());
failure_record::record_and_point(root, what, &report);
}
static INTEGRITY_FAILURE_ROUNDS: RoundCounter = RoundCounter::new();
fn record_integrity_failure(
name: &'static str,
conn: &Connection,
e: &anyhow::Error,
root: Option<&Path>,
) {
let further_rounds = INTEGRITY_FAILURE_ROUNDS.prior_rounds(name);
if further_rounds > 0 {
warn!(
error = %e,
db = %name,
further_rounds,
"Database integrity check failed again — already recorded on its first failure",
);
return;
}
error!(error = %e, db = %name, "Database integrity check failed");
record_store_failure(
FailureKind::RuntimeIntegrityFailure,
"runtime integrity failure",
name,
conn,
e,
root,
);
}
async fn build_failure_report(
name: &'static str,
failures: &[(usize, anyhow::Error)],
repair: Option<&TicketTitleFtsRuntimeRepair>,
window: &FailureWindow,
conn: &Connection,
) -> FailureReport {
let mut report = FailureReport::new(FailureKind::CheckpointFailure)
.store(name)
.environment(
failures
.iter()
.any(|(_, e)| crate::db::is_actionable_signal(e)),
);
let (_, first) = failures
.first()
.expect("a failure report is built only after a failed attempt");
let cause = format!("{first:#}");
report = report.extra(format!("checkpoint error: {cause}"));
for (attempt_no, e) in failures.iter().skip(1) {
report = report.extra(format!("attempt {attempt_no} error: {e:#}"));
}
if window.first_failure != cause {
report = report.extra(format!(
"window first failure error: {}",
window.first_failure
));
}
report = report.extra(window.summary());
report = report.extra(
"retry policy: a transient engine refusal cannot be told from a persistent one except by \
whether a later attempt succeeds — a persistent refusal exhausts this window and the \
service still stops, by design",
);
let repair = repair.expect("the round runs the repair before it can report");
report = report.extra(format!("repair outcome: {}", repair.summary()));
report = report.extra(
match AssertUnwindSafe(conn.quick_check_problems())
.catch_unwind()
.await
{
Ok(Ok(problems)) if problems.is_empty() => "quick_check: ok".to_string(),
Ok(Ok(problems)) => format!("quick_check problems: {}", problems.join("; ")),
Ok(Err(e)) => format!("quick_check error: {e:#}"),
Err(_) => "quick_check: probe panicked".to_string(),
},
);
with_store_file_state(report, conn.db_path())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::db::checkpoint_cause;
use crate::db::test_support::{
build_empty_main_file_store, build_zero_page_count_store, build_zero_page_count_wal_store,
fixture_rows, foreign_client_is_refused, fts_corruption_ddl, insert_fts_ticket, page_count,
};
use crate::db::wal_guard::{ShapeDefect, StoreShape};
#[tokio::test]
async fn noop_when_no_stores() {
checkpoint_all_databases().await;
periodic_checkpoint_and_verify().await;
}
async fn temp_store(name: &'static str) -> (tempfile::TempDir, Connection) {
let tmp = tempfile::TempDir::new().unwrap();
let conn = crate::db::open_with_schema(
&crate::db::store_db_path(tmp.path(), name),
"CREATE TABLE plain (id INTEGER PRIMARY KEY);",
)
.await
.unwrap();
(tmp, conn)
}
async fn a_store_with_a_corrupt_fts_index(root: &Path) -> Connection {
let conn = crate::db::open_consolidated_store(root).await.unwrap();
insert_fts_ticket(&conn, "t-1", "Important bug fix one").await;
insert_fts_ticket(&conn, "t-2", "Another relevant thing").await;
conn.execute_batch(&fts_corruption_ddl()).await.unwrap();
conn
}
fn blocked_checkpoint_error() -> anyhow::Error {
anyhow::anyhow!("Unexpected result from PRAGMA wal_checkpoint").context(format!(
"engine cause: {}",
checkpoint_cause::BLOCKED_REASON
))
}
const TEST_PAUSES: [Duration; CHECKPOINT_ATTEMPTS - 1] =
[Duration::ZERO; CHECKPOINT_ATTEMPTS - 1];
fn window_open(conn: &Connection) -> bool {
FAILURE_WINDOWS
.lock()
.unwrap_poison()
.contains_key(&failure_window_key(conn))
}
fn window_rounds() -> usize {
usize::try_from(CHECKPOINT_FAILURE_MIN_ROUNDS)
.expect("the failure window's round floor fits a usize")
}
async fn failing_rounds(
name: &'static str,
conn: &Connection,
root: &Path,
rounds: usize,
error: &str,
retry_pauses: [Duration; CHECKPOINT_ATTEMPTS - 1],
span: Duration,
) -> usize {
let error = error.to_string();
rounds_with_answers(name, conn, root, rounds, retry_pauses, span, move |_| {
let error = error.clone();
async move { Attempted::Ran(Err(anyhow::anyhow!(error))) }
})
.await
}
async fn rounds_with_answers<A, Fut>(
name: &'static str,
conn: &Connection,
root: &Path,
rounds: usize,
retry_pauses: [Duration; CHECKPOINT_ATTEMPTS - 1],
span: Duration,
mut answer: A,
) -> usize
where
A: FnMut(usize) -> Fut,
Fut: Future<Output = Attempted> + Send,
{
let mut attempts = 0usize;
for _ in 0..rounds {
let mut attempt_no = 0usize;
periodic_checkpoint_inner(name, conn, Some(root), retry_pauses, span, || {
attempt_no += 1;
attempts += 1;
answer(attempt_no)
})
.await;
}
attempts
}
#[tokio::test]
async fn exit_checkpoint_failure_is_recorded() {
let (tmp, conn) = temp_store("core").await;
record_store_failure(
FailureKind::ExitCheckpointFailure,
"exit checkpoint failure",
"core",
&conn,
&anyhow::anyhow!("no space left on device"),
Some(tmp.path()),
);
let body = std::fs::read_to_string(tmp.path().join("error.log")).unwrap();
for needle in [
"MahBot exit checkpoint failure",
"store: core",
failure_record::ENVIRONMENT_CAUSE,
"reason: no space left on device",
"artifact state: wal_size=",
] {
assert!(
body.contains(needle),
"error.log must contain {needle:?}: {body}"
);
}
assert!(
body.contains(&format!("db path: {}", conn.db_path().display())),
"the block must name the file the connection has open: {body}"
);
}
#[tokio::test]
async fn the_failure_record_carries_the_engines_own_reason() {
use tracing_subscriber::layer::SubscriberExt;
let _guard = tracing::subscriber::set_default(
tracing_subscriber::registry().with(checkpoint_cause::CauseCaptureLayer),
);
let (tmp, conn) = temp_store("core").await;
let sink = checkpoint_cause::CauseSink::new();
sink.scoped(async {
tracing::debug!(
target: "turso_core::vdbe::execute",
"PRAGMA wal_checkpoint failed: engine-side detail"
);
})
.await;
let e = sink.attach(anyhow::anyhow!(
"Unexpected result from PRAGMA wal_checkpoint"
));
record_store_failure(
FailureKind::ExitCheckpointFailure,
"exit checkpoint failure",
"core",
&conn,
&e,
Some(tmp.path()),
);
let body = std::fs::read_to_string(tmp.path().join("error.log")).unwrap();
assert!(
body.contains(
"reason: engine cause: PRAGMA wal_checkpoint failed: engine-side detail: \
Unexpected result from PRAGMA wal_checkpoint"
),
"the reason line must be the full chain — the engine's own reason and the \
product's text as separate links: {body}"
);
}
#[tokio::test]
async fn integrity_failure_is_recorded_once_then_counted() {
let name = "integrity_probe";
let (tmp, conn) = temp_store(name).await;
for _ in 0..3 {
record_integrity_failure(
name,
&conn,
&anyhow::anyhow!("the integrity check cannot run"),
Some(tmp.path()),
);
}
let body = std::fs::read_to_string(tmp.path().join("error.log")).unwrap();
assert_eq!(
body.matches("MahBot runtime integrity failure").count(),
1,
"only the first failing round may file a block: {body}"
);
for needle in [
"store: integrity_probe",
"db path:",
"reason: the integrity check cannot run",
] {
assert!(
body.contains(needle),
"error.log must contain {needle:?}: {body}"
);
}
}
#[tokio::test]
async fn the_report_names_the_measured_file_and_its_artifact_state() {
let (_tmp, conn) = temp_store("report_probe").await;
let report = store_failure_report(
FailureKind::ExitCheckpointFailure,
"report_probe",
&anyhow::anyhow!("disk gone"),
conn.db_path(),
)
.render();
for needle in [
"MahBot exit checkpoint failure",
"store: report_probe",
"reason: disk gone",
"artifact state: wal_size=",
] {
assert!(
report.contains(needle),
"the report must contain {needle:?}: {report}"
);
}
assert!(
report.contains(&format!("db path: {}", conn.db_path().display())),
"the report must name the file the connection has open: {report}"
);
}
#[tokio::test]
#[serial_test::serial(drain)] async fn failed_checkpoint_with_broken_fts_repairs_and_retries() {
crate::shutdown::drain_clear();
let tmp = tempfile::TempDir::new().unwrap();
let conn = a_store_with_a_corrupt_fts_index(tmp.path()).await;
assert!(
!crate::db::is_fts_index(&conn, crate::db::TICKETS_FTS_INDEX_NAME).await,
"index must be a btree before the recovery"
);
let before: i64 = conn
.query_row("SELECT COUNT(*) FROM tickets", (), |r| r.get::<i64>(0))
.await
.unwrap();
let mut attempt_no = 0;
periodic_checkpoint_inner(
"core",
&conn,
Some(tmp.path()),
TEST_PAUSES,
Duration::ZERO,
|| {
attempt_no += 1;
let n = attempt_no;
let conn = conn.clone();
async move {
Attempted::Ran(if n == 1 {
Err(anyhow::anyhow!("injected checkpoint failure"))
} else {
conn.checkpoint_ungated().await
})
}
},
)
.await;
assert_eq!(
attempt_no, 2,
"the failed first attempt must be retried exactly once after the repair"
);
assert!(
!crate::shutdown::is_draining(),
"a checkpoint completed by the retry must not begin the drain"
);
assert!(
!tmp.path().join("error.log").exists(),
"a checkpoint completed by the retry must not file a failure block"
);
assert!(
crate::db::is_fts_index(&conn, crate::db::TICKETS_FTS_INDEX_NAME).await,
"FTS index must be restored after the recovery"
);
let matched: String = conn
.query_row(
"SELECT id FROM tickets WHERE title MATCH ?1 LIMIT 1",
crate::db::params![crate::db::sanitize_fts_query("Important bug fix one")],
|r| r.get::<String>(0),
)
.await
.unwrap();
assert_eq!(
matched, "t-1",
"MATCH must find the known ticket after the rebuild"
);
let after: i64 = conn
.query_row("SELECT COUNT(*) FROM tickets", (), |r| r.get::<i64>(0))
.await
.unwrap();
assert_eq!(
after, before,
"ticket rows must be untouched by the FTS rebuild"
);
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_rebuilt_index_does_not_exempt_a_persistently_failing_round() {
crate::shutdown::drain_clear();
let tmp = tempfile::TempDir::new().unwrap();
let conn = a_store_with_a_corrupt_fts_index(tmp.path()).await;
let store = &conn;
let rounds = window_rounds();
let total_attempts = rounds_with_answers(
"core",
store,
tmp.path(),
rounds,
TEST_PAUSES,
Duration::ZERO,
|attempt_no| async move {
if attempt_no == 1 {
store.execute_batch(&fts_corruption_ddl()).await.unwrap();
return Attempted::Ran(Err(anyhow::anyhow!("injected checkpoint failure")));
}
Attempted::Ran(Err(anyhow::anyhow!("injected persistent failure")))
},
)
.await;
assert_eq!(
total_attempts,
CHECKPOINT_ATTEMPTS * rounds,
"a persistent failure must run every attempt of every round before the window decides"
);
let body = std::fs::read_to_string(tmp.path().join("error.log")).unwrap();
assert!(
body.contains("injected checkpoint failure"),
"the original checkpoint error must be in the report: {body}"
);
assert!(
body.contains("injected persistent failure"),
"the later attempts' errors must be in the report: {body}"
);
assert!(
body.contains("repair outcome: rebuilt after detection"),
"the report must name what the round's repair returned: {body}"
);
assert!(
crate::shutdown::is_draining(),
"a rebuilt index must not exempt a round in which no attempt completed a checkpoint"
);
crate::shutdown::drain_clear();
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_rebuilt_index_does_not_exempt_a_partial_retry() {
crate::shutdown::drain_clear();
let tmp = tempfile::TempDir::new().unwrap();
let conn = a_store_with_a_corrupt_fts_index(tmp.path()).await;
let rounds = window_rounds();
let total_attempts = rounds_with_answers(
"core",
&conn,
tmp.path(),
rounds,
TEST_PAUSES,
Duration::ZERO,
|attempt_no| {
let answer = if attempt_no == 1 {
Attempted::Ran(Err(anyhow::anyhow!("injected checkpoint failure")))
} else {
Attempted::Ran(Ok(CheckpointOutcome {
busy: false,
log_frames: 5,
checkpointed_frames: 1,
}))
};
async move { answer }
},
)
.await;
assert_eq!(
total_attempts,
CHECKPOINT_ATTEMPTS * rounds,
"a partially folded attempt must not cut a failing round's budget short"
);
let body = std::fs::read_to_string(tmp.path().join("error.log")).unwrap();
assert!(
body.contains("checkpoint error: injected checkpoint failure"),
"the genuine failure must still be recorded with its real cause: {body}"
);
assert!(
crate::shutdown::is_draining(),
"a rebuilt index must not exempt a round in which a genuine failure had no completion"
);
crate::shutdown::drain_clear();
}
async fn a_round_cured_by_its_last_attempt(conn: &Connection, root: &Path) {
let mut attempt_no = 0;
periodic_checkpoint_inner(
"core",
conn,
Some(root),
TEST_PAUSES,
Duration::ZERO,
|| {
attempt_no += 1;
let n = attempt_no;
async move {
Attempted::Ran(if n == CHECKPOINT_ATTEMPTS {
Ok(CheckpointOutcome {
busy: false,
log_frames: 0,
checkpointed_frames: 0,
})
} else {
Err(anyhow::anyhow!("injected checkpoint failure"))
})
}
},
)
.await;
assert_eq!(
attempt_no, CHECKPOINT_ATTEMPTS,
"a completion on the last attempt leaves no attempt to follow it"
);
assert!(
!crate::shutdown::is_draining(),
"a completed checkpoint after a failure must keep the service serving"
);
assert!(
!root.join("error.log").exists(),
"a round cured by its own retry must not file a failure block"
);
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_completed_retry_keeps_serving_whatever_the_repair_returned() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
assert!(
matches!(
crate::db::repair_ticket_title_fts_runtime(&conn).await,
TicketTitleFtsRuntimeRepair::NotApplicable
),
"this store must classify as a store with nothing to rebuild"
);
a_round_cured_by_its_last_attempt(&conn, tmp.path()).await;
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_failed_repair_does_not_stop_a_round_its_retry_completes() {
crate::shutdown::drain_clear();
let tmp = tempfile::TempDir::new().unwrap();
let conn = crate::db::open_with_schema(
&crate::db::store_db_path(tmp.path(), "core"),
"CREATE TABLE tickets (id INTEGER PRIMARY KEY);\
CREATE INDEX idx_tickets_title_fts ON tickets (id);",
)
.await
.unwrap();
assert!(
matches!(
crate::db::repair_ticket_title_fts_runtime(&conn).await,
TicketTitleFtsRuntimeRepair::Failed(_)
),
"this store must classify as a store whose repair fails"
);
a_round_cured_by_its_last_attempt(&conn, tmp.path()).await;
}
#[tokio::test]
#[serial_test::serial(drain)] async fn checkpoint_failure_without_fts_store_writes_error_log_and_drains() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
let rounds = window_rounds();
let attempts = failing_rounds(
"core",
&conn,
tmp.path(),
rounds,
"no space left on device",
TEST_PAUSES,
Duration::ZERO,
)
.await;
assert_eq!(
attempts,
CHECKPOINT_ATTEMPTS * rounds,
"every store gets every attempt of every failing round, indexed or not"
);
let body = std::fs::read_to_string(tmp.path().join("error.log")).unwrap();
assert!(
body.contains("checkpoint error: no space left on device"),
"the checkpoint error must be in the report: {body}"
);
assert!(
body.contains(failure_record::ENVIRONMENT_CAUSE),
"a resource-caused checkpoint failure must be marked environment-caused: {body}"
);
assert!(
crate::shutdown::is_draining(),
"persistent failure must begin the graceful drain"
);
crate::shutdown::drain_clear();
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_healthy_index_store_gets_every_attempt_and_stops() {
crate::shutdown::drain_clear();
let tmp = tempfile::TempDir::new().expect("temp dir");
let conn = crate::db::open_consolidated_store(tmp.path())
.await
.expect("open the consolidated store");
insert_fts_ticket(&conn, "t-1", "Important bug fix one").await;
let rounds = window_rounds();
let attempts = failing_rounds(
"core",
&conn,
tmp.path(),
rounds,
"injected checkpoint failure",
TEST_PAUSES,
Duration::ZERO,
)
.await;
assert_eq!(
attempts,
CHECKPOINT_ATTEMPTS * rounds,
"an indexed store gets every attempt of every failing round too",
);
let body = std::fs::read_to_string(tmp.path().join("error.log")).expect("error.log");
assert!(
body.contains("repair outcome: healthy"),
"the report must name what the round's repair returned: {body}"
);
assert!(
crate::shutdown::is_draining(),
"a persistent failure on an indexed store must still stop the service",
);
crate::shutdown::drain_clear();
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_failing_round_spaces_its_retries_by_the_production_schedule() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
let schedule = CHECKPOINT_RETRY_PAUSES.map(|pause| pause / 1000);
let started = tokio::time::Instant::now();
let attempts = failing_rounds(
"core",
&conn,
tmp.path(),
1,
"injected checkpoint failure",
schedule,
Duration::ZERO,
)
.await;
assert_eq!(
attempts, CHECKPOINT_ATTEMPTS,
"a failing round runs every attempt the schedule has pauses for",
);
assert!(
started.elapsed() >= schedule.into_iter().sum::<Duration>(),
"the round must spend the schedule's own pause before each retry",
);
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_partially_folded_journal_is_normal_and_ends_the_round() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
let mut attempts = 0usize;
periodic_checkpoint_inner(
"core",
&conn,
Some(tmp.path()),
TEST_PAUSES,
Duration::ZERO,
|| {
attempts += 1;
async {
Attempted::Ran(Ok(CheckpointOutcome {
busy: false,
log_frames: 5,
checkpointed_frames: 1,
}))
}
},
)
.await;
assert_eq!(
attempts, 1,
"a partially folded journal must not be retried"
);
assert!(
!crate::shutdown::is_draining(),
"a partially folded journal must not begin the drain"
);
assert!(
!tmp.path().join("error.log").exists(),
"a partially folded journal must not file a failure block"
);
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_busy_store_is_retried_and_a_still_busy_round_keeps_serving() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
let mut attempts = 0usize;
periodic_checkpoint_inner(
"core",
&conn,
Some(tmp.path()),
TEST_PAUSES,
Duration::ZERO,
|| {
attempts += 1;
async { Attempted::Ran(Err(blocked_checkpoint_error())) }
},
)
.await;
assert_eq!(
attempts, CHECKPOINT_ATTEMPTS,
"a blocked checkpoint must be retried up to the round's budget"
);
assert!(
!crate::shutdown::is_draining(),
"a round that stayed busy must keep the service serving"
);
assert!(
!tmp.path().join("error.log").exists(),
"a round that stayed busy must not file a failure block"
);
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_busy_attempt_before_a_failure_leaves_the_cause_line_unmarked() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
rounds_with_answers(
"core",
&conn,
tmp.path(),
window_rounds(),
TEST_PAUSES,
Duration::ZERO,
|attempt_no| {
let answer = if attempt_no == 1 {
Attempted::Ran(Err(blocked_checkpoint_error()))
} else {
Attempted::Ran(Err(anyhow::anyhow!("injected checkpoint failure")))
};
async move { answer }
},
)
.await;
let body = std::fs::read_to_string(tmp.path().join("error.log")).unwrap();
assert!(
body.lines()
.any(|line| line == "checkpoint error: injected checkpoint failure"),
"a busy first attempt leaves no line, so the cause line stays unmarked: {body}"
);
assert!(
body.contains("attempt 3 error: injected checkpoint failure"),
"the later failure must keep its own labelled line: {body}"
);
assert!(
crate::shutdown::is_draining(),
"a genuine failure with no completion must stop the service"
);
crate::shutdown::drain_clear();
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_failure_followed_by_busy_attempts_still_stops() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
let rounds = window_rounds();
let attempts = rounds_with_answers(
"core",
&conn,
tmp.path(),
rounds,
TEST_PAUSES,
Duration::ZERO,
|attempt_no| {
let answer = if attempt_no == 1 {
Attempted::Ran(Err(anyhow::anyhow!("injected checkpoint failure")))
} else {
Attempted::Ran(Err(blocked_checkpoint_error()))
};
async move { answer }
},
)
.await;
assert_eq!(
attempts,
CHECKPOINT_ATTEMPTS * rounds,
"a busy attempt must not cut a genuine failure's budget short — every attempt of every \
failing round must run"
);
let body = std::fs::read_to_string(tmp.path().join("error.log")).unwrap();
assert!(
body.contains("checkpoint error: injected checkpoint failure"),
"the round's real failure must be the report's reason: {body}"
);
assert!(
crate::shutdown::is_draining(),
"a genuine failure must still stop the service despite later busy attempts"
);
crate::shutdown::drain_clear();
}
#[tokio::test]
#[serial_test::serial(drain)] async fn the_report_numbers_each_failure_after_the_first() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
rounds_with_answers(
"core",
&conn,
tmp.path(),
window_rounds(),
TEST_PAUSES,
Duration::ZERO,
|attempt_no| {
let answer = match attempt_no {
1 => Attempted::Ran(Err(anyhow::anyhow!("injected first failure"))),
2 => Attempted::Ran(Err(blocked_checkpoint_error())),
_ => Attempted::Ran(Err(anyhow::anyhow!("injected third failure"))),
};
async move { answer }
},
)
.await;
let body = std::fs::read_to_string(tmp.path().join("error.log")).unwrap();
assert!(
body.lines()
.any(|line| line == "checkpoint error: injected first failure"),
"attempt 1's cause line keeps its long-standing byte-for-byte shape: {body}"
);
assert!(
body.contains("attempt 3 error: injected third failure"),
"the later error must carry the attempt that produced it: {body}"
);
assert!(
!body.contains("attempt 2 error:"),
"the busy/partial attempt produced no error to record: {body}"
);
assert!(
crate::shutdown::is_draining(),
"the accumulated failure must still stop the service"
);
crate::shutdown::drain_clear();
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_panicking_attempt_is_a_failed_attempt_and_stops_after_the_window() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
let rounds = window_rounds();
let attempts = rounds_with_answers(
"core",
&conn,
tmp.path(),
rounds,
TEST_PAUSES,
Duration::ZERO,
|_| async {
Attempted::Ran(guarded_checkpoint(async { panic!("injected attempt panic") }).await)
},
)
.await;
assert_eq!(
attempts,
CHECKPOINT_ATTEMPTS * rounds,
"a panic is a failed attempt, so every attempt of every round must run"
);
assert!(
crate::shutdown::is_draining(),
"a panic must not escape the round without beginning the drain"
);
let body = std::fs::read_to_string(tmp.path().join("error.log")).unwrap();
assert!(
body.contains("checkpoint attempt panicked"),
"the panic must be filed as the failed attempt's reason: {body}"
);
assert!(
body.contains("injected attempt panic"),
"the panicking attempt's message must be in the report: {body}"
);
crate::shutdown::drain_clear();
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_refused_shrink_is_not_a_failure_and_not_an_attempt() {
crate::shutdown::drain_clear();
let name = "refused_shrink_probe";
let tmp = tempfile::TempDir::new().unwrap();
let db_path = crate::db::store_db_path(tmp.path(), name);
build_zero_page_count_store(&db_path).await;
let conn = Connection::open(&db_path)
.await
.expect("open the fixture store");
let path = tmp.path();
let conn_ref = &conn;
let mut attempts = 0usize;
periodic_checkpoint_inner(name, &conn, Some(path), TEST_PAUSES, Duration::ZERO, || {
attempts += 1;
checkpoint_attempt(name, conn_ref, true, Some(path))
})
.await;
assert_eq!(attempts, 1, "a refusal must not be retried");
let body = std::fs::read_to_string(tmp.path().join("error.log"))
.expect("the refusal must be filed in the round's root");
assert!(
body.contains("MahBot store shrink refused"),
"the gate must refuse a shrink on a zero-page-count store and record it: {body}"
);
assert!(
!body.contains("MahBot checkpoint failure"),
"a refused shrink must not file a checkpoint failure: {body}"
);
assert!(
!crate::shutdown::is_draining(),
"a refused shrink must not begin the drain"
);
}
#[tokio::test]
#[serial_test::serial(drain)] async fn an_exit_round_refusal_is_recorded_and_is_not_an_exit_failure() {
crate::shutdown::drain_clear();
let name = "exit_refused_shrink_probe";
let tmp = tempfile::TempDir::new().unwrap();
let db_path = crate::db::store_db_path(tmp.path(), name);
build_zero_page_count_store(&db_path).await;
let conn = Connection::open(&db_path)
.await
.expect("open the fixture store");
let before = std::fs::read(&db_path).unwrap();
exit_checkpoint(name, &conn, true, Some(tmp.path())).await;
let body = std::fs::read_to_string(tmp.path().join("error.log"))
.expect("the refusal must be filed in the round's root");
assert!(
body.contains("MahBot store shrink refused"),
"the exit round must record a refused shrink: {body}"
);
assert!(
!body.contains("MahBot exit checkpoint failure"),
"a refused shrink is not an exit-round checkpoint failure: {body}"
);
assert!(
!crate::shutdown::is_draining(),
"the exit round never drains — the process is already exiting"
);
assert_eq!(
std::fs::read(&db_path).unwrap(),
before,
"the refused reclaiming checkpoint must not shrink the store"
);
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_refusal_on_a_retry_still_stops_on_the_rounds_real_failure() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
let rounds = window_rounds();
let attempts = rounds_with_answers(
"core",
&conn,
tmp.path(),
rounds,
TEST_PAUSES,
Duration::ZERO,
|attempt_no| {
let answer = if attempt_no == 1 {
Attempted::Ran(Err(anyhow::anyhow!("injected checkpoint failure")))
} else {
Attempted::ShrinkRefused
};
async move { answer }
},
)
.await;
assert_eq!(
attempts,
2 * rounds,
"every failing round must stop at the refusal instead of retrying forever — the total \
is the failure and the refusal of each round"
);
let body = std::fs::read_to_string(tmp.path().join("error.log")).unwrap();
assert!(
body.contains("checkpoint error: injected checkpoint failure"),
"the report's reason must be the round's real failure: {body}"
);
assert!(
body.contains("MahBot checkpoint failure"),
"the round's stop must be filed as a checkpoint failure — a refusal is \
not one and is never written as one: {body}"
);
assert!(
crate::shutdown::is_draining(),
"the round's real failure must still stop the service"
);
crate::shutdown::drain_clear();
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_single_failing_round_warns_and_never_stops() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
let attempts = failing_rounds(
"core",
&conn,
tmp.path(),
1,
"injected checkpoint failure",
TEST_PAUSES,
Duration::ZERO,
)
.await;
assert_eq!(
attempts, CHECKPOINT_ATTEMPTS,
"one failing round must still make its whole attempt budget"
);
assert!(
!crate::shutdown::is_draining(),
"one failing round must never stop the service"
);
assert!(
window_open(&conn),
"one failing round must open the store's failure window"
);
assert!(
!tmp.path().join("error.log").exists(),
"a window that is not exhausted must file no failure block"
);
}
#[tokio::test]
#[serial_test::serial(drain)] async fn the_second_consecutive_failing_round_stops_with_the_windows_cumulative_facts() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
let rounds = window_rounds();
let attempts = failing_rounds(
"core",
&conn,
tmp.path(),
rounds,
"injected checkpoint failure",
TEST_PAUSES,
Duration::ZERO,
)
.await;
assert_eq!(
attempts,
CHECKPOINT_ATTEMPTS * rounds,
"a persistent failure must run every attempt of every round before the window decides"
);
assert!(
crate::shutdown::is_draining(),
"an exhausted window must stop the service"
);
let body = std::fs::read_to_string(tmp.path().join("error.log")).expect("error.log");
for needle in [
"MahBot checkpoint failure",
"checkpoint error: injected checkpoint failure",
"retry policy:",
] {
assert!(
body.contains(needle),
"error.log must contain {needle:?}: {body}"
);
}
assert!(
!body.contains("window first failure error:"),
"the opening cause is not repeated when it is the deciding round's own: {body}"
);
let summary = format!(
"failure window: {} attempts (blocked and partial ones included) over {} consecutive \
failing rounds",
CHECKPOINT_ATTEMPTS * rounds,
rounds,
);
assert!(
body.contains(&summary),
"the record must carry the window's cumulative attempts and rounds — {summary:?}: {body}"
);
assert!(
!window_open(&conn),
"the exhausted window is spent: it must not be left for a later round to re-file"
);
crate::shutdown::drain_clear();
}
#[tokio::test]
#[serial_test::serial(drain)] async fn the_window_records_the_cause_that_opened_it() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
let (retry_pauses, span) = (TEST_PAUSES, Duration::ZERO);
failing_rounds(
"core",
&conn,
tmp.path(),
1,
"injected first failure",
retry_pauses,
span,
)
.await;
failing_rounds(
"core",
&conn,
tmp.path(),
1,
"injected second failure",
retry_pauses,
span,
)
.await;
assert!(
crate::shutdown::is_draining(),
"the exhausted window must stop the service"
);
let body = std::fs::read_to_string(tmp.path().join("error.log")).expect("error.log");
for needle in [
"checkpoint error: injected second failure",
"window first failure error: injected first failure",
] {
assert!(
body.contains(needle),
"error.log must contain {needle:?}: {body}"
);
}
crate::shutdown::drain_clear();
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_completed_checkpoint_closes_the_window_and_the_next_failure_starts_a_new_one() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
let (retry_pauses, span) = (TEST_PAUSES, Duration::ZERO);
failing_rounds(
"core",
&conn,
tmp.path(),
1,
"injected checkpoint failure",
retry_pauses,
span,
)
.await;
assert!(
window_open(&conn),
"the failing round must open the store's window"
);
periodic_checkpoint_inner("core", &conn, Some(tmp.path()), retry_pauses, span, || {
let conn = conn.clone();
async move { Attempted::Ran(conn.checkpoint_ungated().await) }
})
.await;
assert!(
!window_open(&conn),
"a completed checkpoint must close the window"
);
assert!(
!crate::shutdown::is_draining(),
"a completion must keep the service serving"
);
assert!(
!tmp.path().join("error.log").exists(),
"a completion must not file a failure block"
);
failing_rounds(
"core",
&conn,
tmp.path(),
1,
"injected checkpoint failure",
retry_pauses,
span,
)
.await;
assert!(
window_open(&conn),
"the next failure must open a fresh window"
);
assert!(
!crate::shutdown::is_draining(),
"a fresh window is one failing round, so the service must keep serving"
);
assert!(
!tmp.path().join("error.log").exists(),
"a fresh, unexhausted window must file no failure block"
);
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_round_that_did_not_fail_closes_the_window() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
let (retry_pauses, span) = (TEST_PAUSES, Duration::ZERO);
failing_rounds(
"core",
&conn,
tmp.path(),
1,
"injected checkpoint failure",
retry_pauses,
span,
)
.await;
assert!(
window_open(&conn),
"the failing round must open the store's window"
);
periodic_checkpoint_inner(
"core",
&conn,
Some(tmp.path()),
retry_pauses,
span,
|| async { Attempted::Ran(Err(blocked_checkpoint_error())) },
)
.await;
assert!(
!window_open(&conn),
"a round that ended without a genuine failure must close the window"
);
failing_rounds(
"core",
&conn,
tmp.path(),
1,
"injected checkpoint failure",
retry_pauses,
span,
)
.await;
assert!(
!crate::shutdown::is_draining(),
"failing rounds separated by a round that did not fail are not consecutive, so the \
fresh window must not stop the service"
);
assert!(
!tmp.path().join("error.log").exists(),
"an unexhausted window must file no failure block"
);
}
#[tokio::test]
#[serial_test::serial(drain)] async fn the_span_floor_alone_blocks_the_stop_the_round_floor_would_decide() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
let rounds = window_rounds() + 1;
failing_rounds(
"core",
&conn,
tmp.path(),
rounds,
"injected checkpoint failure",
TEST_PAUSES,
Duration::from_secs(3600),
)
.await;
assert!(
!crate::shutdown::is_draining(),
"rounds past the round floor inside an unelapsed span must keep the service serving"
);
assert!(
window_open(&conn),
"the window must stay open until it has spanned both floors"
);
assert!(
!tmp.path().join("error.log").exists(),
"a window that has not spanned both floors must file no failure block"
);
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_window_that_outlives_its_span_stops_on_the_next_failing_round() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
let (retry_pauses, span) = (TEST_PAUSES, Duration::from_millis(100));
failing_rounds(
"core",
&conn,
tmp.path(),
1,
"injected checkpoint failure",
retry_pauses,
span,
)
.await;
assert!(
!crate::shutdown::is_draining(),
"one failing round must never stop the service, whatever the span"
);
tokio::time::sleep(Duration::from_millis(150)).await;
failing_rounds(
"core",
&conn,
tmp.path(),
1,
"injected checkpoint failure",
retry_pauses,
span,
)
.await;
assert!(
crate::shutdown::is_draining(),
"the window spans both floors once the span has really elapsed"
);
crate::shutdown::drain_clear();
}
#[tokio::test]
#[serial_test::serial(drain)] async fn a_round_cut_short_by_a_drain_decides_nothing() {
crate::shutdown::drain_clear();
let (tmp, conn) = temp_store("core").await;
let (mut retry_pauses, span) = (TEST_PAUSES, Duration::ZERO);
retry_pauses[0] = Duration::from_millis(10);
crate::shutdown::drain_begin();
let attempts = failing_rounds(
"core",
&conn,
tmp.path(),
1,
"injected checkpoint failure",
retry_pauses,
span,
)
.await;
assert_eq!(
attempts, 1,
"the drain must cut the spaced retry short after the round's first attempt"
);
assert!(
!window_open(&conn),
"an interrupted round is not a failing one, so the window must stay untouched"
);
assert!(
!tmp.path().join("error.log").exists(),
"a round that decided nothing must file no failure block"
);
crate::shutdown::drain_clear();
}
fn assert_refusal_recorded(root: &Path) {
let body = std::fs::read_to_string(root.join("error.log"))
.expect("the refusal must be filed in the round's root");
assert!(
body.contains("MahBot store shrink refused"),
"the gate must refuse the shrink and record it: {body}"
);
for forbidden in [
"MahBot checkpoint failure",
"MahBot exit checkpoint failure",
] {
assert!(
!body.contains(forbidden),
"a refused shrink must not read as {forbidden:?}: {body}"
);
}
assert!(
!crate::shutdown::is_draining(),
"a refused shrink must not begin the drain",
);
}
fn assert_store_untouched_by_the_refusal(
db_path: &Path,
before_main: &[u8],
before_wal: &[u8],
verdict: StoreShape,
) {
foreign_client_is_refused(db_path);
assert_eq!(
crate::db::wal_guard::classify_store_shape(db_path),
verdict,
"a refused shrink must not change what counts as a damaged store",
);
assert_eq!(
std::fs::read(db_path).expect("read the fixture main file"),
before_main,
"a refused shrink must not change a byte of the store's main file",
);
assert_eq!(
std::fs::read(crate::db::wal_path(db_path)).expect("read the fixture journal"),
before_wal,
"a refused shrink must not change a byte of the store's journal",
);
}
async fn refused_state_lets_no_shrink_through(tmp: &Path, db_path: &Path, state: &str) {
let before_main = std::fs::read(db_path).expect("read the fixture main file");
let before_wal = std::fs::read(crate::db::wal_path(db_path)).expect("read the journal");
assert!(
!before_main.is_empty(),
"the fixture must hold data in its main file",
);
let verdict = crate::db::wal_guard::classify_store_shape(db_path);
assert_eq!(
verdict,
StoreShape::Present,
"the gate must add no file-level defect to a store it refuses",
);
let conn = Connection::open(db_path)
.await
.expect("open the fixture store");
assert_eq!(page_count(&conn).await, 0, "{state} must declare no pages");
assert_eq!(
fixture_rows(&conn).await,
1,
"the fixture's rows must stay readable",
);
periodic_checkpoint("core", &conn, true, Some(tmp)).await;
assert_refusal_recorded(tmp);
assert_eq!(
fixture_rows(&conn).await,
1,
"the refused round must leave the store serving",
);
assert_store_untouched_by_the_refusal(db_path, &before_main, &before_wal, verdict);
conn.checkpoint_ungated()
.await
.expect("run the ungated reclaiming checkpoint");
let main_len = std::fs::metadata(db_path)
.expect("stat the fixture main file")
.len();
assert_eq!(
main_len, 0,
"an ungated reclaiming checkpoint must truncate the main file to nothing, \
left {main_len} bytes",
);
}
#[tokio::test]
async fn a_healthy_store_still_runs_its_reclaiming_checkpoint() {
let tmp = tempfile::TempDir::new().expect("temp dir");
let db_path = crate::db::store_db_path(tmp.path(), "core");
let conn =
crate::db::open_with_schema(&db_path, "CREATE TABLE plain (id INTEGER PRIMARY KEY);")
.await
.expect("open the healthy store");
conn.execute("INSERT INTO plain (id) VALUES (1)", ())
.await
.expect("commit a row so the journal holds frames");
let wal = crate::db::wal_path(&db_path);
let before = crate::db::wal_guard::stat_size(&wal).expect("stat the store's journal");
assert!(
before > crate::db::wal_guard::WAL_HEADER_BYTES,
"the healthy store must have frames to reclaim, found {before} bytes",
);
let Attempted::Ran(Ok(outcome)) =
checkpoint_attempt("core", &conn, true, Some(tmp.path())).await
else {
panic!("a healthy store must still run its reclaiming checkpoint");
};
assert!(
outcome.is_complete(),
"the reclaiming checkpoint must complete on a healthy store",
);
let after = crate::db::wal_guard::stat_size(&wal).expect("stat the reclaimed journal");
assert_eq!(
after, 0,
"the reclaiming mode must still reclaim the journal, left {after} bytes",
);
}
#[tokio::test]
#[serial_test::serial(drain)] async fn empty_main_file_with_a_journal_is_refused_at_boot_and_never_reaches_a_round() {
crate::shutdown::drain_clear();
let tmp = tempfile::TempDir::new().expect("temp dir");
let db_path = crate::db::store_db_path(tmp.path(), "core");
let wal = crate::db::wal_path(&db_path);
build_empty_main_file_store(&db_path).await;
let journal_bytes = std::fs::metadata(&wal)
.expect("stat the fixture journal")
.len();
assert!(
journal_bytes > crate::db::wal_guard::WAL_HEADER_BYTES,
"the fixture must leave committed frames in its journal, found \
{journal_bytes} bytes",
);
assert_eq!(
crate::db::wal_guard::classify_store_shape(&db_path),
StoreShape::Unusable(ShapeDefect::ZeroBytes),
"an empty main file must stay a boot refusal",
);
let conn = Connection::open(&db_path)
.await
.expect("open the fixture store");
assert_eq!(
page_count(&conn).await,
0,
"the engine must answer no pages for the empty main file",
);
assert_eq!(
std::fs::metadata(&wal).expect("stat the journal").len(),
0,
"the engine must discard the orphan journal as it opens the empty store",
);
periodic_checkpoint("core", &conn, true, Some(tmp.path())).await;
assert!(
!tmp.path().join("error.log").exists(),
"a round with nothing left to destroy files no refusal",
);
}
#[tokio::test]
#[serial_test::serial(drain)] async fn zero_page_count_store_is_refused_end_to_end() {
crate::shutdown::drain_clear();
let tmp = tempfile::TempDir::new().expect("temp dir");
let db_path = crate::db::store_db_path(tmp.path(), "core");
build_zero_page_count_store(&db_path).await;
let state = "the fixture's page-1 header";
refused_state_lets_no_shrink_through(tmp.path(), &db_path, state).await;
}
#[tokio::test]
#[serial_test::serial(drain)] async fn zero_page_count_in_the_journal_is_refused_end_to_end() {
crate::shutdown::drain_clear();
let tmp = tempfile::TempDir::new().expect("temp dir");
let db_path = crate::db::store_db_path(tmp.path(), "core");
build_zero_page_count_wal_store(&db_path).await;
let state = "the journal's newest page-1 image";
refused_state_lets_no_shrink_through(tmp.path(), &db_path, state).await;
}
}