use actl_core::state::{SignalPaths, unix_ms};
const PROBE_INTERVAL_MS: u64 = 5000;
const PROBE_GIVE_UP_MS: u64 = 10 * 60 * 1000;
const AUTO_RESUME_BUDGET: u32 = 2;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Verdict {
Retriable,
AlreadyDone,
Pending,
}
struct ProbeState {
last_ms: u64,
verdict: Verdict,
}
static PROBE: std::sync::Mutex<Option<(String, ProbeState)>> = std::sync::Mutex::new(None);
pub fn verdict(before_ok: Option<bool>, verify_ok: Option<bool>) -> Verdict {
match (before_ok, verify_ok) {
(_, Some(true)) => Verdict::AlreadyDone,
(Some(true), _) => Verdict::Retriable,
_ => Verdict::Pending,
}
}
fn newest_probe_target(paths: &SignalPaths) -> Option<String> {
let now = unix_ms();
let dir = std::fs::read_dir(paths.dir.join("flows")).ok()?;
let mut best: Option<(u64, String)> = None;
for entry in dir.flatten() {
let id = entry.file_name().into_string().ok()?;
let Some(state) = super::flow_ui::read(paths, &id) else {
continue;
};
if state["status"] != "paused" || state["reason"] != "action_failed" {
continue;
}
let updated = state["updated_ms"].as_u64()?;
if now.saturating_sub(updated) < PROBE_GIVE_UP_MS
&& best.as_ref().is_none_or(|(t, _)| updated > *t)
{
best = Some((updated, id));
}
}
best.map(|(_, id)| id)
}
pub fn tick(
paths: &SignalPaths,
bound: Option<&str>,
status: Option<&str>,
reason: Option<&str>,
paused_since_ms: u64,
plan: Option<&serde_json::Value>,
) -> Option<(String, Verdict)> {
let eligible = bound.is_some()
&& status == Some("paused")
&& reason == Some("action_failed")
&& plan
.map(|p| p["steps"].as_array().is_some_and(|s| !s.is_empty()))
.unwrap_or(false)
&& unix_ms().saturating_sub(paused_since_ms) < PROBE_GIVE_UP_MS;
let Some(id) = bound
.filter(|_| eligible)
.map(str::to_string)
.or_else(|| newest_probe_target(paths))
else {
if let Ok(mut guard) = PROBE.lock() {
*guard = None;
}
return None;
};
let id = id.as_str();
let now = unix_ms();
let mut guard = PROBE.lock().ok()?;
let entry = guard.get_or_insert_with(|| {
(
id.to_string(),
ProbeState {
last_ms: 0,
verdict: Verdict::Pending,
},
)
});
if entry.0 != id {
*entry = (
id.to_string(),
ProbeState {
last_ms: 0,
verdict: Verdict::Pending,
},
);
}
let state = &mut entry.1;
if now.saturating_sub(state.last_ms) < PROBE_INTERVAL_MS {
return Some((id.to_string(), state.verdict));
}
state.last_ms = now;
let plan_value = match plan {
Some(p) if bound == Some(id) => p.clone(),
_ => super::ui_text::plan(paths, id)?,
};
let flow = super::flow_ui::read(paths, id)?;
let steps = flow["steps"].as_array()?;
let plan_steps = plan_value["steps"].as_array()?;
let index = steps
.iter()
.position(|s| s["status"] != "completed")
.or_else(|| steps.len().checked_sub(1))?;
let run_probe = |invocation: Option<&serde_json::Value>| -> Option<bool> {
let invocation = invocation?;
let mut args = vec![invocation["command"].as_str()?.to_string()];
if let Some(list) = invocation["args"].as_array() {
args.extend(list.iter().filter_map(|a| a.as_str().map(str::to_string)));
}
run_actl(args).ok()
};
let before_ok = run_probe(plan_steps.get(index).and_then(|s| s.get("before")));
let verify_ok = run_probe(plan_steps.get(index).and_then(|s| s.get("verify")));
state.verdict = verdict(before_ok, verify_ok);
record(
paths,
id,
"repair_probe",
&serde_json::json!({
"step_index": index,
"before_ok": before_ok,
"verify_ok": verify_ok,
"verdict": format!("{:?}", state.verdict).to_lowercase(),
}),
);
Some((id.to_string(), state.verdict))
}
pub fn record(paths: &SignalPaths, id: &str, event: &str, data: &serde_json::Value) {
use std::io::Write;
let line = serde_json::json!({
"event": event, "ts_ms": unix_ms(), "data": data,
});
let Ok(mut file) = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(paths.dir.join("flows").join(id).join("probe.jsonl"))
else {
return;
};
let _ = writeln!(file, "{line}");
}
pub fn consume_auto_resume(paths: &SignalPaths, id: &str, declared: bool, v: Verdict) -> bool {
if !declared || v != Verdict::Retriable {
return false;
}
let used = std::fs::read_to_string(paths.dir.join("flows").join(id).join("probe.jsonl"))
.map(|text| {
text.lines()
.filter(|l| l.contains("\"event\":\"auto_resume\""))
.count()
})
.unwrap_or(0);
(used as u32) < AUTO_RESUME_BUDGET
}
fn run_actl(args: Vec<String>) -> Result<bool, ()> {
use std::os::windows::process::CommandExt;
let exe = std::env::current_exe()
.map_err(|_| ())?
.with_file_name("actl.exe");
let output = std::process::Command::new(exe)
.args(&args)
.env("ACTL_SILENT_PROBE", "1")
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::null())
.creation_flags(0x08000000)
.output()
.map_err(|_| ())?;
let body: serde_json::Value = serde_json::from_slice(&output.stdout).map_err(|_| ())?;
Ok(output.status.success() && body["ok"] == true)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn verdict_prioritizes_already_done_over_retriable() {
assert_eq!(verdict(Some(true), Some(true)), Verdict::AlreadyDone);
assert_eq!(verdict(Some(false), Some(true)), Verdict::AlreadyDone);
assert_eq!(verdict(Some(true), None), Verdict::Retriable);
assert_eq!(verdict(Some(true), Some(false)), Verdict::Retriable);
assert_eq!(verdict(None, None), Verdict::Pending);
assert_eq!(verdict(Some(false), Some(false)), Verdict::Pending);
}
}