use std::env;
use std::sync::Arc;
use std::sync::atomic::Ordering;
use std::thread::{self, JoinHandle};
use std::time::Duration;
use crate::{SHUTDOWN_FLAG, ServerState};
const DEFAULT_NAPTIME_MS: u64 = 1000;
const MIN_NAPTIME_MS: u64 = 10;
const SHUTDOWN_POLL_CAP: Duration = Duration::from_millis(50);
pub(crate) fn naptime_from_env() -> Duration {
let ms = env::var("SPG_AUTOVACUUM_NAPTIME_MS")
.ok()
.and_then(|s| s.parse::<u64>().ok())
.filter(|&n| n >= MIN_NAPTIME_MS)
.unwrap_or(DEFAULT_NAPTIME_MS);
Duration::from_millis(ms)
}
pub(crate) fn spawn(state: Arc<ServerState>) -> JoinHandle<()> {
let naptime = naptime_from_env();
thread::Builder::new()
.name("spg-autovacuum".into())
.spawn(move || run(&state, naptime))
.expect("spawn autovacuum thread")
}
fn run(state: &ServerState, naptime: Duration) {
let mut slept = Duration::ZERO;
loop {
if SHUTDOWN_FLAG.load(Ordering::Acquire) {
break;
}
if slept < naptime {
let remaining = naptime.saturating_sub(slept);
let chunk = remaining.min(SHUTDOWN_POLL_CAP);
thread::sleep(chunk);
slept += chunk;
continue;
}
slept = Duration::ZERO;
let Ok(mut engine) = state.engine.write() else {
break;
};
let vacuumed = engine.autovacuum_tick();
drop(engine);
if vacuumed > 0 {
state
.metrics
.autovacuum_ticks
.fetch_add(vacuumed as u64, Ordering::Relaxed);
}
}
}