use std::path::Path;
use std::time::Duration;
use stackless_core::paths::Paths;
use stackless_core::state::{ReapAttempt, ReapDecision, Store};
use stackless_core::types::TcpPort;
use tokio::time::{self, MissedTickBehavior};
const TICK: Duration = Duration::from_secs(60);
pub async fn run(paths: Paths, proxy_port: TcpPort) {
let mut interval = time::interval(TICK);
interval.set_missed_tick_behavior(MissedTickBehavior::Delay);
loop {
interval.tick().await;
tick(&paths, proxy_port).await;
}
}
async fn tick(paths: &Paths, proxy_port: TcpPort) {
let worklist = plan_reaps(paths);
let exe = std::env::current_exe().map_err(|err| format!("cannot resolve binary path: {err}"));
for instance in worklist {
let outcome = match &exe {
Ok(path) => run_down(&instance, path, paths, proxy_port).await,
Err(err) => Err(err.clone()),
};
record_outcome(paths, &instance, outcome);
}
if let Ok(store) = Store::open_with_paths(paths) {
gc_tombstones(&store, paths);
}
}
fn plan_reaps(paths: &Paths) -> Vec<String> {
let Ok(store) = Store::open_with_paths(paths) else {
return Vec::new();
};
let expired = store.expired_instances().unwrap_or_default();
let now = Store::now_secs();
expired
.into_iter()
.filter(|instance| {
let lock_held = store.lock_holder_alive(instance).unwrap_or(false);
let prior = store.reap_attempt(instance).ok().flatten();
matches!(
ReapDecision::decide(now, lock_held, prior.as_ref()),
ReapDecision::Reap
)
})
.collect()
}
fn record_outcome(paths: &Paths, instance: &str, outcome: Result<(), String>) {
let Ok(store) = Store::open_with_paths(paths) else {
return;
};
match outcome {
Ok(()) => {
let _ = store.clear_reap_failure(instance);
}
Err(reason) => {
let _ = store.record_reap_failure(instance, &reason);
let attempts = store
.reap_attempt(instance)
.ok()
.flatten()
.map(|a| a.attempts)
.unwrap_or(1);
eprintln!(
"stackless reaper: reap of {instance:?} failed: {reason} \
(attempt {attempts}, retrying in {}s)",
ReapAttempt::backoff_after(attempts).as_secs()
);
}
}
}
async fn run_down(
instance: &str,
executable: &Path,
paths: &Paths,
proxy_port: TcpPort,
) -> Result<(), String> {
let port = proxy_port.get().to_string();
let output = tokio::process::Command::new(executable)
.args(["down", instance, "--json"])
.arg("--state-dir")
.arg(paths.state_dir())
.arg("--proxy-port")
.arg(&port)
.output()
.await
.map_err(|err| format!("cannot spawn `down`: {err}"))?;
if output.status.success() {
return Ok(());
}
let stdout = String::from_utf8_lossy(&output.stdout);
let reason = stdout
.lines()
.last()
.map(str::trim)
.filter(|line| !line.is_empty())
.map(str::to_owned)
.unwrap_or_else(|| format!("`down` exited with {}", output.status));
Err(reason)
}
fn gc_tombstones(store: &Store, paths: &Paths) {
for instance in store.gc_due_tombstones().unwrap_or_default() {
let logs = paths.logs_dir(&instance);
if logs.exists() {
let _ = std::fs::remove_dir_all(&logs);
}
if let Err(err) = store.delete_instance(&instance) {
eprintln!("stackless reaper: GC of tombstone {instance:?} failed: {err}");
}
}
}
pub async fn tick_once(paths: &Paths, proxy_port: TcpPort) {
tick(paths, proxy_port).await;
}
#[cfg(test)]
mod tests {
use super::*;
use stackless_core::state::TOMBSTONE_GC_WINDOW;
use std::collections::BTreeMap;
use std::time::SystemTime;
fn temp_store() -> (tempfile::TempDir, Paths, Store) {
let dir = tempfile::tempdir().expect("tempdir");
let paths = Paths::new(dir.path());
let store = Store::open_with_paths(&paths).expect("open");
(dir, paths, store)
}
const DEF: &str = "[stack]\nname = \"t\"\n[services.web]\nsource = { repo = \"https://example.invalid/x\", ref = \"main\" }\nhealth = { path = \"/\" }\n[services.web.mock]\nrun = \"true\"\n";
#[test]
fn reaper_skips_an_instance_holding_its_lock() {
let (_dir, _paths, store) = temp_store();
store
.create_instance("held", "mock", DEF, &BTreeMap::new(), "", false)
.expect("create");
store
.renew_lease("held", Duration::from_secs(0))
.expect("renew");
let _claim = store.claim_lock("held", "up").expect("claim");
assert_eq!(store.expired_instances().expect("expired"), vec!["held"]);
let lock_held = store.lock_holder_alive("held").expect("alive");
assert!(lock_held);
assert_eq!(
ReapDecision::decide(Store::now_secs(), lock_held, None),
ReapDecision::SkipLocked
);
}
#[test]
fn gc_removes_only_tombstones_past_the_window() {
let (_dir, paths, store) = temp_store();
store
.create_instance("old", "mock", DEF, &BTreeMap::new(), "", false)
.expect("create old");
store
.create_instance("recent", "mock", DEF, &BTreeMap::new(), "", false)
.expect("create recent");
store.tombstone_instance("old").expect("tombstone old");
store
.tombstone_instance("recent")
.expect("tombstone recent");
let now = SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.expect("clock")
.as_secs() as i64;
let stale = now - TOMBSTONE_GC_WINDOW.as_secs() as i64 - 1;
store
.conn_for_tests()
.execute(
"UPDATE instances SET tombstoned_at = ?1 WHERE name = 'old'",
[stale],
)
.expect("backdate");
assert_eq!(store.gc_due_tombstones().expect("due"), vec!["old"]);
gc_tombstones(&store, &paths);
assert!(store.instance("old").expect("q").is_none());
assert!(store.instance("recent").expect("q").is_some());
}
}