use futures_util::future::{FutureExt, join_all};
use std::future::Future;
use std::panic::AssertUnwindSafe;
use std::path::Path;
use tracing::{debug, error, info, warn};
use crate::db::{CheckpointOutcome, Connection, TicketTitleFtsRuntimeRepair};
const WAL_CHECKPOINT_CAP_BYTES: u64 = 32 * 1024 * 1024;
const DEFAULT_CHECKPOINT_MIN_FREE_BYTES: u64 = 64 * 1024 * 1024;
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)
}
#[cfg(unix)]
fn available_free_bytes(path: &Path) -> u64 {
use std::os::unix::ffi::OsStrExt;
let Ok(c_path) = std::ffi::CString::new(path.as_os_str().as_bytes()) else {
return 0;
};
let mut stats = unsafe { std::mem::zeroed::<libc::statvfs>() };
if unsafe { libc::statvfs(c_path.as_ptr(), std::ptr::addr_of_mut!(stats)) } != 0 {
return 0;
}
u64::from(stats.f_bavail).saturating_mul(stats.f_frsize)
}
#[cfg(not(unix))]
fn available_free_bytes(_path: &Path) -> u64 {
u64::MAX
}
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(Debug, Clone, Copy, PartialEq, Eq)]
enum CheckpointRound {
Exit,
Periodic,
}
async fn checkpoint_stores(round: CheckpointRound) {
let root = crate::config::CONFIG.try_storage_root();
let truncate_gate = root.as_deref().is_none_or(truncate_allowed);
let verify = matches!(round, CheckpointRound::Periodic);
let policy = match round {
CheckpointRound::Exit => CheckpointPolicy::Truncate,
CheckpointRound::Periodic => CheckpointPolicy::periodic(),
};
for_each_store(|name, conn| {
let root = root.clone();
async move {
let status = root
.as_deref()
.map(|r| crate::db::wal_guard::inspect_store(r, name));
let truncate = match policy {
CheckpointPolicy::Truncate => truncate_gate,
CheckpointPolicy::PassiveCapped(cap) => {
status.as_ref().is_some_and(|s| s.wal_size > cap) && truncate_gate
}
};
let outcome = if truncate {
conn.checkpoint().await
} else {
conn.checkpoint_passive().await
};
match outcome {
Ok(o) if o.is_complete() => debug!(
db = %name,
log = o.log_frames,
checkpointed = o.checkpointed_frames,
"Database WAL checkpointed",
),
Ok(o) => warn!(
db = %name,
busy = o.busy,
log = o.log_frames,
checkpointed = o.checkpointed_frames,
"Checkpoint busy or partial — WAL frames left uncheckpointed",
),
Err(e) => {
warn!(error = %e, db = %name, "Failed to checkpoint database WAL");
if matches!(round, CheckpointRound::Periodic) {
recover_failed_checkpoint(name, conn, &e, truncate, root.as_deref()).await;
}
}
}
if verify {
match conn.quick_check().await {
Ok(()) => debug!(db = %name, "Database integrity check passed"),
Err(e) => error!(error = %e, db = %name, "Database integrity check failed"),
}
}
}
})
.await;
}
async fn recover_failed_checkpoint(
name: &str,
conn: &Connection,
error: &anyhow::Error,
truncate: bool,
root: Option<&Path>,
) {
let retry = async move {
if truncate {
conn.checkpoint().await
} else {
conn.checkpoint_passive().await
}
};
recover_failed_checkpoint_inner(name, conn, error, root, retry).await;
}
async fn recover_failed_checkpoint_inner(
name: &str,
conn: &Connection,
error: &anyhow::Error,
root: Option<&Path>,
retry: impl Future<Output = anyhow::Result<CheckpointOutcome>>,
) -> bool {
let repair = crate::db::repair_ticket_title_fts_on_failed_checkpoint(conn).await;
let retried = match &repair {
TicketTitleFtsRuntimeRepair::Rebuilt(_) => {
Some(match AssertUnwindSafe(retry).catch_unwind().await {
Ok(result) => result,
Err(panic) => Err(anyhow::anyhow!(
"retry checkpoint panicked: {}",
crate::util::panic_message(&panic)
)),
})
}
_ => None,
};
if let Some(Ok(o)) = &retried {
info!(
db = %name,
repair = ?repair,
complete = o.is_complete(),
checkpointed = o.checkpointed_frames,
"Checkpoint recovered after FTS repair — continuing"
);
return true;
}
let report = build_failure_report(
name,
error,
&repair,
retried.as_ref().and_then(|r| r.as_ref().err()),
conn,
root,
)
.await;
match root {
Some(root) => match write_checkpoint_error_log(root, &report) {
Ok(path) => {
error!(
db = %name,
log = %path.display(),
"Checkpoint failure persists after repair — error.log written, initiating graceful shutdown"
);
}
Err(_) => {
error!(
db = %name,
"Checkpoint failure persists after repair — error.log write FAILED, initiating graceful shutdown"
);
}
},
None => {
error!(
db = %name,
"storage root unresolvable — error.log not written, initiating graceful shutdown"
);
}
}
crate::shutdown::drain_begin();
false
}
async fn build_failure_report(
name: &str,
error: &anyhow::Error,
repair: &TicketTitleFtsRuntimeRepair,
retry_error: Option<&anyhow::Error>,
conn: &Connection,
root: Option<&Path>,
) -> String {
use std::fmt::Write;
let mut body = String::new();
let _ = writeln!(
body,
"MahBot checkpoint failure — {}",
chrono::Utc::now().to_rfc3339()
);
let _ = writeln!(body, "store: {name}");
match root {
Some(r) => {
let _ = writeln!(
body,
"db path: {}",
crate::db::store_db_path(r, name).display()
);
}
None => {
let _ = writeln!(body, "db path: unresolvable (storage root unavailable)");
}
}
let _ = writeln!(body, "checkpoint error: {error:#}");
if let Some(re) = retry_error {
let _ = writeln!(body, "retry error: {re:#}");
}
let _ = writeln!(body, "repair outcome: {}", repair.summary());
match AssertUnwindSafe(conn.quick_check_problems())
.catch_unwind()
.await
{
Ok(Ok(problems)) if problems.is_empty() => {
let _ = writeln!(body, "quick_check: ok");
}
Ok(Ok(problems)) => {
let _ = writeln!(body, "quick_check problems: {}", problems.join("; "));
}
Ok(Err(e)) => {
let _ = writeln!(body, "quick_check error: {e:#}");
}
Err(_) => {
let _ = writeln!(body, "quick_check: probe panicked");
}
}
if let Some(r) = root {
let status = std::panic::catch_unwind(AssertUnwindSafe(|| {
crate::db::wal_guard::inspect_store(r, name)
}));
match status {
Ok(status) => {
let _ = writeln!(
body,
"artifact state: store={} class={:?} wal_size={} has_stale_tshm={}",
status.store, status.class, status.wal_size, status.has_stale_tshm
);
}
Err(_) => {
let _ = writeln!(body, "artifact state: inspection panicked");
}
}
}
body
}
fn write_checkpoint_error_log(root: &Path, report: &str) -> std::io::Result<std::path::PathBuf> {
use std::io::Write;
let path = root.join("error.log");
let mut file = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(&path)?;
file.write_all(format!("{report}\n").as_bytes())?;
Ok(path)
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn noop_when_no_stores() {
checkpoint_all_databases().await;
periodic_checkpoint_and_verify().await;
}
async fn insert_fts_ticket(conn: &Connection, id: &str, title: &str) {
conn.execute(
"INSERT INTO tickets (id, title, description, workspace_name, created_at, updated_at) \
VALUES (?1, ?2, 'desc', 'ws', ?3, ?3)",
crate::db::params![id.to_string(), title.to_string(), crate::db::now()],
)
.await
.unwrap();
}
#[tokio::test]
async fn failed_checkpoint_with_broken_fts_repairs_and_retries() {
let tmp = tempfile::TempDir::new().unwrap();
let conn = crate::db::open_consolidated_store(tmp.path())
.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(
"DROP INDEX idx_tickets_title_fts; \
CREATE INDEX idx_tickets_title_fts ON tickets(title);",
)
.await
.unwrap();
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 retry = conn.checkpoint();
let ok = recover_failed_checkpoint_inner(
"core",
&conn,
&anyhow::anyhow!("injected checkpoint failure"),
Some(tmp.path()),
retry,
)
.await;
assert!(
ok,
"the checkpoint must be retried and succeed after the repair"
);
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]
#[expect(clippy::await_holding_lock)] async fn repair_ran_but_checkpoint_still_fails_shuts_down() {
let _lock = crate::util::test::retry_tests_lock();
crate::shutdown::drain_clear();
let tmp = tempfile::TempDir::new().unwrap();
let conn = crate::db::open_consolidated_store(tmp.path())
.await
.unwrap();
insert_fts_ticket(&conn, "t-1", "Important bug fix one").await;
conn.execute_batch(
"DROP INDEX idx_tickets_title_fts; \
CREATE INDEX idx_tickets_title_fts ON tickets(title);",
)
.await
.unwrap();
let retry = async {
Err::<crate::db::CheckpointOutcome, anyhow::Error>(anyhow::anyhow!(
"injected persistent failure"
))
};
let ok = recover_failed_checkpoint_inner(
"core",
&conn,
&anyhow::anyhow!("injected checkpoint failure"),
Some(tmp.path()),
retry,
)
.await;
assert!(!ok, "a persistent retry failure must terminate the service");
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"
);
assert!(
body.contains("injected persistent failure"),
"the retry error must be in the report"
);
assert!(
body.contains("quick_check"),
"the report must carry a quick_check section"
);
assert!(
crate::shutdown::is_draining(),
"persistent failure must begin the graceful drain"
);
crate::shutdown::drain_clear();
}
#[tokio::test]
#[expect(clippy::await_holding_lock)] async fn checkpoint_failure_without_fts_store_writes_error_log_and_drains() {
let _lock = crate::util::test::retry_tests_lock();
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 plain (id INTEGER PRIMARY KEY);",
)
.await
.unwrap();
let retry = async { panic!("retry must not be polled for a non-FTS store") };
let ok = recover_failed_checkpoint_inner(
"core",
&conn,
&anyhow::anyhow!("injected checkpoint failure"),
Some(tmp.path()),
retry,
)
.await;
assert!(
!ok,
"a non-FTS store must still terminate on a checkpoint failure"
);
let body = std::fs::read_to_string(tmp.path().join("error.log")).unwrap();
assert!(
body.contains("injected checkpoint failure"),
"the checkpoint error must be in the report"
);
assert!(
crate::shutdown::is_draining(),
"persistent failure must begin the graceful drain"
);
crate::shutdown::drain_clear();
}
}