use std::sync::{Arc, Mutex};
use super::model::{ProcStatus, Session, SessionLifecycle, Store, JOB_FAIL_STREAK_CAP};
pub const SUPERVISOR_INTERVAL_SECS: u64 = 15;
fn job_backoff_secs(restarts_done: u32, salt: u64) -> u64 {
let initial = std::env::var("SCSH_JOB_BACKOFF_INITIAL_SECS")
.ok()
.and_then(|s| s.trim().parse().ok())
.filter(|n| *n >= 1)
.unwrap_or(5 * 60);
crate::failure::RetryPolicy {
max_retries: u32::MAX,
budget_secs: u64::MAX,
backoff_initial_secs: initial,
backoff_cap_secs: (60 * 60).max(initial),
signature_cap: 1,
}
.backoff_delay_secs(restarts_done, salt)
}
fn salt_of(id: &str) -> u64 {
use std::hash::{Hash, Hasher};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
id.hash(&mut hasher);
hasher.finish()
}
fn job_failure_signature(s: &Session) -> String {
s.procs
.iter()
.find(|p| p.status == ProcStatus::Fail && !s.proc_is_superseded(p))
.map(|p| format!("{}|{}", p.skill_name.as_deref().unwrap_or(&p.label), p.fail_reason.as_deref().unwrap_or("")))
.unwrap_or_else(|| "(no failed proc)".to_string())
}
pub fn schedule_pass(store: &mut Store, now: u64) -> Vec<String> {
let mut dirty = Vec::new();
let ids: Vec<String> = store.sessions.keys().cloned().collect();
for id in ids {
let Some(s) = store.sessions.get_mut(&id) else { continue };
if !s.supervisor.supervised()
|| s.supervisor.gave_up.is_some()
|| s.supervisor.restarted_as.is_some()
|| s.supervisor.next_retry_at.is_some()
|| s.parent_session.is_some()
|| s.profile.is_none()
|| s.repo == super::server::IMAGE_BUILDS_REPO
|| s.repo == crate::daemon::INTERNAL_REPO
{
continue;
}
if s.lifecycle_status(now) != SessionLifecycle::Failed {
continue;
}
let attempt = s.supervisor.attempt();
let max = s.supervisor.retries;
if attempt > max {
let why = format!("retries budget exhausted ({max} restarts)");
s.supervisor.gave_up = Some(why.clone());
crate::failure::log_session_proc(&id, "supervisor_gave_up", "(supervisor)", &why);
dirty.push(id);
continue;
}
let signature = job_failure_signature(s);
let streak =
if s.supervisor.fail_signature.as_deref() == Some(signature.as_str()) { s.supervisor.fail_streak + 1 } else { 1 };
s.supervisor.fail_signature = Some(signature.clone());
s.supervisor.fail_streak = streak;
if streak > JOB_FAIL_STREAK_CAP {
let why = format!("{streak} consecutive runs failed identically at {signature}");
s.supervisor.gave_up = Some(why.clone());
crate::failure::log_session_proc(&id, "supervisor_gave_up", "(supervisor)", &why);
dirty.push(id);
continue;
}
let delay = job_backoff_secs(attempt.saturating_sub(1), salt_of(&id));
s.supervisor.next_retry_at = Some(now + delay);
crate::failure::log_session_proc(
&id,
"supervisor_scheduled",
"(supervisor)",
&format!("restart {attempt}/{max} scheduled after failure at {signature}; starting in {delay}s"),
);
dirty.push(id);
}
dirty
}
pub fn fire_due(store: &Arc<Mutex<Store>>, now: u64) -> Vec<String> {
let due: Vec<(String, bool)> = {
let guard = store.lock().unwrap_or_else(|e| e.into_inner());
guard
.sessions
.values()
.filter(|s| s.supervisor.supervised() && s.supervisor.gave_up.is_none() && s.supervisor.restarted_as.is_none())
.filter(|s| s.supervisor.next_retry_at.is_some_and(|at| now >= at))
.filter(|s| s.lifecycle_status(now) == SessionLifecycle::Failed)
.map(|s| (s.id.clone(), s.kind.as_deref() == Some("workflow")))
.collect()
};
let mut dirty = Vec::new();
for (id, workflow) in due {
let mode = if workflow { "resume" } else { "scratch" };
let body = format!("{{\"session\":{},\"mode\":\"{mode}\"}}", crate::json::quote(&id));
let (status, out, _) = super::server::jobs_restart_response(&body, store);
let mut guard = store.lock().unwrap_or_else(|e| e.into_inner());
let Some(s) = guard.sessions.get_mut(&id) else { continue };
if status == 200 {
s.supervisor.next_retry_at = None;
crate::failure::log_session_proc(
&id,
"supervisor_restart",
"(supervisor)",
&format!(
"restart {} of {} continued as {} ({mode})",
s.supervisor.attempt(),
s.supervisor.retries,
s.supervisor.restarted_as.as_deref().unwrap_or("?")
),
);
} else if status == 409 {
s.supervisor.next_retry_at = Some(now + 60);
} else {
let why = format!("restart refused (HTTP {status}): {out}");
s.supervisor.gave_up = Some(why.clone());
s.supervisor.next_retry_at = None;
crate::failure::log_session_proc(&id, "supervisor_gave_up", "(supervisor)", &why);
}
dirty.push(id);
}
dirty
}
#[cfg(test)]
mod tests {
use super::*;
use crate::daemon::model::{DaemonMode, ProcKind, ProcRecord, SupervisorState, DEFAULT_JOB_RETRIES};
fn failed_supervised_session(id: &str, now: u64) -> Session {
Session {
id: id.into(),
started_at: now - 100,
ended_at: Some(now - 10),
profile: Some("greet".into()),
kind: Some("workflow".into()),
repo: "/r".into(),
branch: "main".into(),
skills: Vec::new(),
procs: vec![ProcRecord {
index: 0,
previous_attempt: None,
label: "claude: fix".into(),
kind: ProcKind::Skill,
status: ProcStatus::Fail,
skill_name: Some("fix".into()),
harness: Some("claude".into()),
model: None,
started_at: Some(now - 50),
note: None,
detail: None,
fail_reason: Some(crate::failure::reason::HARNESS_NONZERO.into()),
elapsed: Some(5.0),
lines: Vec::new(),
container_name: None,
container_runtime: None,
cast_path: None,
diff_path: None,
skill_source: None,
route: None,
result_path: None,
annotate_target: None,
}],
last_seen_at: now - 10,
client_connected: false,
run_pid: None,
workflow: None,
parent_session: None,
supervisor: SupervisorState::fresh(DEFAULT_JOB_RETRIES),
}
}
#[test]
fn schedule_pass_schedules_failed_supervised_jobs_with_backoff() {
let now = 1_000_000;
let mut store = Store::new(DaemonMode::Persistent, 7274, now);
store.insert_session("aaaaaa".into(), failed_supervised_session("aaaaaa", now));
let dirty = schedule_pass(&mut store, now);
assert_eq!(dirty, vec!["aaaaaa".to_string()]);
let s = &store.sessions["aaaaaa"];
let at = s.supervisor.next_retry_at.expect("scheduled");
assert!(at >= now + 4 * 60 && at <= now + 6 * 60, "5m±20%, got +{}s", at - now);
assert_eq!(s.supervisor.fail_signature.as_deref(), Some("fix|harness_nonzero_exit"));
assert_eq!(s.supervisor.fail_streak, 1);
assert!(schedule_pass(&mut store, now + 1).is_empty());
}
#[test]
fn schedule_pass_ignores_zero_budget_cancelled_and_settled_jobs() {
let now = 1_000_000;
let mut store = Store::new(DaemonMode::Persistent, 7274, now);
let mut opted_out = failed_supervised_session("attend", now);
opted_out.supervisor = Default::default();
store.insert_session("attend".into(), opted_out);
let mut cancelled = failed_supervised_session("cancel", now);
cancelled.procs[0].fail_reason = Some(crate::failure::reason::FORCE_STOPPED.into());
store.insert_session("cancel".into(), cancelled);
let mut chained = failed_supervised_session("chained", now);
chained.supervisor.restarted_as = Some("newone".into());
store.insert_session("chained".into(), chained);
let mut done = failed_supervised_session("gaveup", now);
done.supervisor.gave_up = Some("budget".into());
store.insert_session("gaveup".into(), done);
assert!(schedule_pass(&mut store, now).is_empty());
}
#[test]
fn breaker_trips_after_identical_failures_and_budget_after_max_attempts() {
let now = 1_000_000;
let mut store = Store::new(DaemonMode::Persistent, 7274, now);
let mut identical = failed_supervised_session("identic", now);
identical.supervisor.fail_signature = Some("fix|harness_nonzero_exit".into());
identical.supervisor.fail_streak = JOB_FAIL_STREAK_CAP;
identical.supervisor.job_attempt = 5;
store.insert_session("identic".into(), identical);
schedule_pass(&mut store, now);
let s = &store.sessions["identic"];
assert!(s.supervisor.gave_up.as_deref().is_some_and(|w| w.contains("identically")), "{:?}", s.supervisor.gave_up);
assert!(s.supervisor.next_retry_at.is_none());
let mut different = failed_supervised_session("differs", now);
different.supervisor.fail_signature = Some("prepare|container_timeout".into());
different.supervisor.fail_streak = JOB_FAIL_STREAK_CAP;
store.insert_session("differs".into(), different);
schedule_pass(&mut store, now);
let s = &store.sessions["differs"];
assert_eq!(s.supervisor.fail_streak, 1, "a new signature resets the streak");
assert!(s.supervisor.next_retry_at.is_some());
let mut spent = failed_supervised_session("spentit", now);
spent.supervisor.job_attempt = DEFAULT_JOB_RETRIES + 1;
store.insert_session("spentit".into(), spent);
schedule_pass(&mut store, now);
let s = &store.sessions["spentit"];
assert!(
s.supervisor.gave_up.as_deref().is_some_and(|w| w.contains("budget exhausted")),
"{:?}",
s.supervisor.gave_up
);
}
#[test]
fn job_backoff_doubles_from_five_minutes_and_caps_at_an_hour() {
for (n, base) in [(0u32, 300u64), (1, 600), (2, 1200), (3, 2400), (4, 3600), (9, 3600)] {
let d = job_backoff_secs(n, 42);
let span = base / 5;
assert!(d >= base - span && d <= base + span, "restart {n}: {d}s outside {base}±{span}");
}
}
}